Skip to content

Commit 22353ee

Browse files
committed
Move the run.json turn-boundary snapshot trigger into run-sink
run.json is session domain state, and src/session/run-sink.ts already owns the turn count and already special-cases inference.done there (clearing a stale runError). The mid-run snapshot trigger lived in src/tui/runner.ts instead, keyed off a second subscription to the same event stream — the exact shape that has already cost this constraint three renderer swaps. Move the cadence into run-sink via an onTurnBoundary callback, so the renderer only owns how to persist a snapshot, not when one is due. The end-to-end test now drives the callback the way production wiring does, instead of calling saveState directly after each sink call, so it actually exercises the trigger rather than simulating it.
1 parent 196344c commit 22353ee

3 files changed

Lines changed: 63 additions & 45 deletions

File tree

‎src/session/run-sink.ts‎

Lines changed: 13 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,17 @@ export type RunSinkArgs = {
2222
// Continues a resumed session's persisted run.json turn count instead of
2323
// restarting the collector at zero.
2424
initialTurnCount?: number;
25+
// Fired at every turn boundary so a caller can persist a mid-run run.json
26+
// snapshot. `inference.done` is the turn boundary every reactor cycle
27+
// guarantees; `reactor.done` fires once, at shutdown, and never between
28+
// turns of a long-lived interactive session. Keying the mid-run snapshot
29+
// off `reactor.done` left turnsUsed frozen at its resume-time value for
30+
// the entire session — a live monorepo session showed turnsUsed: 0 with
31+
// dozens of turns already in the turns log. This cadence lives here,
32+
// alongside the turn count it reports, rather than in a second
33+
// subscription to the same event stream in a renderer: the renderer has
34+
// already been swapped out from under this constraint three times.
35+
onTurnBoundary?: () => void;
2536
};
2637

