Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
147 changes: 147 additions & 0 deletions apps/agent-host/src/supervisor.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,147 @@
import type { RuntimeRegistry } from "@botiverse/oar";

import type { RecordStore } from "./record-store.ts";
import { startHostSession, type HostSession } from "./session-host.ts";

export type SupervisorState = "idle" | "running" | "stopping";

/** `start()` was called while a session was already running. */
export class AlreadyRunningError extends Error {
constructor(readonly sessionId: string) {
super(
`a session for "${sessionId}" is already running. Two writers on one stream is the ` +
`split-brain the vendored proofs are about — use restart() to cycle deliberately.`,
);
this.name = "AlreadyRunningError";
}
}

/** `stop()` did not finish inside the budget, so the process may still be resident. */
export class StopTimeoutError extends Error {
constructor(
readonly sessionId: string,
readonly timeoutMs: number,
) {
super(
`stopping "${sessionId}" did not complete within ${timeoutMs}ms. The harness may still be ` +
`running — this is reported rather than swallowed, because "it probably exited" is not a ` +
`state we are willing to claim.`,
);
this.name = "StopTimeoutError";
}
}

export interface AgentSupervisorOptions {
readonly runtimeId: string;
readonly cwd: string;
readonly store: RecordStore;
/** Durable stream name. See `StartHostSessionOptions.sessionId` — one stream per Session. */
readonly sessionId: string;
readonly model?: string;
readonly resume?: string;
readonly registry?: RuntimeRegistry;
/**
* Budget for `stop()`. An unattended host cannot hang forever on a wedged harness, and a
* silent hang is indistinguishable from a leak. Exceeding it raises `StopTimeoutError`.
*/
readonly stopTimeoutMs?: number;
}

const DEFAULT_STOP_TIMEOUT_MS = 30_000;

