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 e85099db1..b50714794 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( @@ -141,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) => { @@ -174,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", @@ -197,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" }], ]); @@ -212,7 +255,7 @@ describe("collectUncollectedTerminals", () => { return taken; }, }; - const reports = collectUncollectedTerminals( + const reports = await collectUncollectedTerminals( mailbox, [{ id: "ghost", description: "from store", report: "store report" }], true, @@ -228,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", @@ -240,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", @@ -250,24 +293,101 @@ describe("collectUncollectedTerminals", () => { ]); }); - test("clips oversized reports", () => { + 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: "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 = await 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("consume false peeks without take", () => { + 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 = await 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", async () => { + const report = "short enough"; + const records = new Map([ + ["w1", { status: "done", report }], + ]); + const store = fakeBlobStore(); + const reports = await 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", async () => { + const original = "e".repeat(FLEET_DRY_REPORT_CHARS + 20); + const records = new Map([ + ["boom", { status: "failed", error: original }], + ]); + const store = fakeBlobStore(); + const reports = await 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", async () => { const records = new Map([ ["w1", { status: "done", report: "ok" }], ]); @@ -282,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" }], @@ -306,7 +498,7 @@ describe("driveOpenTasksAfterFleetDry", () => { }, }; const sent: string[] = []; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ previousRunning: 1, running: 0, openTasks: [openTask], @@ -329,6 +521,42 @@ describe("driveOpenTasksAfterFleetDry", () => { expect(records.get("w1")?.collected).toBe(true); }); + 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 = await 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, @@ -369,7 +597,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" }], ]); @@ -385,7 +613,7 @@ describe("driveOpenTasksAfterFleetDry", () => { }, }; const sent: string[] = []; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ previousRunning: 0, running: 0, deferredDryEdge: true, @@ -403,7 +631,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" }], ]); @@ -418,7 +646,7 @@ describe("driveOpenTasksAfterFleetDry", () => { return taken; }, }; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ previousRunning: 1, running: 0, openTasks: [openTask], @@ -453,7 +681,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, @@ -488,7 +716,7 @@ describe("driveOpenTasksAfterFleetDry", () => { await Promise.resolve(); return false; }; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ deferredDryEdge: true, openTasks: [openTask], parentProcessing: false, @@ -523,7 +751,7 @@ describe("driveOpenTasksAfterFleetDry", () => { new Promise((resolve) => { resolveSend = resolve; }); - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ deferredDryEdge: true, openTasks: [openTask], parentProcessing: false, @@ -539,7 +767,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" }], ]); @@ -555,7 +783,7 @@ describe("driveOpenTasksAfterFleetDry", () => { }, }; let failures = 0; - const driven = driveOpenTasksAfterFleetDry({ + const driven = await driveOpenTasksAfterFleetDry({ previousRunning: 1, running: 0, openTasks: [openTask], @@ -573,9 +801,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], @@ -608,8 +836,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); @@ -617,7 +845,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 b02c1979d..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. */ @@ -68,10 +72,48 @@ 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}`; +} + +async function clipField( + text: string | undefined, + agentId: string, + field: "report" | "error", + writeBlob?: FleetDryBlobWriter, +): Promise { 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); + 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({ + maxChars: FLEET_DRY_REPORT_CHARS, + remaining: text.length - keptLen, + fullLength: text.length, + contentType: "text/plain", + ...(uri !== undefined ? { uri } : {}), + }), + ); } export function projectMailboxRecord( @@ -115,11 +157,22 @@ export function takeAndProjectMailboxRecord( return projectMailboxRecord(id, taken, lane); } -function clipCollectedReport( +async function clipCollectedReport( report: CollectedWorkerReport, -): CollectedWorkerReport { - const clippedReport = clipField(report.report); - const clippedError = clipField(report.error); + writeBlob?: FleetDryBlobWriter, +): Promise { + const clippedReport = await clipField( + report.report, + report.agent_id, + "report", + writeBlob, + ); + const clippedError = await clipField( + report.error, + report.agent_id, + "error", + writeBlob, + ); return { ...report, ...(clippedReport !== undefined ? { report: clippedReport } : {}), @@ -127,11 +180,12 @@ function clipCollectedReport( }; } -export function collectUncollectedTerminals( +export async function collectUncollectedTerminals( mailbox: FleetDryMailbox | undefined, lanes: readonly FleetDryLane[], consume: boolean, -): CollectedWorkerReport[] { + writeBlob?: FleetDryBlobWriter, +): Promise { if (mailbox === undefined) return []; const byId = new Map(lanes.map((lane) => [lane.id, lane])); const reports: CollectedWorkerReport[] = []; @@ -144,7 +198,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(await clipCollectedReport(projected, writeBlob)); } return reports; } @@ -180,10 +234,11 @@ export function driveOpenTasksAfterFleetDry(args: { deferredDryEdge?: boolean; mailbox: FleetDryMailbox | undefined; lanes: readonly FleetDryLane[]; + writeBlob?: FleetDryBlobWriter; beginSystemContinuation: (prompt: string) => void; send: (prompt: string) => unknown; onSendFailure?: () => void; -}): boolean { +}): boolean | Promise { const tasks = [...args.openTasks]; if ( !shouldDriveOpenTasks({ @@ -196,7 +251,19 @@ export function driveOpenTasksAfterFleetDry(args: { ) { return false; } - const reports = collectUncollectedTerminals(args.mailbox, args.lanes, false); + return driveOpenTasksAfterFleetDrySpill(args, tasks); +} + +async function driveOpenTasksAfterFleetDrySpill( + args: Parameters[0], + tasks: Task[], +): Promise { + const reports = await 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/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 def5e8e03..9c1bcaaa2 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; - return driveOpenTasksAfterFleetDry({ + const storage = state.currentStorage; + const driven = 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); }, @@ -214,14 +221,24 @@ export function wirePostStartup( sessionBridge.abortSystemContinuation(); }, }); + if (driven === false) return false; + void driven; + return true; }); 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); }, @@ -230,6 +247,9 @@ export function wirePostStartup( sessionBridge.abortSystemContinuation({ rearmDry: false }); }, }); + if (driven === false) return false; + void driven; + return true; }); sessionBridge.setWaitYieldWake(() => { services.subAgentSessions.wake();