2738
export type RunSink = {
@@ -69,7 +80,7 @@ export function resolveExecRunStatus(args: {
6980
}
7081

7182
export function createRunSink(args: RunSinkArgs): RunSink {
72-
const { emitter, hookManager, onTurnComplete, initialTurnCount } = args;
83+
const { emitter, hookManager, onTurnComplete, initialTurnCount, onTurnBoundary } = args;
7384

7485
function hasConfiguredHooks(): boolean {
7586
return hookManager.getStatuses().length > 0;
@@ -114,6 +125,7 @@ export function createRunSink(args: RunSinkArgs): RunSink {
114125
// would mark a recovered successful send as failed.
115126
if (event.type === "inference.done") {
116127
runError = undefined;
128+
onTurnBoundary?.();
117129
}
118130
if (event.type === "reactor.error") {
119131
const data = event.data as { error: string };

‎src/tui/runner.ts‎

Lines changed: 5 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -201,18 +201,6 @@ export function resolveExitCode(args: ResolveExitCodeArgs): number {
201201
return 0;
202202
}
203203

204-
// `inference.done` is the turn boundary every reactor cycle guarantees;
205-
// `reactor.done` fires once, at shutdown, and never between turns of a
206-
// long-lived interactive session. Keying the mid-run run.json snapshot off
207-
// `reactor.done` left turnsUsed frozen at its resume-time value for the
208-
// entire session — a live monorepo session showed turnsUsed: 0 with dozens
209-
// of turns already in the turns log. The terminal write on close still goes
210-
// through writeRunSnapshot directly with the real final status, so this
211-
// only needs to cover progress snapshots taken while the run is live.
212-
function isRunSnapshotTurnBoundary(eventType: string): boolean {
213-
return eventType === "inference.done";
214-
}
215-
216204
/** One-line transcript block when resume history fails to load. */
217205
export function resumeTranscriptLoadErrorBlock(err: unknown): {
218206
type: "error";
@@ -1358,6 +1346,11 @@ export async function runTUI(initialConfig: Config): Promise<number> {
13581346
duration_ms: ctx.durationMs,
13591347
});
13601348
},
1349+
// persistRunSnapshot is defined below but not invoked until the stream
1350+
// starts consuming events, well after this closure captures it.
1351+
onTurnBoundary: () => {
1352+
void persistRunSnapshot("running");
1353+
},
13611354
});
13621355

13631356
// MCP servers connected so far, keyed by name so a reconnect after a failure
@@ -1411,9 +1404,6 @@ export async function runTUI(initialConfig: Config): Promise<number> {
14111404
const streamSink = (event: Parameters<typeof runSink.sink>[0]): void => {
14121405
runSink.sink(event);
14131406
cycleRecorder.handleEvent(event);
1414-
if (isRunSnapshotTurnBoundary(event.type)) {
1415-
void persistRunSnapshot("running");
1416-
}
14171407
};
14181408

14191409
// Tool count before any MCP server connects; a reload is only worthwhile if

‎tests/unit/session/run-state-e2e.test.ts‎

Lines changed: 45 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -12,10 +12,11 @@ import { saveState, loadState, type RunState } from "../../../src/session/state.
1212

1313
// End-to-end coverage for the run.json turn-boundary snapshot fix (CL-5534):
1414
// createRunSink, saveState, and loadState run for real against a temp
15-
// session directory — nothing mocked. This is what now covers
16-
// isRunSnapshotTurnBoundary's observable effect, since that predicate was
17-
// un-exported from src/tui/runner.ts as a testability-only surface with a
18-
// single production caller.
15+
// session directory — nothing mocked. The snapshot cadence itself now lives
16+
// in createRunSink's onTurnBoundary callback (moved out of src/tui/runner.ts
17+
// so a future renderer swap cannot silently drop it again), so these tests
18+
// drive that callback exactly as production wiring does: nothing here calls
19+
// saveState directly from the turn loop, only from inside onTurnBoundary.
1920

2021
const noopHookManager = { dispatchPostTurn: () => undefined, getStatuses: () => [] };
2122

@@ -48,18 +49,22 @@ describe("run.json turn-boundary snapshots — end to end", () => {
4849
const home = mkdtempSync(join(tmpdir(), "corbits-run-state-home-"));
4950
const sessionId = generateSessionId();
5051
try {
51-
const runSink = createRunSink({ emitter: new EventEmitter(), hookManager: noopHookManager });
52+
const writes: Promise<void>[] = [];
53+
const runSink = createRunSink({
54+
emitter: new EventEmitter(),
55+
hookManager: noopHookManager,
56+
onTurnBoundary: () => {
57+
writes.push(saveState(cwd, sessionId, baseState({ status: "running" }, runSink.getTurnCount()), home));
58+
},
59+
});
5260

53-
const observed: number[] = [];
5461
for (let turn = 1; turn <= 4; turn++) {
5562
runSink.sink(inferenceDone());
56-
await saveState(cwd, sessionId, baseState({ status: "running" }, runSink.getTurnCount()), home);
57-
const onDisk = await loadState(cwd, sessionId, home);
58-
expect(onDisk).not.toBeNull();
59-
observed.push(onDisk!.turnsUsed);
6063
}
64+
await Promise.all(writes);
6165

62-
expect(observed).toEqual([1, 2, 3, 4]);
66+
const onDisk = await loadState(cwd, sessionId, home);
67+
expect(onDisk?.turnsUsed).toBe(4);
6368

6469
runSink.sink({ type: "reactor.done", data: {} } as unknown as ReactorEmittedEvent);
6570
await saveState(
@@ -82,16 +87,21 @@ describe("run.json turn-boundary snapshots — end to end", () => {
8287
const home = mkdtempSync(join(tmpdir(), "corbits-run-state-home-"));
8388
const sessionId = generateSessionId();
8489
try {
85-
const runSink = createRunSink({ emitter: new EventEmitter(), hookManager: noopHookManager });
86-
87-
// Fire all 20 turns and their snapshot writes back to back, with no
88-
// await between them — the per-session writeChains promise chain in
89-
// src/session/state.ts is what keeps these ordered rather than the
90-
// caller awaiting each one before starting the next.
9190
const writes: Promise<void>[] = [];
91+
const runSink = createRunSink({
92+
emitter: new EventEmitter(),
93+
hookManager: noopHookManager,
94+
onTurnBoundary: () => {
95+
writes.push(saveState(cwd, sessionId, baseState({ status: "running" }, runSink.getTurnCount()), home));
96+
},
97+
});
98+
99+
// Fire all 20 turns back to back, with no await between them — the
100+
// per-session writeChains promise chain in src/session/state.ts is what
101+
// keeps the resulting writes ordered, not the caller awaiting each one
102+
// before starting the next.
92103
for (let turn = 1; turn <= 20; turn++) {
93104
runSink.sink(inferenceDone());
94-
writes.push(saveState(cwd, sessionId, baseState({ status: "running" }, runSink.getTurnCount()), home));
95105
}
96106
await Promise.all(writes);
97107

@@ -109,18 +119,24 @@ describe("run.json turn-boundary snapshots — end to end", () => {
109119
const home = mkdtempSync(join(tmpdir(), "corbits-run-state-home-"));
110120
const sessionId = generateSessionId();
111121
try {
112-
const runSink = createRunSink({ emitter: new EventEmitter(), hookManager: noopHookManager });
113-
runSink.sink(inferenceDone());
114-
115-
// Issue a "running" progress snapshot but do not await it before
116-
// issuing the terminal "done" write right behind it — this models a
122+
let runningWrite: Promise<void> | undefined;
123+
const runSink = createRunSink({
124+
emitter: new EventEmitter(),
125+
hookManager: noopHookManager,
126+
onTurnBoundary: () => {
127+
runningWrite = saveState(
128+
cwd,
129+
sessionId,
130+
baseState({ status: "running" }, runSink.getTurnCount()),
131+
home,
132+
);
133+
},
134+
});
135+
136+
// The turn-boundary snapshot fires and is left un-awaited before the
137+
// terminal "done" write follows right behind it — this models a
117138
// straggler turn-boundary snapshot racing the close-out write.
118-
const runningWrite = saveState(
119-
cwd,
120-
sessionId,
121-
baseState({ status: "running" }, runSink.getTurnCount()),
122-
home,
123-
);
139+
runSink.sink(inferenceDone());
124140
runSink.sink({ type: "reactor.done", data: {} } as unknown as ReactorEmittedEvent);
125141
const doneWrite = saveState(
126142
cwd,

0 commit comments

Comments
 (0)