From 8ee6b9add3f1f5714607ec0e1f11e7984ff4d7d4 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Thu, 10 Sep 2026 01:07:08 -0700 Subject: [PATCH 1/4] Spill oversized fleet-dry reports to a retrievable URI --- src/subagent/fleet-dry-drive.test.ts | 176 ++++++++++++++++++++++++--- src/subagent/fleet-dry-drive.ts | 103 +++++++++++++++- src/tui/runner/wiring.ts | 7 ++ 3 files changed, 266 insertions(+), 20 deletions(-) diff --git a/src/subagent/fleet-dry-drive.test.ts b/src/subagent/fleet-dry-drive.test.ts index e85099db1..c73020b16 100644 --- a/src/subagent/fleet-dry-drive.test.ts +++ b/src/subagent/fleet-dry-drive.test.ts @@ -1,4 +1,6 @@ +import { createBlobReader } from "@intx/types/runtime"; import { describe, expect, test } from "bun:test"; +import type { Task } from "../agent/tasks.js"; import { createFleetMailbox } from "./agent-fleet.js"; import { buildFleetDryContinuationPrompt, @@ -6,15 +8,52 @@ import { driveOpenTasksAfterFleetDry, FLEET_DRY_CONTINUATION_PREFIX, FLEET_DRY_REPORT_CHARS, + fleetDrySpillKey, shouldDriveOpenTasks, type FleetDryMailbox, type FleetDryMailboxRecord, } from "./fleet-dry-drive.js"; import { createSubAgentSessionStore } from "./session-store.js"; -import type { Task } from "../agent/tasks.js"; const openTask: Task = { id: "t1", title: "keep going", status: "todo" }; +function peekMailbox( + records: Map, +): FleetDryMailbox { + return { + ids: () => [...records.keys()], + peek: (id) => records.get(id), + take: (id) => records.get(id), + }; +} + +function fakeBlobStore() { + const blobs = new Map(); + return { + blobs, + writeBlob: (key: string, bytes: Uint8Array, contentType: string) => { + blobs.set(key, { bytes, contentType }); + }, + readBlob: async (key: string) => { + const entry = blobs.get(key); + if (entry === undefined) throw new Error(`Blob not found: ${key}`); + return entry.bytes; + }, + }; +} + +function reportsJSONFromPrompt(prompt: string): unknown { + const header = + "Collected worker reports (already collected — do not call wait_agents for these agent_ids):\n"; + const start = prompt.indexOf(header); + expect(start).toBeGreaterThanOrEqual(0); + const jsonStart = start + header.length; + const jsonEnd = prompt.indexOf("\n", jsonStart); + return JSON.parse( + prompt.slice(jsonStart, jsonEnd === -1 ? undefined : jsonEnd), + ); +} + describe("shouldDriveOpenTasks", () => { test("is true only on wentDry && open tasks && !parentProcessing", () => { expect( @@ -250,21 +289,94 @@ describe("collectUncollectedTerminals", () => { ]); }); - test("clips oversized reports", () => { + test("clips oversized reports with an honest not-retrievable notice when no writer is provided", () => { + const original = "x".repeat(FLEET_DRY_REPORT_CHARS + 40); const records = new Map([ - [ - "big", - { status: "done", report: "x".repeat(FLEET_DRY_REPORT_CHARS + 40) }, - ], + ["big", { status: "done", report: original }], ]); - const mailbox: FleetDryMailbox = { - ids: () => [...records.keys()], - peek: (id) => records.get(id), - take: (id) => records.get(id), - }; - const reports = collectUncollectedTerminals(mailbox, [], true); - expect(reports[0]?.report?.length).toBe(FLEET_DRY_REPORT_CHARS); - expect(reports[0]?.report?.endsWith("…")).toBe(true); + const reports = collectUncollectedTerminals(peekMailbox(records), [], true); + const clipped = reports[0]?.report ?? ""; + expect(clipped.length).toBeLessThanOrEqual(FLEET_DRY_REPORT_CHARS); + expect(clipped.length).toBeLessThan(original.length); + expect(clipped).toContain("[output truncated"); + expect(clipped).toContain("NOT retrievable"); + expect(clipped).not.toContain("tool-output:///"); + expect(clipped.endsWith("…") && !clipped.includes("truncated")).toBe(false); + }); + + test("spills oversized reports to a tool-output URI when a writer is provided", async () => { + const original = `head-${"x".repeat(FLEET_DRY_REPORT_CHARS)}TAIL-MARKER`; + const records = new Map([ + ["big", { status: "done", report: original }], + ]); + const store = fakeBlobStore(); + const reports = collectUncollectedTerminals( + peekMailbox(records), + [], + true, + store.writeBlob, + ); + const clipped = reports[0]?.report ?? ""; + expect(clipped.length).toBeLessThanOrEqual(FLEET_DRY_REPORT_CHARS); + expect(clipped).not.toContain("TAIL-MARKER"); + expect(clipped).toContain("[output truncated"); + expect(clipped).not.toContain("NOT retrievable"); + const key = fleetDrySpillKey("big", "report"); + const uri = `tool-output:///${key}`; + expect(clipped).toContain(uri); + expect(clipped).toContain("read_file"); + const entry = store.blobs.get(key); + expect(entry).toBeDefined(); + expect(new TextDecoder().decode(entry?.bytes ?? new Uint8Array())).toBe( + original, + ); + + const recovered = new TextDecoder().decode( + await createBlobReader(store).read(uri), + ); + expect(recovered).toBe(original); + + const parsed = reportsJSONFromPrompt( + buildFleetDryContinuationPrompt([openTask], reports), + ); + expect(parsed).toEqual(reports); + }); + + test("leaves under-budget reports unchanged even when a writer is provided", () => { + const report = "short enough"; + const records = new Map([ + ["w1", { status: "done", report }], + ]); + const store = fakeBlobStore(); + const reports = collectUncollectedTerminals( + peekMailbox(records), + [], + true, + store.writeBlob, + ); + expect(reports[0]?.report).toBe(report); + expect(store.blobs.size).toBe(0); + }); + + test("spills oversized error fields under a distinct key", () => { + const original = "e".repeat(FLEET_DRY_REPORT_CHARS + 20); + const records = new Map([ + ["boom", { status: "failed", error: original }], + ]); + const store = fakeBlobStore(); + const reports = collectUncollectedTerminals( + peekMailbox(records), + [], + true, + store.writeBlob, + ); + const clipped = reports[0]?.error ?? ""; + const key = fleetDrySpillKey("boom", "error"); + expect(clipped).toContain(`tool-output:///${key}`); + expect(store.blobs.has(fleetDrySpillKey("boom", "report"))).toBe(false); + expect( + new TextDecoder().decode(store.blobs.get(key)?.bytes ?? new Uint8Array()), + ).toBe(original); }); test("consume false peeks without take", () => { @@ -329,6 +441,42 @@ describe("driveOpenTasksAfterFleetDry", () => { expect(records.get("w1")?.collected).toBe(true); }); + test("dry+open continuation JSON includes a spill URI for oversized reports", () => { + const original = `head-${"x".repeat(FLEET_DRY_REPORT_CHARS)}TAIL-MARKER`; + const records = new Map([ + ["big", { status: "done", report: original }], + ]); + const store = fakeBlobStore(); + const sent: string[] = []; + const driven = driveOpenTasksAfterFleetDry({ + previousRunning: 1, + running: 0, + openTasks: [openTask], + parentProcessing: false, + mailbox: peekMailbox(records), + lanes: [], + writeBlob: store.writeBlob, + beginSystemContinuation: (prompt) => { + sent.push(prompt); + }, + send: () => undefined, + }); + expect(driven).toBe(true); + const parsed = reportsJSONFromPrompt(sent[0] ?? ""); + expect(Array.isArray(parsed)).toBe(true); + const report = (parsed as { report?: string }[])[0]?.report ?? ""; + expect(report).toContain( + `tool-output:///${fleetDrySpillKey("big", "report")}`, + ); + expect(report).not.toContain("TAIL-MARKER"); + expect( + new TextDecoder().decode( + store.blobs.get(fleetDrySpillKey("big", "report"))?.bytes ?? + new Uint8Array(), + ), + ).toBe(original); + }); + test("dry+terminal, live+open, and parentProcessing only skip", () => { const noop = { mailbox: undefined, diff --git a/src/subagent/fleet-dry-drive.ts b/src/subagent/fleet-dry-drive.ts index b02c1979d..a6fe230a2 100644 --- a/src/subagent/fleet-dry-drive.ts +++ b/src/subagent/fleet-dry-drive.ts @@ -68,10 +68,83 @@ function isPromiseLike(value: unknown): value is Promise { return typeof value === "object" && value !== null && "then" in value; } -function clipField(text: string | undefined): string | undefined { +/** Session blob-store write; same shape as ContextStore.writeBlob. */ +export type FleetDryBlobWriter = ( + key: string, + bytes: Uint8Array, + contentType: string, +) => void | Promise; + +/** Distinct from leisure `{callId}:full` so a reactor size-cap cannot clobber this spill. */ +export function fleetDrySpillKey( + agentId: string, + field: "report" | "error", +): string { + return `fleet-dry:${agentId}:${field}`; +} + +function truncationNotice(args: { + maxChars: number; + remaining: number; + fullLength: number; + uri?: string; +}): string { + const { maxChars, remaining, fullLength, uri } = args; + if (uri === undefined) { + return ( + `\n[output truncated at ${maxChars.toLocaleString()} chars — ` + + `${remaining.toLocaleString()} chars discarded, NOT retrievable ` + + `(no blob store is configured; re-running gives the same cut). ` + + `Use offset/limit or a narrower query.]` + ); + } + return ( + `\n[output truncated at ${maxChars.toLocaleString()} chars — ` + + `${remaining.toLocaleString()} more chars omitted here. The full result ` + + `(${fullLength.toLocaleString()} chars, text/plain) is saved at ${uri}` + + ` — use read_file with that URI (offset/limit supported) to see the rest.]` + ); +} + +function truncateWithReservedNotice( + text: string, + maxChars: number, + buildNotice: (keptLen: number) => string, +): string { + let keptLen = maxChars; + for (let i = 0; i < 8; i++) { + const notice = buildNotice(keptLen); + const total = keptLen + notice.length; + if (total <= maxChars) return text.slice(0, keptLen) + notice; + keptLen -= total - maxChars; + if (keptLen < 0) keptLen = 0; + } + const notice = buildNotice(keptLen); + return (text.slice(0, keptLen) + notice).slice(0, maxChars); +} + +function clipField( + text: string | undefined, + agentId: string, + field: "report" | "error", + writeBlob?: FleetDryBlobWriter, +): string | undefined { if (text === undefined) return undefined; if (text.length <= FLEET_DRY_REPORT_CHARS) return text; - return `${text.slice(0, FLEET_DRY_REPORT_CHARS - 1).trimEnd()}…`; + let uri: string | undefined; + if (writeBlob !== undefined) { + const key = fleetDrySpillKey(agentId, field); + uri = `tool-output:///${key}`; + void writeBlob(key, new TextEncoder().encode(text), "text/plain"); + } + return truncateWithReservedNotice(text, FLEET_DRY_REPORT_CHARS, (keptLen) => + truncationNotice({ + maxChars: FLEET_DRY_REPORT_CHARS, + remaining: text.length - keptLen, + fullLength: text.length, + ...(uri !== undefined ? { uri } : {}), + }), + ); } export function projectMailboxRecord( @@ -117,9 +190,20 @@ export function takeAndProjectMailboxRecord( function clipCollectedReport( report: CollectedWorkerReport, + writeBlob?: FleetDryBlobWriter, ): CollectedWorkerReport { - const clippedReport = clipField(report.report); - const clippedError = clipField(report.error); + const clippedReport = clipField( + report.report, + report.agent_id, + "report", + writeBlob, + ); + const clippedError = clipField( + report.error, + report.agent_id, + "error", + writeBlob, + ); return { ...report, ...(clippedReport !== undefined ? { report: clippedReport } : {}), @@ -131,6 +215,7 @@ export function collectUncollectedTerminals( mailbox: FleetDryMailbox | undefined, lanes: readonly FleetDryLane[], consume: boolean, + writeBlob?: FleetDryBlobWriter, ): CollectedWorkerReport[] { if (mailbox === undefined) return []; const byId = new Map(lanes.map((lane) => [lane.id, lane])); @@ -144,7 +229,7 @@ export function collectUncollectedTerminals( ? takeAndProjectMailboxRecord(mailbox, id, byId.get(id)) : projectMailboxRecord(id, peeked, byId.get(id)); if (projected === undefined) continue; - reports.push(clipCollectedReport(projected)); + reports.push(clipCollectedReport(projected, writeBlob)); } return reports; } @@ -180,6 +265,7 @@ export function driveOpenTasksAfterFleetDry(args: { deferredDryEdge?: boolean; mailbox: FleetDryMailbox | undefined; lanes: readonly FleetDryLane[]; + writeBlob?: FleetDryBlobWriter; beginSystemContinuation: (prompt: string) => void; send: (prompt: string) => unknown; onSendFailure?: () => void; @@ -196,7 +282,12 @@ export function driveOpenTasksAfterFleetDry(args: { ) { return false; } - const reports = collectUncollectedTerminals(args.mailbox, args.lanes, false); + const reports = collectUncollectedTerminals( + args.mailbox, + args.lanes, + false, + args.writeBlob, + ); const prompt = buildFleetDryContinuationPrompt(tasks, reports); const takeReports = (): void => { for (const report of reports) { diff --git a/src/tui/runner/wiring.ts b/src/tui/runner/wiring.ts index def5e8e03..83697ac17 100644 --- a/src/tui/runner/wiring.ts +++ b/src/tui/runner/wiring.ts @@ -200,12 +200,19 @@ export function wirePostStartup( sessionBridge.setDryOpenTaskDriver(() => { const send = state.sendWithAttemptIdentity; if (send === undefined) return false; + const storage = state.currentStorage; return driveOpenTasksAfterFleetDry({ deferredDryEdge: true, openTasks: services.directorHolder.instance?.getTasks() ?? [], parentProcessing: false, mailbox: services.toolset.fleetRecords, lanes: services.subAgentSessions.list(), + ...(storage !== null + ? { + writeBlob: (key, bytes, contentType) => + storage.writeBlob(key, bytes, contentType), + } + : {}), beginSystemContinuation: (prompt) => { sessionBridge.beginSystemContinuation(prompt); }, From 930ec12a7e9eafc9ad0ea0f2ddc58de339219553 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Thu, 10 Sep 2026 02:04:28 -0700 Subject: [PATCH 2/4] Await fleet-dry spill writes before naming the URI --- src/subagent/fleet-dry-drive.test.ts | 224 +++++++++++++++++++++++---- src/subagent/fleet-dry-drive.ts | 37 +++-- src/tui/runner/wiring.ts | 5 +- 3 files changed, 220 insertions(+), 46 deletions(-) diff --git a/src/subagent/fleet-dry-drive.test.ts b/src/subagent/fleet-dry-drive.test.ts index c73020b16..b2b43c344 100644 --- a/src/subagent/fleet-dry-drive.test.ts +++ b/src/subagent/fleet-dry-drive.test.ts @@ -180,7 +180,7 @@ describe("buildFleetDryContinuationPrompt", () => { }); describe("collectUncollectedTerminals", () => { - test("take()s terminals and leaves live / awaiting_director / already-collected", () => { + test("take()s terminals and leaves live / awaiting_director / already-collected", async () => { const sessions = createSubAgentSessionStore(); const mailbox = createFleetMailbox(sessions); const start = (id: string, description: string) => { @@ -213,7 +213,11 @@ describe("collectUncollectedTerminals", () => { }), ).toBe(true); - const reports = collectUncollectedTerminals(mailbox, sessions.list(), true); + const reports = await collectUncollectedTerminals( + mailbox, + sessions.list(), + true, + ); expect(reports.map((r) => r.agent_id).sort()).toEqual(["done", "fail"]); expect(reports.find((r) => r.agent_id === "done")).toEqual({ agent_id: "done", @@ -236,7 +240,7 @@ describe("collectUncollectedTerminals", () => { expect(mailbox.peek("coll")?.collected).toBe(true); }); - test("fills report/error from the session-store lane when the mailbox snapshot is empty", () => { + test("fills report/error from the session-store lane when the mailbox snapshot is empty", async () => { const records = new Map([ ["ghost", { status: "done" }], ]); @@ -251,7 +255,7 @@ describe("collectUncollectedTerminals", () => { return taken; }, }; - const reports = collectUncollectedTerminals( + const reports = await collectUncollectedTerminals( mailbox, [{ id: "ghost", description: "from store", report: "store report" }], true, @@ -267,7 +271,7 @@ describe("collectUncollectedTerminals", () => { expect(records.get("ghost")?.collected).toBe(true); }); - test("projects mailbox stopReason as stop_reason", () => { + test("projects mailbox stopReason as stop_reason", async () => { const records = new Map([ [ "w1", @@ -279,7 +283,7 @@ describe("collectUncollectedTerminals", () => { peek: (id) => records.get(id), take: (id) => records.get(id), }; - expect(collectUncollectedTerminals(mailbox, [], true)).toEqual([ + expect(await collectUncollectedTerminals(mailbox, [], true)).toEqual([ { agent_id: "w1", status: "interrupted", @@ -289,12 +293,16 @@ describe("collectUncollectedTerminals", () => { ]); }); - test("clips oversized reports with an honest not-retrievable notice when no writer is provided", () => { + test("clips oversized reports with an honest not-retrievable notice when no writer is provided", async () => { const original = "x".repeat(FLEET_DRY_REPORT_CHARS + 40); const records = new Map([ ["big", { status: "done", report: original }], ]); - const reports = collectUncollectedTerminals(peekMailbox(records), [], true); + const reports = await collectUncollectedTerminals( + peekMailbox(records), + [], + true, + ); const clipped = reports[0]?.report ?? ""; expect(clipped.length).toBeLessThanOrEqual(FLEET_DRY_REPORT_CHARS); expect(clipped.length).toBeLessThan(original.length); @@ -310,7 +318,7 @@ describe("collectUncollectedTerminals", () => { ["big", { status: "done", report: original }], ]); const store = fakeBlobStore(); - const reports = collectUncollectedTerminals( + const reports = await collectUncollectedTerminals( peekMailbox(records), [], true, @@ -342,13 +350,13 @@ describe("collectUncollectedTerminals", () => { expect(parsed).toEqual(reports); }); - test("leaves under-budget reports unchanged even when a writer is provided", () => { + test("leaves under-budget reports unchanged even when a writer is provided", async () => { const report = "short enough"; const records = new Map([ ["w1", { status: "done", report }], ]); const store = fakeBlobStore(); - const reports = collectUncollectedTerminals( + const reports = await collectUncollectedTerminals( peekMailbox(records), [], true, @@ -358,13 +366,13 @@ describe("collectUncollectedTerminals", () => { expect(store.blobs.size).toBe(0); }); - test("spills oversized error fields under a distinct key", () => { + test("spills oversized error fields under a distinct key", async () => { const original = "e".repeat(FLEET_DRY_REPORT_CHARS + 20); const records = new Map([ ["boom", { status: "failed", error: original }], ]); const store = fakeBlobStore(); - const reports = collectUncollectedTerminals( + const reports = await collectUncollectedTerminals( peekMailbox(records), [], true, @@ -379,7 +387,7 @@ describe("collectUncollectedTerminals", () => { ).toBe(original); }); - test("consume false peeks without take", () => { + test("consume false peeks without take", async () => { const records = new Map([ ["w1", { status: "done", report: "ok" }], ]); @@ -394,14 +402,86 @@ describe("collectUncollectedTerminals", () => { return taken; }, }; - const reports = collectUncollectedTerminals(mailbox, [], false); + const reports = await collectUncollectedTerminals(mailbox, [], false); expect(reports).toEqual([{ agent_id: "w1", status: "done", report: "ok" }]); expect(records.get("w1")?.collected).not.toBe(true); }); + + test("names a tool-output URI only after an async writeBlob resolves", async () => { + const original = `head-${"x".repeat(FLEET_DRY_REPORT_CHARS)}TAIL-MARKER`; + const records = new Map([ + ["big", { status: "done", report: original }], + ]); + const store = fakeBlobStore(); + let resolveWrite: (() => void) | undefined; + const writeBlob = (key: string, bytes: Uint8Array, contentType: string) => + new Promise((resolve) => { + resolveWrite = () => { + store.writeBlob(key, bytes, contentType); + resolve(); + }; + }); + const reportsP = collectUncollectedTerminals( + peekMailbox(records), + [], + true, + writeBlob, + ); + let settled = false; + void reportsP.then(() => { + settled = true; + }); + await Promise.resolve(); + expect(settled).toBe(false); + expect(store.blobs.size).toBe(0); + resolveWrite?.(); + const reports = await reportsP; + expect(settled).toBe(true); + const clipped = reports[0]?.report ?? ""; + const key = fleetDrySpillKey("big", "report"); + const uri = `tool-output:///${key}`; + expect(clipped).toContain(uri); + expect(clipped).not.toContain("NOT retrievable"); + const recovered = new TextDecoder().decode( + await createBlobReader(store).read(uri), + ); + expect(recovered).toBe(original); + }); + + test("rejected writeBlob yields NOT retrievable with no URI and no unhandled rejection", async () => { + const original = "x".repeat(FLEET_DRY_REPORT_CHARS + 40); + const records = new Map([ + ["big", { status: "done", report: original }], + ]); + const rejections: unknown[] = []; + const onUnhandled = (reason: unknown) => { + rejections.push(reason); + }; + process.on("unhandledRejection", onUnhandled); + try { + const reports = await collectUncollectedTerminals( + peekMailbox(records), + [], + true, + async () => { + throw new Error("disk full"); + }, + ); + await Promise.resolve(); + const clipped = reports[0]?.report ?? ""; + expect(clipped.length).toBeLessThanOrEqual(FLEET_DRY_REPORT_CHARS); + expect(clipped).toContain("[output truncated"); + expect(clipped).toContain("NOT retrievable"); + expect(clipped).not.toContain("tool-output:///"); + expect(rejections).toEqual([]); + } finally { + process.off("unhandledRejection", onUnhandled); + } + }); }); describe("driveOpenTasksAfterFleetDry", () => { - test("dry+open collects, begins continuation, then sends", () => { + test("dry+open collects, begins continuation, then sends", async () => { const order: string[] = []; const records = new Map([ ["w1", { status: "done", report: "ok", description: "lane" }], @@ -418,7 +498,7 @@ describe("driveOpenTasksAfterFleetDry", () => { }, }; const sent: string[] = []; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ previousRunning: 1, running: 0, openTasks: [openTask], @@ -441,14 +521,14 @@ describe("driveOpenTasksAfterFleetDry", () => { expect(records.get("w1")?.collected).toBe(true); }); - test("dry+open continuation JSON includes a spill URI for oversized reports", () => { + test("dry+open continuation JSON includes a spill URI for oversized reports", async () => { const original = `head-${"x".repeat(FLEET_DRY_REPORT_CHARS)}TAIL-MARKER`; const records = new Map([ ["big", { status: "done", report: original }], ]); const store = fakeBlobStore(); const sent: string[] = []; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ previousRunning: 1, running: 0, openTasks: [openTask], @@ -477,6 +557,86 @@ describe("driveOpenTasksAfterFleetDry", () => { ).toBe(original); }); + test("send waits for async writeBlob before naming the URI", async () => { + const original = `head-${"x".repeat(FLEET_DRY_REPORT_CHARS)}TAIL-MARKER`; + const records = new Map([ + ["big", { status: "done", report: original }], + ]); + const store = fakeBlobStore(); + let resolveWrite: (() => void) | undefined; + const writeBlob = (key: string, bytes: Uint8Array, contentType: string) => + new Promise((resolve) => { + resolveWrite = () => { + store.writeBlob(key, bytes, contentType); + resolve(); + }; + }); + const sent: string[] = []; + const drivenP = driveOpenTasksAfterFleetDry({ + previousRunning: 1, + running: 0, + openTasks: [openTask], + parentProcessing: false, + mailbox: peekMailbox(records), + lanes: [], + writeBlob, + beginSystemContinuation: (prompt) => { + sent.push(prompt); + }, + send: () => undefined, + }); + await Promise.resolve(); + expect(sent).toEqual([]); + resolveWrite?.(); + expect(await drivenP).toBe(true); + const uri = `tool-output:///${fleetDrySpillKey("big", "report")}`; + expect(sent[0]).toContain(uri); + expect(sent[0]).not.toContain("TAIL-MARKER"); + const recovered = new TextDecoder().decode( + await createBlobReader(store).read(uri), + ); + expect(recovered).toBe(original); + }); + + test("rejected writeBlob send is NOT retrievable and names no URI", async () => { + const original = "x".repeat(FLEET_DRY_REPORT_CHARS + 40); + const records = new Map([ + ["big", { status: "done", report: original }], + ]); + const rejections: unknown[] = []; + const onUnhandled = (reason: unknown) => { + rejections.push(reason); + }; + process.on("unhandledRejection", onUnhandled); + try { + const sent: string[] = []; + const driven = await driveOpenTasksAfterFleetDry({ + previousRunning: 1, + running: 0, + openTasks: [openTask], + parentProcessing: false, + mailbox: peekMailbox(records), + lanes: [], + writeBlob: async () => { + throw new Error("disk full"); + }, + beginSystemContinuation: (prompt) => { + sent.push(prompt); + }, + send: () => undefined, + }); + await Promise.resolve(); + expect(driven).toBe(true); + const parsed = reportsJSONFromPrompt(sent[0] ?? ""); + const report = (parsed as { report?: string }[])[0]?.report ?? ""; + expect(report).toContain("NOT retrievable"); + expect(report).not.toContain("tool-output:///"); + expect(rejections).toEqual([]); + } finally { + process.off("unhandledRejection", onUnhandled); + } + }); + test("dry+terminal, live+open, and parentProcessing only skip", () => { const noop = { mailbox: undefined, @@ -517,7 +677,7 @@ describe("driveOpenTasksAfterFleetDry", () => { ).toBe(false); }); - test("deferred dry edge after parentProcessing still collects and sends", () => { + test("deferred dry edge after parentProcessing still collects and sends", async () => { const records = new Map([ ["w1", { status: "done", report: "ok" }], ]); @@ -533,7 +693,7 @@ describe("driveOpenTasksAfterFleetDry", () => { }, }; const sent: string[] = []; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ previousRunning: 0, running: 0, deferredDryEdge: true, @@ -551,7 +711,7 @@ describe("driveOpenTasksAfterFleetDry", () => { expect(records.get("w1")?.collected).toBe(true); }); - test("send failure after take leaves reports waitable", () => { + test("send failure after take leaves reports waitable", async () => { const records = new Map([ ["w1", { status: "done", report: "ok" }], ]); @@ -566,7 +726,7 @@ describe("driveOpenTasksAfterFleetDry", () => { return taken; }, }; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ previousRunning: 1, running: 0, openTasks: [openTask], @@ -601,7 +761,7 @@ describe("driveOpenTasksAfterFleetDry", () => { await Promise.resolve(); throw new Error("agentProxy.send failed"); }; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ deferredDryEdge: true, openTasks: [openTask], parentProcessing: false, @@ -636,7 +796,7 @@ describe("driveOpenTasksAfterFleetDry", () => { await Promise.resolve(); return false; }; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ deferredDryEdge: true, openTasks: [openTask], parentProcessing: false, @@ -671,7 +831,7 @@ describe("driveOpenTasksAfterFleetDry", () => { new Promise((resolve) => { resolveSend = resolve; }); - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ deferredDryEdge: true, openTasks: [openTask], parentProcessing: false, @@ -687,7 +847,7 @@ describe("driveOpenTasksAfterFleetDry", () => { expect(records.get("w1")?.collected).toBe(true); }); - test("sync send false returns false, calls onSendFailure, and leaves mailbox uncollected", () => { + test("sync send false returns false, calls onSendFailure, and leaves mailbox uncollected", async () => { const records = new Map([ ["w1", { status: "done", report: "ok" }], ]); @@ -703,7 +863,7 @@ describe("driveOpenTasksAfterFleetDry", () => { }, }; let failures = 0; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ previousRunning: 1, running: 0, openTasks: [openTask], @@ -721,9 +881,9 @@ describe("driveOpenTasksAfterFleetDry", () => { expect(records.get("w1")?.collected).not.toBe(true); }); - test("sync send throw calls onSendFailure", () => { + test("sync send throw calls onSendFailure", async () => { let failures = 0; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ previousRunning: 1, running: 0, openTasks: [openTask], @@ -756,8 +916,8 @@ describe("driveOpenTasksAfterFleetDry", () => { failures += 1; }, }); - expect(driven).toBe(true); expect(failures).toBe(0); + expect(await driven).toBe(true); await Promise.resolve(); await Promise.resolve(); expect(failures).toBe(1); @@ -765,7 +925,7 @@ describe("driveOpenTasksAfterFleetDry", () => { test("TUI send rejection calls onSendFailure", async () => { let failures = 0; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ deferredDryEdge: true, openTasks: [openTask], parentProcessing: false, diff --git a/src/subagent/fleet-dry-drive.ts b/src/subagent/fleet-dry-drive.ts index a6fe230a2..29cb2ffe9 100644 --- a/src/subagent/fleet-dry-drive.ts +++ b/src/subagent/fleet-dry-drive.ts @@ -123,19 +123,23 @@ function truncateWithReservedNotice( return (text.slice(0, keptLen) + notice).slice(0, maxChars); } -function clipField( +async function clipField( text: string | undefined, agentId: string, field: "report" | "error", writeBlob?: FleetDryBlobWriter, -): string | undefined { +): Promise { if (text === undefined) return undefined; if (text.length <= FLEET_DRY_REPORT_CHARS) return text; let uri: string | undefined; if (writeBlob !== undefined) { const key = fleetDrySpillKey(agentId, field); - uri = `tool-output:///${key}`; - void writeBlob(key, new TextEncoder().encode(text), "text/plain"); + try { + await writeBlob(key, new TextEncoder().encode(text), "text/plain"); + uri = `tool-output:///${key}`; + } catch { + // Rejected write: same honest cut as a missing writer — no URI. + } } return truncateWithReservedNotice(text, FLEET_DRY_REPORT_CHARS, (keptLen) => truncationNotice({ @@ -188,17 +192,17 @@ export function takeAndProjectMailboxRecord( return projectMailboxRecord(id, taken, lane); } -function clipCollectedReport( +async function clipCollectedReport( report: CollectedWorkerReport, writeBlob?: FleetDryBlobWriter, -): CollectedWorkerReport { - const clippedReport = clipField( +): Promise { + const clippedReport = await clipField( report.report, report.agent_id, "report", writeBlob, ); - const clippedError = clipField( + const clippedError = await clipField( report.error, report.agent_id, "error", @@ -211,12 +215,12 @@ function clipCollectedReport( }; } -export function collectUncollectedTerminals( +export async function collectUncollectedTerminals( mailbox: FleetDryMailbox | undefined, lanes: readonly FleetDryLane[], consume: boolean, writeBlob?: FleetDryBlobWriter, -): CollectedWorkerReport[] { +): Promise { if (mailbox === undefined) return []; const byId = new Map(lanes.map((lane) => [lane.id, lane])); const reports: CollectedWorkerReport[] = []; @@ -229,7 +233,7 @@ export function collectUncollectedTerminals( ? takeAndProjectMailboxRecord(mailbox, id, byId.get(id)) : projectMailboxRecord(id, peeked, byId.get(id)); if (projected === undefined) continue; - reports.push(clipCollectedReport(projected, writeBlob)); + reports.push(await clipCollectedReport(projected, writeBlob)); } return reports; } @@ -269,7 +273,7 @@ export function driveOpenTasksAfterFleetDry(args: { beginSystemContinuation: (prompt: string) => void; send: (prompt: string) => unknown; onSendFailure?: () => void; -}): boolean { +}): boolean | Promise { const tasks = [...args.openTasks]; if ( !shouldDriveOpenTasks({ @@ -282,7 +286,14 @@ export function driveOpenTasksAfterFleetDry(args: { ) { return false; } - const reports = collectUncollectedTerminals( + return driveOpenTasksAfterFleetDrySpill(args, tasks); +} + +async function driveOpenTasksAfterFleetDrySpill( + args: Parameters[0], + tasks: Task[], +): Promise { + const reports = await collectUncollectedTerminals( args.mailbox, args.lanes, false, diff --git a/src/tui/runner/wiring.ts b/src/tui/runner/wiring.ts index 83697ac17..42b7e7c85 100644 --- a/src/tui/runner/wiring.ts +++ b/src/tui/runner/wiring.ts @@ -201,7 +201,7 @@ export function wirePostStartup( const send = state.sendWithAttemptIdentity; if (send === undefined) return false; const storage = state.currentStorage; - return driveOpenTasksAfterFleetDry({ + const driven = driveOpenTasksAfterFleetDry({ deferredDryEdge: true, openTasks: services.directorHolder.instance?.getTasks() ?? [], parentProcessing: false, @@ -221,6 +221,9 @@ export function wirePostStartup( sessionBridge.abortSystemContinuation(); }, }); + if (driven === false) return false; + void driven; + return true; }); sessionBridge.setMailboxMailDriver(() => { const send = state.sendWithAttemptIdentity; From 6b4719e1cb08642bc6d2df5ea1c87a05d6ca60a7 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Thu, 10 Sep 2026 07:07:29 -0700 Subject: [PATCH 3/4] Reuse leisure truncation notices for fleet-dry spills --- src/plugins/result-truncation-plugin.ts | 4 +- src/subagent/fleet-dry-drive.test.ts | 80 ------------------------- src/subagent/fleet-dry-drive.ts | 45 ++------------ 3 files changed, 7 insertions(+), 122 deletions(-) diff --git a/src/plugins/result-truncation-plugin.ts b/src/plugins/result-truncation-plugin.ts index 66dbcabbe..7a0600d0d 100644 --- a/src/plugins/result-truncation-plugin.ts +++ b/src/plugins/result-truncation-plugin.ts @@ -55,7 +55,7 @@ export function spillBlobKey(callId: string): string { return `${callId}:full`; } -function truncationNotice(args: { +export function truncationNotice(args: { maxChars: number; remaining: number; fullLength: number; @@ -92,7 +92,7 @@ function truncationNotice(args: { * remaining/fullLength (and optional absolutePath), so shrink kept until the * assembled result fits. */ -function truncateWithReservedNotice( +export function truncateWithReservedNotice( text: string, maxChars: number, buildNotice: (keptLen: number) => string, diff --git a/src/subagent/fleet-dry-drive.test.ts b/src/subagent/fleet-dry-drive.test.ts index b2b43c344..b50714794 100644 --- a/src/subagent/fleet-dry-drive.test.ts +++ b/src/subagent/fleet-dry-drive.test.ts @@ -557,86 +557,6 @@ describe("driveOpenTasksAfterFleetDry", () => { ).toBe(original); }); - test("send waits for async writeBlob before naming the URI", async () => { - const original = `head-${"x".repeat(FLEET_DRY_REPORT_CHARS)}TAIL-MARKER`; - const records = new Map([ - ["big", { status: "done", report: original }], - ]); - const store = fakeBlobStore(); - let resolveWrite: (() => void) | undefined; - const writeBlob = (key: string, bytes: Uint8Array, contentType: string) => - new Promise((resolve) => { - resolveWrite = () => { - store.writeBlob(key, bytes, contentType); - resolve(); - }; - }); - const sent: string[] = []; - const drivenP = driveOpenTasksAfterFleetDry({ - previousRunning: 1, - running: 0, - openTasks: [openTask], - parentProcessing: false, - mailbox: peekMailbox(records), - lanes: [], - writeBlob, - beginSystemContinuation: (prompt) => { - sent.push(prompt); - }, - send: () => undefined, - }); - await Promise.resolve(); - expect(sent).toEqual([]); - resolveWrite?.(); - expect(await drivenP).toBe(true); - const uri = `tool-output:///${fleetDrySpillKey("big", "report")}`; - expect(sent[0]).toContain(uri); - expect(sent[0]).not.toContain("TAIL-MARKER"); - const recovered = new TextDecoder().decode( - await createBlobReader(store).read(uri), - ); - expect(recovered).toBe(original); - }); - - test("rejected writeBlob send is NOT retrievable and names no URI", async () => { - const original = "x".repeat(FLEET_DRY_REPORT_CHARS + 40); - const records = new Map([ - ["big", { status: "done", report: original }], - ]); - const rejections: unknown[] = []; - const onUnhandled = (reason: unknown) => { - rejections.push(reason); - }; - process.on("unhandledRejection", onUnhandled); - try { - const sent: string[] = []; - const driven = await driveOpenTasksAfterFleetDry({ - previousRunning: 1, - running: 0, - openTasks: [openTask], - parentProcessing: false, - mailbox: peekMailbox(records), - lanes: [], - writeBlob: async () => { - throw new Error("disk full"); - }, - beginSystemContinuation: (prompt) => { - sent.push(prompt); - }, - send: () => undefined, - }); - await Promise.resolve(); - expect(driven).toBe(true); - const parsed = reportsJSONFromPrompt(sent[0] ?? ""); - const report = (parsed as { report?: string }[])[0]?.report ?? ""; - expect(report).toContain("NOT retrievable"); - expect(report).not.toContain("tool-output:///"); - expect(rejections).toEqual([]); - } finally { - process.off("unhandledRejection", onUnhandled); - } - }); - test("dry+terminal, live+open, and parentProcessing only skip", () => { const noop = { mailbox: undefined, diff --git a/src/subagent/fleet-dry-drive.ts b/src/subagent/fleet-dry-drive.ts index 29cb2ffe9..d7bba4c3b 100644 --- a/src/subagent/fleet-dry-drive.ts +++ b/src/subagent/fleet-dry-drive.ts @@ -5,6 +5,10 @@ */ import { hasActiveTasks, type Task } from "../agent/tasks.js"; +import { + truncationNotice, + truncateWithReservedNotice, +} from "../plugins/result-truncation-plugin.js"; import { isLiveWaitStatus, type WaitJSONStatus } from "./lifecycle.js"; /** Enough of a lane report for a parent continuation; traces stay on disk. */ @@ -83,46 +87,6 @@ export function fleetDrySpillKey( return `fleet-dry:${agentId}:${field}`; } -function truncationNotice(args: { - maxChars: number; - remaining: number; - fullLength: number; - uri?: string; -}): string { - const { maxChars, remaining, fullLength, uri } = args; - if (uri === undefined) { - return ( - `\n[output truncated at ${maxChars.toLocaleString()} chars — ` + - `${remaining.toLocaleString()} chars discarded, NOT retrievable ` + - `(no blob store is configured; re-running gives the same cut). ` + - `Use offset/limit or a narrower query.]` - ); - } - return ( - `\n[output truncated at ${maxChars.toLocaleString()} chars — ` + - `${remaining.toLocaleString()} more chars omitted here. The full result ` + - `(${fullLength.toLocaleString()} chars, text/plain) is saved at ${uri}` + - ` — use read_file with that URI (offset/limit supported) to see the rest.]` - ); -} - -function truncateWithReservedNotice( - text: string, - maxChars: number, - buildNotice: (keptLen: number) => string, -): string { - let keptLen = maxChars; - for (let i = 0; i < 8; i++) { - const notice = buildNotice(keptLen); - const total = keptLen + notice.length; - if (total <= maxChars) return text.slice(0, keptLen) + notice; - keptLen -= total - maxChars; - if (keptLen < 0) keptLen = 0; - } - const notice = buildNotice(keptLen); - return (text.slice(0, keptLen) + notice).slice(0, maxChars); -} - async function clipField( text: string | undefined, agentId: string, @@ -146,6 +110,7 @@ async function clipField( maxChars: FLEET_DRY_REPORT_CHARS, remaining: text.length - keptLen, fullLength: text.length, + contentType: "text/plain", ...(uri !== undefined ? { uri } : {}), }), ); From 877922f6aa2af57d090546b2c0de389b3a5b7762 Mon Sep 17 00:00:00 2001 From: Sawyer Date: Thu, 10 Sep 2026 07:48:19 -0700 Subject: [PATCH 4/4] Await collectUncollectedTerminals in mailbox mail drive --- src/subagent/mailbox-mail-drive.test.ts | 35 ++++++++++++++----------- src/subagent/mailbox-mail-drive.ts | 17 ++++++++++-- src/tui/runner/wiring.ts | 12 ++++++++- 3 files changed, 46 insertions(+), 18 deletions(-) diff --git a/src/subagent/mailbox-mail-drive.test.ts b/src/subagent/mailbox-mail-drive.test.ts index db49f99e0..4917cf79f 100644 --- a/src/subagent/mailbox-mail-drive.test.ts +++ b/src/subagent/mailbox-mail-drive.test.ts @@ -48,14 +48,14 @@ describe("buildMailboxMailPrompt", () => { }); describe("driveMailboxMail", () => { - test("idle parent with one terminal drives even while siblings run", () => { + test("idle parent with one terminal drives even while siblings run", async () => { const records = new Map([ ["done", { status: "done", report: "ok", description: "lane" }], ["live", { status: "running" }], ]); const order: string[] = []; const sent: string[] = []; - const driven = driveMailboxMail({ + const driven = await driveMailboxMail({ parentProcessing: false, mailbox: mapMailbox(records), lanes: [], @@ -77,12 +77,12 @@ describe("driveMailboxMail", () => { expect(records.get("live")?.collected).not.toBe(true); }); - test("fail path is the same terminal collect", () => { + test("fail path is the same terminal collect", async () => { const records = new Map([ ["fail", { status: "failed", error: "boom" }], ]); const sent: string[] = []; - const driven = driveMailboxMail({ + const driven = await driveMailboxMail({ parentProcessing: false, mailbox: mapMailbox(records), lanes: [], @@ -97,7 +97,7 @@ describe("driveMailboxMail", () => { expect(records.get("fail")?.collected).toBe(true); }); - test("parentProcessing or empty mailbox is a no-op", () => { + test("parentProcessing or empty mailbox is a no-op", async () => { const records = new Map([ ["done", { status: "done", report: "ok" }], ]); @@ -118,7 +118,7 @@ describe("driveMailboxMail", () => { }), ).toBe(false); expect( - driveMailboxMail({ + await driveMailboxMail({ parentProcessing: false, mailbox: mapMailbox(new Map()), lanes: [], @@ -128,7 +128,7 @@ describe("driveMailboxMail", () => { expect(records.get("done")?.collected).not.toBe(true); }); - test("already-collected terminals are not driven again", () => { + test("already-collected terminals are not driven again", async () => { const sessions = createSubAgentSessionStore(); const mailbox = createFleetMailbox(sessions); const session = sessions.start({ @@ -141,7 +141,7 @@ describe("driveMailboxMail", () => { sessions.complete("coll", "already taken"); mailbox.take("coll"); expect( - driveMailboxMail({ + await driveMailboxMail({ parentProcessing: false, mailbox, lanes: sessions.list(), @@ -155,11 +155,11 @@ describe("driveMailboxMail", () => { ).toBe(false); }); - test("send failure leaves reports waitable", () => { + test("send failure leaves reports waitable", async () => { const records = new Map([ ["w1", { status: "done", report: "ok" }], ]); - const driven = driveMailboxMail({ + const driven = await driveMailboxMail({ parentProcessing: false, mailbox: mapMailbox(records), lanes: [], @@ -176,7 +176,7 @@ describe("driveMailboxMail", () => { const records = new Map([ ["w1", { status: "done", report: "ok" }], ]); - const driven = driveMailboxMail({ + const driven = await driveMailboxMail({ parentProcessing: false, mailbox: mapMailbox(records), lanes: [], @@ -193,26 +193,31 @@ describe("driveMailboxMail", () => { const records = new Map([ ["w1", { status: "done", report: "ok" }], ]); - const driven = driveMailboxMail({ + let resolveSend: ((ok: boolean) => void) | undefined; + const driven = await driveMailboxMail({ parentProcessing: false, mailbox: mapMailbox(records), lanes: [], beginSystemContinuation: () => undefined, - send: () => Promise.resolve(true), + send: () => + new Promise((resolve) => { + resolveSend = resolve; + }), }); expect(driven).toBe(true); expect(records.get("w1")?.collected).not.toBe(true); + resolveSend?.(true); await Promise.resolve(); expect(records.get("w1")?.collected).toBe(true); }); - test("awaiting_director is not mailbox mail", () => { + test("awaiting_director is not mailbox mail", async () => { const records = new Map([ ["ask", { status: "awaiting_director" }], ["live", { status: "running" }], ]); expect( - driveMailboxMail({ + await driveMailboxMail({ parentProcessing: false, mailbox: mapMailbox(records), lanes: [], diff --git a/src/subagent/mailbox-mail-drive.ts b/src/subagent/mailbox-mail-drive.ts index f9dbc29e7..7ac1be076 100644 --- a/src/subagent/mailbox-mail-drive.ts +++ b/src/subagent/mailbox-mail-drive.ts @@ -9,6 +9,7 @@ import { isLiveWaitStatus } from "./lifecycle.js"; import { collectUncollectedTerminals, type CollectedWorkerReport, + type FleetDryBlobWriter, type FleetDryLane, type FleetDryMailbox, } from "./fleet-dry-drive.js"; @@ -46,12 +47,24 @@ export function driveMailboxMail(args: { parentProcessing: boolean; mailbox: FleetDryMailbox | undefined; lanes: readonly FleetDryLane[]; + writeBlob?: FleetDryBlobWriter; beginSystemContinuation: (prompt: string) => void; send: (prompt: string) => unknown; onSendFailure?: () => void; -}): boolean { +}): boolean | Promise { if (args.parentProcessing) return false; - const reports = collectUncollectedTerminals(args.mailbox, args.lanes, false); + return driveMailboxMailAfterCollect(args); +} + +async function driveMailboxMailAfterCollect( + args: Parameters[0], +): Promise { + const reports = await collectUncollectedTerminals( + args.mailbox, + args.lanes, + false, + args.writeBlob, + ); if (reports.length === 0) return false; const prompt = buildMailboxMailPrompt(reports); const takeReports = (): void => { diff --git a/src/tui/runner/wiring.ts b/src/tui/runner/wiring.ts index 42b7e7c85..9c1bcaaa2 100644 --- a/src/tui/runner/wiring.ts +++ b/src/tui/runner/wiring.ts @@ -228,10 +228,17 @@ export function wirePostStartup( sessionBridge.setMailboxMailDriver(() => { const send = state.sendWithAttemptIdentity; if (send === undefined) return false; - return driveMailboxMail({ + const storage = state.currentStorage; + const driven = driveMailboxMail({ parentProcessing: sessionBridge.turn.isProcessing, mailbox: services.toolset.fleetRecords, lanes: services.subAgentSessions.list(), + ...(storage !== null + ? { + writeBlob: (key, bytes, contentType) => + storage.writeBlob(key, bytes, contentType), + } + : {}), beginSystemContinuation: (prompt) => { sessionBridge.beginSystemContinuation(prompt); }, @@ -240,6 +247,9 @@ export function wirePostStartup( sessionBridge.abortSystemContinuation({ rearmDry: false }); }, }); + if (driven === false) return false; + void driven; + return true; }); sessionBridge.setWaitYieldWake(() => { services.subAgentSessions.wake();