import { clearCompletionSidecar, validateOperation, type CompletionOperation, } from "./completion-sidecar.ts"; /** Parent-facing data prepared once for one Completion operation. */ export type CompletionDeliveryPayload> = Readonly; /** Destination signal for intentional lifecycle suppression. */ export class CompletionDeliverySuppressedError extends Error { constructor(message = "Completion delivery suppressed by lifecycle cancellation.") { super(message); this.name = "CompletionDeliverySuppressedError"; } } export type CompletionDeliveryDestination = | ((payload: CompletionDeliveryPayload) => Value | PromiseLike) | { accept(payload: CompletionDeliveryPayload): Value | PromiseLike; }; export interface CompletionDeliveryManagerOptions { /** Kept as an extension seam for callers that construct the manager. */ acceptingDeliveries?: boolean; } export interface CompletionDeliveryRequestOptions { signal?: AbortSignal; } export interface CompletionDeliveryRequest extends CompletionDeliveryRequestOptions { operation: CompletionOperation; payload: CompletionDeliveryPayload; destination: CompletionDeliveryDestination; } export type CompletionDeliveryResult = | { status: "delivered"; value: Value } | { status: "suppressed" } | { status: "cancelled" } | { status: "failed"; error: string }; interface DeliveryState { operation: CompletionOperation; payload?: CompletionDeliveryPayload; destination?: CompletionDeliveryDestination; status: "pending" | "delivering" | "delivered" | "suppressed" | "cancelled" | "failed"; value?: Value; error?: string; settled: boolean; resolve: (result: CompletionDeliveryResult) => void; promise: Promise>; terminalResult?: CompletionDeliveryResult; removeAbortListener?: () => void; } function validatePayload(payload: unknown): asserts payload is Record { if (payload == null || typeof payload !== "object" || Array.isArray(payload)) { throw new Error("Completion delivery requires a valid delivery payload."); } } /** Freeze the complete parent-facing payload before its one delivery attempt. */ function freezePayload(payload: Payload): CompletionDeliveryPayload { const seen = new WeakSet(); const freeze = (value: unknown): void => { if ( value === null || (typeof value !== "object" && typeof value !== "function") || seen.has(value) ) return; seen.add(value); for (const key of Reflect.ownKeys(value)) { const descriptor = Object.getOwnPropertyDescriptor(value, key); if (!descriptor) continue; if ("get" in descriptor || "set" in descriptor) { throw new Error("Completion delivery payload cannot contain accessor properties."); } freeze(descriptor.value); } Object.freeze(value); }; freeze(payload); return payload as CompletionDeliveryPayload; } function operationKey(operation: CompletionOperation): string { // A Completion operation ID is allocated once per Fresh attempt and is the // sole logical identity for delivery. Artifact paths locate protocol data // but must not create a second delivery claim. return operation.operationId; } /** * Parent-wide owner of one-shot Completion delivery. * * The manager receives stable parent-facing data and never reads Child * sessions or supervises Native agents. Each operation gets one destination * attempt during the owning parent runtime; failures are returned explicitly * and are not retried. */ export class CompletionDeliveryManager { private readonly states = new Map>(); private acceptingDeliveries: boolean; constructor(options: CompletionDeliveryManagerOptions = {}) { this.acceptingDeliveries = options.acceptingDeliveries ?? true; } deliver( request: CompletionDeliveryRequest, ): Promise> { const { operation, payload, destination } = request; const deliveryOptions: CompletionDeliveryRequestOptions = { signal: request.signal }; let stablePayload: CompletionDeliveryPayload; try { validateOperation(operation); validatePayload(payload); stablePayload = freezePayload(payload); if ( typeof destination !== "function" && (destination == null || typeof destination.accept !== "function") ) { throw new Error("Completion delivery requires a destination adapter."); } } catch (error) { return Promise.reject(error); } if (!this.acceptingDeliveries) return Promise.resolve({ status: "suppressed" }); const key = operationKey(operation); const existing = this.states.get(key) as DeliveryState | undefined; if (existing) { if (existing.terminalResult) return Promise.resolve(existing.terminalResult); return existing.promise; } let resolve!: (result: CompletionDeliveryResult) => void; const promise = new Promise>((resolvePromise) => { resolve = resolvePromise; }); const state: DeliveryState = { operation: { ...operation }, payload: stablePayload, destination: destination as CompletionDeliveryDestination, status: "pending", settled: false, resolve, promise, }; this.states.set(key, state); if (deliveryOptions.signal) { const onAbort = () => this.cancel(operation); state.removeAbortListener = () => deliveryOptions.signal?.removeEventListener("abort", onAbort); if (deliveryOptions.signal.aborted) { this.cancel(operation); return promise; } deliveryOptions.signal.addEventListener("abort", onAbort, { once: true }); } void this.startAttempt(state); return promise; } /** Suppress pending and in-flight delivery for one Completion operation. */ suppress(operation: CompletionOperation): void { this.terminate(operation, { status: "suppressed" }); } /** Cancel a foreground delivery without turning it into a Child failure. */ cancel(operation: CompletionOperation): void { this.terminate(operation, { status: "cancelled" }); } /** Suppress every operation still tracked by this parent runtime. */ suppressAll(): void { for (const state of this.states.values()) this.suppress(state.operation); } /** Close the parent delivery boundary during permanent shutdown. */ shutdown(): void { this.acceptingDeliveries = false; this.suppressAll(); } /** Reopen delivery for /reload within the same parent runtime. */ reopenForReload(): void { this.acceptingDeliveries = true; } status(operation: CompletionOperation): DeliveryState["status"] | undefined { validateOperation(operation); return this.states.get(operationKey(operation))?.status; } private terminate( operation: CompletionOperation, result: { status: "suppressed" } | { status: "cancelled" }, ): void { validateOperation(operation); const state = this.states.get(operationKey(operation)); if (!state || state.terminalResult) return; state.status = result.status; state.terminalResult = result; this.clearDeliveryInputs(state); state.removeAbortListener?.(); state.removeAbortListener = undefined; this.settle(state, result); } private async startAttempt( state: DeliveryState, ): Promise { if (state.status !== "pending" || !state.destination || !state.payload) return; state.status = "delivering"; const destination = state.destination; const payload = state.payload; try { // A signal may cancel synchronously after the handoff is reserved but // before this microtask invokes the destination. await Promise.resolve(); if (state.status !== "delivering") return; const value = typeof destination === "function" ? await destination(payload) : await destination.accept(payload); if (state.status === "delivering") this.accept(state, value); } catch (error) { if (error instanceof CompletionDeliverySuppressedError) { this.terminate(state.operation, { status: "suppressed" }); return; } this.fail(state, error); } } private accept(state: DeliveryState, value: Value): void { if (state.status !== "delivering") return; state.value = value; state.status = "delivered"; state.terminalResult = { status: "delivered", value }; // Successful handoff consumes this operation's evidence. A failure to // remove an already-delivered sidecar must not change the parent result. try { clearCompletionSidecar(state.operation); } catch { // Best-effort artifact cleanup; there is no cross-runtime recovery path. } this.clearDeliveryInputs(state); state.removeAbortListener?.(); state.removeAbortListener = undefined; this.settle(state, state.terminalResult); } private fail(state: DeliveryState, error: unknown): void { if (state.terminalResult) return; const message = error instanceof Error ? error.message : String(error); state.error = message; state.status = "failed"; state.terminalResult = { status: "failed", error: message }; this.clearDeliveryInputs(state); state.removeAbortListener?.(); state.removeAbortListener = undefined; this.settle(state, state.terminalResult); } private clearDeliveryInputs(state: DeliveryState): void { state.payload = undefined; state.destination = undefined; } private settle( state: DeliveryState, result: CompletionDeliveryResult, ): void { if (state.settled) return; state.settled = true; state.resolve(result); } } export function createCompletionDeliveryManager( options: CompletionDeliveryManagerOptions = {}, ): CompletionDeliveryManager { return new CompletionDeliveryManager(options); }