import { spawn } from "node:child_process";
const ADAPTER = "omp-spt";
export function decodeBody(body) {
return body
.replaceAll("
", "\n")
.replaceAll("<", "<")
.replaceAll(">", ">")
.replaceAll(""", '"')
.replaceAll("&", "&");
}
function protocolError(message) {
const error = new Error(`invalid spt EVENT stream: ${message}`);
error.code = "SPT_PROTOCOL_ERROR";
return error;
}
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) {
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);
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: parsed.attributes.from,
body: decodeBody(raw.slice(openEnd + 1, close)),
envelope: raw.slice(start, end),
});
}
cursor = end;
if (events.length >= maxEvents) return { events, rest: raw.slice(cursor) };
}
}
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 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) => {
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", 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);
}
}
});
}
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 requestedShutdownCommandTimeoutMs =
overrides.shutdownCommandTimeoutMs ?? DEFAULT_SHUTDOWN_COMMAND_TIMEOUT_MS;
const terminationWindowMs = killGraceMs + killForceMs;
const maxShutdownCommandTimeoutMs = Math.max(
1,
Math.floor((shutdownBudgetMs - 3 * terminationWindowMs) / 2),
);
const shutdownCommandTimeoutMs = Math.max(
1,
Math.min(requestedShutdownCommandTimeoutMs, maxShutdownCommandTimeoutMs),
);
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);
return function ompSpt(pi) {
const id = env.SPT_ENDPOINT_ID?.trim();
if (!id) return;
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 api 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;
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 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) {
desiredState = "idle";
try {
await setState("idle");
} catch (stateError) {
await failClosed(
"omp-spt could not restore idle state after a failed submission",
stateError,
);
}
}
}
// [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();
}
}
// [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?.();
}
function startListener() {
if (stopping) return;
const args = ["api", "--adapter", ADAPTER, "listen", id, "--session-id", sid];
if (env.OMP_SPT_SUBNET) args.push("--subnet", env.OMP_SPT_SUBNET);
let child;
try {
child = spawnProcess(env.OMP_SPT_SPT_BIN || "spt", args, {
stdio: ["ignore", "pipe", "pipe"],
windowsHide: true,
});
} catch (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 api 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 api listen exited ${status}`), true);
});
ui?.setStatus("omp-spt", `spt:${id}`);
}
// [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();