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:
@@ -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);
|
||||
|
||||
@@ -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") {
|
||||
|
||||
Reference in New Issue
Block a user