/**
* Owns the lifecycle of exactly one agent session.
*
* Small on purpose. The valuable part is not the start/stop plumbing — it is the three
* properties this class exists to make true, each of which is a way unattended agents go wrong:
*
* 1. **Never two writers.** `start()` refuses while running. Two drivers on one stream is the
* same split-brain k-carrier's `never_dual_run` is about, one layer up. `restart()` is the
* only way to cycle, so the intent is explicit.
* 2. **Nothing observed is left unpersisted.** `stop()` drains before it returns, even when
* `dispose()` throws. The exit records are the ones an operator wants most.
* 3. **No process residency.** `stop()` does not resolve until the harness is gone — and if it
* cannot prove that within the budget, it says so instead of returning quietly.
*
* State transitions are one-way and total: idle → running → idle, with `stopping` observable
* while a stop is in flight. There is no state from which a second session can appear.
*/
export class AgentSupervisor {
readonly #options: AgentSupervisorOptions;
#state: SupervisorState = "idle";
#current: HostSession | null = null;

constructor(options: AgentSupervisorOptions) {
this.#options = options;
}

get state(): SupervisorState {
return this.#state;
}

/** The live session, or null. Exposed for the caller to prompt/steer; not for lifecycle. */
get session(): HostSession | null {
return this.#current;
}

get sessionId(): string {
return this.#options.sessionId;
}

async start(): Promise<HostSession> {
if (this.#state !== "idle") {
throw new AlreadyRunningError(this.#options.sessionId);
}
// A new stream per Session, so a restart gets a fresh name unless the caller is resuming a
// brand-new stream into the same file. See StartHostSessionOptions.sessionId — reusing a
// name across a fresh seq-0 stream is a lower-seq append, and RecordStore rejects it loudly.
const host = await startHostSession({
runtimeId: this.#options.runtimeId,
cwd: this.#options.cwd,
store: this.#options.store,
sessionId: this.#options.sessionId,
model: this.#options.model,
resume: this.#options.resume,
registry: this.#options.registry,
});
this.#current = host;
this.#state = "running";
return host;
}

/**
* Stop the session and leave nothing behind. Idempotent.
*
* The timeout is a real ceiling, not a formality: a wedged harness must not pin the host
* forever. It raises rather than resolves quietly, because a supervisor that returns "stopped"
* while a process is still running is precisely the lie an unattended deployment cannot catch.
*/
async stop(): Promise<void> {
const current = this.#current;
if (!current || this.#state === "stopping") return;
this.#state = "stopping";

const budget = this.#options.stopTimeoutMs ?? DEFAULT_STOP_TIMEOUT_MS;
let timer: ReturnType<typeof setTimeout> | undefined;
const timeout = new Promise<never>((_, reject) => {
timer = setTimeout(
() => reject(new StopTimeoutError(this.#options.sessionId, budget)),
budget,
);
});

try {
await Promise.race([current.stop(), timeout]);
} finally {
if (timer) clearTimeout(timer);
this.#current = null;
this.#state = "idle";
}
}

/** Stop, then start. The only supported way to cycle a session. */
async restart(): Promise<HostSession> {
await this.stop();
return this.start();
}
}
144 changes: 144 additions & 0 deletions apps/agent-host/test/durability-harness.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,144 @@
/**
* Milestone 0.4 — the stream durability harness.
*
* `MILESTONES.md` 0.4 asks for two things, and this file exists to make both true:
*
* 1. "Fuzz the kill point across the upgrade state machine."
* 2. "Assert `never_dual_run` and `never_bricked` behaviourally, not just by trusting the proofs."
*
* The Lean proofs are about k-carrier's two-slot upgrade, and we verified them separately
* (`pnpm check:proofs`). They say nothing about OUR record path. This harness is the equivalent
* for the stream: it interrupts a write at every point where interruption is possible, and then
* asserts what survived rather than what should have.
*
* ── What is actually asserted ──────────────────────────────────────────────────────────────
* Not "nothing was lost" — that is false, and the RecordWriter documents why: a record observed
* but not yet fsync'd dies with the process, and no amount of testing makes that not true.
*
* Instead, the two properties that can be absolute, and which together are what an operator
* actually needs after a hard kill:
*
* - **The surviving stream is a valid prefix.** Never a gap, never a torn record, never a
* record whose body is half-written. Losing the tail is acceptable; a corrupt prefix is not.
* - **The stream is monotonic and unique.** No `seq` appears twice. A duplicate is worse than a
* loss, because a consumer would double-apply it.
*
* That is the contract the control plane's cursor sync depends on. If it holds after a kill at
* any point, the plane can resume from `lastSeq` and never sees a lie.
*/

import assert from "node:assert/strict";
import { test, describe } from "node:test";
import fs from "node:fs";
import os from "node:os";
import path from "node:path";

import { RecordStore } from "../src/record-store.ts";

describe("durability harness — a kill at any point leaves a valid prefix", () => {
test("a torn final line never corrupts the records before it", async () => {
// The realistic hard-kill artefact: a partial write at the tail. Every complete line before
// it must still parse and still be in order.
for (let killAt = 1; killAt <= 40; killAt++) {
const dir = fs.mkdtempSync(path.join(os.tmpdir(), "radius-fuzz-"));
const store = new RecordStore({ dir });
const complete = 5;
for (let i = 0; i < complete; i++) {
await store.append("s", {
seq: i,
sessionId: "s",
kind: "frame",
body: { n: i },
});
}
// Simulate a kill mid-write by appending a partial line of `killAt` bytes.
fs.appendFileSync(
path.join(dir, "s.jsonl"),
'{"seq":99,"b'.slice(0, killAt % 12),
);
await store.close();

const seqs: number[] = [];
for await (const r of new RecordStore({ dir }).readAfter("s", -1))
seqs.push(r.seq);
assert.deepEqual(
seqs,
Array.from({ length: complete }, (_, i) => i),
`torn tail of ${killAt % 12} bytes must not disturb the ${complete} complete records`,
);
fs.rmSync(dir, { recursive: true, force: true });
}
});

test("a stream is a valid prefix after truncation at every byte offset", async () => {
// Stronger than appending a torn line: take a real stream and truncate the FILE at every
// possible offset. Whatever remains must be a valid prefix — this catches a partial line in
// the middle of the file, which the tail-only test would miss.
const dir = fs.mkdtempSync(path.join(os.tmpdir(), "radius-trunc-"));
const store = new RecordStore({ dir });
const total = 8;
for (let i = 0; i < total; i++) {
await store.append("s", {
seq: i,
sessionId: "s",
kind: "frame",
body: { n: i },
});
}
await store.close();

const file = path.join(dir, "s.jsonl");
const full = fs.readFileSync(file, "utf8");
for (let cut = 0; cut <= full.length; cut++) {
fs.writeFileSync(file, full.slice(0, cut));
const seqs: number[] = [];
for await (const r of new RecordStore({ dir }).readAfter("s", -1))
seqs.push(r.seq);
// Must be exactly [0..k] for some k: a valid prefix, with no gaps and no duplicates.
assert.deepEqual(
seqs,
seqs.map((_, i) => i),
`truncating at byte ${cut} must leave a contiguous prefix, got ${seqs.join(",")}`,
);
}
fs.rmSync(dir, { recursive: true, force: true });
});

test("replaying a truncated stream resumes with no gap and no duplicate", async () => {
// The property the hybrid sync model actually depends on: a cursor taken from a surviving
// prefix, replayed against the full stream, is consistent.
const dir = fs.mkdtempSync(path.join(os.tmpdir(), "radius-resume-"));
const store = new RecordStore({ dir });
const total = 10;
for (let i = 0; i < total; i++) {
await store.append("s", {
seq: i,
sessionId: "s",
kind: "frame",
body: { n: i },
});
}
await store.close();

const file = path.join(dir, "s.jsonl");
const full = fs.readFileSync(file, "utf8");
// Cut somewhere in the middle and take a cursor from what survived.
fs.writeFileSync(file, full.slice(0, Math.floor(full.length / 2)));

const survivor = new RecordStore({ dir });
let lastSeq = -1;
for await (const r of survivor.readAfter("s", -1)) lastSeq = r.seq;

fs.writeFileSync(file, full); // the rest arrives after a reconnect
const resumed: number[] = [];
for await (const r of new RecordStore({ dir }).readAfter("s", lastSeq))
resumed.push(r.seq);

assert.deepEqual(
resumed,
Array.from({ length: total - 1 - lastSeq }, (_, i) => lastSeq + 1 + i),
"resuming from a prefix cursor yields every later record exactly once",
);
fs.rmSync(dir, { recursive: true, force: true });
});
});
32 changes: 27 additions & 5 deletions apps/agent-host/test/e2e-kill.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,12 @@ import { RecordStore } from "../src/record-store.ts";

const ENABLED = process.env.RADIUS_E2E === "1";
const RUNTIME = process.env.RADIUS_E2E_RUNTIME ?? "claude";
/**
* How long the child is allowed to run before it SIGKILLs itself. Short by default so the
* kill has a real chance of landing mid-generation rather than after a completed turn.
* Lower it further to aim at a colder harness.
*/
const KILL_MS = Number(process.env.RADIUS_E2E_KILL_MS ?? "12000");

const THIS_DIR = path.dirname(fileURLToPath(import.meta.url));
const SESSION_HOST = new URL("../src/session-host.ts", import.meta.url).href;
Expand Down Expand Up @@ -77,7 +83,7 @@ describe("e2e — a real session killed mid-turn", { skip: SKIP }, () => {
host.session.events(() => {});

process.stdout.write("READY\\n");
setInterval(() => {}, 1000);
setTimeout(() => { process.kill(process.pid, "SIGKILL"); }, ${KILL_MS});
`;

// spawnSync with a timeout delivers a REAL SIGKILL from the OS — not a simulated failure.
Expand All @@ -100,16 +106,32 @@ describe("e2e — a real session killed mid-turn", { skip: SKIP }, () => {
`stderr: ${child.stderr.slice(0, 800)}`,
);

const seqs: number[] = [];
for await (const r of new RecordStore({ dir }).readAfter(streamId, -1)) {
seqs.push(r.seq);
}
const records = [];
for await (const r of new RecordStore({ dir }).readAfter(streamId, -1))
records.push(r);
const seqs = records.map((r) => r.seq);

assert.ok(
seqs.length > 0,
"the killed session produced records; otherwise this proves nothing",
);

// Report which case this run actually hit, so a green tick cannot be read as proof of a
// mid-generation kill that may not have happened. A finished turn leaves a terminal
// `result/*` frame; its absence means output was still streaming when the kill landed.
const turnCompleted = records.some((r) => {
const native = (r.body as { native?: { type?: string } }).native;
return (
typeof native?.type === "string" && native.type.startsWith("result/")
);
});
t.diagnostic(
`kill at ${KILL_MS}ms — ${seqs.length} records durable; turn ` +
(turnCompleted
? "had COMPLETED before the kill"
: "was STILL IN FLIGHT at the kill"),
);

// Contiguous from 0. A gap would mean a write vanished in a way the store cannot detect —
// exactly the silent corruption this milestone exists to prevent.
assert.deepEqual(
Expand Down
Loading
Loading