diff --git a/adapter/strings/omp-spt.mjs b/adapter/strings/omp-spt.mjs index e6ef979..8c8fc7d 100644 --- a/adapter/strings/omp-spt.mjs +++ b/adapter/strings/omp-spt.mjs @@ -11,190 +11,1138 @@ export function decodeBody(body) { .replaceAll("&", "&"); } -function attribute(tag, name) { - const marker = ` ${name}="`; - const start = tag.indexOf(marker); - if (start < 0) return undefined; - const valueStart = start + marker.length; - const end = tag.indexOf('"', valueStart); - return end < 0 ? undefined : tag.slice(valueStart, end); +function protocolError(message) { + const error = new Error(`invalid spt EVENT stream: ${message}`); + error.code = "SPT_PROTOCOL_ERROR"; + return error; } -export function drainEvents(raw) { +function parseEventTag(tag) { + if (!tag.startsWith("]*)"/.exec(tag.slice(cursor)); + if (!match) return { error: protocolError("malformed EVENT attributes") }; + const [, name, value] = match; + if (Object.hasOwn(attributes, name)) { + return { error: protocolError(`duplicate EVENT ${name} attribute`) }; + } + attributes[name] = value; + cursor += match[0].length; + } + if (!attributes.type) return { error: protocolError("missing EVENT type attribute") }; + if (attributes.type === "msg" && !attributes.from) { + return { error: protocolError("missing EVENT from attribute") }; + } + return { attributes }; +} + +export function drainEvents(raw, options = {}) { + const maxEvents = options.maxEvents ?? Number.POSITIVE_INFINITY; + const maxFrameChars = options.maxFrameChars ?? DEFAULT_LISTENER_BUFFER_LIMIT; const events = []; let cursor = 0; while (true) { const start = raw.indexOf("", start); - if (openEnd < 0) return { events, rest: raw.slice(start) }; + if (openEnd < 0) { + if (raw.length - start > maxFrameChars) { + return { + error: protocolError(`EVENT frame exceeded ${maxFrameChars} characters`), + events, + rest: raw.slice(start), + }; + } + return { events, rest: raw.slice(start) }; + } + if (openEnd + 1 - start > maxFrameChars) { + return { + error: protocolError(`EVENT frame exceeded ${maxFrameChars} characters`), + events, + rest: raw.slice(start), + }; + } const close = raw.indexOf("", openEnd + 1); - if (close < 0) return { events, rest: raw.slice(start) }; - const tag = raw.slice(start, openEnd); - if (attribute(tag, "type") === "msg") { + const nested = raw.indexOf("= 0 && (close < 0 || nested < close)) { + return { + error: protocolError("nested EVENT before closing the current frame"), + events, + rest: raw.slice(start), + }; + } + if (close < 0) { + if (raw.length - start > maxFrameChars) { + return { + error: protocolError(`EVENT frame exceeded ${maxFrameChars} characters`), + events, + rest: raw.slice(start), + }; + } + return { events, rest: raw.slice(start) }; + } + const end = close + "".length; + if (end - start > maxFrameChars) { + return { + error: protocolError(`EVENT frame exceeded ${maxFrameChars} characters`), + events, + rest: raw.slice(start), + }; + } + const parsed = parseEventTag(raw.slice(start, openEnd)); + if (parsed.error) return { error: parsed.error, events, rest: raw.slice(start) }; + if (parsed.attributes.type === "msg") { events.push({ - from: attribute(tag, "from"), + from: parsed.attributes.from, body: decodeBody(raw.slice(openEnd + 1, close)), + envelope: raw.slice(start, end), }); } - cursor = close + "".length; + cursor = end; + if (events.length >= maxEvents) return { events, rest: raw.slice(cursor) }; } } -export function extractReply(messages) { - const assistant = [...(messages ?? [])].reverse().find((message) => message?.role === "assistant"); - if (!assistant) return ""; - if (typeof assistant.content === "string") return assistant.content; - return (assistant.content ?? []) +function messageText(message) { + if (typeof message?.content === "string") return message.content; + return (message?.content ?? []) .filter((part) => part?.type === "text" && typeof part.text === "string") .map((part) => part.text) .join(""); } +export function extractReply(messages, afterUserMessage) { + const allMessages = messages ?? []; + let start = 0; + if (afterUserMessage !== undefined) { + const userIndex = allMessages.findLastIndex((message) => { + if (message?.role !== "user") return false; + const text = messageText(message); + return text === afterUserMessage || text.startsWith(`${afterUserMessage}\n\n message?.role === "assistant"); + return assistant ? messageText(assistant) : ""; +} + function firstLine(text) { return text.split(/\r?\n/).find((line) => line.trim()) ?? ""; } -function runSpt(args, input) { +function errorSummary(error) { + const detail = error instanceof Error ? error.message : String(error); + return firstLine(detail).trim() || "unknown error"; +} + +function failureMessage(reason, error) { + const detail = error === undefined ? "" : `: ${errorSummary(error)}`; + return `[omp-spt] ${reason}${detail}`; +} + +function senderStub(sender) { + const escaped = sender + .replaceAll("&", "&") + .replaceAll('"', """) + .replaceAll("<", "<") + .replaceAll(">", ">"); + return ``; +} + +function injectEnvelope(messages, item) { + const index = messages.findLastIndex( + (message) => message?.role === "user" && messageText(message) === item.stub, + ); + if (index < 0) return messages; + const original = messages[index]; + const content = + typeof original.content === "string" + ? `${original.content}\n\n${item.envelope}` + : [...(original.content ?? []), { type: "text", text: `\n\n${item.envelope}` }]; + const injected = [...messages]; + injected[index] = { ...original, content }; + return injected; +} + +const DEFAULT_COMMAND_TIMEOUT_MS = 15_000; +const DEFAULT_KILL_GRACE_MS = 100; +const DEFAULT_KILL_FORCE_MS = 100; +const DEFAULT_LISTENER_BUFFER_LIMIT = 256 * 1024; +const DEFAULT_ACCEPTED_QUEUE_LIMIT = 128; +const DEFAULT_ACCEPTED_BYTES_LIMIT = 1024 * 1024; +const DEFAULT_SHUTDOWN_BUDGET_MS = 1_800; +const DEFAULT_SHUTDOWN_COMMAND_TIMEOUT_MS = 300; + +function childExited(child) { + return ( + (child.exitCode !== undefined && child.exitCode !== null) || + (child.signalCode !== undefined && child.signalCode !== null) + ); +} + +function waitForChildExit(child, timeoutMs, setTimer, clearTimer) { + if (childExited(child)) return Promise.resolve(true); + return new Promise((resolve) => { + let timer; + let finished = false; + const finish = (exited) => { + if (finished) return; + finished = true; + if (timer !== undefined) clearTimer(timer); + child.off("close", onClose); + resolve(exited); + }; + const onClose = () => finish(true); + child.once("close", onClose); + timer = setTimer(() => finish(false), timeoutMs); + timer?.unref?.(); + if (childExited(child)) finish(true); + }); +} + +async function terminateChild(child, label, options) { + const { clearTimer, forceMs, graceMs, setTimer } = options; + if (childExited(child)) return; + + const gracefulExit = waitForChildExit(child, graceMs, setTimer, clearTimer); + let killError; + try { + child.kill(); + } catch (error) { + killError = error; + } + if (await gracefulExit) return; + + const forcedExit = waitForChildExit(child, forceMs, setTimer, clearTimer); + try { + child.kill("SIGKILL"); + } catch (error) { + killError ??= error; + } + if (await forcedExit) return; + + child.stdin?.destroy?.(); + child.stdout?.destroy?.(); + child.stderr?.destroy?.(); + child.unref?.(); + const detail = killError === undefined ? "" : `: ${errorSummary(killError)}`; + throw new Error(`${label} did not exit after forced termination${detail}`); +} + +function commandLabel(args) { + const command = args[0] === "api" ? args[3] : args[0]; + return `spt ${command ?? "command"}`; +} + +export function runSpt(args, input, overrides = {}) { + const spawnProcess = overrides.spawnProcess ?? spawn; + const setTimer = overrides.setTimeout ?? globalThis.setTimeout; + const clearTimer = overrides.clearTimeout ?? globalThis.clearTimeout; + const timeoutMs = overrides.commandTimeoutMs ?? DEFAULT_COMMAND_TIMEOUT_MS; + const graceMs = overrides.killGraceMs ?? DEFAULT_KILL_GRACE_MS; + const forceMs = overrides.killForceMs ?? DEFAULT_KILL_FORCE_MS; + const env = overrides.env ?? process.env; + const label = commandLabel(args); + const signal = overrides.signal; + if (signal?.aborted) { + return Promise.reject( + signal.reason instanceof Error ? signal.reason : new Error(`${label} aborted`), + ); + } + return new Promise((resolve, reject) => { - const child = spawn(process.env.OMP_SPT_SPT_BIN || "spt", args, { - stdio: [input === undefined ? "ignore" : "pipe", "pipe", "pipe"], - windowsHide: true, - }); + let child; + try { + child = spawnProcess(env.OMP_SPT_SPT_BIN || "spt", args, { + stdio: [input === undefined ? "ignore" : "pipe", "pipe", "pipe"], + windowsHide: true, + }); + } catch (error) { + reject(error); + return; + } + let output = ""; + let settled = false; + let terminating = false; + let stdinFinished = input === undefined; + let timeoutTimer; + const onAbort = () => { + const error = + signal.reason instanceof Error ? signal.reason : new Error(`${label} aborted`); + void terminateAndReject(error); + }; + const onStdout = (chunk) => (output += chunk); + const onStderr = (chunk) => (output += chunk); + const cleanup = () => { + if (timeoutTimer !== undefined) clearTimer(timeoutTimer); + child.stdout.off("data", onStdout); + child.stderr.off("data", onStderr); + child.stdin?.off("finish", onStdinFinish); + child.off("close", onClose); + signal?.removeEventListener("abort", onAbort); + }; + const settle = (error) => { + if (settled) return; + settled = true; + cleanup(); + if (error === undefined) resolve(output.trim()); + else reject(error); + }; + const terminateAndReject = async (error) => { + if (settled || terminating) return; + terminating = true; + if (timeoutTimer !== undefined) { + clearTimer(timeoutTimer); + timeoutTimer = undefined; + } + try { + await terminateChild(child, label, { + clearTimer, + forceMs, + graceMs, + setTimer, + }); + } catch (terminationError) { + error = new Error(`${errorSummary(error)}; ${errorSummary(terminationError)}`, { + cause: error, + }); + } + settle(error); + }; + const onStdinError = (error) => { + void terminateAndReject(error); + }; + const onStdinFinish = () => { + stdinFinished = true; + }; + const onClose = (code, signal) => { + if (terminating || settled) return; + if (code === 0 && stdinFinished) { + settle(); + return; + } + const status = signal ? `signal ${signal}` : `exit ${code}`; + const detail = firstLine(output); + const suffix = detail ? `: ${detail}` : ""; + if (code === 0) { + void terminateAndReject(new Error(`${label} exited before stdin completed${suffix}`)); + } else { + void terminateAndReject(new Error(`${label} ${status}${suffix}`)); + } + }; + child.stdout.setEncoding("utf8"); child.stderr.setEncoding("utf8"); - child.stdout.on("data", (chunk) => (output += chunk)); - child.stderr.on("data", (chunk) => (output += chunk)); - child.on("error", reject); - child.on("close", (code) => { - if (code === 0) resolve(output.trim()); - else reject(new Error(`spt exited ${code}: ${firstLine(output)}`)); - }); - if (input !== undefined) child.stdin.end(input); + child.stdout.on("data", onStdout); + child.stderr.on("data", onStderr); + child.on("error", onStdinError); + child.on("close", onClose); + if (input !== undefined) { + child.stdin.on("error", onStdinError); + child.stdin.once("finish", onStdinFinish); + } + signal?.addEventListener("abort", onAbort, { once: true }); + if (signal?.aborted) { + onAbort(); + return; + } + timeoutTimer = setTimer(() => { + timeoutTimer = undefined; + void terminateAndReject(new Error(`${label} timed out after ${timeoutMs}ms`)); + }, timeoutMs); + timeoutTimer?.unref?.(); + if (input !== undefined) { + try { + child.stdin.end(input); + } catch (error) { + void terminateAndReject(error); + } + } }); } -// [impl->REQ-OMP-NATIVE-TUI] -export default function ompSpt(pi) { - const id = process.env.SPT_ENDPOINT_ID?.trim(); - if (!id) return; - - let sid; - let token; - let listener; - let listenerBuffer = ""; - let agentActive = false; - let current; - let stopping = false; - let ui; - const queue = []; - - const logError = (message, error) => { - pi.logger.error(message, { error: String(error) }); - ui?.notify(`${message}: ${error}`, "error"); - }; +export function createOmpSpt(overrides = {}) { + const spawnProcess = overrides.spawnProcess ?? spawn; + const setTimer = overrides.setTimeout ?? globalThis.setTimeout; + const clearTimer = overrides.clearTimeout ?? globalThis.clearTimeout; + const env = overrides.env ?? process.env; + const killGraceMs = overrides.killGraceMs ?? DEFAULT_KILL_GRACE_MS; + const killForceMs = overrides.killForceMs ?? DEFAULT_KILL_FORCE_MS; + const customRunSptCommand = overrides.runSptCommand; + const commandTimeoutMs = overrides.commandTimeoutMs ?? DEFAULT_COMMAND_TIMEOUT_MS; + const runSptCommand = + customRunSptCommand ?? + ((args, input, options = {}) => + runSpt(args, input, { + clearTimeout: clearTimer, + commandTimeoutMs: options.timeoutMs ?? commandTimeoutMs, + env, + killForceMs, + killGraceMs, + setTimeout: setTimer, + signal: options.signal, + spawnProcess, + })); + const shutdownBudgetMs = + overrides.shutdownBudgetMs ?? DEFAULT_SHUTDOWN_BUDGET_MS; + const shutdownCommandTimeoutMs = + overrides.shutdownCommandTimeoutMs ?? DEFAULT_SHUTDOWN_COMMAND_TIMEOUT_MS; + const listenerBufferLimit = + overrides.listenerBufferLimit ?? DEFAULT_LISTENER_BUFFER_LIMIT; + const acceptedQueueLimit = + overrides.acceptedQueueLimit ?? DEFAULT_ACCEPTED_QUEUE_LIMIT; + const acceptedBytesLimit = + overrides.acceptedBytesLimit ?? DEFAULT_ACCEPTED_BYTES_LIMIT; + const restartDelaysMs = [...(overrides.restartDelaysMs ?? [250, 1000, 4000])]; + const outcomeRetryDelaysMs = [ + ...(overrides.outcomeRetryDelaysMs ?? [250, 1000, 4000]), + ]; + const sessionEndRetryDelaysMs = [ + ...(overrides.sessionEndRetryDelaysMs ?? [250, 1000]), + ]; + const listenerStableMs = + overrides.listenerStableMs === false ? undefined : (overrides.listenerStableMs ?? 30_000); - async function setState(state) { - if (!sid) return; - const auth = token ? ["--token", token] : ["--session-id", sid]; - await runSpt(["api", "--adapter", ADAPTER, "state", state, id, ...auth]); - } + return function ompSpt(pi) { + const id = env.SPT_ENDPOINT_ID?.trim(); + if (!id) return; - function dispatchNext() { - if (stopping || agentActive || current || queue.length === 0) return; - current = queue.shift(); - try { - pi.sendUserMessage(current.body); - } catch (error) { - logError("omp-spt could not submit an inbound message", error); + let sid; + let token; + let listener; + let listenerBuffer = ""; + let listenerRestartCount = 0; + let listenerStableTimer; + let restartTimer; + let dispatchTimer; + let bindPromise; + let agentActive = false; + let desiredState = "idle"; + let dispatching = false; + let turnCompletionPromise; + let listenerTerminationPromise; + let shutdownMode = false; + let shutdownDeadlineExpired = false; + const activeCommands = new Map(); + const retryWaiters = new Set(); + let current; + let stopping = false; + let ui; + let runtimeCtx; + let endpointState; + let stateOperation = Promise.resolve(); + let endPromise; + let fatalPromise; + let teardownPromise; + let acceptedBytes = 0; + let overflowItem; + const queue = []; + + const logError = (message, error) => { + pi.logger.error(message, { error: errorSummary(error) }); + ui?.notify(`${message}: ${errorSummary(error)}`, "error"); + }; + + function runCommand(args, input, options = {}) { + if (shutdownDeadlineExpired) { + return Promise.reject(new Error("omp-spt shutdown deadline expired")); + } + const timeoutMs = + options.timeoutMs ?? (shutdownMode ? shutdownCommandTimeoutMs : commandTimeoutMs); + const controller = new AbortController(); + activeCommands.set(controller, { args, abortTimer: undefined }); + let command; + try { + command = Promise.resolve( + runSptCommand(args, input, { + signal: controller.signal, + timeoutMs, + }), + ); + } catch (error) { + activeCommands.delete(controller); + return Promise.reject(error); + } + if (customRunSptCommand) { + const rawCommand = command; + command = new Promise((resolve, reject) => { + let timer; + let finished = false; + const finish = (error, value) => { + if (finished) return; + finished = true; + if (timer !== undefined) clearTimer(timer); + controller.signal.removeEventListener("abort", onAbort); + if (error === undefined) resolve(value); + else reject(error); + }; + const onAbort = () => + finish( + controller.signal.reason instanceof Error + ? controller.signal.reason + : new Error(`${commandLabel(args)} aborted`), + ); + controller.signal.addEventListener("abort", onAbort, { once: true }); + if (shutdownMode) { + timer = setTimer( + () => + controller.abort( + new Error(`${commandLabel(args)} timed out after ${timeoutMs}ms`), + ), + timeoutMs, + ); + timer?.unref?.(); + } + rawCommand.then( + (value) => finish(undefined, value), + (error) => finish(error), + ); + }); + } + return command.finally(() => { + const active = activeCommands.get(controller); + if (active?.abortTimer !== undefined) clearTimer(active.abortTimer); + activeCommands.delete(controller); + }); + } + + function abortActiveCommands(reason, allowBindGrace = false) { + for (const [controller, active] of activeCommands) { + const isBind = active.args[0] === "api" && active.args[3] === "bind"; + if (allowBindGrace && isBind && active.abortTimer === undefined) { + active.abortTimer = setTimer( + () => controller.abort(reason), + shutdownCommandTimeoutMs, + ); + active.abortTimer?.unref?.(); + continue; + } + controller.abort(reason); + } + } + + function waitForRetry(delay) { + if (shutdownMode) return Promise.resolve(); + return new Promise((resolve) => { + let timer; + const finish = () => { + if (timer !== undefined) clearTimer(timer); + retryWaiters.delete(finish); + resolve(); + }; + retryWaiters.add(finish); + timer = setTimer(finish, delay); + timer?.unref?.(); + }); + } + + function enterShutdownMode() { + if (shutdownMode) return; + shutdownMode = true; + const reason = new Error("omp-spt command interrupted for bounded shutdown"); + abortActiveCommands(reason, true); + for (const finish of [...retryWaiters]) finish(); + } + + function authArgs() { + if (!token) throw new Error("bind did not return an authentication token"); + return ["--token", token]; + } + + function setState(state) { + if (!sid || !token || stopping) return Promise.resolve(); + const operation = stateOperation.catch(() => {}).then(async () => { + if (endpointState === state || stopping) return; + await runCommand(["api", "--adapter", ADAPTER, "state", state, id, ...authArgs()]); + endpointState = state; + }); + stateOperation = operation; + return operation; + } + + async function syncDesiredState() { + if (!bindPromise) return; + await bindPromise; + while (!stopping && endpointState !== desiredState) { + await setState(desiredState); + } + } + + function endSession() { + if (!sid || !token) return Promise.resolve(); + if (!endPromise) { + const operation = (async () => { + await stateOperation.catch(() => {}); + await runCommand(["api", "--adapter", ADAPTER, "session-end", id, ...authArgs()]); + endpointState = undefined; + })(); + endPromise = operation; + void operation.catch(() => { + if (endPromise === operation) endPromise = undefined; + }); + } + return endPromise; + } + + async function endSessionWithRetry() { + for (let attempt = 0; ; attempt += 1) { + try { + await endSession(); + return; + } catch (error) { + if (shutdownMode || attempt >= sessionEndRetryDelaysMs.length) throw error; + const delay = sessionEndRetryDelaysMs[attempt]; + pi.logger.error( + `omp-spt session teardown failed; retrying ${ + attempt + 1 + }/${sessionEndRetryDelaysMs.length} in ${delay}ms`, + { error: errorSummary(error) }, + ); + await waitForRetry(delay); + } + } + } + + // [impl->REQ-OMP-EXTENSION-CUSTODY] + function settleItem(item, payload) { + if (!item) return Promise.resolve(); + if (item.outcomePromise) return item.outcomePromise; + item.settling = true; + item.outcomePromise = (async () => { + if (!item.from) throw new Error("missing EVENT from attribute"); + for (let attempt = 0; ; attempt += 1) { + try { + await runCommand(["send", item.from, "--from", id], payload); + item.settled = true; + return; + } catch (error) { + if (shutdownMode || attempt >= outcomeRetryDelaysMs.length) throw error; + const delay = outcomeRetryDelaysMs[attempt]; + pi.logger.error( + `omp-spt could not send the outcome to ${item.from}; retrying ${ + attempt + 1 + }/${outcomeRetryDelaysMs.length} in ${delay}ms`, + { error: errorSummary(error) }, + ); + await waitForRetry(delay); + } + } + })(); + return item.outcomePromise; + } + + function releaseItem(item) { + if (!item?.accounted) return; + item.accounted = false; + acceptedBytes -= item.acceptedBytes; + } + + function beginListenerTermination(child, label) { + if (listenerTerminationPromise) return listenerTerminationPromise; + const operation = (async () => { + try { + await terminateChild(child, label, { + clearTimer, + forceMs: killForceMs, + graceMs: killGraceMs, + setTimer, + }); + } catch (error) { + pi.logger.error("omp-spt could not reap the listener", { + error: errorSummary(error), + }); + } + })(); + listenerTerminationPromise = operation; + void operation.then(() => { + if (listenerTerminationPromise === operation) listenerTerminationPromise = undefined; + }); + return operation; + } + + async function stopResources() { + if (dispatchTimer !== undefined) { + clearTimer(dispatchTimer); + dispatchTimer = undefined; + } + if (restartTimer !== undefined) { + clearTimer(restartTimer); + restartTimer = undefined; + } + if (listenerStableTimer !== undefined) { + clearTimer(listenerStableTimer); + listenerStableTimer = undefined; + } + const child = listener; + listener = undefined; + listenerBuffer = ""; + if (child) { + await beginListenerTermination(child, "spt ready listener"); + } else { + await listenerTerminationPromise; + } + } + + async function settlePendingItem(item, reason) { + try { + if (item.outcomePromise && !item.settled) { + let existingError; + try { + await item.outcomePromise; + } catch (error) { + existingError = error; + } + if (item.settled) return; + if (existingError && !shutdownMode) { + logError( + `omp-spt could not return custody to ${item.from ?? "unknown"}`, + existingError, + ); + return; + } + item.outcomePromise = undefined; + item.settling = false; + } + try { + await settleItem(item, failureMessage(reason)); + } catch (error) { + logError(`omp-spt could not return custody to ${item.from ?? "unknown"}`, error); + } + } finally { + releaseItem(item); + } + } + + async function failPending(reason) { + const pending = current ? [current, ...queue] : [...queue]; + if (overflowItem) pending.push(overflowItem); current = undefined; - setTimeout(dispatchNext, 0); + queue.length = 0; + overflowItem = undefined; + for (let index = 0; index < pending.length; index += 1) { + if (shutdownMode) { + await Promise.all( + pending.slice(index).map((item) => settlePendingItem(item, reason)), + ); + return; + } + await settlePendingItem(pending[index], reason); + } } - } - function startListener() { - const args = ["ready", id]; - if (process.env.OMP_SPT_SUBNET) args.push("--subnet", process.env.OMP_SPT_SUBNET); - listener = spawn(process.env.OMP_SPT_SPT_BIN || "spt", args, { - stdio: ["ignore", "pipe", "pipe"], - windowsHide: true, - }); - listener.stdout.setEncoding("utf8"); - listener.stderr.setEncoding("utf8"); - listener.stdout.on("data", (chunk) => { - listenerBuffer += chunk; - const drained = drainEvents(listenerBuffer); - listenerBuffer = drained.rest; - for (const event of drained.events) queue.push(event); - dispatchNext(); - }); - listener.stderr.on("data", (chunk) => pi.logger.debug("omp-spt listener", { output: chunk.trim() })); - listener.on("error", (error) => logError("omp-spt listener failed", error)); - listener.on("close", (code) => { + function teardownSession(pendingReason) { + if (!teardownPromise) { + stopping = true; + const operation = (async () => { + await stopResources(); + await failPending(pendingReason); + await bindPromise?.catch(() => {}); + await endSessionWithRetry(); + })(); + teardownPromise = operation; + void operation.catch(() => { + if (teardownPromise === operation) teardownPromise = undefined; + }); + } + return teardownPromise; + } + + async function shutdownWithinBudget(pendingReason) { + enterShutdownMode(); + const teardown = teardownSession(pendingReason); + let budgetTimer; + const expired = new Promise((resolve) => { + budgetTimer = setTimer(() => { + budgetTimer = undefined; + shutdownDeadlineExpired = true; + const error = new Error( + `omp-spt shutdown exceeded its ${shutdownBudgetMs}ms budget`, + ); + abortActiveCommands(error); + for (const finish of [...retryWaiters]) finish(); + resolve(false); + }, shutdownBudgetMs); + budgetTimer?.unref?.(); + }); + const completed = teardown.then( + () => true, + (error) => { + logError("omp-spt session teardown failed", error); + return true; + }, + ); + const finished = await Promise.race([completed, expired]); + if (budgetTimer !== undefined) clearTimer(budgetTimer); + if (!finished) { + pi.logger.error("omp-spt bounded shutdown expired", { + error: `${shutdownBudgetMs}ms budget exhausted`, + }); + } + } + + // [impl->REQ-OMP-LISTENER-FAIL-CLOSED] + async function failClosed(message, error) { + if (stopping && shutdownMode) return teardownPromise ?? Promise.resolve(); + if (fatalPromise) return fatalPromise; + fatalPromise = (async () => { + ui?.setStatus("omp-spt", "spt failed"); + logError(message, error); + try { + await teardownSession("endpoint stopped before your message could complete"); + } catch (teardownError) { + logError("omp-spt session teardown failed", teardownError); + } + runtimeCtx?.shutdown(); + })(); + return fatalPromise; + } + + function scheduleDispatch() { + if ( + stopping || + agentActive || + dispatching || + current || + queue.length === 0 || + dispatchTimer !== undefined + ) { + return; + } + dispatchTimer = setTimer(() => { + dispatchTimer = undefined; + void dispatchNext().catch((error) => { + if (!stopping) return failClosed("omp-spt dispatch failed", error); + }); + }, 0); + dispatchTimer?.unref?.(); + } + + async function rejectItem(item, reason, error) { + if (stopping) return; + logError(`omp-spt ${reason}`, error); + try { + await settleItem(item, failureMessage(reason, error)); + } catch (outcomeError) { + if (stopping) return; + await failClosed(`omp-spt could not send the outcome to ${item.from ?? "unknown"}`, outcomeError); + return; + } + if (current === item) { + current = undefined; + releaseItem(item); + } if (!stopping) { - ui?.setStatus("omp-spt", "spt offline"); - ui?.notify(`omp-spt listener exited (${code})`, "error"); + desiredState = "idle"; + try { + await setState("idle"); + } catch (stateError) { + await failClosed( + "omp-spt could not restore idle state after a failed submission", + stateError, + ); + } } - }); - } + } - pi.on("session_start", async (_event, ctx) => { - ui = ctx.ui; - sid = ctx.sessionManager.getSessionId(); - try { - const bindArgs = ["api", "--adapter", ADAPTER, "bind", id, "--set-session-id", sid]; - if (process.env.OMP_SPT_SUBNET) bindArgs.push("--subnet", process.env.OMP_SPT_SUBNET); - const bound = await runSpt(bindArgs); - token = bound.match(/\btoken=([^\s]+)/)?.[1]; - await setState("idle"); - startListener(); - ctx.ui.setStatus("omp-spt", `spt:${id}`); - } catch (error) { - ctx.ui.setStatus("omp-spt", "spt bind failed"); - logError(`omp-spt could not bind ${id}`, error); + // [impl->REQ-OMP-EXTENSION-CUSTODY] + async function dispatchNext() { + if (stopping || agentActive || dispatching || current || queue.length === 0) return; + dispatching = true; + const item = queue.shift(); + current = item; + try { + try { + desiredState = "busy"; + await setState("busy"); + } catch (error) { + await rejectItem(item, "could not accept your message", error); + return; + } + if (stopping) return; + if (agentActive) { + if (current === item) current = undefined; + queue.unshift(item); + return; + } + item.stub = senderStub(item.from ?? "unknown"); + item.submitted = true; + try { + pi.sendUserMessage(item.stub); + } catch (error) { + item.submitted = false; + await rejectItem(item, "could not submit your message to OMP", error); + } + } finally { + dispatching = false; + scheduleDispatch(); + } } - }); - pi.on("agent_start", async () => { - agentActive = true; - try { - await setState("busy"); - } catch (error) { - logError("omp-spt could not mark the endpoint busy", error); + // [impl->REQ-OMP-LISTENER-FAIL-CLOSED] + function handleListenerDeath(reason) { + listenerBuffer = ""; + if (listenerStableTimer !== undefined) { + clearTimer(listenerStableTimer); + listenerStableTimer = undefined; + } + if (stopping || restartTimer !== undefined) return; + if (listenerRestartCount >= restartDelaysMs.length) { + void failClosed("omp-spt listener restart budget exhausted", reason); + return; + } + const attempt = listenerRestartCount + 1; + const delay = restartDelaysMs[listenerRestartCount]; + listenerRestartCount = attempt; + const message = `omp-spt listener stopped; restarting ${attempt}/${restartDelaysMs.length} in ${delay}ms`; + pi.logger.error(message, { error: errorSummary(reason) }); + ui?.setStatus("omp-spt", `spt reconnecting (${attempt}/${restartDelaysMs.length})`); + ui?.notify(message, "warning"); + restartTimer = setTimer(() => { + restartTimer = undefined; + startListener(); + }, delay); + restartTimer?.unref?.(); } - }); - pi.on("agent_end", async (event) => { - agentActive = false; - const completed = current; - current = undefined; - if (completed?.from) { - const reply = extractReply(event.messages) || "[omp-spt] turn ended without an assistant response."; + function startListener() { + if (stopping) return; + const args = ["ready", id]; + if (env.OMP_SPT_SUBNET) args.push("--subnet", env.OMP_SPT_SUBNET); + let child; try { - await runSpt(["send", completed.from, "--from", id], reply); + child = spawnProcess(env.OMP_SPT_SPT_BIN || "spt", args, { + stdio: ["ignore", "pipe", "pipe"], + windowsHide: true, + }); } catch (error) { - logError(`omp-spt could not reply to ${completed.from}`, error); + handleListenerDeath(error); + return; } + listener = child; + listenerBuffer = ""; + let dead = false; + const died = (reason, alreadyExited) => { + if (dead) return; + dead = true; + if (listener === child) listener = undefined; + if (stopping || alreadyExited) { + handleListenerDeath(reason); + return; + } + const termination = beginListenerTermination(child, "dead spt ready listener"); + void termination.then(() => handleListenerDeath(reason)); + }; + if (listenerStableMs !== undefined) { + listenerStableTimer = setTimer(() => { + listenerStableTimer = undefined; + if (listener === child && !stopping) listenerRestartCount = 0; + }, listenerStableMs); + listenerStableTimer?.unref?.(); + } + child.stdout.setEncoding("utf8"); + child.stderr.setEncoding("utf8"); + child.stdout.on("data", (chunk) => { + if (listener !== child || stopping) return; + listenerBuffer += String(chunk); + if (listenerBuffer.length > listenerBufferLimit) { + void failClosed( + "omp-spt listener protocol corruption", + protocolError( + `EVENT buffer exceeded ${listenerBufferLimit} characters without a complete drain`, + ), + ); + return; + } + while (!stopping) { + const drained = drainEvents(listenerBuffer, { + maxEvents: 1, + maxFrameChars: listenerBufferLimit, + }); + listenerBuffer = drained.rest; + if (drained.error) { + void failClosed("omp-spt listener protocol corruption", drained.error); + return; + } + if (drained.events.length === 0) break; + const event = drained.events[0]; + const acceptedCount = queue.length + (current ? 1 : 0); + const eventBytes = Buffer.byteLength(event.envelope, "utf8"); + if ( + acceptedCount >= acceptedQueueLimit || + acceptedBytes + eventBytes > acceptedBytesLimit + ) { + overflowItem = event; + void failClosed( + "omp-spt inbound custody capacity exceeded", + new Error( + `accepted queue limit is ${acceptedQueueLimit} messages and ${acceptedBytesLimit} bytes`, + ), + ); + return; + } + event.acceptedBytes = eventBytes; + event.accounted = true; + acceptedBytes += eventBytes; + queue.push(event); + } + if (!agentActive && !current && !dispatching) void dispatchNext(); + }); + child.stderr.on("data", (chunk) => + pi.logger.debug("omp-spt listener", { output: String(chunk).trim() }), + ); + child.on("error", (error) => died(error, false)); + child.on("close", (code, signal) => { + const status = signal ? `signal ${signal}` : code; + died(new Error(`spt ready exited ${status}`), true); + }); + ui?.setStatus("omp-spt", `spt:${id}`); } - try { - await setState("idle"); - } catch (error) { - logError("omp-spt could not mark the endpoint idle", error); - } - setTimeout(dispatchNext, 0); - }); - pi.on("session_shutdown", async () => { - stopping = true; - listener?.kill(); - ui?.setStatus("omp-spt", undefined); - if (!sid) return; - const auth = token ? ["--token", token] : ["--session-id", sid]; - try { - await runSpt(["api", "--adapter", ADAPTER, "session-end", id, ...auth]); - } catch (error) { - pi.logger.error("omp-spt session teardown failed", { error: String(error) }); + // [impl->REQ-OMP-NATIVE-TUI] + pi.on("session_start", async (_event, ctx) => { + runtimeCtx = ctx; + ui = ctx.ui; + sid = ctx.sessionManager.getSessionId(); + const bindArgs = ["api", "--adapter", ADAPTER, "bind", id, "--set-session-id", sid]; + if (env.OMP_SPT_SUBNET) bindArgs.push("--subnet", env.OMP_SPT_SUBNET); + bindPromise = (async () => { + const bound = await runCommand(bindArgs); + token = bound.match(/\btoken=([^\s]+)/)?.[1]; + if (!token) throw new Error("spt bind response did not include token="); + })(); + try { + await bindPromise; + } catch (error) { + if (stopping) return; + ui.setStatus("omp-spt", "spt bind failed"); + await failClosed(`omp-spt could not bind ${id}`, error); + return; + } + if (stopping) { + try { + await teardownSession("OMP session shut down before initialization completed"); + } catch (error) { + logError("omp-spt session teardown failed", error); + } + return; + } + try { + await syncDesiredState(); + if (stopping) { + await teardownSession("OMP session shut down before initialization completed"); + return; + } + ui.setStatus("omp-spt", `spt:${id}`); + startListener(); + } catch (error) { + if (stopping) { + logError("omp-spt session teardown failed", error); + return; + } + ui.setStatus("omp-spt", "spt bind failed"); + await failClosed(`omp-spt could not bind ${id}`, error); + } + }); + + // [impl->REQ-OMP-SESSION-IMMUTABLE] + const blockSessionChange = (description, ctx) => { + ctx.ui.notify( + `omp-spt blocked the in-TUI ${description}; end this SPT session first`, + "warning", + ); + return { cancel: true }; + }; + + pi.on("session_before_switch", (event, ctx) => + blockSessionChange(`${event.reason} session switch`, ctx), + ); + pi.on("session_before_branch", (_event, ctx) => + blockSessionChange("session branch", ctx), + ); + + // [impl->REQ-OMP-MESSAGE-CONTEXT] + pi.on("context", (event) => { + if (!current?.submitted || current.settling) return; + const messages = injectEnvelope(event.messages, current); + if (messages !== event.messages) return { messages }; + }); + + async function completeTurn(event) { + agentActive = false; + desiredState = "idle"; + if (stopping) return; + const completed = current; + if (completed?.submitted && !completed.settled) { + const reply = extractReply(event.messages, completed.stub); + try { + await settleItem( + completed, + reply || failureMessage("turn ended without an assistant response"), + ); + } catch (error) { + if (stopping) return; + await failClosed( + `omp-spt could not send the outcome to ${completed.from ?? "unknown"}`, + error, + ); + return; + } + } + if (current === completed) { + current = undefined; + releaseItem(completed); + } + try { + await setState("idle"); + } catch (error) { + if (stopping) return; + await failClosed("omp-spt could not mark the endpoint idle", error); + return; + } + scheduleDispatch(); } - }); + + pi.on("agent_start", async () => { + if (stopping) return; + turnCompletionPromise = undefined; + agentActive = true; + desiredState = "busy"; + try { + await syncDesiredState(); + } catch (error) { + if (!stopping) await failClosed("omp-spt could not mark the endpoint busy", error); + } + }); + + // [impl->REQ-OMP-EXTENSION-CUSTODY] + pi.on("agent_end", (event) => { + turnCompletionPromise ??= completeTurn(event); + return turnCompletionPromise; + }); + + pi.on("session_stop", async (event) => { + turnCompletionPromise ??= completeTurn(event); + await turnCompletionPromise; + }); + + pi.on("session_shutdown", async (_event, ctx) => { + runtimeCtx ??= ctx; + ui?.setStatus("omp-spt", undefined); + await shutdownWithinBudget("OMP session shut down before your message could complete"); + }); + }; } + +export default createOmpSpt(); diff --git a/tests/omp-extension.mjs b/tests/omp-extension.mjs index 2ccf034..6bde45a 100644 --- a/tests/omp-extension.mjs +++ b/tests/omp-extension.mjs @@ -1,25 +1,1266 @@ import assert from "node:assert/strict"; -import { decodeBody, drainEvents, extractReply } from "../adapter/strings/omp-spt.mjs"; +import { EventEmitter } from "node:events"; +import { + createOmpSpt, + decodeBody, + drainEvents, + extractReply, + runSpt, +} from "../adapter/strings/omp-spt.mjs"; +const flush = () => new Promise((resolve) => setImmediate(resolve)); +function deferred() { + let resolve; + let reject; + const promise = new Promise((resolvePromise, rejectPromise) => { + resolve = resolvePromise; + reject = rejectPromise; + }); + return { promise, reject, resolve }; +} + +class FakeStream extends EventEmitter { + setEncoding(encoding) { + this.encoding = encoding; + } + + end(input) { + this.input = input; + if (this.onEnd?.(input) === false) return; + this.emit("finish"); + } +} + +class FakeChild extends EventEmitter { + constructor(options = {}) { + super(); + this.stdin = new FakeStream(); + this.stdout = new FakeStream(); + this.stderr = new FakeStream(); + this.exitCode = null; + this.signalCode = null; + this.kills = 0; + this.killSignals = []; + this.onKill = options.onKill; + } + + close(code = 0, signal = null) { + this.exitCode = code; + this.signalCode = signal; + this.emit("close", code, signal); + } + + kill(signal = "SIGTERM") { + this.kills += 1; + this.killSignals.push(signal); + const handled = this.onKill?.(signal, this); + if (handled !== undefined) return handled; + this.close(null, signal); + return true; + } +} + +class FakeClock { + constructor() { + this.nextId = 1; + this.timers = new Map(); + } + + setTimeout(fn, delay) { + const handle = { id: this.nextId++, unref() {} }; + this.timers.set(handle, { fn, delay }); + return handle; + } + + clearTimeout(handle) { + this.timers.delete(handle); + } + + delays() { + return [...this.timers.values()].map(({ delay }) => delay); + } + + async runNext(expectedDelay) { + const entry = [...this.timers.entries()].sort((left, right) => left[1].delay - right[1].delay)[0]; + assert.ok(entry, `expected a ${expectedDelay}ms timer`); + const [handle, timer] = entry; + assert.equal(timer.delay, expectedDelay); + this.timers.delete(handle); + timer.fn(); + await flush(); + } +} + +function createHarness(options = {}) { + const handlers = new Map(); + const calls = []; + const children = []; + const submitted = []; + const statuses = []; + const notifications = []; + const errors = []; + const debug = []; + const clock = new FakeClock(); + let shutdowns = 0; + + const ui = { + notify(message, type) { + notifications.push({ message, type }); + }, + setStatus(key, text) { + statuses.push({ key, text }); + }, + }; + const ctx = { + ui, + sessionManager: { getSessionId: () => options.sessionId ?? "session-1" }, + shutdown() { + shutdowns += 1; + }, + }; + const pi = { + logger: { + error(message, details) { + errors.push({ message, details }); + }, + debug(message, details) { + debug.push({ message, details }); + }, + }, + on(name, handler) { + const registered = handlers.get(name) ?? []; + registered.push(handler); + handlers.set(name, registered); + }, + sendUserMessage(content) { + submitted.push(content); + options.onSubmit?.(content); + }, + }; + const runSptCommand = async (args, input) => { + const call = { args: [...args], input }; + calls.push(call); + const overridden = await options.onRun?.(call); + if (overridden !== undefined) return overridden; + if (args[0] === "api" && args[3] === "bind") { + return options.bindOutput ?? "BOUND endpoint token=token-123"; + } + return ""; + }; + const spawnProcess = (binary, args, spawnOptions) => { + const child = new FakeChild(); + child.binary = binary; + child.args = [...args]; + child.spawnOptions = spawnOptions; + options.onSpawn?.(child); + children.push(child); + return child; + }; + const extension = createOmpSpt({ + env: { + SPT_ENDPOINT_ID: options.id ?? "omp-agent", + OMP_SPT_SUBNET: options.subnet, + OMP_SPT_SPT_BIN: "spt-test", + }, + acceptedBytesLimit: options.acceptedBytesLimit, + acceptedQueueLimit: options.acceptedQueueLimit, + restartDelaysMs: options.restartDelaysMs ?? [5, 10], + outcomeRetryDelaysMs: options.outcomeRetryDelaysMs ?? [], + sessionEndRetryDelaysMs: options.sessionEndRetryDelaysMs ?? [], + shutdownBudgetMs: options.shutdownBudgetMs, + shutdownCommandTimeoutMs: options.shutdownCommandTimeoutMs, + listenerStableMs: options.listenerStableMs ?? false, + killForceMs: options.killForceMs ?? 4, + killGraceMs: options.killGraceMs ?? 3, + listenerBufferLimit: options.listenerBufferLimit, + runSptCommand, + spawnProcess, + setTimeout: clock.setTimeout.bind(clock), + clearTimeout: clock.clearTimeout.bind(clock), + }); + extension(pi); + + async function emit(name, event = {}) { + let result; + for (const handler of handlers.get(name) ?? []) { + const returned = await handler({ type: name, ...event }, ctx); + if (returned !== undefined) result = returned; + } + return result; + } + + return { + calls, + children, + clock, + ctx, + debug, + emit, + errors, + handlers, + notifications, + statuses, + submitted, + get shutdowns() { + return shutdowns; + }, + }; +} + +function commandCalls(harness, command) { + return harness.calls.filter((call) => call.args[0] === command); +} + +function stateCalls(harness) { + return harness.calls.filter((call) => call.args[0] === "api" && call.args[3] === "state"); +} + +async function testParsingAndReplies() { + assert.equal(decodeBody('a<b>
"c" & &lt;'), 'a\n"c" & <'); + + const partialEnvelope = 'hello
wo'; + const partial = drainEvents(`noise${partialEnvelope}`); + assert.deepEqual(partial.events, []); + assert.equal(partial.rest, partialEnvelope); + + const envelope = `${partial.rest}rld
`; + const complete = drainEvents(`${envelope}skip`); + assert.deepEqual(complete.events, [ + { from: "doyle", body: "hello\nworld", envelope }, + ]); + assert.equal(complete.rest, ""); + + const truncatedA = + 'truncatedvalid'; + const nested = drainEvents(truncatedA); + assert.deepEqual(nested.events, []); + assert.equal(nested.rest, truncatedA); + assert.match(nested.error.message, /nested EVENT/); + assert.match( + drainEvents('missing sender').error.message, + /missing EVENT from/, + ); + assert.match( + drainEvents('bad attrs').error.message, + /malformed EVENT attributes/, + ); + assert.match( + drainEvents('0123456789', { + maxFrameChars: 32, + }).error.message, + /EVENT frame exceeded/, + ); + + assert.equal( + extractReply([ + { role: "assistant", content: [{ type: "text", text: "first" }] }, + { role: "toolResult", content: [] }, + { + role: "assistant", + content: [ + { type: "text", text: "final " }, + { type: "text", text: "answer" }, + ], + }, + ]), + "final answer", + ); + assert.equal( + extractReply( + [ + { role: "assistant", content: "stale answer" }, + { role: "user", content: '' }, + ], + '', + ), + "", + ); +} + +// [unit->REQ-OMP-EXTENSION-CUSTODY] +// [unit->REQ-OMP-SESSION-IMMUTABLE] +// [unit->REQ-OMP-MESSAGE-CONTEXT] // [unit->REQ-OMP-NATIVE-TUI] -assert.equal(decodeBody('a<b>
"c" & &lt;'), 'a\n"c" & <'); - -const partial = drainEvents('noisehello
wo'); -assert.deepEqual(partial.events, []); -assert.equal(partial.rest, 'hello
wo'); - -const complete = drainEvents(`${partial.rest}rld
skip`); -assert.deepEqual(complete.events, [{ from: "doyle", body: "hello\nworld" }]); -assert.equal(complete.rest, ""); - -assert.equal( - extractReply([ - { role: "assistant", content: [{ type: "text", text: "first" }] }, - { role: "toolResult", content: [] }, - { role: "assistant", content: [{ type: "text", text: "final " }, { type: "text", text: "answer" }] }, - ]), - "final answer", -); -assert.equal(extractReply([{ role: "user", content: "hello" }]), ""); +async function testLifecycleCustodyAndContext() { + const harness = createHarness({ subnet: "mesh-a" }); + assert.deepEqual([...harness.handlers.keys()], [ + "session_start", + "session_before_switch", + "session_before_branch", + "context", + "agent_start", + "agent_end", + "session_stop", + "session_shutdown", + ]); + + await harness.emit("session_start"); + assert.deepEqual(harness.calls[0], { + args: [ + "api", + "--adapter", + "omp-spt", + "bind", + "omp-agent", + "--set-session-id", + "session-1", + "--subnet", + "mesh-a", + ], + input: undefined, + }); + assert.deepEqual(harness.calls[1].args, [ + "api", + "--adapter", + "omp-spt", + "state", + "idle", + "omp-agent", + "--token", + "token-123", + ]); + assert.equal(harness.children.length, 1); + assert.equal(harness.children[0].binary, "spt-test"); + assert.deepEqual(harness.children[0].args, ["ready", "omp-agent", "--subnet", "mesh-a"]); + + for (const reason of ["new", "resume", "fork", "handoff"]) { + assert.deepEqual(await harness.emit("session_before_switch", { reason }), { cancel: true }); + assert.ok( + harness.notifications.some(({ message }) => + message.includes(`${reason} session switch`), + ), + ); + } + assert.deepEqual(await harness.emit("session_before_branch"), { cancel: true }); + assert.ok( + harness.notifications.some(({ message }) => message.includes("session branch")), + ); + + const aliceEnvelope = + 'hello<world
line
'; + const bobEnvelope = 'second'; + harness.children[0].stdout.emit("data", `${aliceEnvelope}${bobEnvelope}`); + await flush(); + assert.deepEqual(harness.submitted, ['']); + assert.deepEqual( + stateCalls(harness).map((call) => call.args[4]), + ["idle", "busy"], + ); + + const originalMessages = [{ role: "user", content: '' }]; + const context = await harness.emit("context", { messages: originalMessages }); + assert.equal(originalMessages[0].content, ''); + assert.equal(context.messages[0].content, `\n\n${aliceEnvelope}`); + + const arrayContext = await harness.emit("context", { + messages: [{ role: "user", content: [{ type: "text", text: '' }] }], + }); + assert.deepEqual(arrayContext.messages[0].content, [ + { type: "text", text: '' }, + { type: "text", text: `\n\n${aliceEnvelope}` }, + ]); + + await harness.emit("agent_start"); + assert.deepEqual( + stateCalls(harness).map((call) => call.args[4]), + ["idle", "busy"], + "agent_start must not duplicate the already-honest busy transition", + ); + await harness.emit("agent_end", { + messages: [ + { role: "user", content: '' }, + { role: "assistant", content: [{ type: "text", text: "alice reply" }] }, + ], + }); + assert.deepEqual(harness.submitted, ['']); + assert.deepEqual(harness.clock.delays(), [0]); + + await harness.clock.runNext(0); + assert.deepEqual(harness.submitted, ['', '']); + await harness.emit("agent_start"); + await harness.emit("agent_end", { + messages: [ + { role: "assistant", content: "previous reply" }, + { role: "user", content: '' }, + ], + }); + + const outcomes = commandCalls(harness, "send"); + assert.equal(outcomes.length, 2); + assert.deepEqual( + outcomes.map((call) => call.args), + [ + ["send", "alice", "--from", "omp-agent"], + ["send", "bob", "--from", "omp-agent"], + ], + ); + assert.equal(outcomes[0].input, "alice reply"); + assert.match(outcomes[1].input, /turn ended without an assistant response/); + assert.deepEqual( + stateCalls(harness).map((call) => call.args[4]), + ["idle", "busy", "idle", "busy", "idle"], + ); + for (const call of [...stateCalls(harness), ...harness.calls.filter((candidate) => candidate.args[3] === "session-end")]) { + assert.deepEqual(call.args.slice(-2), ["--token", "token-123"]); + } + + await harness.emit("session_shutdown"); + const ended = harness.calls.filter((call) => call.args[3] === "session-end"); + assert.equal(ended.length, 1); + assert.deepEqual(ended[0].args.slice(-2), ["--token", "token-123"]); + assert.equal(harness.children[0].kills, 1); + assert.deepEqual(harness.clock.delays(), []); +} + +async function testDeferredBindLifecycleSerialization() { + const busyBind = deferred(); + const busyHarness = createHarness({ + onRun(call) { + if (call.args[0] === "api" && call.args[3] === "bind") return busyBind.promise; + }, + }); + const startingBusy = busyHarness.emit("session_start"); + await flush(); + const becomingBusy = busyHarness.emit("agent_start"); + await flush(); + assert.deepEqual(stateCalls(busyHarness), []); + assert.equal(busyHarness.children.length, 0); + + busyBind.resolve("BOUND endpoint token=token-busy"); + await Promise.all([startingBusy, becomingBusy]); + assert.deepEqual( + stateCalls(busyHarness).map((call) => call.args[4]), + ["busy"], + "agent_start before bind completion must suppress the stale idle publication", + ); + assert.equal(busyHarness.children.length, 1); + await busyHarness.emit("session_shutdown"); + assert.deepEqual(busyHarness.clock.delays(), []); + + const shutdownBind = deferred(); + const shutdownHarness = createHarness({ + onRun(call) { + if (call.args[0] === "api" && call.args[3] === "bind") return shutdownBind.promise; + }, + }); + const startingShutdown = shutdownHarness.emit("session_start"); + await flush(); + const shuttingDown = shutdownHarness.emit("session_shutdown"); + await flush(); + assert.equal(shutdownHarness.children.length, 0); + assert.equal( + shutdownHarness.calls.filter((call) => call.args[3] === "session-end").length, + 0, + ); + + shutdownBind.resolve("BOUND endpoint token=token-shutdown"); + await Promise.all([startingShutdown, shuttingDown]); + assert.deepEqual(stateCalls(shutdownHarness), []); + assert.equal(shutdownHarness.children.length, 0); + assert.equal( + shutdownHarness.calls.filter((call) => call.args[3] === "session-end").length, + 1, + ); + assert.ok( + !shutdownHarness.statuses.some(({ text }) => text === "spt:omp-agent"), + "bind completion after shutdown must not restore live status", + ); + assert.deepEqual(shutdownHarness.clock.delays(), []); +} + +// [unit->REQ-OMP-EXTENSION-CUSTODY] +// [unit->REQ-OMP-LISTENER-FAIL-CLOSED] +async function testOutcomeSendRetriesAndExhaustion() { + let retryAttempts = 0; + const retryHarness = createHarness({ + outcomeRetryDelaysMs: [5, 10], + onRun(call) { + if (call.args[0] === "send" && call.args[1] === "retry") { + retryAttempts += 1; + if (retryAttempts < 3) throw new Error(`outcome failure ${retryAttempts}`); + } + }, + }); + await retryHarness.emit("session_start"); + retryHarness.children[0].stdout.emit( + "data", + 'work', + ); + await flush(); + await retryHarness.emit("agent_start"); + const retryEnding = retryHarness.emit("agent_end", { + messages: [ + { role: "user", content: '' }, + { role: "assistant", content: "eventual outcome" }, + ], + }); + await flush(); + assert.equal(commandCalls(retryHarness, "send").length, 1); + assert.deepEqual(retryHarness.clock.delays(), [5]); + await retryHarness.clock.runNext(5); + assert.equal(commandCalls(retryHarness, "send").length, 2); + assert.deepEqual(retryHarness.clock.delays(), [10]); + await retryHarness.clock.runNext(10); + await retryEnding; + assert.equal(commandCalls(retryHarness, "send").length, 3); + assert.equal(commandCalls(retryHarness, "send").at(-1).input, "eventual outcome"); + assert.equal(retryHarness.shutdowns, 0); + await retryHarness.emit("session_shutdown"); + assert.deepEqual(retryHarness.clock.delays(), []); + + const exhaustedHarness = createHarness({ + outcomeRetryDelaysMs: [7], + onRun(call) { + if (call.args[0] === "send") throw new Error("outcome channel unavailable"); + }, + }); + await exhaustedHarness.emit("session_start"); + const exhaustedListener = exhaustedHarness.children[0]; + exhaustedListener.stdout.emit( + "data", + 'work', + ); + await flush(); + await exhaustedHarness.emit("agent_start"); + const exhaustedEnding = exhaustedHarness.emit("agent_end", { + messages: [ + { role: "user", content: '' }, + { role: "assistant", content: "undeliverable outcome" }, + ], + }); + await flush(); + assert.equal(commandCalls(exhaustedHarness, "send").length, 1); + assert.deepEqual(exhaustedHarness.clock.delays(), [7]); + await exhaustedHarness.clock.runNext(7); + await exhaustedEnding; + + assert.equal(commandCalls(exhaustedHarness, "send").length, 2); + assert.deepEqual( + stateCalls(exhaustedHarness).map((call) => call.args[4]), + ["idle", "busy"], + "exhausted custody must never be advertised idle", + ); + assert.equal(exhaustedHarness.shutdowns, 1); + assert.equal(exhaustedListener.kills, 1); + assert.equal( + exhaustedHarness.calls.filter((call) => call.args[3] === "session-end").length, + 1, + ); + assert.ok( + exhaustedHarness.errors.some(({ message }) => + message.includes("could not send the outcome to exhausted"), + ), + ); + await exhaustedHarness.emit("session_shutdown"); + assert.deepEqual(exhaustedHarness.clock.delays(), []); +} + +// [unit->REQ-OMP-EXTENSION-CUSTODY] +async function testSubmissionFailureAdvancesQueue() { + const harness = createHarness({ + onSubmit(content) { + if (content === '') throw new Error("OMP prompt flow rejected input"); + }, + }); + await harness.emit("session_start"); + harness.children[0].stdout.emit( + "data", + 'onetwo', + ); + await flush(); + + const firstOutcome = commandCalls(harness, "send"); + assert.equal(firstOutcome.length, 1); + assert.deepEqual(firstOutcome[0].args, ["send", "broken", "--from", "omp-agent"]); + assert.match(firstOutcome[0].input, /could not submit your message to OMP/); + assert.deepEqual(harness.clock.delays(), [0]); + + await harness.clock.runNext(0); + assert.deepEqual(harness.submitted, ['', '']); + await harness.emit("agent_start"); + await harness.emit("agent_end", { + messages: [ + { role: "user", content: '' }, + { role: "assistant", content: "next reply" }, + ], + }); + const outcomes = commandCalls(harness, "send"); + assert.equal(outcomes.length, 2); + assert.deepEqual( + outcomes.map((call) => call.args[1]), + ["broken", "next"], + ); + assert.equal(outcomes[1].input, "next reply"); + await harness.emit("session_shutdown"); + assert.deepEqual(harness.clock.delays(), []); +} + +async function testFailedIdleRecoveryFailsClosed() { + let idleCalls = 0; + const harness = createHarness({ + onSubmit() { + throw new Error("OMP prompt flow rejected input"); + }, + onRun(call) { + if (call.args[0] === "api" && call.args[3] === "state" && call.args[4] === "idle") { + idleCalls += 1; + if (idleCalls === 2) throw new Error("state channel unavailable"); + } + }, + }); + await harness.emit("session_start"); + harness.children[0].stdout.emit("data", 'one'); + await flush(); + + assert.equal(commandCalls(harness, "send").length, 1); + assert.match(commandCalls(harness, "send")[0].input, /could not submit your message to OMP/); + assert.equal(harness.shutdowns, 1); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.ok( + harness.errors.some(({ message }) => + message.includes("could not restore idle state after a failed submission"), + ), + ); + assert.deepEqual(harness.clock.delays(), []); +} + +// [unit->REQ-OMP-LISTENER-FAIL-CLOSED] +async function testListenerRestartExhaustion() { + const harness = createHarness({ restartDelaysMs: [5, 10] }); + await harness.emit("session_start"); + const first = harness.children[0]; + first.emit("close", 7); + assert.deepEqual(harness.clock.delays(), [5]); + + await harness.clock.runNext(5); + const second = harness.children[1]; + second.stdout.emit("data", 'half'); + second.emit("error", new Error("listener crashed")); + second.emit("close", 8); + await flush(); + assert.deepEqual(harness.clock.delays(), [10], "error plus close schedules one restart"); + + await harness.clock.runNext(10); + const third = harness.children[2]; + third.emit("close", 9); + await flush(); + + assert.equal(harness.children.length, 3); + assert.equal(harness.shutdowns, 1); + assert.ok( + harness.errors.some(({ message }) => message.includes("listener restart budget exhausted")), + ); + assert.ok( + harness.notifications.some(({ message, type }) => + type === "error" && message.includes("listener restart budget exhausted"), + ), + ); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.deepEqual(harness.clock.delays(), []); + + await harness.emit("session_shutdown"); + assert.equal( + harness.calls.filter((call) => call.args[3] === "session-end").length, + 1, + "fatal shutdown and lifecycle shutdown share one teardown", + ); +} + +async function testListenerStableIntervalResetsRetries() { + const harness = createHarness({ + listenerStableMs: 20, + restartDelaysMs: [5, 10], + }); + await harness.emit("session_start"); + await harness.emit("agent_start"); + assert.deepEqual(harness.clock.delays(), [20]); + + harness.children[0].emit("close", 1); + assert.deepEqual(harness.clock.delays(), [5]); + await harness.clock.runNext(5); + + const shortLived = harness.children[1]; + shortLived.stdout.emit( + "data", + 'a parsed event is not stability', + ); + shortLived.emit("close", 2); + assert.deepEqual( + harness.clock.delays(), + [10], + "a parsed event followed by an immediate crash remains in the consecutive crash loop", + ); + await harness.clock.runNext(10); + + const stable = harness.children[2]; + assert.deepEqual(harness.clock.delays(), [20]); + await harness.clock.runNext(20); + stable.emit("close", 3); + assert.deepEqual( + harness.clock.delays(), + [5], + "a listener surviving the stable interval resets the next retry to attempt one", + ); + assert.match( + harness.notifications.filter(({ type }) => type === "warning").at(-1).message, + /restarting 1\/2 in 5ms/, + ); + + await harness.clock.runNext(5); + assert.equal(harness.children.length, 4); + await harness.emit("session_shutdown"); + assert.deepEqual(harness.clock.delays(), []); +} + +// [unit->REQ-OMP-EXTENSION-CUSTODY] +// [unit->REQ-OMP-LISTENER-FAIL-CLOSED] +async function testFatalTeardownAwaitsInFlightOutcome() { + let releaseOutcome; + const outcomeGate = new Promise((resolve) => { + releaseOutcome = resolve; + }); + const harness = createHarness({ + restartDelaysMs: [], + onRun(call) { + if (call.args[0] === "send" && call.args[1] === "slow") return outcomeGate; + }, + }); + await harness.emit("session_start"); + harness.children[0].stdout.emit("data", 'work'); + await flush(); + await harness.emit("agent_start"); + const ending = harness.emit("agent_end", { + messages: [ + { role: "user", content: '' }, + { role: "assistant", content: "finished" }, + ], + }); + await flush(); + assert.equal(commandCalls(harness, "send").length, 1); + + harness.children[0].emit("close", 11); + await flush(); + assert.equal(harness.shutdowns, 0, "fatal teardown must join the sender outcome"); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 0); + + const shutdown = harness.emit("session_shutdown"); + await flush(); + assert.equal( + harness.calls.filter((call) => call.args[3] === "session-end").length, + 0, + "concurrent lifecycle shutdown must join fatal custody teardown", + ); + for (const reason of ["new", "resume", "fork", "handoff"]) { + assert.deepEqual( + await harness.emit("session_before_switch", { reason }), + { cancel: true }, + `teardown must keep blocking the ${reason} switch while custody is pending`, + ); + } + assert.deepEqual( + await harness.emit("session_before_branch"), + { cancel: true }, + "teardown must keep blocking branches while custody is pending", + ); + + releaseOutcome(); + await Promise.all([ending, shutdown]); + await flush(); + assert.equal(harness.shutdowns, 1); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.deepEqual(harness.clock.delays(), []); +} + +async function testSessionEndRetriesAfterTransientFailure() { + const firstEnd = deferred(); + let endAttempts = 0; + const harness = createHarness({ + restartDelaysMs: [], + sessionEndRetryDelaysMs: [5], + onRun(call) { + if (call.args[0] === "api" && call.args[3] === "session-end") { + endAttempts += 1; + if (endAttempts === 1) return firstEnd.promise; + } + }, + }); + await harness.emit("session_start"); + harness.children[0].emit("close", 12); + await flush(); + assert.equal( + harness.calls.filter((call) => call.args[3] === "session-end").length, + 1, + ); + + firstEnd.reject(new Error("transient teardown failure")); + await flush(); + assert.equal(harness.shutdowns, 0, "fatal close must wait for the bounded teardown retry"); + assert.deepEqual(harness.clock.delays(), [5]); + assert.ok( + harness.errors.some( + ({ message, details }) => + message.includes("session teardown failed; retrying") && + details.error.includes("transient teardown failure"), + ), + ); + + await harness.clock.runNext(5); + await flush(); + const shutdown = harness.emit("session_shutdown"); + await shutdown; + await flush(); + assert.equal( + harness.calls.filter((call) => call.args[3] === "session-end").length, + 2, + "normal fatal teardown may retry before bounded lifecycle shutdown begins", + ); + assert.equal(harness.shutdowns, 1); + assert.deepEqual(harness.clock.delays(), []); +} + +async function testHumanBusyFailureFailsClosed() { + const harness = createHarness({ + onRun(call) { + if (call.args[0] === "api" && call.args[3] === "state" && call.args[4] === "busy") { + throw new Error("state channel unavailable"); + } + }, + }); + await harness.emit("session_start"); + await harness.emit("agent_start"); + + assert.equal(harness.shutdowns, 1); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.ok( + harness.errors.some(({ message }) => message.includes("could not mark the endpoint busy")), + ); + assert.equal(harness.children[0].kills, 1); + assert.deepEqual(harness.clock.delays(), []); +} + +// [unit->REQ-OMP-EXTENSION-CUSTODY] +// [unit->REQ-OMP-LISTENER-FAIL-CLOSED] +async function testShutdownReapsAndFailsQueuedCustody() { + const harness = createHarness(); + await harness.emit("session_start"); + const listener = harness.children[0]; + listener.stdout.emit( + "data", + 'onetwo', + ); + await flush(); + await harness.emit("agent_start"); + await harness.emit("agent_end", { + messages: [ + { role: "user", content: '' }, + { role: "assistant", content: "done" }, + ], + }); + assert.deepEqual(harness.clock.delays(), [0]); + + await harness.emit("session_shutdown"); + assert.equal(listener.kills, 1); + assert.deepEqual(harness.clock.delays(), []); + assert.deepEqual( + commandCalls(harness, "send").map((call) => call.args[1]), + ["first", "queued"], + ); + assert.match(commandCalls(harness, "send")[1].input, /OMP session shut down/); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.deepEqual(harness.statuses.at(-1), { key: "omp-spt", text: undefined }); + assert.equal(harness.shutdowns, 0, "normal lifecycle shutdown must not recursively shut down OMP"); +} +async function testRunSptRejectsStdinErrorsAndHungCommands() { + const epipeClock = new FakeClock(); + const epipeChild = new FakeChild(); + const epipe = Object.assign(new Error("write EPIPE"), { code: "EPIPE" }); + epipeChild.stdin.onEnd = () => { + epipeChild.stdin.emit("error", epipe); + return false; + }; + await assert.rejects( + runSpt(["send", "peer", "--from", "omp-agent"], "reply", { + clearTimeout: epipeClock.clearTimeout.bind(epipeClock), + commandTimeoutMs: 20, + killForceMs: 4, + killGraceMs: 3, + setTimeout: epipeClock.setTimeout.bind(epipeClock), + spawnProcess: () => epipeChild, + }), + (error) => error === epipe && error.code === "EPIPE", + ); + assert.deepEqual(epipeChild.killSignals, ["SIGTERM"]); + assert.deepEqual(epipeClock.delays(), []); + + const fastExitClock = new FakeClock(); + const fastExitChild = new FakeChild(); + fastExitChild.stdin.onEnd = () => { + fastExitChild.close(0); + return false; + }; + await assert.rejects( + runSpt(["send", "peer", "--from", "omp-agent"], "reply", { + clearTimeout: fastExitClock.clearTimeout.bind(fastExitClock), + commandTimeoutMs: 20, + killForceMs: 4, + killGraceMs: 3, + setTimeout: fastExitClock.setTimeout.bind(fastExitClock), + spawnProcess: () => fastExitChild, + }), + /exited before stdin completed/, + ); + assert.deepEqual(fastExitChild.killSignals, []); + assert.deepEqual(fastExitClock.delays(), []); + + const commandCases = [ + ["api", "--adapter", "omp-spt", "bind", "omp-agent"], + ["send", "peer", "--from", "omp-agent"], + ["api", "--adapter", "omp-spt", "state", "idle", "omp-agent"], + ["api", "--adapter", "omp-spt", "session-end", "omp-agent"], + ]; + for (const args of commandCases) { + const clock = new FakeClock(); + const child = new FakeChild(); + const pending = runSpt(args, args[0] === "send" ? "outcome" : undefined, { + clearTimeout: clock.clearTimeout.bind(clock), + commandTimeoutMs: 7, + killForceMs: 4, + killGraceMs: 3, + setTimeout: clock.setTimeout.bind(clock), + spawnProcess: () => child, + }); + const rejected = assert.rejects(pending, /timed out after 7ms/); + assert.deepEqual(clock.delays(), [7]); + await clock.runNext(7); + await rejected; + assert.deepEqual(child.killSignals, ["SIGTERM"]); + assert.deepEqual(clock.delays(), []); + } +} + +// [unit->REQ-OMP-LISTENER-FAIL-CLOSED] +async function testListenerTerminationEscalatesAndReaps() { + const harness = createHarness({ + killForceMs: 4, + killGraceMs: 3, + onSpawn(child) { + child.onKill = (signal) => { + if (signal === "SIGKILL") child.close(null, signal); + return true; + }; + }, + }); + await harness.emit("session_start"); + const listener = harness.children[0]; + const shutdown = harness.emit("session_shutdown"); + await flush(); + assert.deepEqual(listener.killSignals, ["SIGTERM"]); + assert.deepEqual(harness.clock.delays().sort((a, b) => a - b), [3, 1_800]); + + await harness.clock.runNext(3); + await shutdown; + assert.deepEqual(listener.killSignals, ["SIGTERM", "SIGKILL"]); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.deepEqual(harness.clock.delays(), []); + + const errorHarness = createHarness({ + killForceMs: 4, + killGraceMs: 3, + onSpawn(child) { + child.onKill = (signal) => { + if (signal === "SIGKILL") child.close(null, signal); + return true; + }; + }, + }); + await errorHarness.emit("session_start"); + const erroredListener = errorHarness.children[0]; + erroredListener.emit("error", new Error("listener pipe failed")); + const concurrentShutdown = errorHarness.emit("session_shutdown"); + await flush(); + assert.deepEqual(erroredListener.killSignals, ["SIGTERM"]); + assert.deepEqual(errorHarness.clock.delays().sort((a, b) => a - b), [3, 1_800]); + await errorHarness.clock.runNext(3); + await concurrentShutdown; + assert.deepEqual(erroredListener.killSignals, ["SIGTERM", "SIGKILL"]); + assert.equal( + errorHarness.calls.filter((call) => call.args[3] === "session-end").length, + 1, + ); + assert.deepEqual(errorHarness.clock.delays(), []); +} + +// [unit->REQ-OMP-EXTENSION-CUSTODY] +// [unit->REQ-OMP-LISTENER-FAIL-CLOSED] +async function testProtocolCorruptionFailsClosed() { + async function failProtocol(payload, expected, options = {}) { + const harness = createHarness({ restartDelaysMs: [], ...options }); + await harness.emit("session_start"); + harness.children[0].stdout.emit("data", payload); + await flush(); + assert.equal(harness.shutdowns, 1); + assert.deepEqual(harness.submitted, []); + assert.deepEqual(commandCalls(harness, "send"), []); + assert.ok( + harness.errors.some( + ({ message, details }) => + message.includes("listener protocol corruption") && expected.test(details.error), + ), + ); + assert.equal(harness.children[0].kills, 1); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.deepEqual(harness.clock.delays(), []); + await harness.emit("session_shutdown"); + return harness; + } + + const truncatedA = + 'truncatedvalid'; + const nested = await failProtocol(truncatedA, /nested EVENT/); + assert.deepEqual( + commandCalls(nested, "send").map((call) => call.args[1]), + [], + "the later valid b frame must not be merged into or consumed as a", + ); + await failProtocol('missing sender', /missing EVENT from/); + await failProtocol('bad attrs', /malformed EVENT attributes/); + await failProtocol('never closes'.padEnd(80, "x"), /buffer exceeded/, { + listenerBufferLimit: 64, + }); +} + +// [unit->REQ-OMP-EXTENSION-CUSTODY] +// [unit->REQ-OMP-LISTENER-FAIL-CLOSED] +async function testInboundQueueOverflowReturnsAcceptedCustody() { + const frames = ["a", "b", "overflow"].map( + (from) => `work`, + ); + const harness = createHarness({ + acceptedQueueLimit: 2, + restartDelaysMs: [], + }); + await harness.emit("session_start"); + harness.children[0].stdout.emit("data", frames.join("")); + await flush(); + + assert.equal(harness.shutdowns, 1); + assert.deepEqual(harness.submitted, []); + assert.deepEqual( + commandCalls(harness, "send").map((call) => call.args[1]), + ["a", "b", "overflow"], + "every accepted item and the capacity-refused item receive an explicit terminal failure", + ); + for (const call of commandCalls(harness, "send")) { + assert.match(call.input, /endpoint stopped before your message could complete/); + } + assert.ok( + harness.errors.some(({ message }) => + message.includes("inbound custody capacity exceeded"), + ), + ); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.deepEqual(harness.clock.delays(), []); + await harness.emit("session_shutdown"); + + const byteFirst = 'x'; + const byteOverflow = '😀'; + const byteHarness = createHarness({ + acceptedBytesLimit: Buffer.byteLength(byteFirst, "utf8") + byteOverflow.length, + acceptedQueueLimit: 10, + restartDelaysMs: [], + }); + await byteHarness.emit("session_start"); + byteHarness.children[0].stdout.emit("data", `${byteFirst}${byteOverflow}`); + await flush(); + assert.deepEqual( + commandCalls(byteHarness, "send").map((call) => call.args[1]), + ["a", "b"], + ); + assert.equal(byteHarness.shutdowns, 1); + assert.deepEqual(byteHarness.clock.delays(), []); + await byteHarness.emit("session_shutdown"); +} + +async function testSessionStopAwaitsOutcomeWithoutEndingEndpoint() { + const firstOutcome = deferred(); + let attempts = 0; + const harness = createHarness({ + outcomeRetryDelaysMs: [5], + onRun(call) { + if (call.args[0] === "send" && call.args[1] === "awaited") { + attempts += 1; + if (attempts === 1) return firstOutcome.promise; + } + }, + }); + await harness.emit("session_start"); + harness.children[0].stdout.emit( + "data", + 'work', + ); + await flush(); + await harness.emit("agent_start"); + const messages = [ + { role: "user", content: '' }, + { role: "assistant", content: "completed outcome" }, + ]; + const agentEnd = harness.emit("agent_end", { messages }); + await flush(); + const sessionStop = harness.emit("session_stop", { messages }); + await flush(); + assert.equal(commandCalls(harness, "send").length, 1); + assert.equal(harness.children[0].kills, 0); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 0); + + firstOutcome.reject(new Error("transient delayed outcome failure")); + await flush(); + assert.deepEqual(harness.clock.delays(), [5]); + await harness.clock.runNext(5); + await Promise.all([agentEnd, sessionStop]); + assert.equal(commandCalls(harness, "send").length, 2); + assert.equal(commandCalls(harness, "send").at(-1).input, "completed outcome"); + assert.equal(harness.children[0].kills, 0, "ordinary session_stop must leave the endpoint live"); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 0); + + await harness.emit("session_shutdown"); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.equal(harness.children[0].kills, 1); + assert.deepEqual(harness.clock.delays(), []); +} + +// [unit->REQ-OMP-EXTENSION-CUSTODY] +// [unit->REQ-OMP-LISTENER-FAIL-CLOSED] +async function testShutdownFallbackStaysBelowHostCap() { + const never = new Promise(() => {}); + const harness = createHarness({ + shutdownBudgetMs: 1_800, + shutdownCommandTimeoutMs: 300, + onRun(call) { + if (call.args[0] === "send" || call.args[3] === "session-end") return never; + }, + }); + await harness.emit("session_start"); + await harness.emit("agent_start"); + harness.children[0].stdout.emit( + "data", + 'work', + ); + await flush(); + + const shutdown = harness.emit("session_shutdown"); + await flush(); + assert.deepEqual( + harness.clock.delays().sort((a, b) => a - b), + [300, 1_800], + "queued custody has a short command timeout inside the 2s host cap", + ); + await harness.clock.runNext(300); + assert.equal(commandCalls(harness, "send").length, 1); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.deepEqual(harness.clock.delays().sort((a, b) => a - b), [300, 1_800]); + await harness.clock.runNext(300); + await shutdown; + + assert.equal(harness.children[0].kills, 1); + assert.ok( + harness.errors.some(({ message }) => + message.includes("could not return custody to queued"), + ), + ); + assert.equal(harness.calls.filter((call) => call.args[3] === "session-end").length, 1); + assert.deepEqual(harness.clock.delays(), []); + assert.ok(1_800 < 2_000); + + const inFlightOutcome = deferred(); + let inFlightAttempts = 0; + const inFlightHarness = createHarness({ + onRun(call) { + if (call.args[0] === "send" && call.args[1] === "in-flight") { + inFlightAttempts += 1; + if (inFlightAttempts === 1) return inFlightOutcome.promise; + } + }, + }); + await inFlightHarness.emit("session_start"); + inFlightHarness.children[0].stdout.emit( + "data", + 'work', + ); + await flush(); + await inFlightHarness.emit("agent_start"); + const ending = inFlightHarness.emit("agent_end", { + messages: [ + { role: "user", content: '' }, + { role: "assistant", content: "answer racing shutdown" }, + ], + }); + await flush(); + assert.equal(commandCalls(inFlightHarness, "send").length, 1); + const inFlightShutdown = inFlightHarness.emit("session_shutdown"); + await Promise.all([ending, inFlightShutdown]); + assert.deepEqual( + commandCalls(inFlightHarness, "send").map((call) => call.args[1]), + ["in-flight", "in-flight"], + ); + assert.match(commandCalls(inFlightHarness, "send")[1].input, /OMP session shut down/); + assert.equal( + inFlightHarness.calls.filter((call) => call.args[3] === "session-end").length, + 1, + ); + assert.equal(inFlightHarness.shutdowns, 0); + assert.deepEqual(inFlightHarness.clock.delays(), []); + inFlightOutcome.resolve(); + await flush(); + assert.equal(commandCalls(inFlightHarness, "send").length, 2); + + const hardCapHarness = createHarness({ + shutdownBudgetMs: 1_800, + shutdownCommandTimeoutMs: 5_000, + onRun(call) { + if (call.args[0] === "send") return never; + }, + }); + await hardCapHarness.emit("session_start"); + await hardCapHarness.emit("agent_start"); + hardCapHarness.children[0].stdout.emit( + "data", + 'work', + ); + await flush(); + const hardCappedShutdown = hardCapHarness.emit("session_shutdown"); + await flush(); + assert.deepEqual( + hardCapHarness.clock.delays().sort((a, b) => a - b), + [1_800, 5_000], + ); + await hardCapHarness.clock.runNext(1_800); + await hardCappedShutdown; + assert.ok( + hardCapHarness.errors.some(({ message }) => + message.includes("bounded shutdown expired"), + ), + ); + assert.equal(hardCapHarness.children[0].kills, 1); + assert.deepEqual(hardCapHarness.clock.delays(), []); +} + +await testParsingAndReplies(); +await testRunSptRejectsStdinErrorsAndHungCommands(); +await testLifecycleCustodyAndContext(); +await testDeferredBindLifecycleSerialization(); +await testOutcomeSendRetriesAndExhaustion(); +await testSubmissionFailureAdvancesQueue(); +await testFailedIdleRecoveryFailsClosed(); +await testListenerRestartExhaustion(); +await testListenerStableIntervalResetsRetries(); +await testFatalTeardownAwaitsInFlightOutcome(); +await testSessionEndRetriesAfterTransientFailure(); +await testHumanBusyFailureFailsClosed(); +await testShutdownReapsAndFailsQueuedCustody(); +await testListenerTerminationEscalatesAndReaps(); +await testProtocolCorruptionFailsClosed(); +await testInboundQueueOverflowReturnsAcceptedCustody(); +await testSessionStopAwaitsOutcomeWithoutEndingEndpoint(); +await testShutdownFallbackStaysBelowHostCap(); console.log("OMP-EXTENSION OK"); [raw output: artifact://671]