import assert from "node:assert/strict"; 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] 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, [ "api", "--adapter", "omp-spt", "listen", "omp-agent", "--session-id", "session-1", "--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(), []); const hungBindHarness = createHarness({ onRun(call) { if (call.args[0] === "api" && call.args[3] === "bind") return new Promise(() => {}); }, }); const hungStart = hungBindHarness.emit("session_start"); await flush(); const hungShutdown = hungBindHarness.emit("session_shutdown"); await flush(); assert.deepEqual(hungBindHarness.clock.delays().sort((a, b) => a - b), [300, 1_800]); await hungBindHarness.clock.runNext(300); await Promise.all([hungStart, hungShutdown]); assert.equal(hungBindHarness.children.length, 0); assert.equal( hungBindHarness.calls.filter((call) => call.args[3] === "session-end").length, 0, "session-end cannot run without a completed bind token", ); assert.deepEqual(hungBindHarness.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 queuedGates = [deferred(), deferred()]; let queuedSendIndex = 0; const concurrentHarness = createHarness({ onRun(call) { if (call.args[0] === "send") { const gate = queuedGates[queuedSendIndex]; queuedSendIndex += 1; return gate.promise; } }, }); await concurrentHarness.emit("session_start"); await concurrentHarness.emit("agent_start"); concurrentHarness.children[0].stdout.emit( "data", 'onetwo', ); await flush(); const concurrentShutdown = concurrentHarness.emit("session_shutdown"); await flush(); assert.deepEqual( commandCalls(concurrentHarness, "send").map((call) => call.args[1]), ["queued-a", "queued-b"], "all pending custody failures must start concurrently", ); assert.deepEqual( concurrentHarness.clock.delays().sort((a, b) => a - b), [300, 300, 1_800], ); for (const gate of queuedGates) gate.resolve(); await concurrentShutdown; assert.equal( concurrentHarness.calls.filter((call) => call.args[3] === "session-end").length, 1, ); assert.deepEqual(concurrentHarness.clock.delays(), []); const busyGate = deferred(); const dispatchHarness = createHarness({ onRun(call) { if (call.args[0] === "api" && call.args[3] === "state" && call.args[4] === "busy") { return busyGate.promise; } }, }); await dispatchHarness.emit("session_start"); dispatchHarness.children[0].stdout.emit( "data", 'work', ); await flush(); assert.deepEqual(dispatchHarness.submitted, []); assert.equal( stateCalls(dispatchHarness).filter((call) => call.args[4] === "busy").length, 1, ); const dispatchShutdown = dispatchHarness.emit("session_shutdown"); await dispatchShutdown; assert.deepEqual(dispatchHarness.submitted, []); assert.deepEqual( commandCalls(dispatchHarness, "send").map((call) => call.args[1]), ["dispatching"], ); assert.match(commandCalls(dispatchHarness, "send")[0].input, /OMP session shut down/); assert.equal( dispatchHarness.calls.filter((call) => call.args[3] === "session-end").length, 1, ); assert.equal(dispatchHarness.shutdowns, 0); assert.deepEqual(dispatchHarness.clock.delays(), []); busyGate.resolve(); await flush(); assert.deepEqual(dispatchHarness.submitted, []); 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, killForceMs: 100, killGraceMs: 100, 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), [600, 1_800], ); await hardCapHarness.clock.runNext(600); await hardCappedShutdown; assert.equal( hardCapHarness.calls.filter((call) => call.args[3] === "session-end").length, 1, "the phase clamp must reserve time for one session-end attempt", ); 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");