fix(advisor): never let a failing advisor stall or abort the primary agent

A broken advisor could hold the primary agent on the per-turn catch-up
gate for its full 30s budget while retrying, and an exception thrown from
onTurnEnd propagated into the primary's turn-end callback.

- waitForCatchup resolves immediately while the advisor is mid-failure
  (new #failing latch, set at the failure catch BEFORE any async hook,
  cleared on the next successful turn or reset/seed).
- Every parked waiter is woken the moment an advisor turn fails.
- The turn-end boundary isolates advisor exceptions per advisor: a
  throwing advisor loses its delta, the primary and sibling advisors
  continue untouched.
- A failed render (poisoned message, formatter bug) restores the delta
  cursor and dedup state, so the delta is re-rendered next turn instead
  of silently lost; the size probe itself is guarded and falls back to
  the deferred renderer.
This commit is contained in:
DarkPhilosophy
2026-07-17 07:24:30 +03:00
parent 1b4c292f8f
commit f4c81434d0
4 changed files with 185 additions and 16 deletions
+1
View File
@@ -13,6 +13,7 @@
- Enriched `/advisor status` to show per-advisor status glyphs, model, spend breakdown, and quota window for every configured advisor (including disabled ones), replacing the previous single-advisor-only summary.
### Fixed
- Fixed a failing advisor stalling the primary agent: the per-turn catch-up gate parked the primary for up to its full 30s budget while a broken advisor (unsupported model, dead endpoint, render bug) retried — and an advisor exception could abort the primary's turn-end outright. A failing advisor now releases parked waiters the moment its turn fails (before any async hook), refuses new parks until a turn succeeds, and the turn-end boundary isolates advisor exceptions completely; a failed render restores the delta cursor so nothing is lost when the advisor recovers.
- Fixed advisors retrying a permanently rejected request forever (e.g. `invalid_request_error: model not supported with this account`): unlike quota exhaustion — which already paused with a notice — this class notified once and silently kept re-attempting every turn, re-building heavy context in a shared daemon. The runtime now hard-stops after a permanent rejection or three consecutive backlog-drop cycles, with a visible notice; an explicit reset (`/new`, config rebuild, restart) re-enables it. `waitForCatchup` resolves immediately while halted so the primary agent is never parked on a runtime that cannot drain.
- Fixed the advisor's delta render freezing the whole process on large transcripts (one agent + one advisor was enough): rendering the transcript slice for the advisor ran synchronously on the event loop, and a post-reset replay of a multi-MB session blocked it for 600ms+ per render. Large deltas now render in size- and count-bounded chunks that yield the event loop, with tool call/result pairing preserved across chunk boundaries via a shared whole-delta result index; small per-turn deltas keep the synchronous fast path.
- `retry.fallbackChains` wildcards now support id-prefixed targets and keys: a chain entry like `"openrouter/google/*"` re-prefixes the failing model's bare id (`google-antigravity/gemini-x` → `openrouter/google/gemini-x`), a plain `"provider/*"` entry falling back *from* an aggregator strips the vendor prefix when the target provider only knows the bare id (`openrouter/google/x` → `google-vertex/x`), and an id-prefixed key (`"openrouter/google/*"`) scopes a chain to that provider's ids under the prefix.
@@ -35,6 +35,14 @@ import {
type WatchdogConfigDoc,
} from "..";
/** Poll until the drain loop reaches the asserted state — waitForCatchup
* releases IMMEDIATELY on advisor failure (the primary must never park on a
* failing advisor), so failure-path tests cannot use it as a settle barrier. */
async function settleUntil(predicate: () => boolean, timeoutMs = 2_000): Promise<void> {
const deadline = Date.now() + timeoutMs;
while (!predicate() && Date.now() < deadline) await Bun.sleep(2);
}
describe("advisor", () => {
describe("advisor system prompt", () => {
it("forbids concrete claims about tool arguments hidden from the advisor transcript", () => {
@@ -1896,6 +1904,94 @@ describe("advisor", () => {
expect(promptInputs).toHaveLength(promptsAtHalt);
});
it("never holds the primary agent on the catch-up gate while the advisor is failing", async () => {
// CRITICAL contract: a broken advisor (wrong model, dead endpoint)
// must not stall the primary agent — not even for one hook. The
// onTurnError hook here NEVER resolves, simulating a wedged host
// callback; a parked waiter must still be released the moment the
// advisor turn fails, and later waits must resolve immediately while
// the advisor is mid-failure.
const agent: AdvisorAgent = {
prompt: async () => {
throw new Error("socket hang up");
},
abort: () => {},
reset: () => {},
state: { messages: [] },
};
const messages: AgentMessage[] = [{ role: "user", content: "aaa", timestamp: 1 } as AgentMessage];
const host: AdvisorRuntimeHost = {
snapshotMessages: () => messages,
enqueueAdvice: () => {},
onTurnError: () => new Promise<undefined>(() => {}),
};
const runtime = new AdvisorRuntime(agent, host, 60_000);
runtime.onTurnEnd(messages);
const started = performance.now();
// Parked with a huge budget: must release on the failure, not the timer.
await runtime.waitForCatchup(60_000, 1);
expect(performance.now() - started).toBeLessThan(2_000);
// While the advisor is mid-failure (retry pending), new waits are free.
const again = performance.now();
await runtime.waitForCatchup(60_000, 1);
expect(performance.now() - again).toBeLessThan(100);
runtime.dispose();
}, 10_000);
it("survives a poisoned message without throwing into the caller or losing the delta", async () => {
// CRITICAL contract: an advisor render failure (throwing getter,
// formatter bug) must neither propagate into the primary agent's
// turn-end callback nor park it on the catch-up gate — and the
// unrendered delta must survive for the next turn.
const promptInputs: string[] = [];
const agent: AdvisorAgent = {
prompt: async input => {
promptInputs.push(input);
},
abort: () => {},
reset: () => {},
state: { messages: [] },
};
const messages: AgentMessage[] = [{ role: "user", content: "aaa", timestamp: 1 } as AgentMessage];
const host: AdvisorRuntimeHost = {
snapshotMessages: () => messages,
enqueueAdvice: () => {},
};
const runtime = new AdvisorRuntime(agent, host, 0);
runtime.onTurnEnd(messages);
await settleUntil(() => promptInputs.length >= 1);
expect(promptInputs).toHaveLength(1);
// Poison: reading `content` throws — during the size probe or render.
const poisoned = {
role: "user",
get content(): string {
throw new Error("poisoned message");
},
timestamp: 2,
} as AgentMessage;
messages.push(poisoned);
expect(() => runtime.onTurnEnd(messages)).not.toThrow();
// A parked primary must not wait out the catch-up budget.
const started = performance.now();
await runtime.waitForCatchup(60_000, 1);
expect(performance.now() - started).toBeLessThan(2_000);
await settleUntil(() => runtime.backlog === 0);
// Replace the poison with a healthy message: the cursor was restored,
// so the next turn re-renders from the failed position.
messages[1] = { role: "user", content: "bbb-recovered", timestamp: 2 } as AgentMessage;
messages.push({ role: "user", content: "ccc", timestamp: 3 } as AgentMessage);
runtime.onTurnEnd(messages);
await settleUntil(() => promptInputs.length >= 2);
expect(promptInputs).toHaveLength(2);
expect(promptInputs[1]).toContain("bbb-recovered");
expect(promptInputs[1]).toContain("ccc");
runtime.dispose();
}, 10_000);
// The live incident shape: ONE agent + ONE advisor froze the whole
// process. The advisor's delta render (formatSessionHistoryMarkdown over
// the transcript slice) ran synchronously on the event loop; a post-reset
@@ -2392,7 +2488,7 @@ describe("advisor", () => {
const runtime = new AdvisorRuntime(agent, host, 1);
runtime.onTurnEnd(messages);
await runtime.waitForCatchup(1000, 1);
await settleUntil(() => promptInputs.length >= 2 && runtime.backlog === 0);
expect(promptInputs).toHaveLength(2);
expect(turnErrors).toHaveLength(1);
@@ -2439,7 +2535,7 @@ describe("advisor", () => {
const runtime = new AdvisorRuntime(agent, host, 1);
runtime.onTurnEnd(messages);
await runtime.waitForCatchup(1000, 1);
await settleUntil(() => failures.length >= 1 && runtime.backlog === 0);
expect(promptInputs).toHaveLength(3);
expect(turnErrors.map(error => (error instanceof Error ? error.message : String(error)))).toEqual([
@@ -2495,7 +2591,7 @@ describe("advisor", () => {
const runtime = new AdvisorRuntime(agent, host, 1);
runtime.onTurnEnd(messages);
await runtime.waitForCatchup(1000, 1);
await settleUntil(() => promptInputs.length >= 2 && runtime.backlog === 0);
expect(promptInputs).toHaveLength(2);
expect(turnErrors).toHaveLength(1);
@@ -2634,7 +2730,7 @@ describe("advisor", () => {
const runtime = new AdvisorRuntime(agent, host, 0);
runtime.onTurnEnd(messages);
await runtime.waitForCatchup(1000, 1);
await settleUntil(() => promptInputs.length >= 1 && runtime.backlog === 0);
expect(promptInputs).toHaveLength(1);
expect(resetCalls).toBe(1);
@@ -2643,7 +2739,7 @@ describe("advisor", () => {
messages.push({ role: "user", content: "bbb", timestamp: 2 } as AgentMessage);
runtime.onTurnEnd(messages);
await runtime.waitForCatchup(1000, 1);
await settleUntil(() => promptInputs.length >= 2 && runtime.backlog === 0);
expect(promptInputs).toHaveLength(2);
expect(lengthsBeforePrompt).toEqual([0, 0]);
@@ -2684,7 +2780,7 @@ describe("advisor", () => {
messages.push({ role: "user", content: "bbb", timestamp: 2 } as AgentMessage);
runtime.onTurnEnd(messages);
rejectFirstPrompt(new AdvisorOutputQuarantinedError("quarantined"));
await runtime.waitForCatchup(1000, 1);
await settleUntil(() => promptInputs.length >= 2 && runtime.backlog === 0);
expect(promptInputs).toHaveLength(2);
expect(promptInputs[1]).toContain("aaa");
@@ -2961,7 +3057,7 @@ describe("advisor", () => {
const runtime = new AdvisorRuntime(agent, host, 0);
runtime.onTurnEnd([{ role: "user", content: "quota-turn", timestamp: 1 } as AgentMessage]);
await runtime.waitForCatchup(1000, 1);
await settleUntil(() => promptInputs.length >= 3 && runtime.backlog === 0);
expect(promptInputs).toHaveLength(3);
expect(hookErrors).toHaveLength(2);
@@ -3178,7 +3274,7 @@ describe("advisor", () => {
{ role: "user", content: "triple-credential", timestamp: 1 } as AgentMessage,
];
runtime.onTurnEnd(messages);
await runtime.waitForCatchup(1000, 1);
await settleUntil(() => promptInputs.length >= 3 && runtime.backlog === 0);
expect(promptInputs).toHaveLength(3);
expect(hookErrors).toHaveLength(2);
+68 -7
View File
@@ -289,6 +289,10 @@ export class AdvisorRuntime {
* explicit {@link reset} (config rebuild, /new, session restart).
*/
#halted = false;
/** True from the moment an advisor turn fails until one succeeds (or an
* explicit reset/seed). While set, {@link waitForCatchup} resolves
* immediately: the primary agent NEVER parks on a failing advisor. */
#failing = false;
#latestMessages?: AgentMessage[];
#waiters: CatchupWaiter[] = [];
/** Bumped by every external {@link reset}/{@link dispose}. A drain iteration
@@ -361,12 +365,38 @@ export class AdvisorRuntime {
// formatted in one synchronous call those block the event loop for
// hundreds of milliseconds, freezing EVERY session hosted by a shared
// daemon.
if (
this.#renderBusy === 0 &&
all.length - this.#lastCount <= RENDER_CHUNK_MESSAGES &&
!deltaExceedsSize(all, this.#lastCount, FAST_RENDER_MAX_CHARS)
) {
const render = this.#renderDelta(all, wip);
let fastPath = false;
try {
fastPath =
this.#renderBusy === 0 &&
all.length - this.#lastCount <= RENDER_CHUNK_MESSAGES &&
!deltaExceedsSize(all, this.#lastCount, FAST_RENDER_MAX_CHARS);
} catch (err) {
// A poisoned message (throwing getter) trips the size probe before
// any state mutates. Route it through the deferred renderer, whose
// catch restores the cursor — never through the caller.
logger.warn("advisor delta size probe failed; deferring render", { err: String(err) });
}
if (fastPath) {
let render: string | null = null;
// The render advances #lastCount/#seenContext before formatting can
// throw; snapshot both so a formatter bug loses NOTHING — the next
// turn re-renders this delta.
const cursorBefore = this.#lastCount;
const seenBefore = [...this.#seenContext];
try {
render = this.#renderDelta(all, wip);
} catch (err) {
// A render bug must never propagate into the primary agent's
// turn-end callback — the advisor skips this delta and stops
// gating the catch-up wait, Luna moves on.
this.#lastCount = cursorBefore;
this.#seenContext.clear();
for (const [key, value] of seenBefore) this.#seenContext.set(key, value);
this.#failing = true;
this.#wakeAllWaiters();
logger.warn("advisor delta render failed", { err: String(err) });
}
if (render) {
this.#pending.push({ text: render, turns: 1, wip });
this.#backlog++;
@@ -388,9 +418,20 @@ export class AdvisorRuntime {
return;
}
let render: string | null = null;
// Snapshot the cursor/dedup state: a formatter bug mid-render must
// lose nothing — the next turn re-renders this delta.
const cursorBefore = this.#lastCount;
const seenBefore = [...this.#seenContext];
try {
render = await this.#renderDeltaChunked(all, wip, epoch);
} catch (err) {
if (!this.disposed && this.#epoch === epoch) {
this.#lastCount = cursorBefore;
this.#seenContext.clear();
for (const [key, value] of seenBefore) this.#seenContext.set(key, value);
this.#failing = true;
this.#wakeAllWaiters();
}
logger.warn("advisor delta render failed", { err: String(err) });
}
if (this.disposed || this.#epoch !== epoch) return;
@@ -423,7 +464,17 @@ export class AdvisorRuntime {
}
waitForCatchup(maxMs: number, threshold: number, signal?: AbortSignal): Promise<void> {
if (this.disposed || signal?.aborted || this.#backlog < threshold || this.#quotaExhausted || this.#halted)
if (
this.disposed ||
signal?.aborted ||
this.#backlog < threshold ||
this.#quotaExhausted ||
this.#halted ||
// An advisor mid-failure/retry must NEVER gate the primary agent:
// its backlog cannot drain until the retry cycle resolves, and the
// primary would otherwise park for the full catch-up budget.
this.#failing
)
return Promise.resolve();
const { promise, resolve } = Promise.withResolvers<void>();
let waiter!: CatchupWaiter;
@@ -487,6 +538,7 @@ export class AdvisorRuntime {
this.#epoch++;
this.#quotaExhausted = false;
this.#halted = false;
this.#failing = false;
this.#droppedBacklogs = 0;
this.#resetAdvisorContext(true, true);
}
@@ -502,6 +554,7 @@ export class AdvisorRuntime {
this.#pending = [];
this.#backlog = 0;
this.#consecutiveFailures = 0;
this.#failing = false;
this.#droppedBacklogs = 0;
this.#failureNotified = false;
this.#seenContext.clear();
@@ -808,10 +861,17 @@ export class AdvisorRuntime {
const turnError = getAdvisorTurnError(this.agent.state.messages.slice(messageSnapshot));
if (turnError) throw turnError;
success = true;
this.#failing = false;
this.#consecutiveFailures = 0;
this.#failureNotified = false;
this.#droppedBacklogs = 0;
} catch (err) {
// Release any parked primary-agent waiters IMMEDIATELY — before
// the async onTurnError hook or any retry sleep — and refuse new
// parks until a turn succeeds. A failing advisor must never hold
// the primary on the catch-up gate.
this.#failing = true;
this.#wakeAllWaiters();
// reset()/dispose() aborts the in-flight prompt; treat it as a
// reset, not a transient failure — drop the stale batch.
if (this.#epoch !== epoch) continue;
@@ -842,6 +902,7 @@ export class AdvisorRuntime {
const retryTurnError = getAdvisorTurnError(this.agent.state.messages.slice(retrySnapshot));
if (retryTurnError) throw retryTurnError;
success = true;
this.#failing = false;
this.#consecutiveFailures = 0;
this.#failureNotified = false;
this.#droppedBacklogs = 0;
@@ -2623,7 +2623,18 @@ export class AgentSession {
this.#advisorPrimaryTurnsCompleted++;
if (this.#advisors.length > 0) {
for (const a of this.#advisors) {
if (!a.runtime.disposed) a.runtime.onTurnEnd(messages, { willContinue: context?.willContinue });
if (a.runtime.disposed) continue;
try {
a.runtime.onTurnEnd(messages, { willContinue: context?.willContinue });
} catch (advisorErr) {
// CRITICAL boundary: NOTHING an advisor does may abort the
// primary agent's turn-end. A throwing advisor loses its
// delta; the primary continues untouched.
logger.warn("advisor onTurnEnd threw; delta dropped", {
advisor: a.name,
err: String(advisorErr),
});
}
}
const syncBacklog = this.settings.get("advisor.syncBacklog");
if (syncBacklog !== "off") {