Skip to content

Commit 5c893fa

Browse files
committed
Continue on non-fatal reactor errors and skip empty checkpoints
reactor.error with fatal:false is a recoverable stream signal; treating it as terminal idle/failed the TUI and run sink while the reactor kept going. Empty managed checkpoints were also creating new commits and risking session junk like partial.jsonl riding along — skip when no allowlisted path differs from HEAD.
1 parent 3391db2 commit 5c893fa

10 files changed

Lines changed: 226 additions & 3 deletions

‎src/agent/reactor-events.test.ts‎

Lines changed: 29 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,11 @@
11
import { describe, expect, test } from "bun:test";
22
import type { ReactorEmittedEvent } from "@intx/inference";
33
import type { ReactorInboundEvent } from "@intx/types/runtime";
4-
import { onReactorShutdown, onTurnBoundary } from "./reactor-events.js";
4+
import {
5+
isReactorErrorFatal,
6+
onReactorShutdown,
7+
onTurnBoundary,
8+
} from "./reactor-events.js";
59

610
// Bare `{ type: string }` literals only prove the string comparison works.
711
// The generic exists so the guards narrow across both `ReactorInboundEvent`
@@ -115,3 +119,27 @@ describe("onReactorShutdown", () => {
115119
expect(shutdowns.length).toBe(1);
116120
});
117121
});
122+
123+
describe("isReactorErrorFatal", () => {
124+
test("fatal:false continues", () => {
125+
expect(
126+
isReactorErrorFatal({ error: "transient write", fatal: false }),
127+
).toBe(false);
128+
});
129+
130+
test("fatal:true stays terminal", () => {
131+
expect(isReactorErrorFatal({ error: "gave up", fatal: true })).toBe(true);
132+
});
133+
134+
test("missing fatal stays terminal", () => {
135+
expect(isReactorErrorFatal({ error: "gave up" })).toBe(true);
136+
});
137+
138+
test("malformed payloads stay terminal", () => {
139+
expect(isReactorErrorFatal(undefined)).toBe(true);
140+
expect(isReactorErrorFatal(null)).toBe(true);
141+
expect(isReactorErrorFatal("boom")).toBe(true);
142+
expect(isReactorErrorFatal({ fatal: "false" })).toBe(true);
143+
expect(isReactorErrorFatal({ fatal: 0 })).toBe(true);
144+
});
145+
});

‎src/agent/reactor-events.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,8 @@
1313
* call sites without re-declaring the union here.
1414
*/
1515

16+
import { type } from "arktype";
17+
1618
/** True when `event` is the turn boundary — fires once per turn, every turn. */
1719
export const onTurnBoundary = <E extends { type: string }>(
1820
event: E,
@@ -24,3 +26,19 @@ export const onReactorShutdown = <E extends { type: string }>(
2426
event: E,
2527
): event is Extract<E, { type: "reactor.done" }> =>
2628
event.type === "reactor.done";
29+
30+
/**
31+
* Explicit non-fatal reactor.error payload. Only `fatal: false` continues;
32+
* missing, malformed, or any other value stays terminal.
33+
*/
34+
const NonFatalReactorErrorData = type({
35+
fatal: "false",
36+
});
37+
38+
/**
39+
* Whether a `reactor.error` payload should terminate the turn/shell.
40+
* Returns false only when the payload explicitly carries `fatal: false`.
41+
*/
42+
export function isReactorErrorFatal(data: unknown): boolean {
43+
return NonFatalReactorErrorData(data) instanceof type.errors;
44+
}

‎src/session/optimized-context-store.test.ts‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -710,6 +710,24 @@ describe("createOptimizedContextStore checkpoint", () => {
710710
}
711711
expect(gitSpawns).toEqual([]);
712712
});
713+
714+
test("skips commit when no managed path differs and never stages partial.jsonl", async () => {
715+
const dir = tempDir();
716+
const store = await createOptimizedContextStore(dir);
717+
await store.writeMetadata(EMPTY_CHECKPOINT_METADATA);
718+
const first = await store.commit({ message: "first managed checkpoint" });
719+
720+
fs.writeFileSync(path.join(dir, "partial.jsonl"), '{"reason":"abort"}\n');
721+
fs.writeFileSync(path.join(dir, "untracked-junk.txt"), "session junk\n");
722+
await store.writeMetadata(EMPTY_CHECKPOINT_METADATA);
723+
724+
const second = await store.commit({ message: "empty managed checkpoint" });
725+
expect(second.hash).toBe(first.hash);
726+
727+
const tree = await gitLsTree(dir);
728+
expect(tree).not.toContain("partial.jsonl");
729+
expect(tree).not.toContain("untracked-junk.txt");
730+
});
713731
});
714732

