Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -162,7 +162,7 @@ The ChatDirector counts consecutive assistant turns that contain tool calls and

#### Sub-agent stall management

`SubAgentDirector` tracks `lastActivityAt`, updated on every real `inference.done` and `tool.done`. Directors are pure `decide(event, ...)` functions with no timer of their own and the reactor has no proactive "idle" event, so a genuinely silent worker (e.g. parked on a long-running background command with nothing else to do) produces no event for the director to react to. `runSubAgent` (`src/subagent/index.ts`) arms an external interval, at `subAgentStallTimeoutMs`, that pings the same content-less continuation channel the compaction governor uses to re-enter an idle reactor (`requestContinuation`). The director only acts on a ping if the elapsed time since `lastActivityAt` has crossed the timeout. A ping can still be delivered while `execute_tools` is in flight; outstanding call ids from the last `inference.done` reset the silence clock and wait rather than recording `stall-nudge` with a huge `silenceMs`. The first stall past the timeout records `stallNudgeAt` and issues one continuation nudge (asking the worker to check on the background work or report status). Later empty pings inside `subAgentStallTimeoutMs` of that instant wait without stopping or treating the ping as activity; salvage fires only once a ping arrives after that grace with still no `tool.done` / turn-boundary reset. That grace is what keeps two queued interval ticks from salvaging hundreds of milliseconds after the nudge. Any real activity clears `stallNudgeAt`, so a worker that is genuinely working through a slow single turn is never penalized. After the worker has already replied with a terminal report (complete envelope or salvage), further empty continuations — idle-compact meter sync or stall pings — return `wait` instead of falling through to `DefaultDirector.infer`; only a non-empty parent message (`resume_agent` / `send_input`) re-opens the brief. A parked `ask_director` is the same class of wait: empty stall pings are dropped (not deferred) for the park duration, compact-continue hops skipped during the park flush after unpark, and `SubAgentDirector` returns `wait` on any empty continuation that still arrives while the ask is pending so a long park cannot start a billable infer.
`SubAgentDirector` tracks `lastActivityAt`, updated on every real `inference.done` and `tool.done`. Directors are pure `decide(event, ...)` functions with no timer of their own and the reactor has no proactive "idle" event, so a genuinely silent worker (e.g. parked on a long-running background command with nothing else to do) produces no event for the director to react to. `runSubAgent` (`src/subagent/index.ts`) arms an external interval, at `subAgentStallTimeoutMs`, that pings the same content-less continuation channel the compaction governor uses to re-enter an idle reactor (`requestContinuation`). The director nudges only once elapsed time since `lastActivityAt` has crossed the timeout. A ping inside that window, or when stall timing is unset, returns `wait` and does not start a model turn or stamp `lastActivityAt`. A ping can still be delivered while `execute_tools` is in flight; outstanding call ids from the last `inference.done` reset the silence clock and wait rather than recording `stall-nudge` with a huge `silenceMs`. The first stall past the timeout records `stallNudgeAt` and issues one continuation nudge (asking the worker to check on the background work or report status). Later empty pings inside `subAgentStallTimeoutMs` of that instant wait without stopping or treating the ping as activity; salvage fires only once a ping arrives after that grace with still no `tool.done` / turn-boundary reset. That grace is what keeps two queued interval ticks from salvaging hundreds of milliseconds after the nudge. Any real activity clears `stallNudgeAt`, so a worker that is genuinely working through a slow single turn is never penalized. After the worker has already replied with a terminal report (complete envelope or salvage), further empty continuations — idle-compact meter sync or stall pings — return `wait` instead of falling through to `DefaultDirector.infer`; only a non-empty parent message (`resume_agent` / `send_input`) re-opens the brief. A parked `ask_director` is the same class of wait: empty stall pings are dropped (not deferred) for the park duration, compact-continue hops skipped during the park flush after unpark, and `SubAgentDirector` returns `wait` on any empty continuation that still arrives while the ask is pending so a long park cannot start a billable infer.

**Intervention log**: every stop and nudge is appended as one JSONL record to `interventions.jsonl` in the firing worker's trace dir (`src/subagent/intervention-log.ts`), carrying the trigger's measured value beside the threshold it crossed, the provider/model/family it fired on, and the run state at that moment (turns used vs budget, tool calls, read/edit counts). A refused parent re-dispatch is recorded on the parent side, where no worker run exists to record it. The parent also appends one `outcome` record per completed dispatch — the salvage kind `classifyBriefSalvage` assigned, or a clean-complete marker, plus the dispatch count — so the log carries dispatch outcomes as well as interventions, and a stop record can later be read alongside what the dispatch it touched actually produced. Writes are fire-and-forget and swallow their own errors — a diagnostic must not be able to fail a run. `scripts/intervention-forensics.ts` aggregates these across local sessions: per-intervention counts by model family, the measured-value distribution against the threshold, two context columns (stops that fired on runs which had already edited files; stops that fired before half the turn budget was spent — neither is a measured false-positive rate, since either is equally consistent with a correct stop or a wrong one), and outcome counts by kind. This exists because every threshold in this tree was set by judgment and four of those judgments were later reverted — a threshold change is expected to cite this data (CL-6938).

