import { Datagram, Message, Request, Response, ResponseInit } from "./Datagram/datagram-types.js" import { CResponse, DatagramShort, isAsyncIterable, isRequest } from "./Datagram/datagram.js" import { Subscription } from "./util/Subscription.js" import { Handler, MessageHandler } from "../Independents/Server.js" export interface SinglePortMachineInit { onRequest?: Handler onMessage?: MessageHandler; } export type MultiPortMachineInit = { ports: { [port: string]: SinglePortMachineInit } } export type MachineInit = SinglePortMachineInit | MultiPortMachineInit function normalizeInit(init: MachineInit): MultiPortMachineInit { if (!isMultiportInit(init)) { if (!init.onRequest && !init.onMessage) { init = { ports: {} } } else { init = { ports: { default: init } } } } return init; } export interface IMachine { tryRequest(request: Request): Subscription | null; tryMessage(message: Message): void; } export class Machine implements IMachine { public init: MultiPortMachineInit; constructor(init: MachineInit) { this.init = normalizeInit(init); //console.log("INIT IS", this.init.ports); } install(init: MachineInit) { init = normalizeInit(init); for (let key in init.ports) { if (this.init.ports.hasOwnProperty(key)) { throw Error(`Cannot install handlers. Conflicting ports ${key}`); } } this.init.ports = { ...this.init.ports, ...init.ports }; } /** * Handle a request and return a promise or return null. * * @param request A request to process * @returns Null if this machine does not handle the request or the response promise */ tryMessage(message: Datagram): Subscription | null { if (message.from == undefined || message.from == null || message.from == "") { console.warn("lacking sender", DatagramShort(message)); throw new Error("lacking sender for datagram in tryMessage") } let machine = this.init.ports[message.port]; if (!machine) { console.warn(`Lacking installed processor on port ${message.port} at host ${message.to} (available ports are ${Object.keys(this.init.ports).join(",")})`); return null; } if (isRequest(message) && machine && machine.onRequest) { //console.log(`Processing request at host ${request.to}`); let ret = processRequestByHandler(machine.onRequest, message) //console.log("Processed response", ret); return ret; } else { if (machine.onMessage) { console.log(`Processing message at host ${message.to}`); processMessageByHandler(machine.onMessage, message as Message) } } return null; } /** * Handle a request and return a promise or return null. * * @param request A request to process * @returns Null if this machine does not handle the request or the response promise */ tryRequest(request: Request): Subscription | null { if (request.from == undefined || request.from == null || request.from == "") { console.warn("lacking sender", request); throw new Error("lacking sender for request in tryrequest") } let machine = this.init.ports[request.port]; if (!machine) { console.warn(`Lacking installed processor on port ${request.port} at host ${request.to} (available ports are ${Object.keys(this.init.ports).join(",")})`); return null; } if (machine) { if (machine.onRequest) { //console.log(`Processing request at host ${request.to}`); let ret = processRequestByHandler(machine.onRequest, request) //console.log("Processed response", ret); return ret; } else if (machine.onMessage) { processMessageByHandler(machine.onMessage, request) let ret = new Subscription(); ret.setCurrent(new CResponse({ from: request.to, to: request.from, reqno: request.reqno, status: 204 })); return ret; } } return null; } } export function isMultiportInit(machine: MachineInit): machine is MultiPortMachineInit { return (machine as MultiPortMachineInit).ports !== undefined } export function processMessageByHandler(handler: MessageHandler, message: Message): void { let subs = new Subscription(); handler(message); } export function processRequestByHandler(handler: Handler, request: Request): Subscription { let subs = new Subscription(); // return new Promise(async (resolve: (value: CResponse) => void) => { let responseInit = handler(request); //console.warn(`Processor returned`, responseInit); let pushResponse = async (ri: ResponseInit) => { if (isAsyncIterable(ri)) { // console.error("subscription result", rio.body); for await (let d of ri) { let rio = CResponse.initToInitObj(ri); let resp = new CResponse({ headers: rio.headers, body: d, status: rio.status, final: false }, request); subs.setCurrent(resp) } } else { let rio = CResponse.initToInitObj(ri); if (isAsyncIterable(rio.body)) { for await (let d of rio.body) { let resp = new CResponse({ headers: rio.headers, body: d, status: rio.status, final: false }, request); subs.setCurrent(resp) } } else { //console.warn("DIRECT RESPONSE IS", JSON.stringify(rio, null, 3)); if (rio.final == undefined) { rio.final = true; } subs.setCurrent(new CResponse(rio, request)); //subs.setCurrent(new CResponse({ headers: rio.headers, body: ri, status: rio.status, final: true }, request)); } } // resolve(new CResponse(rio, request)); } if (responseInit instanceof Promise) { responseInit.then((resp: ResponseInit) => { pushResponse(resp); }) } else { pushResponse(responseInit); } return subs; }