715733
describe("createSessionStores", () => {

‎src/session/optimized-context-store.ts‎

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -419,6 +419,24 @@ function extraCommitPaths(paths: readonly string[]): string[] {
419419
return paths.filter((filepath) => !VENDOR_COMMIT_ROOT_FILES.has(filepath));
420420
}
421421

422+
/**
423+
* True when any allowlisted managed path differs between HEAD and the
424+
* worktree. Callers must pass only managed paths — never "." and never
425+
* session junk such as partial.jsonl.
426+
*/
427+
async function managedPathsDiffer(
428+
dir: string,
429+
filepaths: readonly string[],
430+
): Promise<boolean> {
431+
if (filepaths.length === 0) return false;
432+
const matrix = await git.statusMatrix({
433+
fs,
434+
dir,
435+
filepaths: [...filepaths],
436+
});
437+
return matrix.some(([, head, workdir]) => head !== workdir);
438+
}
439+
422440
export interface SessionStores {
423441
storage: ContextStore;
424442
audit: AuditStore;
@@ -671,6 +689,25 @@ export async function createSessionStores(
671689
);
672690
const extraPaths = [...new Set([...add, ...remove])];
673691

692+
const managedFilepaths = [
693+
...VENDOR_COMMIT_ROOT_FILES,
694+
TOOL_OUTPUT_DIR,
695+
EVIDENCE_ARCHIVE_DIR,
696+
...extraPaths,
697+
];
698+
if (!(await managedPathsDiffer(dir, managedFilepaths))) {
699+
const [head] = await base.log(1);
700+
if (head !== undefined) {
701+
pendingBlobFilepaths.clear();
702+
pendingSegmentPaths.clear();
703+
if (stagedRewrite !== null) {
704+
liveTurnRefs = stagedRewrite;
705+
unpublishedRewrite = null;
706+
}
707+
return head;
708+
}
709+
}
710+
674711
try {
675712
for (const filepath of add) {
676713
await git.add({ fs, dir, filepath });

‎src/session/run-sink.test.ts‎

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -532,4 +532,47 @@ describe("createRunSink", () => {
532532
toolCallCount: 0,
533533
});
534534
});
535+
536+
test("non-fatal reactor.error does not sticky-fail the run", () => {
537+
const runSink = createRunSink({
538+
emitter: new EventEmitter(),
539+
hookManager: stubHookManager([]),
540+
});
541+
542+
runSink.sink(
543+
event("reactor.error", {
544+
error: "transient checkpoint write",
545+
fatal: false,
546+
}),
547+
);
548+
549+
expect(runSink.getRunError()).toBeUndefined();
550+
expect(runSink.getStatus()).not.toBe("failed");
551+
});
552+
553+
test("fatal reactor.error sticky-fails the run", () => {
554+
const runSink = createRunSink({
555+
emitter: new EventEmitter(),
556+
hookManager: stubHookManager([]),
557+
});
558+
559+
runSink.sink(
560+
event("reactor.error", { error: "reactor gave up", fatal: true }),
561+
);
562+
563+
expect(runSink.getRunError()).toBe("reactor gave up");
564+
expect(runSink.getStatus()).toBe("failed");
565+
});
566+
567+
test("reactor.error without fatal sticky-fails the run", () => {
568+
const runSink = createRunSink({
569+
emitter: new EventEmitter(),
570+
hookManager: stubHookManager([]),
571+
});
572+
573+
runSink.sink(event("reactor.error", { error: "reactor gave up" }));
574+
575+
expect(runSink.getRunError()).toBe("reactor gave up");
576+
expect(runSink.getStatus()).toBe("failed");
577+
});
535578
});

‎src/session/run-sink.ts‎

Lines changed: 5 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,10 @@ import type { EventEmitter } from "node:events";
22
import type { ReactorEmittedEvent } from "@intx/inference";
33
import type { LastCycleSource, TokenUsage } from "@intx/types/runtime";
44
import { createPerfReactorObserver } from "../perf/reactor-spans.js";
5-
import { onTurnBoundary } from "../agent/reactor-events.js";
5+
import {
6+
isReactorErrorFatal,
7+
onTurnBoundary,
8+
} from "../agent/reactor-events.js";
69
import {
710
createTurnContextCollector,
811
type LifecycleHookManager,
@@ -191,7 +194,7 @@ export function createRunSink(args: RunSinkArgs): RunSink {
191194
runError = undefined;
192195
onTurnBoundarySnapshot?.();
193196
}
194-
if (event.type === "reactor.error") {
197+
if (event.type === "reactor.error" && isReactorErrorFatal(event.data)) {
195198
const data = event.data as { error: string };
196199
runError = data.error;
197200
}

‎src/tui/stream-event-map.test.ts‎

Lines changed: 33 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -97,6 +97,39 @@ describe("mapProductionEvent", () => {
9797
]);
9898
});
9999

100+
test("non-fatal reactor.error paints the error without idling", () => {
101+
expect(
102+
mapProductionEvent({
103+
type: "reactor.error",
104+
data: { error: "transient checkpoint write", fatal: false },
105+
}),
106+
).toEqual([{ type: "error", message: "transient checkpoint write" }]);
107+
});
108+
109+
test("fatal reactor.error paints the error and idles", () => {
110+
expect(
111+
mapProductionEvent({
112+
type: "reactor.error",
113+
data: { error: "gave up", fatal: true },
114+
}),
115+
).toEqual([
116+
{ type: "error", message: "gave up" },
117+
{ type: "run", state: "idle" },
118+
]);
119+
});
120+
121+
test("reactor.error without fatal paints the error and idles", () => {
122+
expect(
123+
mapProductionEvent({
124+
type: "reactor.error",
125+
data: { error: "gave up" },
126+
}),
127+
).toEqual([
128+
{ type: "error", message: "gave up" },
129+
{ type: "run", state: "idle" },
130+
]);
131+
});
132+
100133
test("connector.reply after deltas is skipped (already painted)", () => {
101134
const ctx = createStreamMapContext();
102135
mapProductionEvent(

‎src/tui/stream-event-map.ts‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import {
99
splitPendingControlTail,
1010
stripTerminalControlSequences,
1111
} from "../util/control-char-strip.js";
12+
import { isReactorErrorFatal } from "../agent/reactor-events.js";
1213
import { terminalProviderFailureMessage } from "../inference-error-message.js";
1314
import {
1415
normalizeInferenceErrorForTerminal,
@@ -575,6 +576,9 @@ function mapEvent(
575576
const error =
576577
typeof data.error === "string" ? data.error : "reactor error";
577578
if (ctx) ctx.hadTextDelta = false;
579+
if (!isReactorErrorFatal(event.data)) {
580+
return [{ type: "error", message: error }];
581+
}
578582
return [
579583
...disarmAttempt(ctx),
580584
{ type: "error", message: error },

‎src/tui/turn-state.test.ts‎

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -508,4 +508,41 @@ describe("repetition tracking", () => {
508508
expect(restarted.repeatingSinceTokenCount).toBeNull();
509509
expect(restarted.streamText).toBe("");
510510
});
511+
512+
test("non-fatal reactor.error leaves the turn running", () => {
513+
const running = fold([
514+
{ type: "inference.start" },
515+
{ type: "inference.text.delta", data: { token: "hi" } },
516+
]);
517+
const continued = turnStateFromEvent(
518+
running,
519+
{
520+
type: "reactor.error",
521+
data: { error: "transient checkpoint write", fatal: false },
522+
},
523+
100,
524+
);
525+
expect(continued.status).toBe("running");
526+
expect(continued.isProcessing).toBe(true);
527+
});
528+
529+
test("fatal reactor.error fails the turn", () => {
530+
const running = fold([{ type: "inference.start" }]);
531+
const failed = turnStateFromEvent(
532+
running,
533+
{ type: "reactor.error", data: { error: "gave up", fatal: true } },
534+
100,
535+
);
536+
expect(failed.status).toBe("failed");
537+
});
538+
539+
test("reactor.error without fatal fails the turn", () => {
540+
const running = fold([{ type: "inference.start" }]);
541+
const failed = turnStateFromEvent(
542+
running,
543+
{ type: "reactor.error", data: { error: "gave up" } },
544+
100,
545+
);
546+
expect(failed.status).toBe("failed");
547+
});
511548
});

‎src/tui/turn-state.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,7 @@
1111

1212
import { type } from "arktype";
1313

14+
import { isReactorErrorFatal } from "../agent/reactor-events.js";
1415
import type { TurnStatus } from "./session-chrome.js";
1516

1617
// Bound on the accumulated stream text kept for the current cycle. Comfortably
@@ -687,6 +688,7 @@ export function turnStateFromEvent(
687688
});
688689

689690
case "reactor.error":
691+
if (!isReactorErrorFatal(event.data)) return state;
690692
return carryBlockedGateCount(state, {
691693
...initialTurnState(nowMs),
692694
status: "failed",

0 commit comments

Comments
 (0)