Skip to content

Commit 98c381c

Browse files
committed
Hold occupancy continuation through a stale connector reply
The occupancy shot re-arms processing during inference.done, so the same turn's connector.reply settled idle and drained queued follow-ups. Hold until the continuation's inference.start. Take mailbox reports only after send actually succeeds on the TUI path.
1 parent d3292b8 commit 98c381c

8 files changed

Lines changed: 203 additions & 37 deletions

File tree

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

Lines changed: 139 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,13 @@ describe("shouldDriveOpenTasks", () => {
7272
deferredDryEdge: true,
7373
}),
7474
).toBe(true);
75+
expect(
76+
shouldDriveOpenTasks({
77+
hasOpenTasks: true,
78+
parentProcessing: false,
79+
deferredDryEdge: true,
80+
}),
81+
).toBe(true);
7582
expect(
7683
shouldDriveOpenTasks({
7784
previousRunning: 0,
@@ -158,7 +165,7 @@ describe("collectUncollectedTerminals", () => {
158165
}),
159166
).toBe(true);
160167

161-
const reports = collectUncollectedTerminals(mailbox, sessions.list());
168+
const reports = collectUncollectedTerminals(mailbox, sessions.list(), true);
162169
expect(reports.map((r) => r.agent_id).sort()).toEqual(["done", "fail"]);
163170
expect(reports.find((r) => r.agent_id === "done")).toEqual({
164171
agent_id: "done",
@@ -194,9 +201,11 @@ describe("collectUncollectedTerminals", () => {
194201
return taken;
195202
},
196203
};
197-
const reports = collectUncollectedTerminals(mailbox, [
198-
{ id: "ghost", description: "from store", report: "store report" },
199-
]);
204+
const reports = collectUncollectedTerminals(
205+
mailbox,
206+
[{ id: "ghost", description: "from store", report: "store report" }],
207+
true,
208+
);
200209
expect(reports).toEqual([
201210
{
202211
agent_id: "ghost",
@@ -217,10 +226,30 @@ describe("collectUncollectedTerminals", () => {
217226
peek: (id) => records.get(id),
218227
take: (id) => records.get(id),
219228
};
220-
const reports = collectUncollectedTerminals(mailbox, []);
229+
const reports = collectUncollectedTerminals(mailbox, [], true);
221230
expect(reports[0]?.report?.length).toBe(FLEET_DRY_REPORT_CHARS);
222231
expect(reports[0]?.report?.endsWith("…")).toBe(true);
223232
});
233+
234+
test("consume false peeks without take", () => {
235+
const records = new Map<string, FleetDryMailboxRecord>([
236+
["w1", { status: "done", report: "ok" }],
237+
]);
238+
const mailbox: FleetDryMailbox = {
239+
ids: () => [...records.keys()],
240+
peek: (id) => records.get(id),
241+
take: (id) => {
242+
const existing = records.get(id);
243+
if (existing === undefined) return undefined;
244+
const taken = { ...existing, collected: true };
245+
records.set(id, taken);
246+
return taken;
247+
},
248+
};
249+
const reports = collectUncollectedTerminals(mailbox, [], false);
250+
expect(reports).toEqual([{ agent_id: "w1", status: "done", report: "ok" }]);
251+
expect(records.get("w1")?.collected).not.toBe(true);
252+
});
224253
});
225254

226255
describe("driveOpenTasksAfterFleetDry", () => {
@@ -368,4 +397,109 @@ describe("driveOpenTasksAfterFleetDry", () => {
368397
expect(driven).toBe(false);
369398
expect(records.get("w1")?.collected).not.toBe(true);
370399
});
400+
401+
test("TUI sendWithAttemptIdentity rejection leaves mailbox uncollected", async () => {
402+
const records = new Map<string, FleetDryMailboxRecord>([
403+
["w1", { status: "done", report: "ok" }],
404+
]);
405+
const mailbox: FleetDryMailbox = {
406+
ids: () => [...records.keys()],
407+
peek: (id) => records.get(id),
408+
take: (id) => {
409+
const existing = records.get(id);
410+
if (existing === undefined) return undefined;
411+
const taken = { ...existing, collected: true };
412+
records.set(id, taken);
413+
return taken;
414+
},
415+
};
416+
const sendWithAttemptIdentity = async (): Promise<boolean> => {
417+
await Promise.resolve();
418+
throw new Error("agentProxy.send failed");
419+
};
420+
const driven = driveOpenTasksAfterFleetDry({
421+
deferredDryEdge: true,
422+
openTasks: [openTask],
423+
parentProcessing: false,
424+
mailbox,
425+
lanes: [],
426+
beginSystemContinuation: () => undefined,
427+
send: () => sendWithAttemptIdentity(),
428+
});
429+
expect(driven).toBe(true);
430+
expect(records.get("w1")?.collected).not.toBe(true);
431+
await Promise.resolve();
432+
await Promise.resolve();
433+
expect(records.get("w1")?.collected).not.toBe(true);
434+
});
435+
436+
test("TUI sendWithAttemptIdentity false after handleSendFailure leaves mailbox uncollected", async () => {
437+
const records = new Map<string, FleetDryMailboxRecord>([
438+
["w1", { status: "done", report: "ok" }],
439+
]);
440+
const mailbox: FleetDryMailbox = {
441+
ids: () => [...records.keys()],
442+
peek: (id) => records.get(id),
443+
take: (id) => {
444+
const existing = records.get(id);
445+
if (existing === undefined) return undefined;
446+
const taken = { ...existing, collected: true };
447+
records.set(id, taken);
448+
return taken;
449+
},
450+
};
451+
const sendWithAttemptIdentity = async (): Promise<boolean> => {
452+
await Promise.resolve();
453+
return false;
454+
};
455+
const driven = driveOpenTasksAfterFleetDry({
456+
deferredDryEdge: true,
457+
openTasks: [openTask],
458+
parentProcessing: false,
459+
mailbox,
460+
lanes: [],
461+
beginSystemContinuation: () => undefined,
462+
send: () => sendWithAttemptIdentity(),
463+
});
464+
expect(driven).toBe(true);
465+
await Promise.resolve();
466+
await Promise.resolve();
467+
expect(records.get("w1")?.collected).not.toBe(true);
468+
});
469+
470+
test("TUI sendWithAttemptIdentity true takes mailbox after send resolves", async () => {
471+
const records = new Map<string, FleetDryMailboxRecord>([
472+
["w1", { status: "done", report: "ok" }],
473+
]);
474+
const mailbox: FleetDryMailbox = {
475+
ids: () => [...records.keys()],
476+
peek: (id) => records.get(id),
477+
take: (id) => {
478+
const existing = records.get(id);
479+
if (existing === undefined) return undefined;
480+
const taken = { ...existing, collected: true };
481+
records.set(id, taken);
482+
return taken;
483+
},
484+
};
485+
let resolveSend: ((ok: boolean) => void) | undefined;
486+
const sendWithAttemptIdentity = (): Promise<boolean> =>
487+
new Promise((resolve) => {
488+
resolveSend = resolve;
489+
});
490+
const driven = driveOpenTasksAfterFleetDry({
491+
deferredDryEdge: true,
492+
openTasks: [openTask],
493+
parentProcessing: false,
494+
mailbox,
495+
lanes: [],
496+
beginSystemContinuation: () => undefined,
497+
send: () => sendWithAttemptIdentity(),
498+
});
499+
expect(driven).toBe(true);
500+
expect(records.get("w1")?.collected).not.toBe(true);
501+
resolveSend?.(true);
502+
await Promise.resolve();
503+
expect(records.get("w1")?.collected).toBe(true);
504+
});
371505
});

src/subagent/fleet-dry-drive.ts

Lines changed: 30 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -46,15 +46,21 @@ export interface CollectedWorkerReport {
4646
}
4747

4848
export function shouldDriveOpenTasks(input: {
49-
previousRunning: number;
50-
running: number;
49+
previousRunning?: number | undefined;
50+
running?: number | undefined;
5151
hasOpenTasks: boolean;
5252
parentProcessing: boolean;
5353
deferredDryEdge?: boolean;
5454
}): boolean {
55-
const wentDry = input.running === 0 && input.previousRunning > 0;
55+
const running = input.running ?? 0;
56+
const previousRunning = input.previousRunning ?? 0;
57+
const wentDry = running === 0 && previousRunning > 0;
5658
const dryEdge = wentDry || input.deferredDryEdge === true;
57-
return dryEdge && input.running === 0 && input.hasOpenTasks && !input.parentProcessing;
59+
return dryEdge && running === 0 && input.hasOpenTasks && !input.parentProcessing;
60+
}
61+
62+
function isPromiseLike(value: unknown): value is Promise<unknown> {
63+
return typeof value === "object" && value !== null && "then" in value;
5864
}
5965

6066
function clipField(text: string | undefined): string | undefined {
@@ -112,7 +118,7 @@ function clipCollectedReport(report: CollectedWorkerReport): CollectedWorkerRepo
112118
export function collectUncollectedTerminals(
113119
mailbox: FleetDryMailbox | undefined,
114120
lanes: readonly FleetDryLane[],
115-
consume = true,
121+
consume: boolean,
116122
): CollectedWorkerReport[] {
117123
if (mailbox === undefined) return [];
118124
const byId = new Map(lanes.map((lane) => [lane.id, lane]));
@@ -151,15 +157,15 @@ export function buildFleetDryContinuationPrompt(
151157
}
152158

153159
export function driveOpenTasksAfterFleetDry(args: {
154-
previousRunning: number;
155-
running: number;
160+
previousRunning?: number | undefined;
161+
running?: number | undefined;
156162
openTasks: readonly Task[];
157163
parentProcessing: boolean;
158164
deferredDryEdge?: boolean;
159165
mailbox: FleetDryMailbox | undefined;
160166
lanes: readonly FleetDryLane[];
161167
beginSystemContinuation: (prompt: string) => void;
162-
send: (prompt: string) => void;
168+
send: (prompt: string) => unknown;
163169
}): boolean {
164170
const tasks = [...args.openTasks];
165171
if (
@@ -175,14 +181,26 @@ export function driveOpenTasksAfterFleetDry(args: {
175181
}
176182
const reports = collectUncollectedTerminals(args.mailbox, args.lanes, false);
177183
const prompt = buildFleetDryContinuationPrompt(tasks, reports);
184+
const takeReports = (): void => {
185+
for (const report of reports) {
186+
args.mailbox?.take(report.agent_id);
187+
}
188+
};
178189
try {
179190
args.beginSystemContinuation(prompt);
180-
args.send(prompt);
191+
const sent = args.send(prompt);
192+
if (isPromiseLike(sent)) {
193+
void sent.then(
194+
(result) => {
195+
if (result !== false) takeReports();
196+
},
197+
() => undefined,
198+
);
199+
return true;
200+
}
201+
if (sent !== false) takeReports();
181202
} catch {
182203
return false;
183204
}
184-
for (const report of reports) {
185-
args.mailbox?.take(report.agent_id);
186-
}
187205
return true;
188206
}

src/subagent/index.ts

Lines changed: 1 addition & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -25,18 +25,7 @@ export {
2525
type FleetWatch,
2626
type PendingAskWake,
2727
} from "./fleet-report.js";
28-
export {
29-
buildFleetDryContinuationPrompt,
30-
collectUncollectedTerminals,
31-
driveOpenTasksAfterFleetDry,
32-
FLEET_DRY_CONTINUATION_PREFIX,
33-
projectMailboxRecord,
34-
shouldDriveOpenTasks,
35-
takeAndProjectMailboxRecord,
36-
type CollectedWorkerReport,
37-
type FleetDryLane,
38-
type FleetDryMailbox,
39-
} from "./fleet-dry-drive.js";
28+
export { driveOpenTasksAfterFleetDry } from "./fleet-dry-drive.js";
4029
export {
4130
EMPTY_THRASH_STATE,
4231
nextThrashState,

src/tui/runner/state.ts

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -222,7 +222,7 @@ export interface RunnerState {
222222
attempt: InferenceAttemptIdentity,
223223
providerFailure: ProviderFailureAttempt,
224224
) => void;
225-
sendWithAttemptIdentity?: (message: InboundMessage) => Promise<void>;
225+
sendWithAttemptIdentity?: (message: InboundMessage) => Promise<boolean>;
226226
sendUserPrompt?: (text: string, pending: readonly PendingImageAttachment[]) => Promise<void>;
227227
dispatchCommand?: (name: string, args: string) => void;
228228
newSession?: () => void;

src/tui/runner/submit.ts

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -255,7 +255,7 @@ export function createSubmitPath(
255255
};
256256
state.handleSendFailure = handleSendFailure;
257257

258-
const sendWithAttemptIdentity = async (message: InboundMessage): Promise<void> => {
258+
const sendWithAttemptIdentity = async (message: InboundMessage): Promise<boolean> => {
259259
const attempt = live.attemptIdentity();
260260
const providerFailure = services.providerFailureAttempts.begin(attempt);
261261
try {
@@ -265,8 +265,10 @@ export function createSubmitPath(
265265
// decision on the correlationId signal channel so the parked run
266266
// resumes.
267267
await services.approvalResume.handle(result);
268+
return true;
268269
} catch (error) {
269270
handleSendFailure(error, attempt, providerFailure);
271+
return false;
270272
} finally {
271273
services.providerFailureAttempts.sendSettled(providerFailure);
272274
}

src/tui/runner/wiring.ts

Lines changed: 1 addition & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -171,8 +171,6 @@ export function wirePostStartup(
171171
const send = state.sendWithAttemptIdentity;
172172
if (send === undefined) return false;
173173
return driveOpenTasksAfterFleetDry({
174-
previousRunning: 0,
175-
running: 0,
176174
deferredDryEdge: true,
177175
openTasks: services.directorHolder.instance?.getTasks() ?? [],
178176
parentProcessing: false,
@@ -181,9 +179,7 @@ export function wirePostStartup(
181179
beginSystemContinuation: (prompt) => {
182180
sessionBridge.beginSystemContinuation(prompt);
183181
},
184-
send: (prompt) => {
185-
void send(buildFleetDryContinuationMessage(prompt));
186-
},
182+
send: (prompt) => send(buildFleetDryContinuationMessage(prompt)),
187183
});
188184
});
189185
const unsubscribeFleetReport = services.subAgentSessions.subscribe(() => {

src/tui/runtime-bridge.test.ts

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1476,6 +1476,17 @@ describe("fleet-dry open-task drive (CL-7540)", () => {
14761476
settleToollessTurn(bridge);
14771477
expect(drives).toBe(1);
14781478
expect(shell.session.run).toBe("busy");
1479+
bridge.submit("when it finishes, summarize", "queue");
1480+
expect(badgeCount(shell.session)).toBe(1);
1481+
port.clear();
1482+
bridge.handle({ type: "connector.reply", data: { content: "" } });
1483+
expect(shell.session.run).toBe("busy");
1484+
expect(badgeCount(shell.session)).toBe(1);
1485+
expect(port.calls.some((c) => c.op === "deliver")).toBe(false);
1486+
settleToollessTurn(bridge);
1487+
expect(drives).toBe(1);
1488+
expect(shell.session.run).toBe("idle");
1489+
expect(badgeCount(shell.session)).toBe(0);
14791490
} finally {
14801491
bridge.dispose();
14811492
shell.dispose();

0 commit comments

Comments
 (0)