Expand Down
215 changes: 211 additions & 4 deletions src/subagent/nudge-director.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -335,10 +335,11 @@ describe("SubAgentDirector tool failure recovery", () => {
expect(resumedTexts).toHaveLength(1);
expect(resumedTexts?.[0]).toContain("A tool call failed");

const later = inferAction(
const later = actions(
await director.decide(messageReceived(""), longState, caps),
);
expect(ephemeralTexts(later)).toBeUndefined();
expect(later).toEqual([{ type: "wait" }]);
expect(later.some((action) => action.type === "infer")).toBe(false);
});

test("recovery nudge appends to ephemeral turns already on the infer", async () => {
Expand Down Expand Up @@ -462,10 +463,11 @@ describe("SubAgentDirector tool failure recovery", () => {
expect(resumedTexts).toHaveLength(1);
expect(resumedTexts?.[0]).toContain("A tool call failed");

const later = inferAction(
const later = actions(
await director.decide(messageReceived(""), state, caps),
);
expect(ephemeralTexts(later)).toBeUndefined();
expect(later).toEqual([{ type: "wait" }]);
expect(later.some((action) => action.type === "infer")).toBe(false);
});

test("successful nudged infer then later overflow does not resurrect recovery", async () => {
Expand Down Expand Up @@ -1530,6 +1532,211 @@ describe("SubAgentDirector ask_director park wait-guard", () => {
});
});

describe("SubAgentDirector idle stall ping", () => {
const STALL_NUDGE_TEXT =
"No activity has been observed for a while. If you are waiting on a " +
"background command, check its status now; otherwise continue working or " +
"write your report.";

test("empty ping inside the stall window waits and does not infer", async () => {
let now = 8_000_000;
const director = new SubAgentDirector(
"system",
[],
undefined,
1_000,
() => now,
);
const caps = createTestCapabilities();

await director.decide(inferenceDoneText("working"), state, caps);

now += 200;
const early = actions(
await director.decide(messageReceived(""), state, caps),
);
expect(early).toEqual([{ type: "wait" }]);
expect(early.some((action) => action.type === "infer")).toBe(false);
expect(early.some((action) => action.type === "checkpoint")).toBe(false);

// The in-window wait must not restart the silence clock. One stall
// timeout from the original activity still nudges, once.
now = 8_000_000 + 1_000;
const nudge = actions(
await director.decide(messageReceived(""), state, caps),
);
expect(nudge).toContainEqual({
type: "checkpoint",
message: "subagent-stall-nudge",
});
expect(ephemeralTexts(inferAction(nudge))).toEqual([STALL_NUDGE_TEXT]);

now += 200;
const grace = actions(
await director.decide(messageReceived(""), state, caps),
);
expect(grace).toEqual([{ type: "wait" }]);

now += 800;
const stopped = actions(
await director.decide(messageReceived(""), state, caps),
);
expect(stopped).toContainEqual({
type: "checkpoint",
message: "subagent-stalled",
});
expect(stopped.some((action) => action.type === "infer")).toBe(false);
expect(stopped.some((action) => action.type === "reply")).toBe(true);
});

test("outstanding post-compact infer still infers on an in-window empty ping", async () => {
let now = 9_000_000;
let continuations = 0;
const director = new SubAgentDirector(
"system",
[],
() => {
continuations++;
},
60_000,
() => now,
);
const caps = createTestCapabilities();

const compact = actions(
await director.decide(overflowError(), state, caps),
);
expect(compact).toEqual([
{
type: "compact",
compactor: "pruning-compactor",
reason: "context-overflow",
},
]);
expect(continuations).toBe(1);

now += 200;
const resumed = actions(
await director.decide(messageReceived(""), state, caps),
);
expect(resumed.some((action) => action.type === "infer")).toBe(true);
expect(resumed.some((action) => action.type === "wait")).toBe(false);
});

test("idle-threshold fold still compacts on an in-window empty ping", async () => {
let now = 10_000_000;
let continuations = 0;
const director = new SubAgentDirector(
"system",
[],
() => {
continuations++;
},
60_000,
() => now,
);
const caps = createTestCapabilities();

await director.decide(inferenceDone(["read-1"]), longState, caps);
await director.decide(toolDone("read-1"), longState, caps);
const complete = actions(
await director.decide(
inferenceDoneText(REPORT_ENVELOPE, 999_999),
longState,
caps,
),
);
expect(complete.some((action) => action.type === "reply")).toBe(true);
expect(continuations).toBe(1);

now += 200;
const folded = actions(
await director.decide(messageReceived(""), longState, caps),
);
expect(folded).toEqual([
{
type: "compact",
compactor: "pruning-compactor",
reason: "context-threshold",
},
]);
expect(continuations).toBe(2);
expect(folded.some((action) => action.type === "infer")).toBe(false);
expect(folded.some((action) => action.type === "wait")).toBe(false);
});

test("cache-ttl recompress still folds on an in-window empty ping", async () => {
// test-model takes the 10-minute default TTL. The stall window is longer
// so the due fold is still an in-window ping, not a stall nudge.
const activityAt = 11_000_000;
const cacheTtlMs = 10 * 60_000;
const stallTimeoutMs = 15 * 60_000;
let now = activityAt;
let continuations = 0;
const director = new SubAgentDirector(
"system",
[],
() => {
continuations++;
},
stallTimeoutMs,
() => now,
);
const caps = createTestCapabilities();

await director.decide(inferenceDoneText("working"), longState, caps);

now += cacheTtlMs + 1;
const folded = actions(
await director.decide(messageReceived(""), longState, caps),
);
expect(folded).toEqual([
{
type: "compact",
compactor: "pruning-compactor",
reason: "cache-ttl-recompress",
},
]);
expect(continuations).toBe(1);
expect(folded.some((action) => action.type === "infer")).toBe(false);
expect(folded.some((action) => action.type === "wait")).toBe(false);

// Meter-only resume of the empty fold. Same clock: still inside the
// stall window, and this wait must not count as activity either.
const resumed = actions(
await director.decide(messageReceived(""), longState, caps),
);
expect(resumed).toEqual([{ type: "wait" }]);
expect(resumed.some((action) => action.type === "infer")).toBe(false);

// The fold must not stamp lastActivityAt or clear stallNudgeAt. One
// stall timeout from the original activity still nudges, once.
now = activityAt + stallTimeoutMs;
const nudge = actions(
await director.decide(messageReceived(""), longState, caps),
);
expect(nudge).toContainEqual({
type: "checkpoint",
message: "subagent-stall-nudge",
});
expect(ephemeralTexts(inferAction(nudge))).toEqual([STALL_NUDGE_TEXT]);
expect(nudge.some((action) => action.type === "wait")).toBe(false);
});

test("no stall timeout waits on an unsolicited empty continuation", async () => {
const director = new SubAgentDirector("system", [], undefined);
const caps = createTestCapabilities();

await director.decide(inferenceDoneText("working"), state, caps);
const ping = actions(
await director.decide(messageReceived(""), state, caps),
);
expect(ping).toEqual([{ type: "wait" }]);
expect(ping.some((action) => action.type === "infer")).toBe(false);
expect(ping.some((action) => action.type === "checkpoint")).toBe(false);
});
});

function stubAdmission(
notes: { provider: string; until: number }[],
): AdmissionQueue {
Expand Down
10 changes: 8 additions & 2 deletions src/subagent/nudge-director.ts
Original file line number Diff line number Diff line change
Expand Up @@ -265,6 +265,7 @@ export class SubAgentDirector extends DefaultDirector {
requestContinuation,
composedPrompt,
toolDefinitions,
now,
);
this.stallTimeoutMs = stallTimeoutMs;
this.now = now;
Expand Down Expand Up @@ -343,6 +344,10 @@ export class SubAgentDirector extends DefaultDirector {

const stallOutcome = this.checkStallPing(event, capabilities);
if (stallOutcome !== null) return stallOutcome;
// Inside the stall window, or with no stall timeout, an empty
// continuation is not a model turn. Leave lastActivityAt and
// stallNudgeAt alone so the next real silence can still nudge.
if (isEmptyContinuation(event)) return capabilities.wait();

// Keep the running local estimate current on every cycle (tool results and
// rewrites included). Arming still happens inside noteInferenceDone, which
Expand Down Expand Up @@ -522,8 +527,9 @@ export class SubAgentDirector extends DefaultDirector {
* stallNudgeAt. Queued pings that arrive inside the stallTimeoutMs grace
* after that nudge neither stop nor restart the grace (and do not count as
* activity). Stop only when a ping arrives after the grace with still no
* activity. Returns null when this event is not a stall check the director
* should act on (let it fall through as an ordinary continuation).
* activity. Returns null when this ping is not yet silence, or when stall
* timing is unconfigured. decide then waits on an empty continuation
* without stamping the silence clock, instead of inferring.
*/
private checkStallPing(
event: ReactorInboundEvent,
Expand Down
Loading