Skip to content

Commit 4292aee

Browse files
committed
Re-arm occupancy send failure and stamp lastSentMessage
A failed occupancy send consumed the dry-open latch and left the turn running. Abort now re-arms, drains, and resets cadence so a later settle can retry; a successful occupancy stamps lastSentMessage so quota retry resends the continuation.
1 parent 83a6c65 commit 4292aee

5 files changed

Lines changed: 294 additions & 3 deletions

File tree

src/subagent/fleet-dry-drive.test.ts

Lines changed: 99 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -502,4 +502,103 @@ describe("driveOpenTasksAfterFleetDry", () => {
502502
await Promise.resolve();
503503
expect(records.get("w1")?.collected).toBe(true);
504504
});
505+
506+
test("sync send false returns false, calls onSendFailure, and leaves mailbox uncollected", () => {
507+
const records = new Map<string, FleetDryMailboxRecord>([
508+
["w1", { status: "done", report: "ok" }],
509+
]);
510+
const mailbox: FleetDryMailbox = {
511+
ids: () => [...records.keys()],
512+
peek: (id) => records.get(id),
513+
take: (id) => {
514+
const existing = records.get(id);
515+
if (existing === undefined) return undefined;
516+
const taken = { ...existing, collected: true };
517+
records.set(id, taken);
518+
return taken;
519+
},
520+
};
521+
let failures = 0;
522+
const driven = driveOpenTasksAfterFleetDry({
523+
previousRunning: 1,
524+
running: 0,
525+
openTasks: [openTask],
526+
parentProcessing: false,
527+
mailbox,
528+
lanes: [],
529+
beginSystemContinuation: () => undefined,
530+
send: () => false,
531+
onSendFailure: () => {
532+
failures += 1;
533+
},
534+
});
535+
expect(driven).toBe(false);
536+
expect(failures).toBe(1);
537+
expect(records.get("w1")?.collected).not.toBe(true);
538+
});
539+
540+
test("sync send throw calls onSendFailure", () => {
541+
let failures = 0;
542+
const driven = driveOpenTasksAfterFleetDry({
543+
previousRunning: 1,
544+
running: 0,
545+
openTasks: [openTask],
546+
parentProcessing: false,
547+
mailbox: undefined,
548+
lanes: [],
549+
beginSystemContinuation: () => undefined,
550+
send: () => {
551+
throw new Error("send failed");
552+
},
553+
onSendFailure: () => {
554+
failures += 1;
555+
},
556+
});
557+
expect(driven).toBe(false);
558+
expect(failures).toBe(1);
559+
});
560+
561+
test("TUI send false after handleSendFailure calls onSendFailure", async () => {
562+
let failures = 0;
563+
const driven = driveOpenTasksAfterFleetDry({
564+
deferredDryEdge: true,
565+
openTasks: [openTask],
566+
parentProcessing: false,
567+
mailbox: undefined,
568+
lanes: [],
569+
beginSystemContinuation: () => undefined,
570+
send: async () => false,
571+
onSendFailure: () => {
572+
failures += 1;
573+
},
574+
});
575+
expect(driven).toBe(true);
576+
expect(failures).toBe(0);
577+
await Promise.resolve();
578+
await Promise.resolve();
579+
expect(failures).toBe(1);
580+
});
581+
582+
test("TUI send rejection calls onSendFailure", async () => {
583+
let failures = 0;
584+
const driven = driveOpenTasksAfterFleetDry({
585+
deferredDryEdge: true,
586+
openTasks: [openTask],
587+
parentProcessing: false,
588+
mailbox: undefined,
589+
lanes: [],
590+
beginSystemContinuation: () => undefined,
591+
send: async () => {
592+
await Promise.resolve();
593+
throw new Error("agentProxy.send failed");
594+
},
595+
onSendFailure: () => {
596+
failures += 1;
597+
},
598+
});
599+
expect(driven).toBe(true);
600+
await Promise.resolve();
601+
await Promise.resolve();
602+
expect(failures).toBe(1);
603+
});
505604
});

src/subagent/fleet-dry-drive.ts

