From 5e8e4fad5f2e2e2c9dbc3f138487f0072c0c1a48 Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Fri, 11 Sep 2026 12:33:52 -0700 Subject: [PATCH 1/2] Stop mailbox-mail occupancy from flooding the send queue Overlapping occupancy flushes each called send before the parent was busy. Raising the depth-16 cap would only delay the same flood: occupancy already batches every uncollected terminal into one wake. --- src/subagent/index.ts | 2 + src/subagent/mailbox-mail-drive.test.ts | 166 +++++++++++++++++++++++- src/subagent/mailbox-mail-drive.ts | 93 +++++++++++-- src/tui/queued-delivery.test.ts | 7 +- src/tui/runner/wiring.ts | 52 ++++---- src/tui/runtime-bridge.test.ts | 8 +- 6 files changed, 276 insertions(+), 52 deletions(-) diff --git a/src/subagent/index.ts b/src/subagent/index.ts index 75ff5fbd4..8e21f4e16 100644 --- a/src/subagent/index.ts +++ b/src/subagent/index.ts @@ -29,7 +29,9 @@ export { export { driveOpenTasksAfterFleetDry } from "./fleet-dry-drive.js"; export { driveMailboxMail, + latchMailboxMailDrive, MAILBOX_MAIL_WAKE_PREFIX, + mailboxMailWakeLine, occupancyShouldYieldWait, } from "./mailbox-mail-drive.js"; export { diff --git a/src/subagent/mailbox-mail-drive.test.ts b/src/subagent/mailbox-mail-drive.test.ts index 4917cf79f..f800523a1 100644 --- a/src/subagent/mailbox-mail-drive.test.ts +++ b/src/subagent/mailbox-mail-drive.test.ts @@ -3,7 +3,9 @@ import { createFleetMailbox } from "./agent-fleet.js"; import { buildMailboxMailPrompt, driveMailboxMail, + latchMailboxMailDrive, MAILBOX_MAIL_WAKE_PREFIX, + mailboxMailWakeLine, occupancyShouldYieldWait, } from "./mailbox-mail-drive.js"; import { createSubAgentSessionStore } from "./session-store.js"; @@ -41,9 +43,8 @@ describe("buildMailboxMailPrompt", () => { expect(prompt.startsWith(MAILBOX_MAIL_WAKE_PREFIX)).toBe(true); expect(prompt).toContain("worker-1"); expect(prompt).toContain("shipped"); - expect(prompt).toContain( - "already collected — do not call wait_agents for these agent_ids", - ); + expect(prompt).toContain(mailboxMailWakeLine()); + expect(prompt).not.toContain("already collected"); }); }); @@ -211,6 +212,80 @@ describe("driveMailboxMail", () => { expect(records.get("w1")?.collected).toBe(true); }); + test("two flushes while send is pending deliver once", async () => { + const records = new Map([ + ["w1", { status: "done", report: "ok" }], + ]); + const mailbox = mapMailbox(records); + const sends: string[] = []; + let resolveSend: ((ok: boolean) => void) | undefined; + const driven = await driveMailboxMail({ + parentProcessing: false, + mailbox, + lanes: [], + beginSystemContinuation: () => undefined, + send: (prompt) => { + sends.push(prompt); + return new Promise((resolve) => { + resolveSend = resolve; + }); + }, + }); + expect(driven).toBe(true); + expect(sends).toHaveLength(1); + expect(records.get("w1")?.collected).not.toBe(true); + expect( + await driveMailboxMail({ + parentProcessing: false, + mailbox, + lanes: [], + beginSystemContinuation: () => { + throw new Error("must not begin"); + }, + send: () => { + throw new Error("must not send"); + }, + }), + ).toBe(false); + expect(sends).toHaveLength(1); + resolveSend?.(true); + await Promise.resolve(); + expect(records.get("w1")?.collected).toBe(true); + }); + + test("failed send can retry once", async () => { + const records = new Map([ + ["w1", { status: "done", report: "ok" }], + ]); + const mailbox = mapMailbox(records); + const sends: string[] = []; + expect( + await driveMailboxMail({ + parentProcessing: false, + mailbox, + lanes: [], + beginSystemContinuation: () => undefined, + send: () => { + throw new Error("send failed"); + }, + }), + ).toBe(false); + expect(records.get("w1")?.collected).not.toBe(true); + expect( + await driveMailboxMail({ + parentProcessing: false, + mailbox, + lanes: [], + beginSystemContinuation: () => undefined, + send: (prompt) => { + sends.push(prompt); + }, + }), + ).toBe(true); + expect(sends).toHaveLength(1); + expect(records.get("w1")?.collected).toBe(true); + }); + test("awaiting_director is not mailbox mail", async () => { const records = new Map([ ["ask", { status: "awaiting_director" }], @@ -232,6 +307,91 @@ describe("driveMailboxMail", () => { }); }); +describe("latchMailboxMailDrive", () => { + test("overlapping flushes send once until the in-flight collect settles", async () => { + const records = new Map([ + ["done", { status: "done", report: "ok", description: "lane" }], + ]); + const sends: string[] = []; + let resolveSend: (() => void) | undefined; + const sent = new Promise((resolve) => { + resolveSend = resolve; + }); + const driver = latchMailboxMailDrive(() => + driveMailboxMail({ + parentProcessing: false, + mailbox: mapMailbox(records), + lanes: [], + beginSystemContinuation: () => undefined, + send: (prompt) => { + sends.push(prompt); + resolveSend?.(); + }, + }), + ); + expect(driver()).toBe(true); + expect(driver()).toBe(false); + expect(driver()).toBe(false); + await sent; + expect(sends).toHaveLength(1); + expect(sends[0]).toContain(MAILBOX_MAIL_WAKE_PREFIX); + }); + + test("a false drive does not latch the next flush", () => { + let calls = 0; + const driver = latchMailboxMailDrive(() => { + calls += 1; + return false; + }); + expect(driver()).toBe(false); + expect(driver()).toBe(false); + expect(calls).toBe(2); + }); + + test("after the in-flight drive settles, a new terminal can send", async () => { + const records = new Map([ + ["first", { status: "done", report: "one" }], + ]); + const mailbox = mapMailbox(records); + const sends: string[] = []; + let sawSend: (() => void) | undefined; + const waitForSend = (): Promise => + new Promise((resolve) => { + sawSend = resolve; + }); + const driver = latchMailboxMailDrive(() => + driveMailboxMail({ + parentProcessing: false, + mailbox, + lanes: [], + beginSystemContinuation: () => undefined, + send: (prompt) => { + sends.push(prompt); + sawSend?.(); + }, + }), + ); + const first = waitForSend(); + expect(driver()).toBe(true); + await first; + expect(sends).toHaveLength(1); + records.set("second", { status: "done", report: "two" }); + const second = waitForSend(); + let retried = false; + for (let i = 0; i < 10; i++) { + await Promise.resolve(); + if (driver()) { + retried = true; + break; + } + } + expect(retried).toBe(true); + await second; + expect(sends).toHaveLength(2); + expect(sends[1]).toContain("second"); + }); +}); + describe("occupancyShouldYieldWait", () => { test("yields on uncollected terminal, fail, or ask; not on live or collected", () => { expect(occupancyShouldYieldWait(undefined)).toBe(false); diff --git a/src/subagent/mailbox-mail-drive.ts b/src/subagent/mailbox-mail-drive.ts index 7ac1be076..95f0722c6 100644 --- a/src/subagent/mailbox-mail-drive.ts +++ b/src/subagent/mailbox-mail-drive.ts @@ -16,10 +16,46 @@ import { export const MAILBOX_MAIL_WAKE_PREFIX = "mailbox mail"; +/** + * First-delivery instruction. Occupancy handed these reports to the parent; + * do not call wait_agents. Must not say "already collected" — that makes the + * first wake look like a replay. + */ +export function mailboxMailWakeLine(): string { + return `${MAILBOX_MAIL_WAKE_PREFIX} — occupancy delivered these worker reports (do not call wait_agents for these agent_ids):`; +} + function isPromiseLike(value: unknown): value is Promise { return typeof value === "object" && value !== null && "then" in value; } +/** + * Ids occupancy has snapshotted and handed to send, but not yet taken. + * A second flush must not start another parent turn for the same reports. + * Failed send clears the set so a later flush can retry. Weak-keyed so a + * mailbox object can go away without a leak. + */ +const deliveringByMailbox = new WeakMap>(); + +function deliveringSet(mailbox: FleetDryMailbox): Set { + let ids = deliveringByMailbox.get(mailbox); + if (ids === undefined) { + ids = new Set(); + deliveringByMailbox.set(mailbox, ids); + } + return ids; +} + +function releaseDelivering( + mailbox: FleetDryMailbox | undefined, + ids: readonly string[], +): void { + if (mailbox === undefined) return; + const delivering = deliveringByMailbox.get(mailbox); + if (delivering === undefined) return; + for (const id of ids) delivering.delete(id); +} + export function occupancyShouldYieldWait( mailbox: FleetDryMailbox | undefined, ): boolean { @@ -37,10 +73,7 @@ export function occupancyShouldYieldWait( export function buildMailboxMailPrompt( reports: readonly CollectedWorkerReport[], ): string { - return [ - `${MAILBOX_MAIL_WAKE_PREFIX} — worker reports (already collected — do not call wait_agents for these agent_ids):`, - JSON.stringify(reports), - ].join("\n"); + return [mailboxMailWakeLine(), JSON.stringify(reports)].join("\n"); } export function driveMailboxMail(args: { @@ -59,20 +92,30 @@ export function driveMailboxMail(args: { async function driveMailboxMailAfterCollect( args: Parameters[0], ): Promise { - const reports = await collectUncollectedTerminals( - args.mailbox, - args.lanes, - false, - args.writeBlob, - ); + const delivering = + args.mailbox !== undefined + ? deliveringSet(args.mailbox) + : new Set(); + const reports = ( + await collectUncollectedTerminals( + args.mailbox, + args.lanes, + false, + args.writeBlob, + ) + ).filter((report) => !delivering.has(report.agent_id)); if (reports.length === 0) return false; + const ids = reports.map((report) => report.agent_id); + for (const id of ids) delivering.add(id); const prompt = buildMailboxMailPrompt(reports); const takeReports = (): void => { - for (const report of reports) { - args.mailbox?.take(report.agent_id); + for (const id of ids) { + args.mailbox?.take(id); } + releaseDelivering(args.mailbox, ids); }; const fail = (): boolean => { + releaseDelivering(args.mailbox, ids); args.onSendFailure?.(); return false; }; @@ -83,10 +126,10 @@ async function driveMailboxMailAfterCollect( void sent.then( (result) => { if (result !== false) takeReports(); - else args.onSendFailure?.(); + else fail(); }, () => { - args.onSendFailure?.(); + fail(); }, ); return true; @@ -98,3 +141,25 @@ async function driveMailboxMailAfterCollect( } return true; } + +/** + * Store-subscribe, stall-poll, and idle-with-fleet settle all flush mailbox + * mail. `driveMailboxMail` does not mark the parent busy until after an + * awaited collect, so overlapping flushes would each call `send()` and fill + * the agent's depth-16 queue. Hold one drive until that promise settles. + */ +export function latchMailboxMailDrive( + drive: () => boolean | Promise, +): () => boolean { + let inFlight = false; + return () => { + if (inFlight) return false; + const driven = drive(); + if (driven === false) return false; + inFlight = true; + void Promise.resolve(driven).finally(() => { + inFlight = false; + }); + return true; + }; +} diff --git a/src/tui/queued-delivery.test.ts b/src/tui/queued-delivery.test.ts index 69b50a548..456340b70 100644 --- a/src/tui/queued-delivery.test.ts +++ b/src/tui/queued-delivery.test.ts @@ -1,4 +1,5 @@ import { describe, expect, test } from "bun:test"; +import { mailboxMailWakeLine } from "../subagent/mailbox-mail-drive.js"; import type { PendingImageAttachment } from "./image-attachments.js"; import { createDeliveryGeneration, @@ -296,15 +297,13 @@ describe("createLeftoverSend", () => { }); leftoverSend("ask_director wake — see @src/foo.ts"); - leftoverSend( - "mailbox mail — worker reports (already collected — do not call wait_agents for these agent_ids):\n[]", - ); + leftoverSend(`${mailboxMailWakeLine()}\n[]`); leftoverSend("please read @src/foo.ts"); await awaitTail(); expect(ingested).toEqual(["please read @src/foo.ts"]); expect(sent).toEqual([ "ask_director wake — see @src/foo.ts", - "mailbox mail — worker reports (already collected — do not call wait_agents for these agent_ids):\n[]", + `${mailboxMailWakeLine()}\n[]`, "please read @src/foo.ts ingested", ]); }); diff --git a/src/tui/runner/wiring.ts b/src/tui/runner/wiring.ts index d9d682891..40fcf09b4 100644 --- a/src/tui/runner/wiring.ts +++ b/src/tui/runner/wiring.ts @@ -23,6 +23,7 @@ import { createFleetWatch, driveOpenTasksAfterFleetDry, driveMailboxMail, + latchMailboxMailDrive, FLEET_REPORT_SETTLE_MS, FLEET_STALL_POLL_MS, liveFleetCount, @@ -246,32 +247,31 @@ export function wirePostStartup( void driven; return true; }); - sessionBridge.setMailboxMailDriver(() => { - const send = state.sendWithAttemptIdentity; - if (send === undefined) return false; - 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); - }, - send: (prompt) => send(buildMailboxMailMessage(prompt)), - onSendFailure: () => { - sessionBridge.abortSystemContinuation({ rearmDry: false }); - }, - }); - if (driven === false) return false; - void driven; - return true; - }); + sessionBridge.setMailboxMailDriver( + latchMailboxMailDrive(() => { + const send = state.sendWithAttemptIdentity; + if (send === undefined) return false; + const storage = state.currentStorage; + return 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); + }, + send: (prompt) => send(buildMailboxMailMessage(prompt)), + onSendFailure: () => { + sessionBridge.abortSystemContinuation({ rearmDry: false }); + }, + }); + }), + ); sessionBridge.setWaitYieldWake(() => { services.subAgentSessions.wake(); }); diff --git a/src/tui/runtime-bridge.test.ts b/src/tui/runtime-bridge.test.ts index d63fdbccc..7efcd6839 100644 --- a/src/tui/runtime-bridge.test.ts +++ b/src/tui/runtime-bridge.test.ts @@ -1,4 +1,5 @@ import { describe, expect, spyOn, test } from "bun:test"; +import { mailboxMailWakeLine } from "../subagent/mailbox-mail-drive.js"; import { defined } from "../../tests/helpers/defined.js"; import { FIXTURE_BUSY_SESSION, @@ -2137,8 +2138,7 @@ describe("mailbox mail occupancy (CL-7518)", () => { const port = createRecordingPort(); const bridge = attachSessionBridge(shell, port); try { - const prompt = - "mailbox mail — worker reports (already collected — do not call wait_agents for these agent_ids):\n[]"; + const prompt = `${mailboxMailWakeLine()}\n[]`; let terminals = 0; let drives = 0; bridge.setMailboxMailDriver(() => { @@ -2366,9 +2366,7 @@ describe("mailbox mail occupancy (CL-7518)", () => { let drives = 0; bridge.handle({ type: "fleet", running: 1 }); settleToollessTurn(bridge); - bridge.beginSystemContinuation( - "mailbox mail — worker reports (already collected — do not call wait_agents for these agent_ids):\n[]", - ); + bridge.beginSystemContinuation(`${mailboxMailWakeLine()}\n[]`); bridge.abortSystemContinuation({ rearmDry: false }); bridge.setMailboxMailDriver(() => { drives += 1; From 57578a26507e00e2df26c0fed35f8a7c2d97814d Mon Sep 17 00:00:00 2001 From: Sawyer Cutler Date: Fri, 11 Sep 2026 12:34:57 -0700 Subject: [PATCH 2/2] Note occupancy mailbox-mail flood and replay for the next release --- CHANGELOG.md | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 060344eee..058674e11 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -36,6 +36,10 @@ parallel copies under `docs/` or `scripts/notes/`. At cut time: rename ### Fixed +- Occupancy injects mailbox mail once when a worker burst finishes: one + in-flight drive, in-flight ids until send succeeds, and wake copy that + does not say the reports were already collected. Overlapping flushes + no longer fill the send queue or replay the same reports as new turns. - Stalled workers get a full `stallTimeoutMs` grace after the first continuation nudge before salvage. Stall pings inside that window wait instead of counting toward escalation. Mailbox mail re-flushes from the