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" & <'), '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");