Skip to content

Commit 877922f

Browse files
committed
Await collectUncollectedTerminals in mailbox mail drive
1 parent 6b4719e commit 877922f

3 files changed

Lines changed: 46 additions & 18 deletions

File tree

‎src/subagent/mailbox-mail-drive.test.ts‎

Lines changed: 20 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -48,14 +48,14 @@ describe("buildMailboxMailPrompt", () => {
4848
});
4949

5050
describe("driveMailboxMail", () => {
51-
test("idle parent with one terminal drives even while siblings run", () => {
51+
test("idle parent with one terminal drives even while siblings run", async () => {
5252
const records = new Map<string, FleetDryMailboxRecord>([
5353
["done", { status: "done", report: "ok", description: "lane" }],
5454
["live", { status: "running" }],
5555
]);
5656
const order: string[] = [];
5757
const sent: string[] = [];
58-
const driven = driveMailboxMail({
58+
const driven = await driveMailboxMail({
5959
parentProcessing: false,
6060
mailbox: mapMailbox(records),
6161
lanes: [],
@@ -77,12 +77,12 @@ describe("driveMailboxMail", () => {
7777
expect(records.get("live")?.collected).not.toBe(true);
7878
});
7979

80-
test("fail path is the same terminal collect", () => {
80+
test("fail path is the same terminal collect", async () => {
8181
const records = new Map<string, FleetDryMailboxRecord>([
8282
["fail", { status: "failed", error: "boom" }],
8383
]);
8484
const sent: string[] = [];
85-
const driven = driveMailboxMail({
85+
const driven = await driveMailboxMail({
8686
parentProcessing: false,
8787
mailbox: mapMailbox(records),
8888
lanes: [],
@@ -97,7 +97,7 @@ describe("driveMailboxMail", () => {
9797
expect(records.get("fail")?.collected).toBe(true);
9898
});
9999

100-
test("parentProcessing or empty mailbox is a no-op", () => {
100+
test("parentProcessing or empty mailbox is a no-op", async () => {
101101
const records = new Map<string, FleetDryMailboxRecord>([
102102
["done", { status: "done", report: "ok" }],
103103
]);
@@ -118,7 +118,7 @@ describe("driveMailboxMail", () => {
118118
}),
119119
).toBe(false);
120120
expect(
121-
driveMailboxMail({
121+
await driveMailboxMail({
122122
parentProcessing: false,
123123
mailbox: mapMailbox(new Map()),
124124
lanes: [],
@@ -128,7 +128,7 @@ describe("driveMailboxMail", () => {
128128
expect(records.get("done")?.collected).not.toBe(true);
129129
});
130130

131-
test("already-collected terminals are not driven again", () => {
131+
test("already-collected terminals are not driven again", async () => {
132132
const sessions = createSubAgentSessionStore();
133133
const mailbox = createFleetMailbox(sessions);
134134
const session = sessions.start({
@@ -141,7 +141,7 @@ describe("driveMailboxMail", () => {
141141
sessions.complete("coll", "already taken");
142142
mailbox.take("coll");
143143
expect(
144-
driveMailboxMail({
144+
await driveMailboxMail({
145145
parentProcessing: false,
146146
mailbox,
147147
lanes: sessions.list(),
@@ -155,11 +155,11 @@ describe("driveMailboxMail", () => {
155155
).toBe(false);
156156
});
157157

158-
test("send failure leaves reports waitable", () => {
158+
test("send failure leaves reports waitable", async () => {
159159
const records = new Map<string, FleetDryMailboxRecord>([
160160
["w1", { status: "done", report: "ok" }],
161161
]);
162-
const driven = driveMailboxMail({
162+
const driven = await driveMailboxMail({
163163
parentProcessing: false,
164164
mailbox: mapMailbox(records),
165165
lanes: [],
@@ -176,7 +176,7 @@ describe("driveMailboxMail", () => {
176176
const records = new Map<string, FleetDryMailboxRecord>([
177177
["w1", { status: "done", report: "ok" }],
178178
]);
179-
const driven = driveMailboxMail({
179+
const driven = await driveMailboxMail({
180180
parentProcessing: false,
181181
mailbox: mapMailbox(records),
182182
lanes: [],
@@ -193,26 +193,31 @@ describe("driveMailboxMail", () => {
193193
const records = new Map<string, FleetDryMailboxRecord>([
194194
["w1", { status: "done", report: "ok" }],
195195
]);
196-
const driven = driveMailboxMail({
196+
let resolveSend: ((ok: boolean) => void) | undefined;
197+
const driven = await driveMailboxMail({
197198
parentProcessing: false,
198199
mailbox: mapMailbox(records),
199200
lanes: [],
200201
beginSystemContinuation: () => undefined,
201-
send: () => Promise.resolve(true),
202+
send: () =>
203+
new Promise((resolve) => {
204+
resolveSend = resolve;
205+
}),
202206
});
203207
expect(driven).toBe(true);
204208
expect(records.get("w1")?.collected).not.toBe(true);
209+
resolveSend?.(true);
205210
await Promise.resolve();
206211
expect(records.get("w1")?.collected).toBe(true);
207212
});
208213

209-
test("awaiting_director is not mailbox mail", () => {
214+
test("awaiting_director is not mailbox mail", async () => {
210215
const records = new Map<string, FleetDryMailboxRecord>([
211216
["ask", { status: "awaiting_director" }],
212217
["live", { status: "running" }],
213218
]);
214219
expect(
215-
driveMailboxMail({
220+
await driveMailboxMail({
216221
parentProcessing: false,
217222
mailbox: mapMailbox(records),
218223
lanes: [],

‎src/subagent/mailbox-mail-drive.ts‎

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import { isLiveWaitStatus } from "./lifecycle.js";
99
import {
1010
collectUncollectedTerminals,
1111
type CollectedWorkerReport,
12+
type FleetDryBlobWriter,
1213
type FleetDryLane,
1314
type FleetDryMailbox,
1415
} from "./fleet-dry-drive.js";
@@ -46,12 +47,24 @@ export function driveMailboxMail(args: {
4647
parentProcessing: boolean;
4748
mailbox: FleetDryMailbox | undefined;
4849
lanes: readonly FleetDryLane[];
50+
writeBlob?: FleetDryBlobWriter;
4951
beginSystemContinuation: (prompt: string) => void;
5052
send: (prompt: string) => unknown;
5153
onSendFailure?: () => void;
52-
}): boolean {
54+
}): boolean | Promise<boolean> {
5355
if (args.parentProcessing) return false;
54-
const reports = collectUncollectedTerminals(args.mailbox, args.lanes, false);
56+
return driveMailboxMailAfterCollect(args);
57+
}
58+
59+
async function driveMailboxMailAfterCollect(
60+
args: Parameters<typeof driveMailboxMail>[0],
61+
): Promise<boolean> {
62+
const reports = await collectUncollectedTerminals(
63+
args.mailbox,
64+
args.lanes,
65+
false,
66+
args.writeBlob,
67+
);
5568
if (reports.length === 0) return false;
5669
const prompt = buildMailboxMailPrompt(reports);
5770
const takeReports = (): void => {

‎src/tui/runner/wiring.ts‎

Lines changed: 11 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -228,10 +228,17 @@ export function wirePostStartup(
228228
sessionBridge.setMailboxMailDriver(() => {
229229
const send = state.sendWithAttemptIdentity;
230230
if (send === undefined) return false;
231-
return driveMailboxMail({
231+
const storage = state.currentStorage;
232+
const driven = driveMailboxMail({
232233
parentProcessing: sessionBridge.turn.isProcessing,
233234
mailbox: services.toolset.fleetRecords,
234235
lanes: services.subAgentSessions.list(),
236+
...(storage !== null
237+
? {
238+
writeBlob: (key, bytes, contentType) =>
239+
storage.writeBlob(key, bytes, contentType),
240+
}
241+
: {}),
235242
beginSystemContinuation: (prompt) => {
236243
sessionBridge.beginSystemContinuation(prompt);
237244
},
@@ -240,6 +247,9 @@ export function wirePostStartup(
240247
sessionBridge.abortSystemContinuation({ rearmDry: false });
241248
},
242249
});
250+
if (driven === false) return false;
251+
void driven;
252+
return true;
243253
});
244254
sessionBridge.setWaitYieldWake(() => {
245255
services.subAgentSessions.wake();

0 commit comments

Comments
 (0)