Skip to content

Commit e5bf142

Browse files
fix(tui): bound the primary turn so a silent stall can't freeze it (#1080)
* fix(tui): bound the primary turn so a silent stall can't freeze it * fix(tui): answer CL-8016 critic nits without changing stall-bound behavior Summary: document turnMarkers as diagnostic-only; drive the stop test through production cancelWorkersForStop; import prod ASK_DEADLINE_MS and pin its value and 2x stall-bound sizing; coalesce expireStaleAsks to one notify per batch; drop dead ?. on FleetMailbox.sessions. Verification: bun run typecheck (exit 0); oxfmt --check and oxlint clean on all five files; bun test --randomize green on wiring.stall-bound (4 pass) plus session-store, agent-fleet, wiring.stall-poll, wiring.ask-wake, ask-director (207 pass, 0 fail). * fix(tui): hang stalled ask-wake resurface on shouldAbortForStall After the silent-turn abort, a second wake-only predicate would duplicate that bound. Un-dedupe armed wakes on interrupt instead. * fix(tui): let occupancy win after a stalled ask-wake abort A stalled armed wake was re-surfacing before mailbox mail, so occupancy lost the next turn. Share one interrupt path: interrupt the hung inference, then mail, then wake only if still idle. Expiring asks must abort that wake, not merely disarm it.
1 parent a7387ba commit e5bf142

9 files changed

Lines changed: 775 additions & 18 deletions

‎src/subagent/agent-fleet.ts‎

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -243,6 +243,20 @@ class FleetMailbox {
243243
);
244244
}
245245

246+
/**
247+
* Stop teardown (CL-8016): drop every mailbox record alongside the store
248+
* wipe, so no stale lane outlives the sessions it pins. The records held
249+
* `pinHeld` refs into the store's refcounts; the store teardown clears
250+
* `pinCounts` wholesale, so dropping the records here keeps both sides in
251+
* agreement with nothing left to unpin. Also resets the parked-ask wake
252+
* fingerprint so a later session reusing an id re-surfaces cleanly.
253+
*/
254+
clear(): void {
255+
this.records.clear();
256+
this.lastSurfacedParkedAsks = undefined;
257+
this.sessions.wake();
258+
}
259+
246260
markProviderFailure(id: string): void {
247261
const existing = this.records.get(id);
248262
if (existing === undefined) return;

‎src/subagent/fleet-report.ask-wake.test.ts‎

Lines changed: 16 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -84,4 +84,20 @@ describe("pendingAskWakeText", () => {
8484
expect(text).toContain("using target a1");
8585
expect(text.toLowerCase()).toContain("worker");
8686
});
87+
88+
test("a re-surface count restates that the earlier wake stalled", () => {
89+
const text = pendingAskWakeText(
90+
{
91+
sessionId: "a1",
92+
agentId: "builder",
93+
description: "Build the thing",
94+
question: "Which port?",
95+
questionId: "q1",
96+
},
97+
{ resurface: 1 },
98+
);
99+
expect(text).toContain("Re-surface 1");
100+
expect(text).toContain("stalled");
101+
expect(text).toContain("send_input");
102+
});
87103
});

‎src/subagent/fleet-report.ts‎

Lines changed: 19 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -174,14 +174,30 @@ export const ASK_DIRECTOR_WAKE_PREFIX = "ask_director wake";
174174
* parent, not as the operator being asked — the parent answers via
175175
* send_input itself and only escalates when it genuinely cannot.
176176
*/
177-
export function pendingAskWakeText(wake: PendingAskWake): string {
178-
return [
177+
export function pendingAskWakeText(
178+
wake: PendingAskWake,
179+
options?: { resurface?: number },
180+
): string {
181+
const lines = [
179182
`${ASK_DIRECTOR_WAKE_PREFIX} — worker ${wake.agentId} (${wake.description}) parked question ${wake.questionId} while this session was not collecting:`,
180183
"",
181184
wake.question,
182185
"",
186+
];
187+
// Escalation for a re-surfaced question (CL-8016): the earlier wake turn
188+
// stalled past the bound and was aborted without an answer, so say so and
189+
// restate the routing — otherwise a second identical wake reads as a
190+
// duplicate rather than as proof the first one never landed.
191+
if (options?.resurface !== undefined && options.resurface > 0) {
192+
lines.push(
193+
`Re-surface ${options.resurface}: the earlier wake turn stalled and was aborted without an answer — reconcile against the live question before replying.`,
194+
"",
195+
);
196+
}
197+
lines.push(
183198
`The worker — not the operator — raised this. Answer it with send_input (soft) using target ${wake.sessionId}; do not relay to the operator unless it genuinely needs them.`,
184-
].join("\n");
199+
);
200+
return lines.join("\n");
185201
}
186202

187203
type Change =

‎src/subagent/lifecycle-tools.ts‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -460,9 +460,12 @@ export function createSendInputTool(deps: LifecycleToolDeps): AgentTool {
460460
: {}),
461461
});
462462
if (!outcome.ok) {
463+
// CL-8016: name the teardown when one is recorded — after a stop the
464+
// session is gone, and a bare status would read as "never existed".
465+
const hint = outcome.hint !== undefined ? ` ${outcome.hint}` : "";
463466
return lifecycleResult(
464467
call.id,
465-
`Error: cannot send_input to "${target}" (status: ${outcome.status}).`,
468+
`Error: cannot send_input to "${target}" (status: ${outcome.status}).${hint}`,
466469
);
467470
}
468471
// CL-7331: an interrupt-with-followup is transitional, not terminal.