Lines changed: 12 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -166,6 +166,7 @@ export function driveOpenTasksAfterFleetDry(args: {
166166
lanes: readonly FleetDryLane[];
167167
beginSystemContinuation: (prompt: string) => void;
168168
send: (prompt: string) => unknown;
169+
onSendFailure?: () => void;
169170
}): boolean {
170171
const tasks = [...args.openTasks];
171172
if (
@@ -186,21 +187,29 @@ export function driveOpenTasksAfterFleetDry(args: {
186187
args.mailbox?.take(report.agent_id);
187188
}
188189
};
190+
const fail = (): boolean => {
191+
args.onSendFailure?.();
192+
return false;
193+
};
189194
try {
190195
args.beginSystemContinuation(prompt);
191196
const sent = args.send(prompt);
192197
if (isPromiseLike(sent)) {
193198
void sent.then(
194199
(result) => {
195200
if (result !== false) takeReports();
201+
else args.onSendFailure?.();
202+
},
203+
() => {
204+
args.onSendFailure?.();
196205
},
197-
() => undefined,
198206
);
199207
return true;
200208
}
201-
if (sent !== false) takeReports();
209+
if (sent === false) return fail();
210+
takeReports();
202211
} catch {
203-
return false;
212+
return fail();
204213
}
205214
return true;
206215
}

src/tui/runner/wiring.ts

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -146,6 +146,9 @@ export function wirePostStartup(
146146
sessionBridge.beginSystemContinuation(prompt);
147147
},
148148
send: (prompt) => send(buildFleetDryContinuationMessage(prompt)),
149+
onSendFailure: () => {
150+
sessionBridge.abortSystemContinuation();
151+
},
149152
});
150153
});
151154
const unsubscribeFleetReport = services.subAgentSessions.subscribe(() => {

src/tui/runtime-bridge.test.ts

Lines changed: 160 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1565,6 +1565,166 @@ describe("fleet-dry open-task drive (CL-7540)", () => {
15651565
{ width: 80, height: 24 },
15661566
);
15671567
});
1568+
1569+
test("occupancy send failure re-arms, clears continuation hold, and drains follow-ups", async () => {
1570+
await withTestRenderer(
1571+
async (h) => {
1572+
const shell = createAppShell(h.renderer, {
1573+
terminal: { columns: 80, rows: 24 },
1574+
wireKeys: false,
1575+
run: "idle",
1576+
});
1577+
const port = createRecordingPort();
1578+
const bridge = attachSessionBridge(shell, port);
1579+
try {
1580+
const prompt = "The fleet has gone dry. Remaining open tasks:\n- t1: keep going (todo)\n";
1581+
let drives = 0;
1582+
bridge.setDryOpenTaskDriver(() => {
1583+
drives += 1;
1584+
bridge.beginSystemContinuation(prompt);
1585+
if (drives === 1) {
1586+
void Promise.resolve().then(() => {
1587+
bridge.abortSystemContinuation();
1588+
});
1589+
}
1590+
return true;
1591+
});
1592+
bridge.submit("dispatch workers", "immediate");
1593+
bridge.handle({ type: "fleet", running: 1 });
1594+
settleToollessTurn(bridge);
1595+
expect(shell.session.run).toBe("busy");
1596+
bridge.handle({ type: "fleet", running: 0 });
1597+
expect(drives).toBe(1);
1598+
expect(shell.session.run).toBe("busy");
1599+
bridge.submit("when it finishes, summarize", "queue");
1600+
expect(badgeCount(shell.session)).toBe(1);
1601+
port.clear();
1602+
await Promise.resolve();
1603+
expect(shell.session.run).toBe("idle");
1604+
expect(badgeCount(shell.session)).toBe(0);
1605+
const deliver = port.calls.find((c) => c.op === "deliver");
1606+
expect(deliver).toEqual({
1607+
op: "deliver",
1608+
item: expect.objectContaining({
1609+
text: "when it finishes, summarize",
1610+
kind: "queue",
1611+
}),
1612+
});
1613+
bridge.handle({ type: "connector.reply", data: { content: "" } });
1614+
expect(shell.session.run).toBe("idle");
1615+
bridge.submit("continue the remaining work", "immediate");
1616+
expect(shell.session.run).toBe("busy");
1617+
settleToollessTurn(bridge);
1618+
expect(drives).toBe(2);
1619+
expect(shell.session.run).toBe("busy");
1620+
} finally {
1621+
bridge.dispose();
1622+
shell.dispose();
1623+
}
1624+
},
1625+
{ width: 80, height: 24 },
1626+
);
1627+
});
1628+
1629+
test("occupancy send abort resets the turn without a following reply", async () => {
1630+
await withTestRenderer(
1631+
async (h) => {
1632+
const shell = createAppShell(h.renderer, {
1633+
terminal: { columns: 80, rows: 24 },
1634+
wireKeys: false,
1635+
run: "idle",
1636+
});
1637+
const port = createRecordingPort();
1638+
const nowMs = 0;
1639+
let tick: (() => void) | undefined;
1640+
const prompt = "The fleet has gone dry. Remaining open tasks:\n- t1: keep going (todo)\n";
1641+
const bridge = attachSessionBridge(shell, port, {
1642+
now: () => nowMs,
1643+
stallNoticeMs: 400,
1644+
stallTimeoutMs: 1_000,
1645+
schedule: (fn) => {
1646+
tick = fn;
1647+
return () => {
1648+
tick = undefined;
1649+
};
1650+
},
1651+
});
1652+
try {
1653+
bridge.setDryOpenTaskDriver(() => {
1654+
bridge.beginSystemContinuation(prompt);
1655+
void Promise.resolve().then(() => {
1656+
bridge.abortSystemContinuation();
1657+
});
1658+
return true;
1659+
});
1660+
bridge.submit("dispatch workers", "immediate");
1661+
bridge.handle({ type: "fleet", running: 1 });
1662+
settleToollessTurn(bridge);
1663+
bridge.handle({ type: "fleet", running: 0 });
1664+
expect(bridge.turn.isProcessing).toBe(true);
1665+
expect(tick).toBeDefined();
1666+
await Promise.resolve();
1667+
expect(bridge.turn.isProcessing).toBe(false);
1668+
expect(bridge.turn.awaitingResponse).toBe(false);
1669+
expect(bridge.turn.status).not.toBe("running");
1670+
expect(shell.lockupPhase).toBeNull();
1671+
expect(tick).toBeUndefined();
1672+
} finally {
1673+
bridge.dispose();
1674+
shell.dispose();
1675+
}
1676+
},
1677+
{ width: 80, height: 24 },
1678+
);
1679+
});
1680+
1681+
test("occupancy continuation is what quota auto-retry resubmits", async () => {
1682+
await withTestRenderer(
1683+
async (h) => {
1684+
const shell = createAppShell(h.renderer, {
1685+
terminal: { columns: 80, rows: 24 },
1686+
wireKeys: false,
1687+
run: "idle",
1688+
});
1689+
const port = createRecordingPort();
1690+
let nowMs = 0;
1691+
let tick: (() => void) | undefined;
1692+
const continuation =
1693+
"The fleet has gone dry. Remaining open tasks:\n- t1: keep going (todo)\n";
1694+
const bridge = attachSessionBridge(shell, port, {
1695+
now: () => nowMs,
1696+
schedule: (fn) => {
1697+
tick = fn;
1698+
return () => {
1699+
tick = undefined;
1700+
};
1701+
},
1702+
});
1703+
try {
1704+
bridge.setDryOpenTaskDriver(() => {
1705+
bridge.beginSystemContinuation(continuation);
1706+
return true;
1707+
});
1708+
bridge.submit("dispatch workers", "immediate");
1709+
bridge.handle({ type: "fleet", running: 1 });
1710+
settleToollessTurn(bridge);
1711+
bridge.handle({ type: "fleet", running: 0 });
1712+
port.clear();
1713+
bridge.handle({
1714+
type: "inference.error",
1715+
data: { error: { category: "quota_exhausted", retryAfterMs: 1_000 } },
1716+
});
1717+
nowMs += 10_000;
1718+
tick?.();
1719+
expect(port.calls).toEqual([{ op: "sendImmediate", text: continuation.trim() }]);
1720+
} finally {
1721+
bridge.dispose();
1722+
shell.dispose();
1723+
}
1724+
},
1725+
{ width: 80, height: 24 },
1726+
);
1727+
});
15681728
});
15691729

