From 5789524b8de3579c8c7a94a72bf4659a232e4a1f Mon Sep 17 00:00:00 2001 From: Aaron Trowbridge Date: Fri, 25 Sep 2026 10:23:42 +0200 Subject: [PATCH] fix: re-dispatch interrupted sessions after a hub engine bounce (#1552) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A hub engine restart (hub-restart.sh's sanctioned systemctl bounce) SIGTERMs the engine and kills every in-flight agent loop — loops live in the engine process's memory — and nothing re-dispatched them once the engine returned (2026-09-24: three campaigns frozen at the bounce instant). The service already sees every in-flight turn (an open chat POST proxied through the engine seam), so the fix is entirely in the extension layer: - track (inflight_journal.ts): the service's dispatch path journals every proxied POST /session/{id}/message (the SDK's session.prompt — start on entry, end on response close, which fires on completion AND on death); best-effort appends only, never a request-path dependency - derive: a pure fold of the journal into the interrupted set (latest unmatched start, within a 12h recency window, capped at 8) - resume: after engine health, the runner fires ONE mechanical resume turn per interrupted session (fire-and-forget — the chat POST holds until the turn completes; boot never wedges), then clears the journal so a subsequent boot cannot double-resume; a failed resume is one log line, never a failed boot - CLI env surface: AMICODE_INFLIGHT_JOURNAL, AMICODE_RESUME_RECENT_MS, AMICODE_RESUME_DISABLED (the test/ops escape hatch) --- .../src/amicode_service/inflight_journal.ts | 189 ++++++ .../extension/src/amicode_service/server.ts | 14 + .../extension/src/amicode_service_runner.ts | 85 +++ .../src/amicode_service_runner_cli.ts | 23 + .../amicode_service_inflight_journal.test.ts | 556 ++++++++++++++++++ 5 files changed, 867 insertions(+) create mode 100644 packages/extension/src/amicode_service/inflight_journal.ts create mode 100644 packages/extension/test/amicode_service_inflight_journal.test.ts diff --git a/packages/extension/src/amicode_service/inflight_journal.ts b/packages/extension/src/amicode_service/inflight_journal.ts new file mode 100644 index 000000000..d07aa4523 --- /dev/null +++ b/packages/extension/src/amicode_service/inflight_journal.ts @@ -0,0 +1,189 @@ +// ============================================================================ +// inflight_journal — #1552 (interrupted-session re-dispatch): the append-only +// JSONL sidecar that makes agent loops durable across an engine bounce. +// +// WHY THIS EXISTS: a running agent loop is an OPEN CHAT POST — the SDK's +// session.prompt(), POST /session/{id}/message with a {parts:[...]} body, +// proxied through this service to the engine, and the response only closes +// when the turn completes (tens of minutes for a long campaign). The loops +// themselves live in the ENGINE process's memory, so the sanctioned deploy +// bounce (hub-restart.sh's systemctl restart → SIGTERM) kills every one of +// them at once, and an engine crash does the same. The service is the one +// layer that already SEES every in-flight turn — this module turns that +// visibility into state the next boot can act on: +// +// track — trackInflightChatTurn: a "start" line when a chat POST enters +// the dispatch seam, an "end" line when the response closes +// (fires on completion AND on client/upstream death — exactly +// the two outcomes that must not be confused with a kill). +// derive — deriveInterruptedSessions (PURE: journal text → interrupted +// set): a session is interrupted iff its latest start has no +// matching end after it AND is within the recency window, +// capped — a corrupted journal can never fan out unbounded work. +// resume — the runner's boot path (amicode_service_runner.ts) reads this +// after engine health, fires ONE mechanical resume turn per +// interrupted session (fire-and-forget), then clears the file. +// +// Because the journal is written THROUGH the request path, it inherits the +// service's never-crash-on-request discipline: every append is best-effort +// (try/catch, a sidecar may never become a dependency), and the derivation +// skips corrupt lines rather than failing the boot. +// +// CONFIG (the runner CLI's ENV SURFACE documents both): +// AMICODE_INFLIGHT_JOURNAL the journal path. Default +// ~/.amico/server/active-sessions.jsonl (the +// parent dirs are created on demand — a fresh +// hub has no ~/.amico/server yet). +// AMICODE_RESUME_RECENT_MS the recency window. Default 12h: an older +// start belongs to a loop the user abandoned — +// resuming yesterday's turn is noise, not +// continuity. +// ============================================================================ +import { appendFileSync, mkdirSync, readFileSync, writeFileSync } from "node:fs"; +import { homedir } from "node:os"; +import { dirname, join } from "node:path"; +import type * as http from "node:http"; + +/** One journal line (JSONL): a chat turn opened ("start") or closed ("end") + * for a session, at the wall-clock millisecond it happened. */ +export interface InflightJournalLine { + kind: "start" | "end"; + sessionID: string; + ts: number; +} + +/** An interrupted session: the sessionID to re-dispatch and the ts of its + * latest unmatched start (the freshest evidence of the killed turn). */ +export interface InterruptedSession { + sessionID: string; + ts: number; +} + +/** #1552's resume cap: at most this many sessions are re-dispatched per boot. + * The journal is a sidecar, not a database — a corrupted or hostile file + * must never fan out unbounded resume turns. */ +export const RESUME_CAP = 8; + +/** #1552's default recency window: 12h. Overnight campaigns are the exact + * case this fix exists for; a start older than this is a stale artifact of + * some much earlier crash, not a loop to resurrect. */ +export const RESUME_RECENT_MS_DEFAULT = 12 * 60 * 60 * 1000; + +/** #1552's mechanical resume nudge — the issue's wording, exactly. The loop + * re-reads its own session ledger/state and continues; the nudge carries + * no new instructions (the campaign's own grammar owns what happens next). */ +export const RESUME_MESSAGE = + "[auto-resume] The engine restarted while your loop was in flight (issue #1552). Re-read your session ledger/state and continue exactly where you left off."; + +/** The journal path (env AMICODE_INFLIGHT_JOURNAL → the ~/.amico default). + * Resolved per call — the late-bound seam convention, so a mid-process env + * change behaves like every other config carrier here. */ +export function inflightJournalPath(): string { + const env = (process.env.AMICODE_INFLIGHT_JOURNAL ?? "").trim(); + if (env !== "") return env; + return join(homedir(), ".amico", "server", "active-sessions.jsonl"); +} + +/** The one tracked route (#1552): the SDK's session.prompt() — + * POST /session/{id}/message, the request whose response resolves only + * when the agent turn completes (verified against the vendored SDK gen at + * the manifest's upstream base: v1 names the path param {id}, v2 names it + * {sessionID} — the URL is the same). Returns the sessionID on a match. + * GETs, SSE, and every other session route observe or manage — they never + * carry a turn — so they stay untracked. */ +export function matchInflightChatTurn(method: string | undefined, pathname: string): string | undefined { + if (method !== "POST") return undefined; + const m = /^\/session\/([^/]+)\/message$/.exec(pathname); + return m === null ? undefined : m[1]; +} + +/** Best-effort append — NEVER throws into the request path (the service's + * founding discipline: a sidecar is a witness, not a dependency). Parent + * dirs are created on demand (~/.amico/server does not pre-exist). */ +function appendJournalLine(path: string, line: InflightJournalLine): void { + try { + mkdirSync(dirname(path), { recursive: true }); + appendFileSync(path, JSON.stringify(line) + "\n"); + } catch { + /* best-effort by contract */ + } +} + +/** The tracking hook (#1552): journal a "start" for `sessionID` now, and a + * matching "end" when the response closes. res "close" is the one event + * that fires on BOTH completion and death (client abort, upstream crash, + * the service's own SIGTERM) — which is exactly the pair that must leave + * the journal in the truthful state: an ended turn is never interrupted, + * a killed turn always is. */ +export function trackInflightChatTurn( + res: http.ServerResponse, + sessionID: string, + journalPath?: string, +): void { + const journal = journalPath ?? inflightJournalPath(); + appendJournalLine(journal, { kind: "start", sessionID, ts: Date.now() }); + res.once("close", () => appendJournalLine(journal, { kind: "end", sessionID, ts: Date.now() })); +} + +/** The derivation (#1552's pure core): journal text → the interrupted set. + * Grouped by sessionID in line order, a session is INTERRUPTED iff its + * latest start has no matching end after it, AND that start is within the + * recency window of `now`; the result is sorted freshest-start-first and + * capped. Corrupt/foreign lines are skipped — the derivation may never + * be the reason a boot fails. */ +export function deriveInterruptedSessions( + journalText: string, + now: number, + opts?: { recentMs?: number; cap?: number }, +): InterruptedSession[] { + const recentMs = opts?.recentMs ?? RESUME_RECENT_MS_DEFAULT; + const cap = opts?.cap ?? RESUME_CAP; + // sessionID → the ts of its latest UNMATCHED start (the only state that + // matters: a matched end deletes the entry, a new start overwrites it). + const open = new Map(); + for (const raw of journalText.split("\n")) { + const text = raw.trim(); + if (text === "") continue; + let line: InflightJournalLine; + try { + line = JSON.parse(text) as InflightJournalLine; + } catch { + continue; // corrupt line: skip, never fatal + } + if (typeof line?.sessionID !== "string" || line.sessionID === "") continue; + if (line.kind === "start" && typeof line.ts === "number") open.set(line.sessionID, line.ts); + else if (line.kind === "end") open.delete(line.sessionID); + } + return [...open.entries()] + .filter(([, ts]) => now - ts <= recentMs) + .sort((a, b) => b[1] - a[1]) // freshest first: the cap keeps the most-recent victims + .slice(0, cap) + .map(([sessionID, ts]) => ({ sessionID, ts })); +} + +/** The boot-side reader: derive from the journal FILE. A missing or + * unreadable journal is a clean boot (nothing was in flight), never a + * failed one — the resume step is continuity sugar, not a dependency. */ +export function readInterruptedSessions( + journalPath: string, + now: number = Date.now(), + opts?: { recentMs?: number; cap?: number }, +): InterruptedSession[] { + try { + return deriveInterruptedSessions(readFileSync(journalPath, "utf8"), now, opts); + } catch { + return []; + } +} + +/** Clear the journal after a re-dispatch: the sessions have been handed + * their resume turns, so a subsequent boot must not double-resume them. + * Best-effort, same as every append. */ +export function clearInflightJournal(journalPath: string): void { + try { + mkdirSync(dirname(journalPath), { recursive: true }); + writeFileSync(journalPath, ""); + } catch { + /* best-effort by contract */ + } +} diff --git a/packages/extension/src/amicode_service/server.ts b/packages/extension/src/amicode_service/server.ts index 4e1724c71..f2285b9f0 100644 --- a/packages/extension/src/amicode_service/server.ts +++ b/packages/extension/src/amicode_service/server.ts @@ -20,6 +20,7 @@ import { mintServerPassword, serverAuthHeader } from "../server_auth"; import { setBindHostname } from "./bind_host"; import { AppShelf, type AppShelfResult } from "./app_shelf"; import { EngineProxy } from "./engine_proxy"; +import { matchInflightChatTurn, trackInflightChatTurn } from "./inflight_journal"; import { HubProxy } from "./hub_proxy"; import type { UpstreamMode } from "./merged_projection"; import { isPublicUiPath } from "./public_ui"; @@ -278,6 +279,19 @@ export class AmicodeServiceServer { return; } if (this.engineProxy) { + // #1552: the in-flight journal seam. dispatch is the one place that + // sees EVERY engine-bound request, so it is where an agent loop (an + // open chat POST — POST /session/{id}/message, the SDK's + // session.prompt()) becomes durable state: the match journals a + // "start" now, and the response's "close" (fires on completion AND on + // client/upstream death) writes the matching "end". The next boot + // reads this journal to re-dispatch the sessions a bounce killed + // mid-loop. Best-effort by contract — never a request-path + // dependency. (The hub runner always routes engine-mode, so this + // seam covers the hub's bounce case; fleet mode is a client of a + // remote hub, not the hub itself.) + const inflight = matchInflightChatTurn(req.method, url.pathname); + if (inflight !== undefined) trackInflightChatTurn(res, inflight); // Streams method/headers/body through to the engine (SSE included); // false = no upstream bound yet → the honest 503 below. if (this.engineProxy.handle(req, res)) return; diff --git a/packages/extension/src/amicode_service_runner.ts b/packages/extension/src/amicode_service_runner.ts index f0baa8a73..5ecbe3e79 100644 --- a/packages/extension/src/amicode_service_runner.ts +++ b/packages/extension/src/amicode_service_runner.ts @@ -41,6 +41,13 @@ import { tmpdir } from "node:os"; import { createServer } from "node:net"; import { join } from "node:path"; import { createAmicodeService } from "./amicode_service"; +import { + RESUME_MESSAGE, + RESUME_RECENT_MS_DEFAULT, + clearInflightJournal, + inflightJournalPath, + readInterruptedSessions, +} from "./amicode_service/inflight_journal"; import type { AmicodeServiceServer } from "./amicode_service/server"; import { mintServerPassword, serverAuthHeader } from "./server_auth"; @@ -100,6 +107,23 @@ export interface AmicodeServiceRunnerOptions { engineUnarmed?: boolean; /** Engine health-wait budget. Default 30_000 (the ServerManager budget). */ healthTimeoutMs?: number; + /** #1552 (interrupted-session re-dispatch): skip the boot-time resume + * entirely — AMICODE_RESUME_DISABLED=1, the test/ops escape hatch. The + * service keeps writing the journal either way; this only stops the + * boot from acting on it. */ + resumeDisabled?: boolean; + /** #1552: the in-flight journal the service's dispatch seam writes and + * this runner reads after engine health. Default: the env-resolved + * production path (AMICODE_INFLIGHT_JOURNAL → + * ~/.amico/server/active-sessions.jsonl — inflightJournalPath(), the + * single source the service itself uses, so the pair can never drift). */ + resumeJournalPath?: string; + /** #1552: the recency window for resume candidates — AMICODE_RESUME_RECENT_MS. + * Default 12h (RESUME_RECENT_MS_DEFAULT). */ + resumeRecentMs?: number; + /** #1552: the fetch used to fire the resume turns — injectable for tests + * (the never-resolving/rejecting legs). */ + resumeFetch?: typeof fetch; /** Log sink (the structural-interface convention — vscode-free). */ log?: (line: string) => void; } @@ -254,6 +278,67 @@ export async function bootAmicodeServiceRunner(opts: AmicodeServiceRunnerOptions } log(`[service-runner] engine up at ${engineUrl}`); + // ── #1552: interrupted-session re-dispatch ────────────────────────────── + // The engine just came back from the bounce/crash that killed every + // in-flight agent loop (loops live in the engine's memory). The journal + // the service wrote says which sessions had an OPEN chat turn at the + // kill instant; each gets ONE mechanical resume turn so the loop + // continues from its own ledger/state — the hub's missing law: it never + // comes back up without resuming the loops it killed. NEVER-WEDGE + // discipline: the chat POST holds until the turn completes (tens of + // minutes), so every fire is fire-and-forget (boot never waits on a + // turn), every outcome is ONE log line, and the whole step is wrapped — + // a resume problem can never fail or stall the boot. + if (opts.resumeDisabled === true) { + log("[service-runner] interrupted-session resume disabled (AMICODE_RESUME_DISABLED) — skipping"); + } else { + try { + const journalPath = opts.resumeJournalPath ?? inflightJournalPath(); + const interrupted = readInterruptedSessions(journalPath, Date.now(), { + recentMs: opts.resumeRecentMs ?? RESUME_RECENT_MS_DEFAULT, + }); + if (interrupted.length > 0) { + log( + `[service-runner] in-flight journal: ${interrupted.length} interrupted session(s) — re-dispatching resume turns (issue #1552)`, + ); + const resumeFetch = opts.resumeFetch ?? fetch; + for (const { sessionID } of interrupted) { + // Fire-and-forget: same chat route + body schema the SDK uses for + // a user prompt (POST /session/{id}/message, {parts:[{type:"text"}]}), + // carrying the engine credential when armed (unarmed hub posture = + // anonymous — no credential exists). + void resumeFetch(`${engineUrl}/session/${encodeURIComponent(sessionID)}/message`, { + method: "POST", + headers: { + "Content-Type": "application/json", + ...(authHeader === undefined ? {} : { Authorization: authHeader }), + }, + body: JSON.stringify({ parts: [{ type: "text", text: RESUME_MESSAGE }] }), + }) + .then((r) => + log(`[service-runner] resume turn for session ${sessionID}: engine answered ${r.status}`), + ) + .catch((err: unknown) => + log( + `[service-runner] resume turn for session ${sessionID} failed: ${ + err instanceof Error ? err.message : String(err) + }`, + ), + ); + } + // Cleared AFTER firing (acceptance-agnostic): the turns are handed + // over, so a subsequent boot cannot double-resume them. + clearInflightJournal(journalPath); + } + } catch (err) { + log( + `[service-runner] interrupted-session re-dispatch skipped: ${ + err instanceof Error ? err.message : String(err) + }`, + ); + } + } + // ── the service: the SAME wiring startAmicodeService performs ──────────── // No fleetActivation is ever passed here (H3): the hub runs the byte- // identical unarmed base service; arming is an extension-host decision. diff --git a/packages/extension/src/amicode_service_runner_cli.ts b/packages/extension/src/amicode_service_runner_cli.ts index dfaa39b29..09a3483c8 100644 --- a/packages/extension/src/amicode_service_runner_cli.ts +++ b/packages/extension/src/amicode_service_runner_cli.ts @@ -28,6 +28,22 @@ // fork's "canonical serves anonymous 200"). The pair // with AMICODE_SERVICE_AUTH=open is the hub posture. // Wins over AMICODE_ENGINE_PASSWORD. +// AMICODE_INFLIGHT_JOURNAL path to the in-flight chat-turn journal +// (#1552, interrupted-session re-dispatch). The +// service's dispatch seam appends a line for every +// proxied chat POST (start on entry, end on response +// close); the runner's boot-time resume reads the +// same file to re-dispatch the sessions an engine +// bounce killed mid-loop. Default +// ~/.amico/server/active-sessions.jsonl. +// AMICODE_RESUME_RECENT_MS the recency window for resume candidates: a +// session whose latest unmatched chat-turn start is +// older than this never resumes (yesterday's +// abandoned loop is noise, not continuity). +// Default 12h (43_200_000). +// AMICODE_RESUME_DISABLED "=1" skips the boot-time resume entirely — +// the test/ops escape hatch (the service keeps +// writing the journal either way). // OPENCODE_DB the canonical pin — passed through to the spawned // engine untouched (the hub's session store). // @@ -82,6 +98,13 @@ async function main(): Promise { engineCwd: (process.env.AMICODE_ENGINE_CWD ?? "").trim() || undefined, enginePassword: (process.env.AMICODE_ENGINE_PASSWORD ?? "").trim() || undefined, engineUnarmed: (process.env.AMICODE_ENGINE_UNARMED ?? "").trim() === "1", + // #1552 (interrupted-session re-dispatch): the same journal path the + // service's tracking writes (AMICODE_INFLIGHT_JOURNAL — the default + // resolution lives in inflightJournalPath, the single source), plus + // the recency window and the disable hatch. + resumeJournalPath: (process.env.AMICODE_INFLIGHT_JOURNAL ?? "").trim() || undefined, + resumeRecentMs: envInt("AMICODE_RESUME_RECENT_MS"), + resumeDisabled: (process.env.AMICODE_RESUME_DISABLED ?? "").trim() === "1", engineEnv: (process.env.OPENCODE_DB ?? "").trim() ? { OPENCODE_DB: process.env.OPENCODE_DB } : undefined, log: (line) => console.log(line), }); diff --git a/packages/extension/test/amicode_service_inflight_journal.test.ts b/packages/extension/test/amicode_service_inflight_journal.test.ts new file mode 100644 index 000000000..387fbfe6d --- /dev/null +++ b/packages/extension/test/amicode_service_inflight_journal.test.ts @@ -0,0 +1,556 @@ +// #1552 (interrupted-session re-dispatch): a hub engine bounce (SIGTERM at +// restart, or a crash) kills every in-flight agent loop — they live in the +// engine process's memory — and nothing re-dispatches them once the engine +// returns. The fix's whole state machine lives in the extension layer: +// track — the service's dispatch seam journals every proxied chat POST +// (start on entry, end on response close — the same seam that +// sees every engine-bound request), +// derive — a PURE function folds the journal into the interrupted set +// (start-without-end, recent enough, capped), +// resume — the runner's boot path, after engine health, fires ONE +// mechanical resume turn per interrupted session (fire-and-forget +// — the chat POST holds until the turn completes, tens of +// minutes; boot must never wedge), then clears the journal. +// +// Groups below, in that order: the pure derivation (the issue's unit-test +// acceptance: interrupted-set from the journal, resume-message shape, the +// cap, the never-wedge failure path), the route match, the append sidecar, +// the dispatch integration (a real proxied chat turn, a client-killed one), +// and the runner's boot-time re-dispatch (the fake-engine idiom of +// amicode_service_runner.test.ts). +import { describe, it, expect, beforeAll, afterAll } from "vitest"; +import { chmodSync, existsSync, mkdtempSync, readFileSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import * as http from "node:http"; +import type { ServerResponse } from "node:http"; +import { AddressInfo } from "node:net"; +import { EventEmitter } from "node:events"; +import { createAmicodeService } from "../src/amicode_service"; +import { bootAmicodeServiceRunner, type AmicodeServiceRunnerBoot } from "../src/amicode_service_runner"; +import { serverAuthHeader } from "../src/server_auth"; +import { + RESUME_CAP, + RESUME_MESSAGE, + RESUME_RECENT_MS_DEFAULT, + clearInflightJournal, + deriveInterruptedSessions, + matchInflightChatTurn, + readInterruptedSessions, + trackInflightChatTurn, +} from "../src/amicode_service/inflight_journal"; + +/** Poll until pred() holds (the fire-and-forget resume POSTs land on the + * event loop AFTER the boot promise resolves — assertions must wait for + * the evidence, not for the boot). */ +async function waitUntil(pred: () => boolean, timeoutMs = 5_000): Promise { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + if (pred()) return; + await new Promise((r) => setTimeout(r, 25)); + } + throw new Error("condition not met within timeout"); +} + +// ── the pure derivation ────────────────────────────────────────────────────── + +describe("#1552 in-flight journal — deriveInterruptedSessions (pure)", () => { + const now = 1_000_000_000_000; + const line = (kind: "start" | "end", sessionID: string, ts: number): string => + JSON.stringify({ kind, sessionID, ts }); + + it("a start without a matching end is the interrupted shape; a matched end closes it", () => { + const interrupted = deriveInterruptedSessions(line("start", "ses-a", now - 5_000) + "\n", now); + expect(interrupted.map((s) => s.sessionID)).toEqual(["ses-a"]); + const completed = deriveInterruptedSessions( + line("start", "ses-a", now - 5_000) + "\n" + line("end", "ses-a", now - 4_000) + "\n", + now, + ); + expect(completed).toEqual([]); + }); + + it("a re-start after an end is interrupted again — the LATEST unmatched start is the one that counts", () => { + const journal = + line("start", "ses-a", now - 30_000) + + "\n" + + line("end", "ses-a", now - 20_000) + + "\n" + + line("start", "ses-a", now - 10_000) + + "\n"; + const out = deriveInterruptedSessions(journal, now); + expect(out).toHaveLength(1); + expect(out[0]).toEqual({ sessionID: "ses-a", ts: now - 10_000 }); + }); + + it("never-resume shapes: only ends, or starts older than the recency window", () => { + expect(deriveInterruptedSessions(line("end", "ses-endonly", now - 1) + "\n", now)).toEqual([]); + const stale = deriveInterruptedSessions( + line("start", "ses-stale", now - RESUME_RECENT_MS_DEFAULT - 1) + "\n", + now, + ); + expect(stale).toEqual([]); + const fresh = deriveInterruptedSessions( + line("start", "ses-fresh", now - RESUME_RECENT_MS_DEFAULT + 1) + "\n", + now, + ); + expect(fresh.map((s) => s.sessionID)).toEqual(["ses-fresh"]); + }); + + it("the cap is 8 and keeps the MOST RECENT starts (a corrupted journal cannot fan out unbounded work)", () => { + expect(RESUME_CAP).toBe(8); + const lines: string[] = []; + for (let i = 0; i < 10; i++) lines.push(line("start", `ses-${i}`, now - i * 1_000)); + const out = deriveInterruptedSessions(lines.join("\n") + "\n", now); + expect(out).toHaveLength(8); + expect(out.map((s) => s.sessionID)).toEqual([ + "ses-0", + "ses-1", + "ses-2", + "ses-3", + "ses-4", + "ses-5", + "ses-6", + "ses-7", + ]); + }); + + it("corrupt or foreign lines are skipped, never fatal; an empty journal is a clean boot", () => { + const journal = + "not json\n" + + JSON.stringify({ kind: "weird", sessionID: "ses-x" }) + + "\n" + + line("start", "ses-good", now - 1_000) + + "\n"; + expect(deriveInterruptedSessions(journal, now).map((s) => s.sessionID)).toEqual(["ses-good"]); + expect(deriveInterruptedSessions("", now)).toEqual([]); + }); + + it("readInterruptedSessions: a missing journal file is a clean boot, not a failure", () => { + expect(readInterruptedSessions(join(tmpdir(), `amicode-no-such-journal-${Date.now()}.jsonl`), now)).toEqual([]); + }); +}); + +// ── the route match ────────────────────────────────────────────────────────── + +describe("#1552 in-flight journal — the chat-turn route match", () => { + it("POST /session/{id}/message is the one tracked route — the SDK's session.prompt (v1 {id}, v2 {sessionID}: the same URL)", () => { + expect(matchInflightChatTurn("POST", "/session/ses_abc123/message")).toBe("ses_abc123"); + }); + + it("GETs, SSE, and every other session route stay untracked — they observe, they don't run a loop", () => { + expect(matchInflightChatTurn("GET", "/session/ses_abc123/message")).toBeUndefined(); // the message LIST + expect(matchInflightChatTurn("GET", "/session/ses_abc123/message/msg_1")).toBeUndefined(); // one message + expect(matchInflightChatTurn("POST", "/session")).toBeUndefined(); // session create + expect(matchInflightChatTurn("POST", "/session/ses_abc123/abort")).toBeUndefined(); + expect(matchInflightChatTurn("POST", "/session/ses_abc123/command")).toBeUndefined(); + expect(matchInflightChatTurn("POST", "/session/ses_abc123/shell")).toBeUndefined(); + expect(matchInflightChatTurn("POST", "/session/ses_abc123/share")).toBeUndefined(); + expect(matchInflightChatTurn("POST", "/session/ses_abc123/summarize")).toBeUndefined(); + expect(matchInflightChatTurn("POST", "/session/ses_abc123/prompt_async")).toBeUndefined(); + expect(matchInflightChatTurn("POST", "/session/ses_abc123")).toBeUndefined(); // session update is PATCH, not this + expect(matchInflightChatTurn("GET", "/event")).toBeUndefined(); // the SSE observer + expect(matchInflightChatTurn("GET", "/session")).toBeUndefined(); + }); +}); + +// ── the append sidecar ─────────────────────────────────────────────────────── + +describe("#1552 in-flight journal — append + clear (trackInflightChatTurn)", () => { + /** A stand-in response: the only surface the hook touches is "once close". */ + const fakeRes = (): ServerResponse => new EventEmitter() as unknown as ServerResponse; + + it("writes a start line on track, and the response's close writes the matching end", () => { + const dir = mkdtempSync(join(tmpdir(), "amicode-inflight-unit-")); + const journal = join(dir, "active-sessions.jsonl"); + const res = fakeRes(); + trackInflightChatTurn(res, "ses-track", journal); + const start = JSON.parse(readFileSync(journal, "utf8")) as { kind: string; sessionID: string }; + expect(start).toMatchObject({ kind: "start", sessionID: "ses-track" }); + res.emit("close"); + const lines = readFileSync(journal, "utf8") + .trim() + .split("\n") + .map((l) => JSON.parse(l) as { kind: string; sessionID: string }); + expect(lines).toHaveLength(2); + expect(lines[1]).toMatchObject({ kind: "end", sessionID: "ses-track" }); + rmSync(dir, { recursive: true, force: true }); + }); + + it("the parent directory is created on demand (~/.amico/server does not pre-exist on a fresh hub)", () => { + const dir = mkdtempSync(join(tmpdir(), "amicode-inflight-mkdir-")); + const journal = join(dir, "nested", "deeper", "active-sessions.jsonl"); + expect(() => trackInflightChatTurn(fakeRes(), "ses-mkdir", journal)).not.toThrow(); + expect(existsSync(journal)).toBe(true); + rmSync(dir, { recursive: true, force: true }); + }); + + it("appends are best-effort: an unwritable journal never throws into the request path", () => { + const dir = mkdtempSync(join(tmpdir(), "amicode-inflight-unwritable-")); + const fileAsParent = join(dir, "a-regular-file"); + writeFileSync(fileAsParent, "not a directory"); + const journal = join(fileAsParent, "active-sessions.jsonl"); // ENOTDIR under any mkdir/append + const res = fakeRes(); + expect(() => trackInflightChatTurn(res, "ses-bad", journal)).not.toThrow(); + expect(() => res.emit("close")).not.toThrow(); // the close-path append is equally best-effort + rmSync(dir, { recursive: true, force: true }); + }); + + it("clearInflightJournal empties the file so the next boot derives nothing (no double-resume)", () => { + const dir = mkdtempSync(join(tmpdir(), "amicode-inflight-clear-")); + const journal = join(dir, "active-sessions.jsonl"); + trackInflightChatTurn(fakeRes(), "ses-clear", journal); + clearInflightJournal(journal); + expect(readFileSync(journal, "utf8")).toBe(""); + expect(readInterruptedSessions(journal)).toEqual([]); + rmSync(dir, { recursive: true, force: true }); + }); + + it("the resume nudge's exact text (the issue's wording — mechanical, no new instructions)", () => { + expect(RESUME_MESSAGE).toBe( + "[auto-resume] The engine restarted while your loop was in flight (issue #1552). Re-read your session ledger/state and continue exactly where you left off.", + ); + }); +}); + +// ── the dispatch seam (a real proxied chat turn through the service) ──────── + +describe("#1552 in-flight journal — the dispatch seam (proxied chat turns are journaled)", () => { + let root: string; + let journal: string; + let engineUrl: string; + let engineSockets: Set; + let engineServer: http.Server; + let service: ReturnType; + let base: string; + const enginePassword = "inflight-engine-mint"; + let engineAuth: string; + let prevEnv: string | undefined; + + beforeAll(async () => { + root = mkdtempSync(join(tmpdir(), "amicode-inflight-dispatch-")); + journal = join(root, "active-sessions.jsonl"); + // The journal path is env-carried (the production config surface — + // AMICODE_INFLIGHT_JOURNAL); pin it for hermetic assertions. + prevEnv = process.env.AMICODE_INFLIGHT_JOURNAL; + process.env.AMICODE_INFLIGHT_JOURNAL = journal; + + // The mock engine: everything answers 200 JSON; the HOLD route never + // answers — an in-flight turn. Sockets are tracked so afterAll can + // destroy held connections (server.close() waits for them otherwise). + engineSockets = new Set(); + engineServer = http.createServer((req, res) => { + if (req.method === "POST" && req.url === "/session/ses-hold/message") return; // never answered + let body = ""; + req.on("data", (c: Buffer) => (body += c)); + req.on("end", () => { + res.writeHead(200, { "content-type": "application/json" }); + res.end(JSON.stringify({ ok: true, mockEngine: true, echo: body })); + }); + }); + engineServer.on("connection", (s) => { + engineSockets.add(s); + s.on("close", () => engineSockets.delete(s)); + }); + await new Promise((r) => engineServer.listen(0, "127.0.0.1", r)); + engineUrl = `http://127.0.0.1:${(engineServer.address() as AddressInfo).port}`; + + service = createAmicodeService({ + password: "service-own-mint", + engine: { password: enginePassword, getUrl: () => engineUrl }, + }); + base = (await service.start()).toString().replace(/\/$/, ""); + engineAuth = serverAuthHeader(enginePassword); + }); + + afterAll(async () => { + await service.stop(); + for (const s of engineSockets) s.destroy(); + await new Promise((r) => engineServer.close(() => r())); + if (prevEnv === undefined) delete process.env.AMICODE_INFLIGHT_JOURNAL; + else process.env.AMICODE_INFLIGHT_JOURNAL = prevEnv; + rmSync(root, { recursive: true, force: true }); + }); + + const journalLines = (): string[] => + readFileSync(journal, "utf8") + .split("\n") + .filter((l) => l.trim() !== ""); + + it("a proxied chat POST journals start + end; the same-route GET and the session list journal nothing", async () => { + const r = await fetch(`${base}/session/ses-ok/message`, { + method: "POST", + headers: { Authorization: engineAuth, "Content-Type": "application/json" }, + body: JSON.stringify({ parts: [{ type: "text", text: "run" }] }), + }); + expect(r.status).toBe(200); + const lines = journalLines().map((l) => JSON.parse(l) as { kind: string; sessionID: string }); + expect(lines).toHaveLength(2); + expect(lines[0]).toMatchObject({ kind: "start", sessionID: "ses-ok" }); + expect(lines[1]).toMatchObject({ kind: "end", sessionID: "ses-ok" }); + // The turn completed — this session is NOT interrupted. + expect(readInterruptedSessions(journal).map((s) => s.sessionID)).toEqual([]); + // Observers ride the same origin without being tracked: the message + // LIST (GET of the same path) and the session list. + await fetch(`${base}/session/ses-ok/message`, { headers: { Authorization: engineAuth } }); + await fetch(`${base}/session`, { headers: { Authorization: engineAuth } }); + expect(journalLines()).toHaveLength(2); + }); + + it("an interrupted proxied chat POST (client death) leaves start-without-end = INTERRUPTED — the bounce shape", async () => { + const controller = new AbortController(); + const pendingTurn = fetch(`${base}/session/ses-hold/message`, { + method: "POST", + headers: { Authorization: engineAuth, "Content-Type": "application/json" }, + body: JSON.stringify({ parts: [{ type: "text", text: "a long research loop" }] }), + signal: controller.signal, + }).catch(() => undefined); // the client abort is the expected outcome — caught HERE, at creation + void pendingTurn; + // The start line lands the instant the request enters dispatch. + await waitUntil(() => journalLines().some((l) => l.includes("ses-hold"))); + expect(readInterruptedSessions(journal).map((s) => s.sessionID)).toEqual(["ses-hold"]); + // The client dies (lid-close, tunnel death — or the SIGTERM that takes + // the service down with it): the response closes, the end line lands. + controller.abort(); + await waitUntil(() => journalLines().some((l) => l.includes('"end"') && l.includes("ses-hold"))); + expect(readInterruptedSessions(journal)).toEqual([]); + await pendingTurn; + }); +}); + +// ── the runner's boot-time re-dispatch (the fake-engine idiom) ─────────────── + +describe("#1552 runner — interrupted-session re-dispatch on boot", () => { + const boots: AmicodeServiceRunnerBoot[] = []; + afterAll(async () => { + for (const b of boots.splice(0)) await b.shutdown().catch(() => undefined); + }); + + /** A stand-in engine binary (the shebang-node fake-engine idiom): answers + * the health GET with 200; records every non-GET request (the resume + * POSTs) to the dump as JSONL {method, url, auth, body}; FAKE_ENGINE_FAIL + * answers 500 to them, FAKE_ENGINE_HOLD never answers (the in-flight + * shape), else 200 with the SDK's SessionPromptResponses shape. */ + function writeResumeEngine(dir: string): { bin: string; dump: string } { + const bin = join(dir, "fake-engine"); + const dump = join(dir, "resume-dump.jsonl"); + writeFileSync( + bin, + `#!/usr/bin/env node +const { createServer } = require("node:http"); +const { appendFileSync } = require("node:fs"); +const port = Number(process.argv[4] ?? 0); +const dump = ${JSON.stringify(dump)}; +createServer((req, res) => { + if (req.method === "GET") { + res.writeHead(200, { "content-type": "text/plain" }); + res.end("fake engine up"); + return; + } + let body = ""; + req.on("data", (c) => (body += c)); + req.on("end", () => { + appendFileSync(dump, JSON.stringify({ method: req.method, url: req.url, auth: req.headers.authorization ?? null, body }) + "\\n"); + if (process.env.FAKE_ENGINE_FAIL) { + res.writeHead(500, { "content-type": "application/json" }); + res.end(JSON.stringify({ ok: false })); + return; + } + if (process.env.FAKE_ENGINE_HOLD) return; // never answer: the turn stays in flight + res.writeHead(200, { "content-type": "application/json" }); + res.end(JSON.stringify({ info: { id: "msg_resume", role: "assistant" }, parts: [] })); + }); +}).listen(port, "127.0.0.1"); +`, + ); + chmodSync(bin, 0o755); + return { bin, dump }; + } + + function writeStubShelf(dir: string): string { + writeFileSync(join(dir, "index.html"), "stub shelf"); + return dir; + } + + function writeJournal(lines: Array<{ kind: "start" | "end"; sessionID: string; ts: number }>): string { + const journal = join(mkdtempSync(join(tmpdir(), "amicode-resume-journal-")), "active-sessions.jsonl"); + writeFileSync(journal, lines.map((l) => JSON.stringify(l)).join("\n") + "\n"); + return journal; + } + + interface ResumeHit { + method: string; + url: string; + auth: string | null; + body: string; + } + const dumpHits = (dump: string): ResumeHit[] => + existsSync(dump) + ? readFileSync(dump, "utf8") + .split("\n") + .filter((l) => l.trim() !== "") + .map((l) => JSON.parse(l) as ResumeHit) + : []; + + it("the interrupted set gets ONE resume POST per session with the exact nudge text, then the journal is cleared", async () => { + const dir = mkdtempSync(join(tmpdir(), "amicode-resume-dispatch-")); + const { bin, dump } = writeResumeEngine(dir); + const now = Date.now(); + const journal = writeJournal([ + { kind: "start", sessionID: "ses-one", ts: now - 60_000 }, + { kind: "start", sessionID: "ses-two", ts: now - 30_000 }, + { kind: "start", sessionID: "ses-done", ts: now - 120_000 }, + { kind: "end", sessionID: "ses-done", ts: now - 110_000 }, + { kind: "start", sessionID: "ses-stale", ts: now - RESUME_RECENT_MS_DEFAULT - 60_000 }, + ]); + const lines: string[] = []; + const boot = await bootAmicodeServiceRunner({ + engineBin: bin, + appDistRoot: writeStubShelf(mkdtempSync(join(tmpdir(), "amicode-resume-shelf-"))), + enginePassword: "resume-engine-mint", + resumeJournalPath: journal, + healthTimeoutMs: 10_000, + servicePort: 0, + enginePort: 0, + log: (l) => lines.push(l), + }); + boots.push(boot); + + // Exactly one POST per interrupted session (fire-and-forget lands + // after the boot resolves — wait for the evidence). + await waitUntil(() => dumpHits(dump).length >= 2); + const hits = dumpHits(dump); + expect(hits).toHaveLength(2); + expect(hits.map((h) => h.url).sort()).toEqual(["/session/ses-one/message", "/session/ses-two/message"]); + for (const hit of hits) { + expect(hit.method).toBe("POST"); + // The POST carries the ENGINE credential (serverAuthHeader(enginePassword)). + expect(hit.auth).toBe(serverAuthHeader("resume-engine-mint")); + // The same body schema the SDK uses for a user prompt, carrying the exact nudge. + const body = JSON.parse(hit.body) as { parts?: Array<{ type?: string; text?: string }> }; + expect(body.parts).toEqual([{ type: "text", text: RESUME_MESSAGE }]); + } + // Third parties untouched: the matched-end session and the stale one never resume. + expect(hits.some((h) => h.url.includes("ses-done"))).toBe(false); + expect(hits.some((h) => h.url.includes("ses-stale"))).toBe(false); + // The journal is cleared — a subsequent boot cannot double-resume. + expect(readFileSync(journal, "utf8")).toBe(""); + // Every resume outcome = one log line. + await waitUntil(() => lines.filter((l) => l.includes("resume turn for session")).length >= 2); + rmSync(dir, { recursive: true, force: true }); + }, 30_000); + + it("fire-and-forget: a chat POST that never answers (the tens-of-minutes turn) never wedges the boot", async () => { + const dir = mkdtempSync(join(tmpdir(), "amicode-resume-hold-")); + const { bin, dump } = writeResumeEngine(dir); + const journal = writeJournal([{ kind: "start", sessionID: "ses-held", ts: Date.now() - 1_000 }]); + // The boot RESOLVES (the await below returns) while the resume turn is + // still in flight — that is the never-wedge contract. + const boot = await bootAmicodeServiceRunner({ + engineBin: bin, + appDistRoot: writeStubShelf(mkdtempSync(join(tmpdir(), "amicode-resume-hold-shelf-"))), + engineEnv: { FAKE_ENGINE_HOLD: "1" }, + resumeJournalPath: journal, + healthTimeoutMs: 10_000, + servicePort: 0, + enginePort: 0, + log: () => undefined, + }); + boots.push(boot); + await waitUntil(() => dumpHits(dump).length >= 1); + expect(dumpHits(dump)[0].url).toBe("/session/ses-held/message"); + expect(readFileSync(journal, "utf8")).toBe(""); + // Teardown destroys the held connection — the dangling fetch rejects + // into its own catch, never the process. + await boot.shutdown(); + rmSync(dir, { recursive: true, force: true }); + }, 30_000); + + it("a FAILED resume turn is one log line and never a failed boot (the service still serves)", async () => { + const dir = mkdtempSync(join(tmpdir(), "amicode-resume-reject-")); + const { bin } = writeResumeEngine(dir); + const journal = writeJournal([{ kind: "start", sessionID: "ses-rej", ts: Date.now() - 1_000 }]); + const lines: string[] = []; + const boot = await bootAmicodeServiceRunner({ + engineBin: bin, + appDistRoot: writeStubShelf(mkdtempSync(join(tmpdir(), "amicode-resume-reject-shelf-"))), + resumeJournalPath: journal, + // Injected for tests: the POST itself rejects (engine unreachable mid-fire). + resumeFetch: () => Promise.reject(new Error("engine refused the resume turn")), + healthTimeoutMs: 10_000, + servicePort: 0, + enginePort: 0, + log: (l) => lines.push(l), + }); + boots.push(boot); + await waitUntil(() => + lines.some((l) => l.includes("resume turn for session ses-rej failed: engine refused the resume turn")), + ); + // The boot is up and serving — a resume failure never wedges or fails it. + const doc = await fetch(`${boot.url}/`, { headers: { Authorization: boot.authHeader } }); + expect(doc.status).toBe(200); + rmSync(dir, { recursive: true, force: true }); + }, 30_000); + + it("a non-2xx engine answer is logged as the outcome, still never a failed boot", async () => { + const dir = mkdtempSync(join(tmpdir(), "amicode-resume-500-")); + const { bin } = writeResumeEngine(dir); + const journal = writeJournal([{ kind: "start", sessionID: "ses-500", ts: Date.now() - 1_000 }]); + const lines: string[] = []; + const boot = await bootAmicodeServiceRunner({ + engineBin: bin, + appDistRoot: writeStubShelf(mkdtempSync(join(tmpdir(), "amicode-resume-500-shelf-"))), + engineEnv: { FAKE_ENGINE_FAIL: "1" }, + resumeJournalPath: journal, + healthTimeoutMs: 10_000, + servicePort: 0, + enginePort: 0, + log: (l) => lines.push(l), + }); + boots.push(boot); + await waitUntil(() => lines.some((l) => l.includes("resume turn for session ses-500") && l.includes("500"))); + rmSync(dir, { recursive: true, force: true }); + }, 30_000); + + it("AMICODE_RESUME_DISABLED=1 (resumeDisabled) skips the step — no POST, journal left intact for the operator", async () => { + const dir = mkdtempSync(join(tmpdir(), "amicode-resume-disabled-")); + const { bin, dump } = writeResumeEngine(dir); + const journal = writeJournal([{ kind: "start", sessionID: "ses-skip", ts: Date.now() - 1_000 }]); + const lines: string[] = []; + const boot = await bootAmicodeServiceRunner({ + engineBin: bin, + appDistRoot: writeStubShelf(mkdtempSync(join(tmpdir(), "amicode-resume-disabled-shelf-"))), + resumeJournalPath: journal, + resumeDisabled: true, + healthTimeoutMs: 10_000, + servicePort: 0, + enginePort: 0, + log: (l) => lines.push(l), + }); + boots.push(boot); + await new Promise((r) => setTimeout(r, 300)); // any (wrongly) fired POST has time to land + expect(dumpHits(dump)).toEqual([]); + expect(readFileSync(journal, "utf8")).not.toBe(""); // untouched — the escape hatch leaves the state alone + expect(lines.some((l) => l.includes("resume disabled"))).toBe(true); + rmSync(dir, { recursive: true, force: true }); + }, 30_000); + + it("the unarmed hub posture sends the resume turn ANONYMOUSLY (no credential exists to carry)", async () => { + const dir = mkdtempSync(join(tmpdir(), "amicode-resume-unarmed-")); + const { bin, dump } = writeResumeEngine(dir); + const journal = writeJournal([{ kind: "start", sessionID: "ses-anon", ts: Date.now() - 1_000 }]); + const boot = await bootAmicodeServiceRunner({ + engineBin: bin, + appDistRoot: writeStubShelf(mkdtempSync(join(tmpdir(), "amicode-resume-unarmed-shelf-"))), + engineUnarmed: true, + resumeJournalPath: journal, + healthTimeoutMs: 10_000, + servicePort: 0, + enginePort: 0, + log: () => undefined, + }); + boots.push(boot); + await waitUntil(() => dumpHits(dump).length >= 1); + expect(dumpHits(dump)[0].auth).toBeNull(); + rmSync(dir, { recursive: true, force: true }); + }, 30_000); +});