diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index e18da782c..96cae85ac 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -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. diff --git a/packages/coding-agent/src/advisor/__tests__/advisor.test.ts b/packages/coding-agent/src/advisor/__tests__/advisor.test.ts index bc1d7d9cb..753409029 100644 --- a/packages/coding-agent/src/advisor/__tests__/advisor.test.ts +++ b/packages/coding-agent/src/advisor/__tests__/advisor.test.ts @@ -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 { + 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(() => {}), + }; + 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); diff --git a/packages/coding-agent/src/advisor/runtime.ts b/packages/coding-agent/src/advisor/runtime.ts index 83a0c1308..3b17bbeb8 100644 --- a/packages/coding-agent/src/advisor/runtime.ts +++ b/packages/coding-agent/src/advisor/runtime.ts @@ -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 { - 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(); 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; diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 5e9930c86..60687b18d 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -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") {