15701730
describe("syncAgentProgress", () => {

src/tui/runtime-bridge.ts

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -219,6 +219,12 @@ export interface SessionBridge {
219219
* sendWithAttemptIdentity with a system mailbox message.
220220
*/
221221
beginSystemContinuation: (text: string) => void;
222+
/**
223+
* Occupancy send failed after beginSystemContinuation. Re-arm the dry-open
224+
* latch, drop the continuation hold, and idle so follow-ups can drain and a
225+
* later settle can take another occupancy shot.
226+
*/
227+
abortSystemContinuation: () => void;
222228
/**
223229
* Occupancy owner for dry+open continuation. Called once from
224230
* settleRunToIdle when a latched fleet-dry edge is still dry. Return true
@@ -1536,12 +1542,26 @@ export function attachSessionBridge(
15361542
const t = text.trim();
15371543
if (t.length === 0) return;
15381544
bag.pendingEchoes.push(t);
1545+
bag.lastSentMessage = t;
15391546
bag.awaitingContinuationInference = true;
15401547
shell.session = setRunState(shell.session, "busy");
15411548
bag.turn = turnStateOnSubmit(bag.turn, now());
15421549
paintChrome(shell);
15431550
paintPhase();
15441551
},
1552+
abortSystemContinuation: () => {
1553+
if (bag.disposed) return;
1554+
bag.awaitingContinuationInference = false;
1555+
bag.pendingDryOpenDrive = true;
1556+
bag.lastSentMessage = "";
1557+
flushOpenRow(shell, bag);
1558+
bag.turnThinking = null;
1559+
shell.inFlightTool = null;
1560+
shell.session = setRunState(shell.session, "idle");
1561+
drainAtBoundary(shell, bag);
1562+
bag.turn = turnStateOnInterrupt(bag.turn, now());
1563+
paintPhase();
1564+
},
15451565
setDryOpenTaskDriver: (driver) => {
15461566
bag.dryOpenTaskDriver = driver;
15471567
},

0 commit comments

Comments
 (0)