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" & <'), '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" & <'), 'a\n"c" & <');
-
-const partial = drainEvents('noisehello
wo');
-assert.deepEqual(partial.events, []);
-assert.equal(partial.rest, 'hello
wo');
-
-const complete = drainEvents(`${partial.rest}rldskip`);
-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");