‎src/subagent/session-store.ts‎

Lines changed: 113 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -290,7 +290,7 @@ export interface SubAgentSessionStore {
290290
},
291291
):
292292
| { ok: true; status: AgentLifecycleStatus }
293-
| { ok: false; status: AgentLifecycleStatus };
293+
| { ok: false; status: AgentLifecycleStatus; hint?: string };
294294
/**
295295
* One pending ask_director per session. `sendInputOne` (soft) resolves it;
296296
* interrupt/settle/close cancel it. Wait JSON projects this, not lifecycle.
@@ -308,6 +308,16 @@ export interface SubAgentSessionStore {
308308
cancelAsk(id: string, reason?: string): boolean;
309309
hasPendingAsk(id: string): boolean;
310310
peekAsk(id: string): { question: string; questionId: string } | undefined;
311+
/**
312+
* Ask deadline (CL-8016): reject every pending ask older than `maxAgeMs`
313+
* with an explicit timeout error naming its question and session, so a
314+
* parked worker whose wake turn stalled can never wait forever. Returns the
315+
* expired descriptors; settlement stays exactly-once through the same
316+
* delete-then-settle path as `resolveAsk`/`cancelAsk`.
317+
*/
318+
expireStaleAsks(
319+
maxAgeMs: number,
320+
): readonly { sessionId: string; questionId: string }[];
311321
/**
312322
* Refcount so wait mailboxes can pin an uncollected result. pruneCompleted
313323
* will not delete a session while its pin count is greater than zero.
@@ -337,10 +347,26 @@ export interface SubAgentSessionStore {
337347
wake(): void;
338348
subscribe(listener: () => void): () => void;
339349
clear(): void;
350+
/**
351+
* Stop teardown (CL-8016): like `clear`, but leaves a tombstone per known
352+
* session so a late `sendInputOne` fails closed naming the teardown instead
353+
* of the bare `not_found` a wipe would give. New ids after teardown still
354+
* report `not_found`; ids the teardown removed name it via `hint`.
355+
*/
356+
teardown(reason?: string): void;
340357
}
341358

342359
export const DEFAULT_CANCEL_REASON = "Cancelled by operator";
343360

361+
/**
362+
* Ask deadline (CL-8016): a parked `ask_director` question older than this
363+
* settles via `expireStaleAsks` with an explicit timeout error, so a worker
364+
* whose wake turn stalled can never wait forever. Twice the stall abort
365+
* bound, so a re-surfaced wake turn gets a full bound cycle to prove the
366+
* parent alive before its questions time out.
367+
*/
368+
export const ASK_DEADLINE_MS = 1_800_000;
369+
344370
const DEFAULT_MAX_COMPLETED = 20;
345371
// CL-7007: sized for fan-out dispatch (dozens of spawn_agent workers), not a
346372
// sidebar list — see maxRetained doc above.
@@ -582,6 +608,10 @@ export function createSubAgentSessionStore(
582608
questionId: string;
583609
resolve: (answer: string) => void;
584610
reject: (reason: unknown) => void;
611+
// Registration clock for the ask deadline (CL-8016): an ask parked
612+
// longer than the bound settles via `expireStaleAsks` instead of
613+
// waiting on a wake turn that may never land.
614+
askedAt: number;
585615
}
586616
>();
587617
const listeners = new Set<() => void>();
@@ -656,7 +686,11 @@ export function createSubAgentSessionStore(
656686
return false;
657687
};
658688

