Skip to content

Commit 270624b

Browse files
fix(tui): flush mailbox mail before ask wake on settle and gate close (#1105)
* fix(tui): flush mailbox mail before ask wake on settle and gate close settleRunToIdle (idle-with-fleet) and gateClosed flushed the ask wake before mailbox mail, letting a re-surface wake steal the turn occupancy should take. Match the #1080 abort-path rule: mail first, wake only if mail did not start a turn (isProcessing guard, deliveredAskWake dedupe — no duplicate wake). Tests: 4 new CL-8061 order/suppression cases in agent-ask-wake.test.ts (failed pre-fix, pass post-fix). Fixes CL-8061 * fix(tui): claim the primary turn for occupancy before ask wake Mailbox mail collect is async, so isProcessing is still false while occupancy already owns the slot. A shared mail-then-wake guard uses the latch claim so ask-wake cannot send on that stack.
1 parent 127ec5a commit 270624b

5 files changed

Lines changed: 339 additions & 24 deletions

File tree

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

Lines changed: 45 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -132,7 +132,7 @@ describe("driveMailboxMail", () => {
132132
}),
133133
).toBe(false);
134134
expect(
135-
await driveMailboxMail({
135+
driveMailboxMail({
136136
parentProcessing: false,
137137
mailbox: mapMailbox(new Map()),
138138
lanes: [],
@@ -315,7 +315,7 @@ describe("driveMailboxMail", () => {
315315
["live", { status: "running" }],
316316
]);
317317
expect(
318-
await driveMailboxMail({
318+
driveMailboxMail({
319319
parentProcessing: false,
320320
mailbox: mapMailbox(records),
321321
lanes: [],
@@ -329,6 +329,28 @@ describe("driveMailboxMail", () => {
329329
).toBe(false);
330330
});
331331

332+
test("does not begin if the parent starts processing during collect", async () => {
333+
const records = new Map<string, FleetDryMailboxRecord>([
334+
["w1", { status: "done", report: "ok" }],
335+
]);
336+
let processing = false;
337+
const driven = driveMailboxMail({
338+
parentProcessing: false,
339+
isParentProcessing: () => processing,
340+
mailbox: mapMailbox(records),
341+
lanes: [],
342+
beginSystemContinuation: () => {
343+
throw new Error("must not begin");
344+
},
345+
send: () => {
346+
throw new Error("must not send");
347+
},
348+
});
349+
processing = true;
350+
expect(await driven).toBe(false);
351+
expect(records.get("w1")?.collected).not.toBe(true);
352+
});
353+
332354
test("not-delivered send leaves the wake retryable", async () => {
333355
const records = new Map<string, FleetDryMailboxRecord>([
334356
["w1", { status: "done", report: "ok" }],
@@ -395,8 +417,10 @@ describe("latchMailboxMailDrive", () => {
395417
}),
396418
);
397419
expect(driver()).toBe(true);
420+
expect(driver.claimed()).toBe(true);
398421
expect(driver()).toBe(false);
399422
expect(driver()).toBe(false);
423+
expect(driver.claimed()).toBe(true);
400424
await sent;
401425
expect(sends).toHaveLength(1);
402426
expect(sends[0]).toContain(MAILBOX_MAIL_WAKE_PREFIX);
@@ -410,9 +434,28 @@ describe("latchMailboxMailDrive", () => {
410434
});
411435
expect(driver()).toBe(false);
412436
expect(driver()).toBe(false);
437+
expect(driver.claimed()).toBe(false);
413438
expect(calls).toBe(2);
414439
});
415440

441+
test("an empty mailbox does not claim the occupancy slot", () => {
442+
const driver = latchMailboxMailDrive(() =>
443+
driveMailboxMail({
444+
parentProcessing: false,
445+
mailbox: mapMailbox(new Map()),
446+
lanes: [],
447+
beginSystemContinuation: () => {
448+
throw new Error("must not begin");
449+
},
450+
send: () => {
451+
throw new Error("must not send");
452+
},
453+
}),
454+
);
455+
expect(driver()).toBe(false);
456+
expect(driver.claimed()).toBe(false);
457+
});
458+
416459
test("after the in-flight drive settles, a new terminal can send", async () => {
417460
const records = new Map<string, FleetDryMailboxRecord>([
418461
["first", { status: "done", report: "one" }],

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

Lines changed: 46 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -109,16 +109,41 @@ export function buildMailboxMailPrompt(
109109
return [mailboxMailWakeLine(), JSON.stringify(reports)].join("\n");
110110
}
111111

112+
function hasDeliverableMailboxMail(
113+
mailbox: FleetDryMailbox | undefined,
114+
): boolean {
115+
if (mailbox === undefined) return false;
116+
const delivering = deliveringByMailbox.get(mailbox);
117+
for (const id of mailbox.ids()) {
118+
if (delivering?.has(id)) continue;
119+
const record = mailbox.peek(id);
120+
if (record === undefined || record.collected === true) continue;
121+
if (isLiveWaitStatus(record.status)) continue;
122+
return true;
123+
}
124+
return false;
125+
}
126+
127+
function parentIsProcessing(args: {
128+
parentProcessing: boolean;
129+
isParentProcessing?: () => boolean;
130+
}): boolean {
131+
return args.isParentProcessing?.() ?? args.parentProcessing;
132+
}
133+
112134
export function driveMailboxMail(args: {
113135
parentProcessing: boolean;
136+
/** Live check after awaited collect; occupancy must not start a turn already claimed. */
137+
isParentProcessing?: () => boolean;
114138
mailbox: FleetDryMailbox | undefined;
115139
lanes: readonly FleetDryLane[];
116140
writeBlob?: FleetDryBlobWriter;
117141
beginSystemContinuation: (prompt: string) => void;
118142
send: (prompt: string) => unknown;
119143
onSendFailure?: () => void;
120144
}): boolean | Promise<boolean> {
121-
if (args.parentProcessing) return false;
145+
if (parentIsProcessing(args)) return false;
146+
if (!hasDeliverableMailboxMail(args.mailbox)) return false;
122147
return driveMailboxMailAfterCollect(args);
123148
}
124149

@@ -138,6 +163,7 @@ async function driveMailboxMailAfterCollect(
138163
)
139164
).filter((report) => !delivering.has(report.agent_id));
140165
if (reports.length === 0) return false;
166+
if (parentIsProcessing(args)) return false;
141167
const ids = reports.map((report) => report.agent_id);
142168
for (const id of ids) delivering.add(id);
143169
const prompt = buildMailboxMailPrompt(reports);
@@ -169,12 +195,28 @@ async function driveMailboxMailAfterCollect(
169195
* mail. `driveMailboxMail` does not mark the parent busy until after an
170196
* awaited collect, so overlapping flushes would each call `send()` and fill
171197
* the agent's depth-16 queue. Hold one drive until that promise settles.
198+
*
199+
* `flush()` is true only when this call starts a drive. `claimed()` is true
200+
* while that drive owns the slot — ask-wake admission uses it so a wake
201+
* cannot send on the same stack, or on a later stack before collect lands.
172202
*/
203+
export type MailboxMailFlush = (() => boolean) & {
204+
readonly claimed: () => boolean;
205+
};
206+
207+
export function mailboxMailDriveClaimed(
208+
flush: (() => boolean) | undefined,
209+
): boolean {
210+
if (flush === undefined) return false;
211+
const claimed = (flush as MailboxMailFlush).claimed;
212+
return typeof claimed === "function" && claimed() === true;
213+
}
214+
173215
export function latchMailboxMailDrive(
174216
drive: () => boolean | Promise<boolean>,
175-
): () => boolean {
217+
): MailboxMailFlush {
176218
let inFlight = false;
177-
return () => {
219+
const flush = (): boolean => {
178220
if (inFlight) return false;
179221
const driven = drive();
180222
if (driven === false) return false;
@@ -184,4 +226,5 @@ export function latchMailboxMailDrive(
184226
});
185227
return true;
186228
};
229+
return Object.assign(flush, { claimed: () => inFlight });
187230
}

0 commit comments

Comments
 (0)