Skip to content

Commit 8a1dc44

Browse files
committed
fix(tui): abort-aware compact lifecycle, visible compaction state
Compact path aborts instead of parking the reactor; interrupt aborts in-flight compact with notice; back-to-back compactions complete unaided. CL-8220.
1 parent caabcf2 commit 8a1dc44

7 files changed

Lines changed: 654 additions & 17 deletions

File tree

Lines changed: 364 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,364 @@
1+
// CL-8220 regression tests: the compact path must always return to dequeue,
2+
// even when the summary call hangs, and back-to-back threshold compactions
3+
// must complete without operator action.
4+
import { describe, expect, test } from "bun:test";
5+
import type {
6+
Compactor,
7+
ConversationTurn,
8+
StrategyContext,
9+
StrategyResult,
10+
} from "@intx/types/runtime";
11+
12+
import {
13+
COMPACTION_ABORTED_REASON,
14+
createCompactionEventNotices,
15+
createCompactionLifecycle,
16+
} from "./compaction-lifecycle.js";
17+
import { createModelSummarizer } from "./summarizer.js";
18+
import type { Telemetry } from "../telemetry/index.js";
19+
20+
const ctx = {} as unknown as StrategyContext;
21+
22+
function turns(count: number): ConversationTurn[] {
23+
return Array.from({ length: count }, (_, i) => ({
24+
role: "user",
25+
content: [{ type: "text", text: `turn ${i}` }],
26+
timestamp: i,
27+
})) as ConversationTurn[];
28+
}
29+
30+
function okResult(
31+
output: ConversationTurn[],
32+
): StrategyResult<ConversationTurn[]> {
33+
return {
34+
output,
35+
record: {
36+
strategy: "inner",
37+
version: "0",
38+
parameters: {},
39+
reason: "folded",
40+
decisions: { summarizedTurnCount: 3 },
41+
},
42+
};
43+
}
44+
45+
function hangingCompactor(): Compactor {
46+
return {
47+
name: "hang",
48+
version: "0",
49+
apply: () =>
50+
new Promise<StrategyResult<ConversationTurn[]>>(() => {
51+
// Never settles on purpose: the abort race must win.
52+
}),
53+
};
54+
}
55+
56+
describe("createCompactionLifecycle", () => {
57+
test("a hanging compact resolves promptly on abort with input unchanged", async () => {
58+
const events: string[] = [];
59+
const tracked = createCompactionLifecycle({
60+
onCompactionStart: () => events.push("start"),
61+
onCompactionEnd: (info) => events.push(`end:${info.aborted}`),
62+
});
63+
const wrapped = tracked.wrapCompactor(hangingCompactor());
64+
const input = turns(10);
65+
const pending = wrapped.apply(input, ctx);
66+
expect(tracked.isCompacting()).toBe(true);
67+
tracked.abortCompaction("operator interrupt");
68+
const result = await pending;
69+
expect(tracked.isCompacting()).toBe(false);
70+
expect(result.output).toBe(input);
71+
expect(result.blobs ?? []).toEqual([]);
72+
expect(result.record.reason).toBe(COMPACTION_ABORTED_REASON);
73+
expect(events).toEqual(["start", "end:true"]);
74+
});
75+
76+
test("abort before start skips the inner run without lifecycle events", async () => {
77+
let calls = 0;
78+
const inner: Compactor = {
79+
name: "inner",
80+
version: "0",
81+
apply: async (input) => {
82+
calls += 1;
83+
return okResult(input);
84+
},
85+
};
86+
const events: string[] = [];
87+
const lifecycle = createCompactionLifecycle({
88+
onCompactionStart: () => events.push("start"),
89+
onCompactionEnd: () => events.push("end"),
90+
});
91+
lifecycle.abortCompaction("operator interrupt");
92+
const input = turns(4);
93+
const result = await lifecycle.wrapCompactor(inner).apply(input, ctx);
94+
expect(calls).toBe(0);
95+
expect(result.output).toBe(input);
96+
expect(result.record.reason).toBe(COMPACTION_ABORTED_REASON);
97+
expect(events).toEqual([]);
98+
expect(lifecycle.isCompacting()).toBe(false);
99+
});
100+
101+
test("a completing compact passes through with end-not-aborted", async () => {
102+
const inner: Compactor = {
103+
name: "inner",
104+
version: "1",
105+
apply: async (input) => okResult(input.slice(-2)),
106+
};
107+
const ends: boolean[] = [];
108+
const lifecycle = createCompactionLifecycle({
109+
onCompactionEnd: (info) => ends.push(info.aborted),
110+
});
111+
const wrapped = lifecycle.wrapCompactor(inner);
112+
expect(wrapped.name).toBe("inner");
113+
const result = await wrapped.apply(turns(8), ctx);
114+
expect(result.output).toHaveLength(2);
115+
expect(result.record.reason).toBe("folded");
116+
expect(ends).toEqual([false]);
117+
expect(lifecycle.isCompacting()).toBe(false);
118+
});
119+
120+
test("a genuine inner failure still propagates (never masked as abort)", async () => {
121+
const inner: Compactor = {
122+
name: "inner",
123+
version: "0",
124+
apply: async () => {
125+
throw new Error("store exploded");
126+
},
127+
};
128+
const ends: boolean[] = [];
129+
const lifecycle = createCompactionLifecycle({
130+
onCompactionEnd: (info) => ends.push(info.aborted),
131+
});
132+
await expect(
133+
lifecycle.wrapCompactor(inner).apply(turns(3), ctx),
134+
).rejects.toThrow("store exploded");
135+
expect(ends).toEqual([false]);
136+
expect(lifecycle.isCompacting()).toBe(false);
137+
});
138+
139+
test("back-to-back compactions complete and the loop keeps going", async () => {
140+
let calls = 0;
141+
const inner: Compactor = {
142+
name: "inner",
143+
version: "0",
144+
apply: async (input) => {
145+
calls += 1;
146+
return okResult(input);
147+
},
148+
};
149+
const lifecycle = createCompactionLifecycle();
150+
const wrapped = lifecycle.wrapCompactor(inner);
151+
const first = turns(12);
152+
const firstResult = await wrapped.apply(first, ctx);
153+
expect(lifecycle.isCompacting()).toBe(false);
154+
const secondResult = await wrapped.apply(firstResult.output, ctx);
155+
expect(lifecycle.isCompacting()).toBe(false);
156+
// A third pass (the post-compact inference turn's threshold re-check)
157+
// still runs: nothing wedged the loop.
158+
await wrapped.apply(secondResult.output, ctx);
159+
expect(calls).toBe(3);
160+
});
161+
162+
test("two compactions, then interrupt-during-compact, then resend works", async () => {
163+
let hang = false;
164+
let releaseHang: (() => void) | undefined;
165+
const inner: Compactor = {
166+
name: "inner",
167+
version: "0",
168+
apply: (input) =>
169+
hang
170+
? new Promise<StrategyResult<ConversationTurn[]>>((resolve) => {
171+
releaseHang = () => resolve(okResult(input));
172+
})
173+
: Promise.resolve(okResult(input)),
174+
};
175+
const lifecycle = createCompactionLifecycle();
176+
const wrapped = lifecycle.wrapCompactor(inner);
177+
// Two threshold compactions complete normally…
178+
await wrapped.apply(turns(10), ctx);
179+
await wrapped.apply(turns(10), ctx);
180+
// …then a third hangs and the operator interrupts mid-compact…
181+
hang = true;
182+
const input = turns(10);
183+
const pending = wrapped.apply(input, ctx);
184+
expect(lifecycle.isCompacting()).toBe(true);
185+
lifecycle.abortCompaction("operator interrupt");
186+
const interrupted = await pending;
187+
expect(interrupted.output).toBe(input);
188+
expect(interrupted.record.reason).toBe(COMPACTION_ABORTED_REASON);
189+
// …the rebuild mints a fresh signal and the resent turn compacts fine.
190+
lifecycle.reset();
191+
hang = false;
192+
const resent = await wrapped.apply(turns(10), ctx);
193+
expect(resent.record.reason).not.toBe(COMPACTION_ABORTED_REASON);
194+
expect(releaseHang).toBeDefined();
195+
});
196+
197+
test("event notices announce the pass and only speak up on abort", () => {
198+
const notices: string[] = [];
199+
const events = createCompactionEventNotices((text) => {
200+
notices.push(text);
201+
});
202+
events.onCompactionStart?.();
203+
expect(notices).toHaveLength(1);
204+
events.onCompactionEnd?.({ aborted: false });
205+
expect(notices).toHaveLength(1);
206+
events.onCompactionEnd?.({ aborted: true });
207+
expect(notices).toHaveLength(2);
208+
});
209+
210+
test("reset mints a fresh signal so the next agent is not pre-aborted", async () => {
211+
const lifecycle = createCompactionLifecycle();
212+
const before = lifecycle.getSignal();
213+
lifecycle.abortCompaction("operator interrupt");
214+
expect(before.aborted).toBe(true);
215+
lifecycle.reset();
216+
const after = lifecycle.getSignal();
217+
expect(after.aborted).toBe(false);
218+
expect(after).not.toBe(before);
219+
});
220+
221+
test("interrupt-during-compact emits one non-failure notice and no telemetry", async () => {
222+
const notices: string[] = [];
223+
const lifecycle = createCompactionLifecycle(
224+
createCompactionEventNotices((text) => {
225+
notices.push(text);
226+
}),
227+
);
228+
const telemetryEvents: string[] = [];
229+
const telemetry: Telemetry = {
230+
enabled: true,
231+
installationId: "test",
232+
capture: (event) => {
233+
telemetryEvents.push(event);
234+
},
235+
captureIntentional: () => false,
236+
flush: async () => undefined,
237+
discard: () => undefined,
238+
};
239+
const summarize = createModelSummarizer({
240+
getSource: () =>
241+
({
242+
id: "test",
243+
provider: "test",
244+
model: "test",
245+
credentialId: "test",
246+
}) as never,
247+
getSignal: () => lifecycle.getSignal(),
248+
telemetry,
249+
onFailure: (text) => {
250+
notices.push(text);
251+
},
252+
complete: (_promptTurns, _source, signal) =>
253+
new Promise<string>((_resolve, reject) => {
254+
signal.addEventListener(
255+
"abort",
256+
() => {
257+
const err = new Error("aborted by lifecycle");
258+
err.name = "AbortError";
259+
reject(err);
260+
},
261+
{ once: true },
262+
);
263+
}),
264+
});
265+
const inner: Compactor = {
266+
name: "inner",
267+
version: "0",
268+
apply: async (input) => {
269+
const text = await summarize(input);
270+
return okResult([
271+
...input.slice(-1),
272+
{
273+
...(input[0] as ConversationTurn),
274+
timestamp: -1,
275+
content: [{ type: "text", text }],
276+
} as ConversationTurn,
277+
]);
278+
},
279+
};
280+
const wrapped = lifecycle.wrapCompactor(inner);
281+
const input = turns(10);
282+
const pending = wrapped.apply(input, ctx);
283+
expect(lifecycle.isCompacting()).toBe(true);
284+
lifecycle.abortCompaction("operator interrupt");
285+
const result = await pending;
286+
expect(result.output).toBe(input);
287+
expect(result.record.reason).toBe(COMPACTION_ABORTED_REASON);
288+
// Exactly the lifecycle's own two notices: start + interrupted. The
289+
// summarizer's "Compaction summary failed … (aborted by lifecycle)"
290+
// failure framing must stay silent on a lifecycle abort.
291+
expect(notices).toEqual([
292+
"Compacting conversation context…",
293+
"Compaction interrupted — keeping prior context.",
294+
]);
295+
expect(notices.some((n) => n.includes("Compaction summary failed"))).toBe(
296+
false,
297+
);
298+
expect(telemetryEvents).toEqual([]);
299+
});
300+
301+
test("a genuine summarizer failure still notifies and emits telemetry", async () => {
302+
const notices: string[] = [];
303+
const captured: {
304+
event: string;
305+
properties?: Record<string, unknown> | undefined;
306+
}[] = [];
307+
const telemetry: Telemetry = {
308+
enabled: true,
309+
installationId: "test",
310+
capture: (event, properties) => {
311+
captured.push({ event, properties });
312+
},
313+
captureIntentional: () => false,
314+
flush: async () => undefined,
315+
discard: () => undefined,
316+
};
317+
const summarize = createModelSummarizer({
318+
getSource: () =>
319+
({
320+
id: "test",
321+
provider: "test",
322+
model: "test",
323+
credentialId: "test",
324+
}) as never,
325+
telemetry,
326+
onFailure: (text) => {
327+
notices.push(text);
328+
},
329+
complete: async () => {
330+
throw new Error("model unreachable");
331+
},
332+
});
333+
await expect(summarize(turns(5))).rejects.toThrow("model unreachable");
334+
expect(notices).toHaveLength(1);
335+
expect(notices[0]).toContain("Compaction summary failed");
336+
const failures = captured.filter((e) => e.event === "summarizer_failure");
337+
expect(failures).toHaveLength(1);
338+
expect(failures[0]?.properties?.["error_kind"]).toBe("failed");
339+
});
340+
341+
test("the summarizer honors the lifecycle signal via getSignal", async () => {
342+
const lifecycle = createCompactionLifecycle();
343+
const seen: AbortSignal[] = [];
344+
const summarize = createModelSummarizer({
345+
getSource: () =>
346+
({
347+
id: "test",
348+
provider: "test",
349+
model: "test",
350+
credentialId: "test",
351+
}) as never,
352+
getSignal: () => lifecycle.getSignal(),
353+
complete: async (_promptTurns, _source, signal) => {
354+
seen.push(signal);
355+
if (signal.aborted) throw new Error("aborted by lifecycle");
356+
return "summary";
357+
},
358+
});
359+
lifecycle.abortCompaction("operator interrupt");
360+
await expect(summarize(turns(5))).rejects.toThrow();
361+
expect(seen).toHaveLength(1);
362+
expect(seen[0]?.aborted).toBe(true);
363+
});
364+
});

0 commit comments

Comments
 (0)