diff --git a/packages/orchestrator/src/index.ts b/packages/orchestrator/src/index.ts index 67ecbe7a..196332d8 100644 --- a/packages/orchestrator/src/index.ts +++ b/packages/orchestrator/src/index.ts @@ -1,3 +1,4 @@ export * from "./config.ts"; export * from "./ipc/client.ts"; export * from "./ipc/protocol.ts"; +export * from "./ipc/server.ts"; diff --git a/packages/orchestrator/src/ipc/protocol.ts b/packages/orchestrator/src/ipc/protocol.ts index 7c275279..ebb2815f 100644 --- a/packages/orchestrator/src/ipc/protocol.ts +++ b/packages/orchestrator/src/ipc/protocol.ts @@ -20,7 +20,14 @@ export interface StatusRequest { instanceId: string; } -export type OrchestratorRequest = SpawnRequest | ListRequest | StopRequest | StatusRequest; +export interface RequestMap { + spawn: SpawnRequest; + list: ListRequest; + stop: StopRequest; + status: StatusRequest; +} + +export type OrchestratorRequest = RequestMap[keyof RequestMap]; export interface InstanceSummary { id: string; @@ -31,7 +38,6 @@ export interface InstanceSummary { } export interface ResponseBase { - type: string; ok: boolean; error?: string; } @@ -56,7 +62,26 @@ export interface StatusResponse extends ResponseBase { instance?: InstanceSummary; } -export type OrchestratorResponse = SpawnResponse | ListResponse | StopResponse | StatusResponse; +export interface ErrorResponse extends ResponseBase { + type: "error"; + ok: false; + error: string; +} + +export interface ResponseMap { + spawn: SpawnResponse; + list: ListResponse; + stop: StopResponse; + status: StatusResponse; +} + +export type OrchestratorResponse = ResponseMap[keyof ResponseMap] | ErrorResponse; + +export type ResponseFor = T extends { type: infer K } + ? K extends keyof ResponseMap + ? ResponseMap[K] | ErrorResponse + : ErrorResponse + : ErrorResponse; export function encodeMessage(message: OrchestratorRequest | OrchestratorResponse): string { return `${JSON.stringify(message)}\n`; diff --git a/packages/orchestrator/src/ipc/server.ts b/packages/orchestrator/src/ipc/server.ts new file mode 100644 index 00000000..59d91b57 --- /dev/null +++ b/packages/orchestrator/src/ipc/server.ts @@ -0,0 +1,107 @@ +import { existsSync, unlinkSync } from "node:fs"; +import { createConnection, createServer, type Server } from "node:net"; +import { getSocketPath } from "../config.ts"; +import { + type ErrorResponse, + encodeMessage, + type OrchestratorRequest, + parseRequestLine, + type ResponseFor, +} from "./protocol.ts"; + +export type IpcRequestHandler = (request: T) => Promise> | ResponseFor; + +export async function startIpcServer(handler: IpcRequestHandler): Promise { + const socketPath = getSocketPath(); + await removeStaleSocketIfNeeded(socketPath); + + const server = createServer((socket) => { + let buffer = ""; + + socket.on("data", async (chunk: Buffer | string) => { + buffer += chunk.toString(); + const newlineIndex = buffer.indexOf("\n"); + if (newlineIndex === -1) { + return; + } + + const line = buffer.slice(0, newlineIndex).trim(); + buffer = buffer.slice(newlineIndex + 1); + if (!line) { + return; + } + + try { + const request = parseRequestLine(line); + const response = await handler(request); + socket.end(encodeMessage(response)); + } catch (error) { + const response: ErrorResponse = { + type: "error", + ok: false, + error: error instanceof Error ? error.message : String(error), + }; + socket.end(encodeMessage(response)); + } + }); + }); + + await new Promise((resolve, reject) => { + server.once("error", reject); + server.listen(socketPath, () => { + server.off("error", reject); + resolve(); + }); + }); + + return server; +} + +async function removeStaleSocketIfNeeded(socketPath: string): Promise { + if (!existsSync(socketPath)) { + return; + } + + const isLive = await isSocketLive(socketPath); + if (isLive) { + throw new Error(`orchestrator is already running: ${socketPath}`); + } + + unlinkSync(socketPath); +} + +async function isSocketLive(socketPath: string): Promise { + return new Promise((resolve, reject) => { + const socket = createConnection(socketPath); + let settled = false; + + const finish = (result: boolean) => { + if (settled) { + return; + } + settled = true; + socket.removeAllListeners(); + socket.destroy(); + resolve(result); + }; + + socket.on("connect", () => finish(true)); + socket.on("error", (error: NodeJS.ErrnoException) => { + if (error.code === "ECONNREFUSED" || error.code === "ENOENT") { + finish(false); + return; + } + if (error.code === "EPIPE" || error.code === "ECONNRESET") { + finish(false); + return; + } + if (settled) { + return; + } + settled = true; + socket.removeAllListeners(); + socket.destroy(); + reject(error); + }); + }); +}