#!/usr/bin/env node import { randomUUID } from "node:crypto"; import { once } from "node:events"; import { fileURLToPath } from "node:url"; import { jsonFingerprint } from "../resource-managers/json.js"; import { loadResourceManagerFile } from "../resource-managers/loader.js"; import { createResultHelpers } from "../resource-managers/results.js"; import type { ChildWorkflowRecord, ResourceManagerEffects, ResourceManagerFollowUpResult, ResourceManagerQueueFollowUpRequest, ResourceManagerRemoveFollowUpRequest, ResourceManagerSettingsChangeRequest, ResourceManagerSettingsChangeResult, ResourceManagerWorkflows, EffectApplication, EffectDefinition, EffectObservation, EffectRecord, ReconcileContext, ReconcileResult, } from "../resource-managers/types.js"; import { parseJson, type JsonValue } from "../state/json.js"; import { errorMessage } from "../workflows/errors.js"; import { RESOURCE_RUNNER_LAUNCH_SCHEMA, encodeResourceRunnerLine, parseResourceRunnerResponse, type ResourceRunnerLaunchEnvelope, type ResourceRunnerMessage, type ResourceRunnerOperation, type ResourceRunnerResponse, } from "./resource-runner-protocol.js"; const STARTUP_ENV = "PI_WORKFLOWS_RESOURCE_RUNNER_LAUNCH"; const MAX_FRAME_BYTES = 1024 * 1024; class ResourceRunnerTransport { private readonly pending = new Map< string, { resolve: (response: ResourceRunnerResponse) => void; reject: (error: Error) => void } >(); private buffered = Buffer.alloc(0); constructor(private readonly launch: ResourceRunnerLaunchEnvelope) { process.stdin.on("data", (chunk: Buffer) => this.onData(chunk)); process.stdin.on("error", (error) => this.failAll(error)); process.stdin.resume(); } async request(operation: ResourceRunnerOperation, payload: JsonValue): Promise { const message: ResourceRunnerMessage = { schema: "pi-workflows.controller-worker-message.v1", launchSchema: this.launch.schema, messageId: randomUUID(), runnerEpoch: this.launch.runnerEpoch, generation: this.launch.generation, operation, payload, }; const response = await this.send(message); if (response.outcome !== "accepted") { throw new Error(response.error ?? `Resource runner ${response.outcome}`); } return response.result as T; } close(): void { process.stdin.pause(); process.stdin.removeAllListeners("data"); process.stdin.removeAllListeners("error"); this.failAll(new Error("Resource runner transport closed")); } private async send(message: ResourceRunnerMessage): Promise { const response = new Promise((resolve, reject) => { this.pending.set(message.messageId, { resolve, reject }); }); if (!process.stdout.write(encodeResourceRunnerLine(message))) { await once(process.stdout, "drain"); } return await response; } private onData(chunk: Buffer): void { this.buffered = Buffer.concat([this.buffered, chunk]); if (this.buffered.byteLength > MAX_FRAME_BYTES && !this.buffered.includes(0x0a)) { this.failAll(new Error("Resource runner response exceeds 1 MiB")); return; } for (;;) { const newline = this.buffered.indexOf(0x0a); if (newline < 0) return; const frame = this.buffered.subarray(0, newline); this.buffered = this.buffered.subarray(newline + 1); if (frame.byteLength === 0) continue; try { const response = parseResourceRunnerResponse(frame); const pending = this.pending.get(response.messageId); if (pending === undefined) throw new Error("Resource runner response is unexpected"); this.pending.delete(response.messageId); pending.resolve(response); } catch (error) { this.failAll(error instanceof Error ? error : new Error(String(error))); } } } private failAll(error: Error): void { for (const pending of this.pending.values()) pending.reject(error); this.pending.clear(); } } class ServerBackedResourceManagerEffects implements ResourceManagerEffects { private used = false; constructor( private readonly transport: ResourceRunnerTransport, private readonly signal: AbortSignal, ) {} async ensure(definition: EffectDefinition): Promise { if (this.used) throw new Error("A reconciliation pass may ensure only one external effect"); this.used = true; const reservation = await this.transport.request<{ record: EffectRecord; created: boolean }>( "effect.reserve", { key: definition.key, kind: definition.kind, requestFingerprint: jsonFingerprint(definition.request), }, ); if (["applied", "rejected", "indeterminate"].includes(reservation.record.state)) { return reservation.record; } if (!reservation.created) { const observation = await definition.observe(this.signal); if (observation.state !== "not_applied") { return await this.settleObservation(definition.key, observation); } } let applied: EffectApplication; try { applied = await definition.apply(this.signal); } catch (error) { applied = { state: "indeterminate", error: errorMessage(error) }; } return await this.transport.request("effect.settle", { key: definition.key, state: applied.state, ...(applied.state === "applied" && applied.externalRef !== undefined ? { externalRef: applied.externalRef } : {}), ...(applied.state !== "applied" && applied.error !== undefined ? { error: applied.error } : {}), }); } private async settleObservation( key: string, observation: EffectObservation, ): Promise { return await this.transport.request("effect.settle", { key, state: observation.state === "applied" ? "applied" : "indeterminate", ...(observation.state === "applied" && observation.externalRef !== undefined ? { externalRef: observation.externalRef } : {}), }); } } class ServerBackedResourceManagerWorkflows implements ResourceManagerWorkflows { constructor(private readonly transport: ResourceRunnerTransport) {} async ensure(request: { requestKey: string; workflow: string; input: unknown; }): Promise { return await this.transport.request( "workflow.ensure", request as JsonValue, ); } async changeSettings( request: ResourceManagerSettingsChangeRequest, ): Promise { return await this.transport.request( "workflow.changeSettings", request as unknown as JsonValue, ); } async queueFollowUp( request: ResourceManagerQueueFollowUpRequest, ): Promise { return await this.transport.request( "workflow.queueFollowUp", request as unknown as JsonValue, ); } async removeFollowUp( request: ResourceManagerRemoveFollowUpRequest, ): Promise { return await this.transport.request( "workflow.removeFollowUp", request as unknown as JsonValue, ); } } export async function runResourceRunner(): Promise { const launch = readLaunchEnvelope(); const transport = new ResourceRunnerTransport(launch); try { const definition = await loadResourceManagerFile(launch.definitionPath); if (definition.name !== launch.resourceManagerName) { throw new Error( `ResourceManager source name changed: expected ${launch.resourceManagerName}, got ${definition.name}`, ); } await transport.request("runner.ready", {}); const abort = new AbortController(); const timeoutMs = definition.timeoutMs ?? launch.timeoutMs; const timer = setTimeout( () => abort.abort(new Error(`Reconciliation timed out after ${timeoutMs}ms`)), timeoutMs, ); try { const context: ReconcileContext = { signal: abort.signal, effects: new ServerBackedResourceManagerEffects(transport, abort.signal), workflows: new ServerBackedResourceManagerWorkflows(transport), ...createResultHelpers(), }; const result = await raceWithAbort( Promise.resolve(definition.reconcile(context, launch.resource)), abort.signal, ); validateResult(result); await transport.request("runner.finished", result as unknown as JsonValue); } catch (error) { await transport.request("runner.failed", { error: errorMessage(error) }); } finally { clearTimeout(timer); } return 0; } finally { transport.close(); } } function readLaunchEnvelope(): ResourceRunnerLaunchEnvelope { const encoded = process.env[STARTUP_ENV]; if (encoded === undefined) throw new Error("Resource runner launch envelope is missing"); const value = parseJson(Buffer.from(encoded, "base64url").toString("utf8")); if ( typeof value !== "object" || value === null || Array.isArray(value) || (value as { schema?: unknown }).schema !== RESOURCE_RUNNER_LAUNCH_SCHEMA ) { throw new Error("Resource runner launch envelope is invalid"); } return value as unknown as ResourceRunnerLaunchEnvelope; } function validateResult(value: unknown): asserts value is ReconcileResult { if (typeof value !== "object" || value === null || Array.isArray(value)) { throw new Error("ResourceManager reconcile result must be an object"); } const result = value as { kind?: unknown; afterMs?: unknown }; if (result.kind !== "settled" && result.kind !== "requeue") { throw new Error("ResourceManager reconcile result kind is invalid"); } if ( result.afterMs !== undefined && (!Number.isSafeInteger(result.afterMs) || (result.afterMs as number) < 0) ) { throw new Error("ResourceManager reconcile afterMs must be a non-negative safe integer"); } } function raceWithAbort(operation: Promise, signal: AbortSignal): Promise { return new Promise((resolve, reject) => { const onAbort = () => reject(signal.reason ?? new Error("Reconciliation aborted")); signal.addEventListener("abort", onAbort, { once: true }); operation.then(resolve, reject).finally(() => signal.removeEventListener("abort", onAbort)); if (signal.aborted) onAbort(); }); } async function main(): Promise { try { process.exitCode = await runResourceRunner(); } catch (error) { process.stderr.write(`${errorMessage(error)}\n`); process.exitCode = 1; } } if (process.argv[1] === fileURLToPath(import.meta.url)) void main();