Skip to content

Commit b9c4d28

Browse files
committed
Spill oversized fleet-dry reports to a retrievable URI
1 parent 2a8cf44 commit b9c4d28

3 files changed

Lines changed: 266 additions & 20 deletions

File tree

‎src/subagent/fleet-dry-drive.test.ts‎

Lines changed: 162 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -1,20 +1,59 @@
1+
import { createBlobReader } from "@intx/types/runtime";
12
import { describe, expect, test } from "bun:test";
3+
import type { Task } from "../agent/tasks.js";
24
import { createFleetMailbox } from "./agent-fleet.js";
35
import {
46
buildFleetDryContinuationPrompt,
57
collectUncollectedTerminals,
68
driveOpenTasksAfterFleetDry,
79
FLEET_DRY_CONTINUATION_PREFIX,
810
FLEET_DRY_REPORT_CHARS,
11+
fleetDrySpillKey,
912
shouldDriveOpenTasks,
1013
type FleetDryMailbox,
1114
type FleetDryMailboxRecord,
1215
} from "./fleet-dry-drive.js";
1316
import { createSubAgentSessionStore } from "./session-store.js";
14-
import type { Task } from "../agent/tasks.js";
1517

1618
const openTask: Task = { id: "t1", title: "keep going", status: "todo" };
1719

