Skip to content
Open
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
189 changes: 189 additions & 0 deletions packages/extension/src/amicode_service/inflight_journal.ts
Original file line number Diff line number Diff line change
@@ -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<string, number>();
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 */
}
}
14 changes: 14 additions & 0 deletions packages/extension/src/amicode_service/server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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;
Expand Down
85 changes: 85 additions & 0 deletions packages/extension/src/amicode_service_runner.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";

Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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.
Expand Down
23 changes: 23 additions & 0 deletions packages/extension/src/amicode_service_runner_cli.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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).
//
Expand Down Expand Up @@ -82,6 +98,13 @@ async function main(): Promise<never> {
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),
});
Expand Down
Loading
Loading