659-
const cancelAskInternal = (id: string, reason: string): boolean => {
689+
const cancelAskInternal = (
690+
id: string,
691+
reason: string,
692+
silent = false,
693+
): boolean => {
660694
const pending = pendingAsks.get(id);
661695
if (pending === undefined) return false;
662696
pendingAsks.delete(id);
@@ -665,7 +699,7 @@ export function createSubAgentSessionStore(
665699
} catch {
666700
// Reject must not throw into settle/interrupt paths.
667701
}
668-
notify();
702+
if (!silent) notify();
669703
return true;
670704
};
671705

@@ -1763,9 +1797,22 @@ export function createSubAgentSessionStore(
17631797
},
17641798
):
17651799
| { ok: true; status: AgentLifecycleStatus }
1766-
| { ok: false; status: AgentLifecycleStatus } {
1800+
| { ok: false; status: AgentLifecycleStatus; hint?: string } {
17671801
const session = sessions.get(id);
1768-
if (session === undefined) return { ok: false, status: "not_found" };
1802+
if (session === undefined) {
1803+
// CL-8016: after a stop teardown the id is gone but the teardown is
1804+
// named — a bare `not_found` would read as "never existed" when the
1805+
// real story is "existed, then the runtime shut down".
1806+
const tombstone = evicted.get(id);
1807+
if (tombstone !== undefined) {
1808+
return {
1809+
ok: false,
1810+
status: tombstone.lifecycleStatus,
1811+
hint: tombstone.hint,
1812+
};
1813+
}
1814+
return { ok: false, status: "not_found" };
1815+
}
17691816
if (session.lifecycle.state !== "running") {
17701817
return { ok: false, status: projectLifecycleStatus(session.lifecycle) };
17711818
}
@@ -1851,7 +1898,7 @@ export function createSubAgentSessionStore(
18511898
if (session === undefined) return false;
18521899
if (session.lifecycle.state !== "running") return false;
18531900
if (pendingAsks.has(id)) return false;
1854-
pendingAsks.set(id, ask);
1901+
pendingAsks.set(id, { ...ask, askedAt: now() });
18551902
mutate(id, () => undefined);
18561903
return true;
18571904
},
@@ -1879,6 +1926,26 @@ export function createSubAgentSessionStore(
18791926
return { question: pending.question, questionId: pending.questionId };
18801927
},
18811928

1929+
expireStaleAsks(
1930+
maxAgeMs: number,
1931+
): readonly { sessionId: string; questionId: string }[] {
1932+
const cutoff = now() - maxAgeMs;
1933+
const expired: { sessionId: string; questionId: string }[] = [];
1934+
for (const [id, pending] of pendingAsks) {
1935+
if (pending.askedAt > cutoff) continue;
1936+
expired.push({ sessionId: id, questionId: pending.questionId });
1937+
cancelAskInternal(
1938+
id,
1939+
`ask_director question ${pending.questionId} for session ${id} expired without an answer after ${maxAgeMs}ms — reply with send_input before the deadline, or not at all`,
1940+
true,
1941+
);
1942+
}
1943+
// One subscriber wake for the batch (none when nothing expired):
1944+
// per-ask notifies would wake N observers for one poll tick.
1945+
if (expired.length > 0) notify();
1946+
return expired;
1947+
},
1948+
18821949
interruptOne(
18831950
id: string,
18841951
): { ok: true } | { ok: false; status: AgentLifecycleStatus } {
@@ -2172,6 +2239,46 @@ export function createSubAgentSessionStore(
21722239
evicted.clear();
21732240
notify();
21742241
},
2242+
2243+
teardown(reason = "Session closed"): void {
2244+
// CL-8016: stop teardown. Cancels the asks with the teardown named
2245+
// (single-resolve preserved through `cancelAskInternal`), releases the
2246+
// handles like `clear`, then leaves a tombstone per removed session so
2247+
// a late `sendInputOne` fails closed naming the teardown instead of a
2248+
// bare `not_found` that would read as "never existed".
2249+
for (const id of pendingAsks.keys())
2250+
cancelAskInternal(id, `ask_director cancelled: ${reason}`);
2251+
for (const id of closeHandles.keys()) releaseHandles(id);
2252+
cancelHandles.clear();
2253+
closeHandles.clear();
2254+
interruptHandles.clear();
2255+
followupHandles.clear();
2256+
deliverHandles.clear();
2257+
for (const id of stashedFollowups.keys())
2258+
dropStashedFollowups(id, reason);
2259+
stashedFollowups.clear();
2260+
for (const session of sessions.values()) {
2261+
evicted.set(session.id, {
2262+
lifecycleStatus: projectLifecycleStatus(session.lifecycle),
2263+
hint: reason,
2264+
});
2265+
}
2266+
if (evicted.size > MAX_EVICTED_TOMBSTONES) {
2267+
const overflow = evicted.size - MAX_EVICTED_TOMBSTONES;
2268+
const keys = evicted.keys();
2269+
for (let i = 0; i < overflow; i++) {
2270+
const oldest = keys.next().value;
2271+
if (oldest === undefined) break;
2272+
evicted.delete(oldest);
2273+
}
2274+
}
2275+
sessions.clear();
2276+
pinCounts.clear();
2277+
runInFlight.clear();
2278+
revisions.clear();
2279+
snapshotCache.clear();
2280+
notify();
2281+
},
21752282
};
21762283
}
21772284

0 commit comments

Comments
 (0)