20+
function peekMailbox(
21+
records: Map<string, FleetDryMailboxRecord>,
22+
): FleetDryMailbox {
23+
return {
24+
ids: () => [...records.keys()],
25+
peek: (id) => records.get(id),
26+
take: (id) => records.get(id),
27+
};
28+
}
29+
30+
function fakeBlobStore() {
31+
const blobs = new Map<string, { bytes: Uint8Array; contentType: string }>();
32+
return {
33+
blobs,
34+
writeBlob: (key: string, bytes: Uint8Array, contentType: string) => {
35+
blobs.set(key, { bytes, contentType });
36+
},
37+
readBlob: async (key: string) => {
38+
const entry = blobs.get(key);
39+
if (entry === undefined) throw new Error(`Blob not found: ${key}`);
40+
return entry.bytes;
41+
},
42+
};
43+
}
44+
45+
function reportsJSONFromPrompt(prompt: string): unknown {
46+
const header =
47+
"Collected worker reports (already collected — do not call wait_agents for these agent_ids):\n";
48+
const start = prompt.indexOf(header);
49+
expect(start).toBeGreaterThanOrEqual(0);
50+
const jsonStart = start + header.length;
51+
const jsonEnd = prompt.indexOf("\n", jsonStart);
52+
return JSON.parse(
53+
prompt.slice(jsonStart, jsonEnd === -1 ? undefined : jsonEnd),
54+
);
55+
}
56+
1857
describe("shouldDriveOpenTasks", () => {
1958
test("is true only on wentDry && open tasks && !parentProcessing", () => {
2059
expect(
@@ -250,21 +289,94 @@ describe("collectUncollectedTerminals", () => {
250289
]);
251290
});
252291

253-
test("clips oversized reports", () => {
292+
test("clips oversized reports with an honest not-retrievable notice when no writer is provided", () => {
293+
const original = "x".repeat(FLEET_DRY_REPORT_CHARS + 40);
254294
const records = new Map<string, FleetDryMailboxRecord>([
255-
[
256-
"big",
257-
{ status: "done", report: "x".repeat(FLEET_DRY_REPORT_CHARS + 40) },
258-
],
295+
["big", { status: "done", report: original }],
259296
]);
260-
const mailbox: FleetDryMailbox = {
261-
ids: () => [...records.keys()],
262-
peek: (id) => records.get(id),
263-
take: (id) => records.get(id),
264-
};
265-
const reports = collectUncollectedTerminals(mailbox, [], true);
266-
expect(reports[0]?.report?.length).toBe(FLEET_DRY_REPORT_CHARS);
267-
expect(reports[0]?.report?.endsWith("…")).toBe(true);
297+
const reports = collectUncollectedTerminals(peekMailbox(records), [], true);
298+
const clipped = reports[0]?.report ?? "";
299+
expect(clipped.length).toBeLessThanOrEqual(FLEET_DRY_REPORT_CHARS);
300+
expect(clipped.length).toBeLessThan(original.length);
301+
expect(clipped).toContain("[output truncated");
302+
expect(clipped).toContain("NOT retrievable");
303+
expect(clipped).not.toContain("tool-output:///");
304+
expect(clipped.endsWith("…") && !clipped.includes("truncated")).toBe(false);
305+
});
306+
307+
test("spills oversized reports to a tool-output URI when a writer is provided", async () => {
308+
const original = `head-${"x".repeat(FLEET_DRY_REPORT_CHARS)}TAIL-MARKER`;
309+
const records = new Map<string, FleetDryMailboxRecord>([
310+
["big", { status: "done", report: original }],
311+
]);
312+
const store = fakeBlobStore();
313+
const reports = collectUncollectedTerminals(
314+
peekMailbox(records),
315+
[],
316+
true,
317+
store.writeBlob,
318+
);
319+
const clipped = reports[0]?.report ?? "";
320+
expect(clipped.length).toBeLessThanOrEqual(FLEET_DRY_REPORT_CHARS);
321+
expect(clipped).not.toContain("TAIL-MARKER");
322+
expect(clipped).toContain("[output truncated");
323+
expect(clipped).not.toContain("NOT retrievable");
324+
const key = fleetDrySpillKey("big", "report");
325+
const uri = `tool-output:///${key}`;
326+
expect(clipped).toContain(uri);
327+
expect(clipped).toContain("read_file");
328+
const entry = store.blobs.get(key);
329+
expect(entry).toBeDefined();
330+
expect(new TextDecoder().decode(entry?.bytes ?? new Uint8Array())).toBe(
331+
original,
332+
);
333+
334+
const recovered = new TextDecoder().decode(
335+
await createBlobReader(store).read(uri),
336+
);
337+
expect(recovered).toBe(original);
338+
339+
const parsed = reportsJSONFromPrompt(
340+
buildFleetDryContinuationPrompt([openTask], reports),
341+
);
342+
expect(parsed).toEqual(reports);
343+
});
344+
345+
test("leaves under-budget reports unchanged even when a writer is provided", () => {
346+
const report = "short enough";
347+
const records = new Map<string, FleetDryMailboxRecord>([
348+
["w1", { status: "done", report }],
349+
]);
350+
const store = fakeBlobStore();
351+
const reports = collectUncollectedTerminals(
352+
peekMailbox(records),
353+
[],
354+
true,
355+
store.writeBlob,
356+
);
357+
expect(reports[0]?.report).toBe(report);
358+
expect(store.blobs.size).toBe(0);
359+
});
360+
361+
test("spills oversized error fields under a distinct key", () => {
362+
const original = "e".repeat(FLEET_DRY_REPORT_CHARS + 20);
363+
const records = new Map<string, FleetDryMailboxRecord>([
364+
["boom", { status: "failed", error: original }],
365+
]);
366+
const store = fakeBlobStore();
367+
const reports = collectUncollectedTerminals(
368+
peekMailbox(records),
369+
[],
370+
true,
371+
store.writeBlob,
372+
);
373+
const clipped = reports[0]?.error ?? "";
374+
const key = fleetDrySpillKey("boom", "error");
375+
expect(clipped).toContain(`tool-output:///${key}`);
376+
expect(store.blobs.has(fleetDrySpillKey("boom", "report"))).toBe(false);
377+
expect(
378+
new TextDecoder().decode(store.blobs.get(key)?.bytes ?? new Uint8Array()),
379+
).toBe(original);
268380
});
269381

270382
test("consume false peeks without take", () => {
@@ -329,6 +441,42 @@ describe("driveOpenTasksAfterFleetDry", () => {
329441
expect(records.get("w1")?.collected).toBe(true);
330442
});
331443

444+
test("dry+open continuation JSON includes a spill URI for oversized reports", () => {
445+
const original = `head-${"x".repeat(FLEET_DRY_REPORT_CHARS)}TAIL-MARKER`;
446+
const records = new Map<string, FleetDryMailboxRecord>([
447+
["big", { status: "done", report: original }],
448+
]);
449+
const store = fakeBlobStore();
450+
const sent: string[] = [];
451+
const driven = driveOpenTasksAfterFleetDry({
452+
previousRunning: 1,
453+
running: 0,
454+
openTasks: [openTask],
455+
parentProcessing: false,
456+
mailbox: peekMailbox(records),
457+
lanes: [],
458+
writeBlob: store.writeBlob,
459+
beginSystemContinuation: (prompt) => {
460+
sent.push(prompt);
461+
},
462+
send: () => undefined,
463+
});
464+
expect(driven).toBe(true);
465+
const parsed = reportsJSONFromPrompt(sent[0] ?? "");
466+
expect(Array.isArray(parsed)).toBe(true);
467+
const report = (parsed as { report?: string }[])[0]?.report ?? "";
468+
expect(report).toContain(
469+
`tool-output:///${fleetDrySpillKey("big", "report")}`,
470+
);
471+
expect(report).not.toContain("TAIL-MARKER");
472+
expect(
473+
new TextDecoder().decode(
474+
store.blobs.get(fleetDrySpillKey("big", "report"))?.bytes ??
475+
new Uint8Array(),
476+
),
477+
).toBe(original);
478+
});
479+
332480
test("dry+terminal, live+open, and parentProcessing only skip", () => {
333481
const noop = {
334482
mailbox: undefined,

‎src/subagent/fleet-dry-drive.ts‎

Lines changed: 97 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -68,10 +68,83 @@ function isPromiseLike(value: unknown): value is Promise<unknown> {
6868
return typeof value === "object" && value !== null && "then" in value;
6969
}
7070

71-
function clipField(text: string | undefined): string | undefined {
71+
/** Session blob-store write; same shape as ContextStore.writeBlob. */
72+
export type FleetDryBlobWriter = (
73+
key: string,
74+
bytes: Uint8Array,
75+
contentType: string,
76+
) => void | Promise<void>;
77+
78+
/** Distinct from leisure `{callId}:full` so a reactor size-cap cannot clobber this spill. */
79+
export function fleetDrySpillKey(
80+
agentId: string,
81+
field: "report" | "error",
82+
): string {
83+
return `fleet-dry:${agentId}:${field}`;
84+
}
85+
86+
function truncationNotice(args: {
87+
maxChars: number;
88+
remaining: number;
89+
fullLength: number;
90+
uri?: string;
91+
}): string {
92+
const { maxChars, remaining, fullLength, uri } = args;
93+
if (uri === undefined) {
94+
return (
95+
`\n[output truncated at ${maxChars.toLocaleString()} chars — ` +
96+
`${remaining.toLocaleString()} chars discarded, NOT retrievable ` +
97+
`(no blob store is configured; re-running gives the same cut). ` +
98+
`Use offset/limit or a narrower query.]`
99+
);
100+
}
101+
return (
102+
`\n[output truncated at ${maxChars.toLocaleString()} chars — ` +
103+
`${remaining.toLocaleString()} more chars omitted here. The full result ` +
104+
`(${fullLength.toLocaleString()} chars, text/plain) is saved at ${uri}` +
105+
` — use read_file with that URI (offset/limit supported) to see the rest.]`
106+
);
107+
}
108+
109+
function truncateWithReservedNotice(
110+
text: string,
111+
maxChars: number,
112+
buildNotice: (keptLen: number) => string,
113+
): string {
114+
let keptLen = maxChars;
115+
for (let i = 0; i < 8; i++) {
116+
const notice = buildNotice(keptLen);
117+
const total = keptLen + notice.length;
118+
if (total <= maxChars) return text.slice(0, keptLen) + notice;
119+
keptLen -= total - maxChars;
120+
if (keptLen < 0) keptLen = 0;
121+
}
122+
const notice = buildNotice(keptLen);
123+
return (text.slice(0, keptLen) + notice).slice(0, maxChars);
124+
}
125+
126+
function clipField(
127+
text: string | undefined,
128+
agentId: string,
129+
field: "report" | "error",
130+
writeBlob?: FleetDryBlobWriter,
131+
): string | undefined {
72132
if (text === undefined) return undefined;
73133
if (text.length <= FLEET_DRY_REPORT_CHARS) return text;
74-
return `${text.slice(0, FLEET_DRY_REPORT_CHARS - 1).trimEnd()}…`;
134+
let uri: string | undefined;
135+
if (writeBlob !== undefined) {
136+
const key = fleetDrySpillKey(agentId, field);
137+
uri = `tool-output:///${key}`;
138+
void writeBlob(key, new TextEncoder().encode(text), "text/plain");
139+
}
140+
return truncateWithReservedNotice(text, FLEET_DRY_REPORT_CHARS, (keptLen) =>
141+
truncationNotice({
142+
maxChars: FLEET_DRY_REPORT_CHARS,
143+
remaining: text.length - keptLen,
144+
fullLength: text.length,
145+
...(uri !== undefined ? { uri } : {}),
146+
}),
147+
);
75148
}
76149

77150
export function projectMailboxRecord(
@@ -117,9 +190,20 @@ export function takeAndProjectMailboxRecord(
117190

118191
function clipCollectedReport(
119192
report: CollectedWorkerReport,
193+
writeBlob?: FleetDryBlobWriter,
120194
): CollectedWorkerReport {
121-
const clippedReport = clipField(report.report);
122-
const clippedError = clipField(report.error);
195+
const clippedReport = clipField(
196+
report.report,
197+
report.agent_id,
198+
"report",
199+
writeBlob,
200+
);
201+
const clippedError = clipField(
202+
report.error,
203+
report.agent_id,
204+
"error",
205+
writeBlob,
206+
);
123207
return {
124208
...report,
125209
...(clippedReport !== undefined ? { report: clippedReport } : {}),
@@ -131,6 +215,7 @@ export function collectUncollectedTerminals(
131215
mailbox: FleetDryMailbox | undefined,
132216
lanes: readonly FleetDryLane[],
133217
consume: boolean,
218+
writeBlob?: FleetDryBlobWriter,
134219
): CollectedWorkerReport[] {
135220
if (mailbox === undefined) return [];
136221
const byId = new Map(lanes.map((lane) => [lane.id, lane]));
@@ -144,7 +229,7 @@ export function collectUncollectedTerminals(
144229
? takeAndProjectMailboxRecord(mailbox, id, byId.get(id))
145230
: projectMailboxRecord(id, peeked, byId.get(id));
146231
if (projected === undefined) continue;
147-
reports.push(clipCollectedReport(projected));
232+
reports.push(clipCollectedReport(projected, writeBlob));
148233
}
149234
return reports;
150235
}
@@ -180,6 +265,7 @@ export function driveOpenTasksAfterFleetDry(args: {
180265
deferredDryEdge?: boolean;
181266
mailbox: FleetDryMailbox | undefined;
182267
lanes: readonly FleetDryLane[];
268+
writeBlob?: FleetDryBlobWriter;
183269
beginSystemContinuation: (prompt: string) => void;
184270
send: (prompt: string) => unknown;
185271
onSendFailure?: () => void;
@@ -196,7 +282,12 @@ export function driveOpenTasksAfterFleetDry(args: {
196282
) {
197283
return false;
198284
}
199-
const reports = collectUncollectedTerminals(args.mailbox, args.lanes, false);
285+
const reports = collectUncollectedTerminals(
286+
args.mailbox,
287+
args.lanes,
288+
false,
289+
args.writeBlob,
290+
);
200291
const prompt = buildFleetDryContinuationPrompt(tasks, reports);
201292
const takeReports = (): void => {
202293
for (const report of reports) {

‎src/tui/runner/wiring.ts‎

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -195,12 +195,19 @@ export function wirePostStartup(
195195
sessionBridge.setDryOpenTaskDriver(() => {
196196
const send = state.sendWithAttemptIdentity;
197197
if (send === undefined) return false;
198+
const storage = state.currentStorage;
198199
return driveOpenTasksAfterFleetDry({
199200
deferredDryEdge: true,
200201
openTasks: services.directorHolder.instance?.getTasks() ?? [],
201202
parentProcessing: false,
202203
mailbox: services.toolset.fleetRecords,
203204
lanes: services.subAgentSessions.list(),
205+
...(storage !== null
206+
? {
207+
writeBlob: (key, bytes, contentType) =>
208+
storage.writeBlob(key, bytes, contentType),
209+
}
210+
: {}),
204211
beginSystemContinuation: (prompt) => {
205212
sessionBridge.beginSystemContinuation(prompt);
206213
},

0 commit comments

Comments
 (0)