import { spawn } from "node:child_process"; import { createHash, randomUUID } from "node:crypto"; import { once } from "node:events"; import fs from "node:fs"; import { createRequire } from "node:module"; import net, { type Socket } from "node:net"; import { fileURLToPath } from "node:url"; import { workflowStatePath } from "../state/database.js"; import { canonicalJson, parseJson, type JsonValue } from "../state/json.js"; import { CLIENT_PROTOCOL_SCHEMA, NdjsonFrameDecoder, assertSocketPathSupported, clientSocketPath, encodeProtocolLine, parseClientMessage, type ClientEvent, type ClientHello, type ClientOperation, type ClientRequest, type ClientResponse, } from "./protocol.js"; import type { ResolvedResourceManagerInitialization, ResolvedSettingsChange, ResolvedWorkflowLaunch, } from "./resolver.js"; import type { WorkflowRunListPage, WorkflowRunSummary, WorkflowRunView } from "./view.js"; const CONNECT_TIMEOUT_MS = 2_000; const START_TIMEOUT_MS = 10_000; const STARTUP_STDERR_LIMIT_BYTES = 64 * 1024; /** Capped exponential reconnect delay with jitter; the first retry is fastest. */ const RECONNECT_BASE_DELAY_MS = 250; const RECONNECT_MAX_DELAY_MS = 10_000; /** After this many failed attempts the client reports a blocker and stops looping. */ const RECONNECT_MAX_ATTEMPTS = 12; const RESOLVER_TIMEOUT_MS = 30_000; const CLIENT_PACKAGE_VERSION = runtimePackageVersion(); type PendingRequest = { resolve: (response: ClientResponse) => void; reject: (error: Error) => void; }; type Subscription = { operation: "view.runs.watch" | "view.run.watch" | "view.session.watch"; runId?: string; payload: JsonValue; listener: (event: ClientEvent) => void; runListGeneration: number; }; /** Bounded reason code and safe message for one failed subscription. */ export type SubscriptionFailure = { reasonCode: "connection_lost" | "reconnect_exhausted" | "projection_failed"; message: string; }; export class WorkflowClientVersionError extends Error { constructor(message: string) { super(message); this.name = "WorkflowClientVersionError"; } } export type WorkflowClientOptions = { clientId?: string; databasePath?: string; serverEntryPath?: string; env?: Record; }; export class WorkflowClient { readonly clientId: string; readonly databasePath: string; readonly endpoint: string; private readonly serverEntryPath: string | undefined; private readonly env: Record | undefined; private socket: Socket | null = null; private connectTask: Promise | null = null; private hello: ClientHello | null = null; private closed = false; private reconnectTimer: ReturnType | null = null; private reconnectAttempts = 0; /** Set only after a connection holds every subscription it asked to restore. */ private subscriptionsRestored = false; private readonly pending = new Map(); private readonly subscriptions = new Map(); constructor(options: WorkflowClientOptions = {}) { this.clientId = options.clientId ?? `client-${randomUUID()}`; this.databasePath = options.databasePath ?? workflowStatePath(); this.endpoint = clientSocketPath(this.databasePath); this.serverEntryPath = options.serverEntryPath; this.env = options.env; } get connectionId(): string | undefined { return this.hello?.connectionId; } get packageVersion(): string | undefined { return this.hello?.packageVersion; } async connect(): Promise { if (this.closed) throw new Error("Workflow client is closed"); if (this.hello !== null && this.socket !== null && !this.socket.destroyed) return this.hello; this.connectTask ??= this.openConnection(); try { return await this.connectTask; } finally { this.connectTask = null; } } async request(options: { operation: ClientOperation; requestId?: string; idempotencyKey?: string; runId?: string; expectedRevision?: number; payload?: JsonValue; signal?: AbortSignal; }): Promise { await this.connect(); return await this.requestConnected(options); } async requestDurable(options: { operation: ClientOperation; requestId?: string; idempotencyKey: string; runId?: string; expectedRevision?: number; payload?: JsonValue; signal?: AbortSignal; }): Promise { try { return await this.request(options); } catch (error) { if ( this.closed || options.signal?.aborted === true || error instanceof WorkflowClientVersionError ) { throw error; } this.resetConnection(); await this.ensureAvailable(); return await this.request({ ...options, requestId: randomUUID() }); } } async getRun(runId: string): Promise { const response = await this.request({ operation: "view.run.get", runId }); if (response.outcome === "notFound") return null; if (response.outcome !== "accepted" || !isWorkflowRunView(response.receipt, runId)) { throw new Error(response.error ?? "Workflow server returned an invalid run view"); } return response.receipt; } async watchRuns( listener: (event: ClientEvent) => void, options: { subscriptionId?: string; limit?: number } = {}, ): Promise<() => Promise> { return await this.subscribe( "view.runs.watch", undefined, { subscriptionId: options.subscriptionId ?? randomUUID(), ...(options.limit === undefined ? {} : { limit: options.limit }), }, listener, ); } async watchRun( runId: string, listener: (event: ClientEvent) => void, options: { subscriptionId?: string; revision?: number } = {}, ): Promise<() => Promise> { return await this.subscribe( "view.run.watch", runId, { subscriptionId: options.subscriptionId ?? randomUUID(), ...(options.revision === undefined ? {} : { revision: options.revision }), }, listener, ); } async watchSession( sessionId: string, listener: (event: ClientEvent) => void, options: { subscriptionId?: string; coordinator?: boolean; nodeCursor?: number } = {}, ): Promise<() => Promise> { return await this.subscribe( "view.session.watch", undefined, { subscriptionId: options.subscriptionId ?? randomUUID(), sessionId, ...(options.coordinator === undefined ? {} : { coordinator: options.coordinator }), ...(options.nodeCursor === undefined ? {} : { nodeCursor: options.nodeCursor }), }, listener, ); } /** * Move the node window of the active session subscription. `null` returns to * the window that follows the node the widget shows as working. The stored * subscription keeps the window, so a reconnect restores it. */ async setSessionNodeWindow(sessionId: string, nodeCursor: number | null): Promise { const entry = [...this.subscriptions.entries()].find( ([, subscription]) => subscription.operation === "view.session.watch" && (subscription.payload as { sessionId?: unknown }).sessionId === sessionId, ); if (entry === undefined) return false; const [subscriptionId, subscription] = entry; const payload = { ...(subscription.payload as Record) }; if (nodeCursor === null) delete payload.nodeCursor; else payload.nodeCursor = nodeCursor; const response = await this.request({ operation: "view.session.window", payload: { subscriptionId, nodeCursor }, }); if (response.outcome !== "accepted" && response.outcome !== "adopted") { throw new Error(response.error ?? `Workflow session window was ${response.outcome}`); } subscription.payload = payload; return true; } async ensureAvailable(): Promise { try { return await this.connect(); } catch (error) { if (error instanceof WorkflowClientVersionError) throw error; // A path the operating system cannot bind is a blocker. Do not spawn a // child that can never serve it. assertSocketPathSupported(this.endpoint); this.resetConnection(); const startupFailure = this.startDetached(); let startupError: Error | undefined; void startupFailure.then((error) => { startupError = error; }); const deadline = Date.now() + START_TIMEOUT_MS; let lastError: unknown; while (Date.now() < deadline) { await delay(50); try { return await this.connect(); } catch (error) { lastError = error; this.resetConnection(); } } throw new Error( `Workflow server did not become ready: ${startupError?.message ?? (lastError instanceof Error ? lastError.message : String(lastError))}`, ); } } async ensureRunning(): Promise { await this.ensureAvailable(); return await this.request({ operation: "server.status" }); } async readContent( runId: string, contentPath: string, ): Promise<{ mediaType: string; content: Buffer; }> { await this.ensureAvailable(); const chunks: Buffer[] = []; let offset = 0; let expectedBytes: number | undefined; let expectedSha256: string | undefined; let mediaType: string | undefined; for (;;) { const response = await this.request({ operation: "view.content", runId, payload: { path: contentPath, offset }, }); if (response.outcome !== "accepted" || !isRecord(response.receipt)) { throw new Error(response.error ?? `Workflow content is unavailable: ${contentPath}`); } const receipt = response.receipt; if ( receipt.path !== contentPath || receipt.offset !== offset || typeof receipt.data !== "string" || typeof receipt.mediaType !== "string" || typeof receipt.sha256 !== "string" || !Number.isSafeInteger(receipt.bytes) || (receipt.bytes as number) < 0 || !Number.isSafeInteger(receipt.nextOffset) || (receipt.nextOffset as number) < offset || typeof receipt.complete !== "boolean" ) { throw new Error("Workflow content receipt is invalid"); } expectedBytes ??= receipt.bytes as number; expectedSha256 ??= receipt.sha256; mediaType ??= receipt.mediaType; if ( expectedBytes !== receipt.bytes || expectedSha256 !== receipt.sha256 || mediaType !== receipt.mediaType ) { throw new Error("Workflow content identity changed during transfer"); } const chunk = Buffer.from(receipt.data, "base64"); if (offset + chunk.byteLength !== receipt.nextOffset) { throw new Error("Workflow content chunk offset is invalid"); } chunks.push(chunk); offset = receipt.nextOffset as number; if (receipt.complete) break; } const content = Buffer.concat(chunks); if (content.byteLength !== expectedBytes) throw new Error("Workflow content is incomplete"); const digest = createHash("sha256").update(content).digest("hex"); if (digest !== expectedSha256) throw new Error("Workflow content digest does not match"); return { mediaType: mediaType as string, content }; } async hydrateContent(runId: string, value: JsonValue): Promise { return await this.hydrateContentValue(runId, value, new Map()); } private async hydrateContentValue( runId: string, value: JsonValue, reads: Map>, ): Promise { if (isEscapedContent(value)) { const escaped = value.$escaped; if (isRecord(escaped)) { const entries = await Promise.all( Object.entries(escaped).map( async ([key, item]) => [key, await this.hydrateContentValue(runId, item as JsonValue, reads)] as const, ), ); return Object.fromEntries(entries) as JsonValue; } return escaped; } if (isContentReference(value)) { const contentPath = value.$artifact.path; let read = reads.get(contentPath); if (read === undefined) { read = this.readContent(runId, contentPath); reads.set(contentPath, read); } const loaded = await read; const digest = createHash("sha256").update(loaded.content).digest("hex"); if ( loaded.mediaType !== value.$artifact.mediaType || loaded.content.byteLength !== value.$artifact.bytes || digest !== value.$artifact.sha256 ) { throw new Error("Workflow content reference does not match its content"); } const decoded = loaded.mediaType === "application/json" ? parseJson(loaded.content.toString("utf8")) : loaded.content.toString("utf8"); return value.$artifact.opaque === true ? decoded : await this.hydrateContentValue(runId, decoded, reads); } if (Array.isArray(value)) { return await Promise.all( value.map(async (item) => await this.hydrateContentValue(runId, item, reads)), ); } if (isRecord(value)) { const entries = await Promise.all( Object.entries(value).map( async ([key, item]) => [key, await this.hydrateContentValue(runId, item as JsonValue, reads)] as const, ), ); return Object.fromEntries(entries) as JsonValue; } return value; } async resolveWorkflow(options: { cwd: string; workflowRef: string; timeoutMs?: number; }): Promise { return (await this.runResolver( { schema: "pi-workflows.resolve-request.v1", cwd: options.cwd, workflowRef: options.workflowRef, }, options.cwd, "pi-workflows.resolved-launch.v1", options.timeoutMs, )) as unknown as ResolvedWorkflowLaunch; } async resolveResourceManagerInitialization(options: { cwd: string; resourceManagerName: string; spec: JsonValue; timeoutMs?: number; }): Promise { return (await this.runResolver( { schema: "pi-workflows.resource-manager-initialization-request.v1", cwd: options.cwd, resourceManagerName: options.resourceManagerName, spec: options.spec, }, options.cwd, "pi-workflows.resolved-resource-manager-initialization.v1", options.timeoutMs, )) as unknown as ResolvedResourceManagerInitialization; } async resolveSettingsChange(options: { cwd: string; workflowRef: string; definitionDigest: string; mountPath: string; current: JsonValue; patch: JsonValue; actorId: string; timeoutMs?: number; }): Promise { return (await this.runResolver( { schema: "pi-workflows.settings-validation-request.v1", cwd: options.cwd, workflowRef: options.workflowRef, definitionDigest: options.definitionDigest, mountPath: options.mountPath, current: options.current, patch: options.patch, actorId: options.actorId, }, options.cwd, "pi-workflows.resolved-settings-change.v1", options.timeoutMs, )) as unknown as ResolvedSettingsChange; } async close(): Promise { if (this.closed) return; this.closed = true; if (this.reconnectTimer !== null) clearTimeout(this.reconnectTimer); this.reconnectTimer = null; this.subscriptions.clear(); const socket = this.socket; this.resetConnection(new Error("Workflow client closed")); if (socket !== null && !socket.destroyed) { socket.end(); await Promise.race([once(socket, "close").then(() => undefined), delay(250)]); socket.destroy(); } } private async subscribe( operation: Subscription["operation"], runId: string | undefined, payload: JsonValue, listener: (event: ClientEvent) => void, ): Promise<() => Promise> { if (!isRecord(payload) || typeof payload.subscriptionId !== "string") { throw new Error("Workflow subscription requires a subscriptionId"); } const subscriptionId = payload.subscriptionId; await this.connect(); this.subscriptions.set(subscriptionId, { operation, ...(runId === undefined ? {} : { runId }), payload, listener, runListGeneration: 0, }); try { const response = await this.requestConnected({ operation, ...(runId === undefined ? {} : { runId }), payload, }); if (response.outcome !== "accepted" && response.outcome !== "adopted") { throw new Error(response.error ?? `Workflow subscription was ${response.outcome}`); } } catch (error) { this.subscriptions.delete(subscriptionId); throw error; } return async () => { if (!this.subscriptions.delete(subscriptionId)) return; if (this.socket !== null && !this.socket.destroyed) { await this.request({ operation: "view.run.unwatch", ...(runId === undefined ? {} : { runId }), payload: { subscriptionId }, }); } }; } private async requestConnected(options: { operation: ClientOperation; requestId?: string; idempotencyKey?: string; runId?: string; expectedRevision?: number; payload?: JsonValue; signal?: AbortSignal; }): Promise { const request: ClientRequest = { schema: CLIENT_PROTOCOL_SCHEMA, type: "request", requestId: options.requestId ?? randomUUID(), clientId: this.clientId, operation: options.operation, idempotencyKey: options.idempotencyKey ?? randomUUID(), ...(options.runId === undefined ? {} : { runId: options.runId }), ...(options.expectedRevision === undefined ? {} : { expectedRevision: options.expectedRevision }), payload: options.payload ?? {}, }; return await this.send(request, options.signal); } private async send(request: ClientRequest, signal?: AbortSignal): Promise { const socket = this.socket; if (socket === null || socket.destroyed) throw new Error("Workflow server is unavailable"); if (signal?.aborted === true) throw abortReason(signal); if (this.pending.has(request.requestId)) { throw new Error(`Workflow request is already pending: ${request.requestId}`); } const response = new Promise((resolve, reject) => { const removeAbort = (): void => signal?.removeEventListener("abort", onAbort); const onAbort = (): void => { if (!this.pending.delete(request.requestId)) return; removeAbort(); reject(abortReason(signal as AbortSignal)); }; this.pending.set(request.requestId, { resolve: (value) => { removeAbort(); resolve(value); }, reject: (error) => { removeAbort(); reject(error); }, }); signal?.addEventListener("abort", onAbort, { once: true }); }); try { if (!socket.write(encodeProtocolLine(request))) await waitForSocketDrain(socket, signal); } catch (error) { const pending = this.pending.get(request.requestId); this.pending.delete(request.requestId); pending?.reject(toError(error)); } return await response; } private async openConnection(): Promise { assertSocketPathSupported(this.endpoint); const socket = net.createConnection(this.endpoint); const decoder = new NdjsonFrameDecoder(); this.socket = socket; let helloResolve!: (hello: ClientHello) => void; let helloReject!: (error: Error) => void; const helloPromise = new Promise((resolve, reject) => { helloResolve = resolve; helloReject = reject; }); let receivedHello = false; socket.on("data", (chunk: Buffer) => { try { for (const frame of decoder.push(chunk)) { const message = parseClientMessage(frame); if (!receivedHello) { if (message.type !== "hello") throw new Error("Workflow server did not send hello first"); if (message.packageVersion !== CLIENT_PACKAGE_VERSION) { throw new WorkflowClientVersionError( `Workflow client version mismatch: server ${message.packageVersion}, client ${CLIENT_PACKAGE_VERSION}. The running workflow server process is from another version. Run "pi-workflows server stop" and retry.`, ); } receivedHello = true; this.hello = message; helloResolve(message); continue; } if (message.type === "response") { const pending = this.pending.get(message.requestId); if (pending === undefined) continue; this.pending.delete(message.requestId); pending.resolve(message); } else if (message.type === "event") { const subscription = this.subscriptions.get(message.subscriptionId); if ( subscription !== undefined && subscription.operation === "view.run.watch" && isRecord(subscription.payload) && isRecord(message.payload) && Number.isSafeInteger(message.payload.revision) ) { subscription.payload = { ...subscription.payload, revision: message.payload.revision as number, }; } if (subscription !== undefined) { void this.deliverSubscriptionEvent(subscription, message).catch(() => { // A stale paged list is replaced by the next subscription snapshot. }); } } else { throw new Error(`Unexpected workflow client message: ${message.type}`); } } } catch (error) { helloReject(toError(error)); socket.destroy(); } }); socket.once("error", (error) => { helloReject(error); }); socket.once("close", () => { if (!receivedHello) helloReject(new Error("Workflow server closed before hello")); if (this.socket === socket) { this.resetConnection(new Error("Workflow server connection closed")); this.scheduleReconnect(); } }); let connectTimer: ReturnType | undefined; try { await Promise.race([ once(socket, "connect"), helloPromise.then(() => undefined), new Promise((_, reject) => { connectTimer = setTimeout( () => reject(new Error("Workflow server connection timed out")), CONNECT_TIMEOUT_MS, ); connectTimer.unref?.(); }), ]); const hello = await Promise.race([ helloPromise, delay(CONNECT_TIMEOUT_MS).then(() => { throw new Error("Workflow server hello timed out"); }), ]); // A successful hello starts a fresh reconnect budget only together with a // restored subscription set. A subscription the server keeps refusing must // not reset the budget, or the client would reconnect forever. await this.restoreSubscriptions(); this.subscriptionsRestored = true; this.reconnectAttempts = 0; return hello; } catch (error) { socket.destroy(); throw error; } finally { if (connectTimer !== undefined) clearTimeout(connectTimer); } } private async deliverSubscriptionEvent( subscription: Subscription, event: ClientEvent, ): Promise { // One accepted view proves the connection works again, so the reconnect // budget resets here as well as after hello. A view from a subscription the // server restored before it refused a later one is not that proof, because // the client still owes work on this connection. if (event.event !== "unavailable" && this.subscriptionsRestored) this.reconnectAttempts = 0; if (subscription.operation !== "view.runs.watch" || event.event !== "runs") { subscription.listener(event); return; } const first = parseWorkflowRunListPage(event.payload); const subscriptionId = requireSubscriptionId(subscription.payload); const generation = subscription.runListGeneration + 1; subscription.runListGeneration = generation; if (first.start !== 0) throw new Error("Workflow run list snapshot must start at zero"); const items: WorkflowRunSummary[] = [...first.items]; let cursor = first.start + first.items.length; while (cursor < first.total) { if ( subscription.runListGeneration !== generation || this.subscriptions.get(subscriptionId) !== subscription ) return; if (items.length === 0) throw new Error("Workflow run list page made no progress"); const response = await this.requestConnected({ operation: "view.runs.page", payload: { cursor, revision: first.revision, ...(isRecord(subscription.payload) && typeof subscription.payload.limit === "number" ? { limit: subscription.payload.limit } : {}), }, }); if (response.outcome === "conflict") return; if (response.outcome !== "accepted") { throw new Error(response.error ?? `Workflow run list page was ${response.outcome}`); } const page = parseWorkflowRunListPage(response.receipt); if ( page.revision !== first.revision || page.total !== first.total || page.start !== cursor || page.items.length === 0 ) { throw new Error("Workflow run list page does not continue the snapshot"); } items.push(...page.items); cursor += page.items.length; } if ( subscription.runListGeneration !== generation || this.subscriptions.get(subscriptionId) !== subscription ) return; subscription.listener({ ...event, payload: items as unknown as JsonValue }); } private async restoreSubscriptions(): Promise { for (const subscription of this.subscriptions.values()) { const response = await this.requestConnected({ operation: subscription.operation, ...(subscription.runId === undefined ? {} : { runId: subscription.runId }), payload: subscription.payload, }); if (response.outcome !== "accepted" && response.outcome !== "adopted") { throw new Error(response.error ?? `Workflow subscription was ${response.outcome}`); } } } private resetConnection(reason = new Error("Workflow server is unavailable")): void { const socket = this.socket; this.socket = null; this.hello = null; this.connectTask = null; this.subscriptionsRestored = false; if (socket !== null && !socket.destroyed) socket.destroy(); for (const pending of this.pending.values()) pending.reject(reason); this.pending.clear(); if (!this.closed) { this.reportSubscriptionFailure( { reasonCode: "connection_lost", message: "Workflow server connection is unavailable." }, false, ); } } /** * Report one bounded failure to every subscriber. The extension keeps its last * view for display and loses command authority until a fresh snapshot arrives. */ private reportSubscriptionFailure(failure: SubscriptionFailure, clear: boolean): void { for (const [subscriptionId, subscription] of this.subscriptions) { try { subscription.listener({ schema: CLIENT_PROTOCOL_SCHEMA, type: "event", subscriptionId, event: "unavailable", payload: { schema: "pi-workflows.subscription-failure.v1", ...failure }, }); } catch { // One renderer cannot block reconnection for other subscriptions. } if (clear) this.subscriptions.delete(subscriptionId); } } private scheduleReconnect(): void { if (this.closed || this.subscriptions.size === 0 || this.reconnectTimer !== null) return; if (this.reconnectAttempts >= RECONNECT_MAX_ATTEMPTS) { // A fast loop would hide the blocker. Report it once and wait for new work. this.reportSubscriptionFailure( { reasonCode: "reconnect_exhausted", // The budget counts a stopped server and a refused restore, so the // blocker must not claim that the server stayed silent. message: `Workflow server did not return this session's view after ${RECONNECT_MAX_ATTEMPTS} reconnect attempts. Start or restart the workflow server, then open or resume a session.`, }, false, ); return; } const capped = Math.min( RECONNECT_MAX_DELAY_MS, RECONNECT_BASE_DELAY_MS * 2 ** this.reconnectAttempts, ); // Full jitter over the upper half of the window keeps one reconnect storm // from repeating on the same instant across sessions. const delayMs = Math.round(capped / 2 + Math.random() * (capped / 2)); this.reconnectAttempts += 1; this.reconnectTimer = setTimeout(() => { this.reconnectTimer = null; void this.connect().catch(() => this.scheduleReconnect()); }, delayMs); this.reconnectTimer.unref?.(); } private async runResolver( request: JsonValue, cwd: string, expectedSchema: string, timeoutMs = RESOLVER_TIMEOUT_MS, ): Promise { const builtEntry = fileURLToPath(new URL("../server/resolver-entry.js", import.meta.url)); const sourceEntry = fileURLToPath(new URL("../server/resolver-entry.ts", import.meta.url)); const args = fs.existsSync(builtEntry) ? [builtEntry] : ["--import", createRequire(import.meta.url).resolve("tsx"), sourceEntry]; const child = spawn(process.execPath, args, { cwd, detached: process.platform !== "win32", stdio: ["pipe", "pipe", "pipe"], env: { ...process.env, ...this.env }, }); let stdout: Buffer = Buffer.alloc(0); let stderr: Buffer = Buffer.alloc(0); let outputError: Error | undefined; const append = (current: Buffer, chunk: Buffer): Buffer => { const next = Buffer.concat([current, chunk]); if (next.byteLength > 1024 * 1024) { outputError = new Error("Workflow resolver output exceeds 1 MiB"); stopProcessGroup(child.pid); return current; } return next; }; child.stdout.on("data", (chunk: Buffer) => { stdout = append(stdout, chunk); }); child.stderr.on("data", (chunk: Buffer) => { stderr = append(stderr, chunk); }); child.stdin.end(canonicalJson(request)); const timeout = setTimeout(() => stopProcessGroup(child.pid), timeoutMs); timeout.unref?.(); const [code, signal] = (await once(child, "exit")) as [number | null, NodeJS.Signals | null]; clearTimeout(timeout); if (outputError !== undefined) throw outputError; if (code !== 0) { const detail = stderr.toString("utf8").trim().slice(0, 2_000); throw new Error( detail || `Workflow resolver exited before completion (code ${code}, signal ${signal})`, ); } const value = parseJson(stdout.toString("utf8").trimEnd()); if (!isRecord(value) || value.schema !== expectedSchema) { throw new Error("Workflow resolver returned an invalid result envelope"); } return value as JsonValue; } private startDetached(): Promise { const builtEntry = fileURLToPath(new URL("../server/server-entry.js", import.meta.url)); const sourceEntry = fileURLToPath(new URL("../server/server-entry.ts", import.meta.url)); const entry = this.serverEntryPath ?? builtEntry; const args = this.serverEntryPath === undefined && !fs.existsSync(builtEntry) ? ["--import", createRequire(import.meta.url).resolve("tsx"), sourceEntry] : [entry]; const child = spawn(process.execPath, [...args, "--database", this.databasePath], { detached: true, stdio: ["ignore", "ignore", "ignore", "pipe"], env: { ...process.env, ...this.env, PI_WORKFLOWS_STARTUP_FD: "3" }, }); const startupDiagnostic = child.stdio[3]; if (startupDiagnostic === null || startupDiagnostic === undefined) { child.unref(); return Promise.resolve(new Error("Workflow server startup diagnostic pipe is unavailable")); } const stderrChunks: Buffer[] = []; let stderrBytes = 0; let stderrTruncated = false; startupDiagnostic.on("data", (chunk: Buffer) => { const remaining = STARTUP_STDERR_LIMIT_BYTES - stderrBytes; if (remaining <= 0) { stderrTruncated = true; return; } const kept = chunk.subarray(0, remaining); stderrChunks.push(kept); stderrBytes += kept.byteLength; if (kept.byteLength < chunk.byteLength) stderrTruncated = true; }); const failure = new Promise((resolve) => { child.once("error", (error) => resolve(error)); startupDiagnostic.once("end", () => { const detail = startupErrorDetail(stderrChunks, stderrTruncated); if (detail.length > 0) resolve(new Error(detail)); }); child.once("close", (code, signal) => { const detail = startupErrorDetail(stderrChunks, stderrTruncated); resolve( new Error( detail.length > 0 ? detail : `Workflow server process exited before becoming ready (code ${code}, signal ${signal})`, ), ); }); }); child.unref(); return failure; } } function startupErrorDetail(chunks: Buffer[], truncated: boolean): string { const detail = Buffer.concat(chunks).toString("utf8").trim(); return detail.length > 0 && truncated ? `${detail}\n[workflow server startup diagnostic truncated]` : detail; } function requireSubscriptionId(value: JsonValue): string { if (!isRecord(value) || typeof value.subscriptionId !== "string") { throw new Error("Workflow subscription requires a subscriptionId"); } return value.subscriptionId; } function parseWorkflowRunListPage(value: unknown): WorkflowRunListPage { if ( !isRecord(value) || value.schema !== "pi-workflows.run-list-page.v1" || typeof value.revision !== "string" || !Number.isSafeInteger(value.start) || !Number.isSafeInteger(value.total) || (value.start as number) < 0 || (value.total as number) < 0 || !Array.isArray(value.items) || !value.items.every( (item) => isRecord(item) && typeof item.runId === "string" && typeof item.workflowName === "string", ) ) { throw new Error("Workflow server returned an invalid run list page"); } return value as unknown as WorkflowRunListPage; } function isWorkflowRunView(value: unknown, runId: string): value is WorkflowRunView { return ( isRecord(value) && value.schema === "pi-workflows.run-view.v1" && value.runId === runId && Number.isSafeInteger(value.revision) && Number.isSafeInteger(value.runRevision) && (value.runRevision as number) >= 0 ); } function abortReason(signal: AbortSignal): Error { return signal.reason instanceof Error ? signal.reason : new Error("Workflow request was cancelled"); } function waitForSocketDrain(socket: Socket, signal?: AbortSignal): Promise { if (signal?.aborted === true) return Promise.reject(abortReason(signal)); if (socket.destroyed) return Promise.reject(new Error("Workflow server connection closed")); return new Promise((resolve, reject) => { const cleanup = (): void => { socket.off("drain", onDrain); socket.off("close", onClose); socket.off("error", onError); signal?.removeEventListener("abort", onAbort); }; const settle = (error?: Error): void => { cleanup(); if (error === undefined) resolve(); else reject(error); }; const onDrain = (): void => settle(); const onClose = (): void => settle(new Error("Workflow server connection closed")); const onError = (error: Error): void => settle(error); const onAbort = (): void => settle(abortReason(signal as AbortSignal)); socket.once("drain", onDrain); socket.once("close", onClose); socket.once("error", onError); signal?.addEventListener("abort", onAbort, { once: true }); if (signal?.aborted === true) onAbort(); else if (socket.destroyed) onClose(); }); } function runtimePackageVersion(): string { const parsed = JSON.parse( fs.readFileSync(new URL("../../package.json", import.meta.url), "utf8"), ) as { version?: unknown }; if (typeof parsed.version !== "string" || parsed.version.length === 0) { throw new Error("Pi Workflows package version is missing"); } return parsed.version; } function stopProcessGroup(pid: number | undefined): void { if (pid === undefined) return; try { if (process.platform !== "win32") process.kill(-pid, "SIGKILL"); else process.kill(pid, "SIGKILL"); } catch { // The resolver has already exited. } } function delay(ms: number): Promise { return new Promise((resolve) => setTimeout(resolve, ms)); } type ContentReference = { $artifact: { path: string; mediaType: string; bytes: number; sha256: string; opaque?: boolean; }; }; function isEscapedContent(value: JsonValue): value is JsonValue & { $escaped: JsonValue } { return isRecord(value) && Object.keys(value).length === 1 && Object.hasOwn(value, "$escaped"); } function isContentReference(value: JsonValue): value is ContentReference { if (!isRecord(value) || Object.keys(value).length !== 1 || !isRecord(value.$artifact)) { return false; } const artifact = value.$artifact; return ( typeof artifact.path === "string" && typeof artifact.mediaType === "string" && Number.isSafeInteger(artifact.bytes) && (artifact.bytes as number) >= 0 && typeof artifact.sha256 === "string" && (artifact.opaque === undefined || typeof artifact.opaque === "boolean") ); } function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null && !Array.isArray(value); } function toError(error: unknown): Error { return error instanceof Error ? error : new Error(String(error)); }