Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions src/subagent/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
166 changes: 163 additions & 3 deletions src/subagent/mailbox-mail-drive.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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");
});
});

Expand Down Expand Up @@ -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<string, FleetDryMailboxRecord>([
["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<boolean>((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<string, FleetDryMailboxRecord>([
["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<string, FleetDryMailboxRecord>([
["ask", { status: "awaiting_director" }],
Expand All @@ -232,6 +307,91 @@ describe("driveMailboxMail", () => {
});
});

describe("latchMailboxMailDrive", () => {
test("overlapping flushes send once until the in-flight collect settles", async () => {
const records = new Map<string, FleetDryMailboxRecord>([
["done", { status: "done", report: "ok", description: "lane" }],
]);
const sends: string[] = [];
let resolveSend: (() => void) | undefined;
const sent = new Promise<void>((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<string, FleetDryMailboxRecord>([
["first", { status: "done", report: "one" }],
]);
const mailbox = mapMailbox(records);
const sends: string[] = [];
let sawSend: (() => void) | undefined;
const waitForSend = (): Promise<void> =>
new Promise<void>((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);
Expand Down
93 changes: 79 additions & 14 deletions src/subagent/mailbox-mail-drive.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<unknown> {
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<FleetDryMailbox, Set<string>>();

function deliveringSet(mailbox: FleetDryMailbox): Set<string> {
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 {
Expand All @@ -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: {
Expand All @@ -59,20 +92,30 @@ export function driveMailboxMail(args: {
async function driveMailboxMailAfterCollect(
args: Parameters<typeof driveMailboxMail>[0],
): Promise<boolean> {
const reports = await collectUncollectedTerminals(
args.mailbox,
args.lanes,
false,
args.writeBlob,
);
const delivering =
args.mailbox !== undefined
? deliveringSet(args.mailbox)
: new Set<string>();
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;
};
Expand All @@ -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;
Expand All @@ -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>,
): () => 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;
};
}
7 changes: 3 additions & 4 deletions src/tui/queued-delivery.test.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand Down Expand Up @@ -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",
]);
});
Expand Down
Loading
Loading