diff --git a/packages/ai/test/gitlab-duo-workflow-provider.test.ts b/packages/ai/test/gitlab-duo-workflow-provider.test.ts index 52a65c1f6..15573ab7e 100644 --- a/packages/ai/test/gitlab-duo-workflow-provider.test.ts +++ b/packages/ai/test/gitlab-duo-workflow-provider.test.ts @@ -291,6 +291,7 @@ describe("GitLab Duo Workflow provider protocol", () => { it("builds startRequest goal as a bare ChatML transcript with tool-run linkage", () => { const patToken = `${"glpat"}-abcdefgh12345678ijkl`; const sessionCookie = "_gitlab_session=0123456789abcdef0123456789abcdef"; + const credentialTokens = [patToken, sessionCookie]; const replayContext: Context = { systemPrompt: [`OMP system instructions: preserve the local tool bridge. token ${patToken}`], diff --git a/packages/ai/test/pi-native-client.test.ts b/packages/ai/test/pi-native-client.test.ts index 930c23166..194c02524 100644 --- a/packages/ai/test/pi-native-client.test.ts +++ b/packages/ai/test/pi-native-client.test.ts @@ -48,30 +48,39 @@ function stalledBody(bytes: Uint8Array[] = []): ReadableStream { } function delayedBody(chunks: Array<{ atMs: number; bytes: Uint8Array }>): ReadableStream { - let active = true; + let closed = false; + const timers: Timer[] = []; + const clearTimers = () => { + closed = true; + for (const timer of timers) clearTimeout(timer); + timers.length = 0; + }; return new ReadableStream({ start(controller) { + const enqueue = (bytes: Uint8Array) => { + if (!closed) controller.enqueue(bytes); + }; for (const chunk of chunks) { - setTimeout(() => { - if (!active) return; - try { - controller.enqueue(chunk.bytes); - } catch {} - }, chunk.atMs); + if (chunk.atMs <= 0) { + enqueue(chunk.bytes); + } else { + timers.push(setTimeout(() => enqueue(chunk.bytes), chunk.atMs)); + } } - setTimeout( - () => { - if (!active) return; - active = false; - try { - controller.close(); - } catch {} - }, - Math.max(...chunks.map(chunk => chunk.atMs)) + 1, + timers.push( + setTimeout( + () => { + if (!closed) { + clearTimers(); + controller.close(); + } + }, + Math.max(...chunks.map(chunk => chunk.atMs)) + 1, + ), ); }, cancel() { - active = false; + clearTimers(); }, }); } diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index 679c0460d..38f400bf8 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -10,6 +10,9 @@ - Added OpenTelemetry log and metric export support alongside existing trace exports, enabling forwarding of centralized-logger events and GenAI-semconv metrics when configured. - Enhanced `retry.fallbackChains` wildcards to support id-prefixed targets and keys, allowing more flexible model fallback routing across different providers. - Added an opt-in per-project model role storage mode with global fallback from the model selector. +- Added per-advisor on/off toggle (`enabled: false` in `WATCHDOG.yml`): advisors stay in the roster but their runtime is never built — they show `○` in `/advisor status` rather than disappearing. Existing configs are backward-compatible (defaults to `true` when absent). +- Colored the status line's advisor `++` badge by roster health (green all running, yellow quota-exhausted, red failed, dim paused); per-advisor glyphs (`●`/`○`/`✕`) show in `/advisor status`. +- Added real provider quota display (usage percent, window, reset timer) to `/advisor status` and the `/advisor configure` preview. ### Changed @@ -20,6 +23,7 @@ - Changed bundled TTSR rules to warn without interrupting generation. - Renamed the system prompt's project-context section wrapper from `` to `` to prevent collisions with the `task` tool's `context` parameter. - Rendered `read xd://` calls in a compact grouped read view instead of a full tool-execution card. +- 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 @@ -70,6 +74,8 @@ - Fixed `/review` aborting entirely when GitHub rejects a pull request's aggregate diff for exceeding the line limit by falling back to the paginated per-file endpoint. - Fixed `/q` + Enter running `/queue` instead of `/quit` by adding an explicit `q` alias to `/quit`. - Fixed Ctrl+L (`app.display.reset`) not refreshing the dark/light theme on certain terminals by issuing a background re-query before repainting. +- 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 pauses with a notice until an explicit reset — 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. ## [17.0.1] - 2026-07-16 diff --git a/packages/coding-agent/src/advisor/__tests__/advisor.test.ts b/packages/coding-agent/src/advisor/__tests__/advisor.test.ts index 4fced58b2..9e5eec54e 100644 --- a/packages/coding-agent/src/advisor/__tests__/advisor.test.ts +++ b/packages/coding-agent/src/advisor/__tests__/advisor.test.ts @@ -36,6 +36,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", () => { @@ -1777,7 +1785,7 @@ describe("advisor", () => { messages.push({ role: "user", content: "overflowing-current-update", timestamp: 3 } as AgentMessage); runtime.onTurnEnd(messages); - await runtime.waitForCatchup(1000, 1); + await settleUntil(() => promptInputs.length >= 2 && runtime.backlog === 0); expect(promptInputs).toHaveLength(2); for (const input of promptInputs) { @@ -1789,7 +1797,7 @@ describe("advisor", () => { messages.push({ role: "user", content: "post-recovery-update", timestamp: 4 } as AgentMessage); runtime.onTurnEnd(messages); - await runtime.waitForCatchup(1000, 1); + await settleUntil(() => promptInputs.length >= 3 && runtime.backlog === 0); expect(promptInputs).toHaveLength(3); expect(promptInputs[2]).toContain("post-recovery-update"); @@ -1860,7 +1868,7 @@ describe("advisor", () => { messages.push({ role: "user", content: "structured-current-update", timestamp: 2 } as AgentMessage); runtime.onTurnEnd(messages); - await runtime.waitForCatchup(1000, 1); + await settleUntil(() => promptInputs.length >= 2 && runtime.backlog === 0); expect(promptInputs).toHaveLength(2); for (const input of promptInputs) { @@ -1917,7 +1925,7 @@ describe("advisor", () => { messages.push({ role: "user", content: "queued-small-update", timestamp: 3 } as AgentMessage); runtime.onTurnEnd(messages); finishSecondAttempt.resolve(); - await runtime.waitForCatchup(1000, 1); + await settleUntil(() => promptInputs.length >= 3 && runtime.backlog === 0); expect(failingAttempts).toBe(2); expect(promptInputs).toHaveLength(3); @@ -2151,6 +2159,381 @@ describe("advisor", () => { expect(failures).toHaveLength(2); }); + it("halts permanently on an invalid_request rejection instead of retrying forever", async () => { + // The runaway observed live: a provider that refuses the configured + // model outright ("not supported ... (code=invalid_request_error)") + // failed 351 turns/hour in a shared daemon, rebuilding heavy context + // every cycle. One drop cycle must latch the runtime off. + const promptInputs: string[] = []; + const failures: unknown[] = []; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + throw new Error( + "Codex error event: The 'gpt-5.3-codex-spark' model is not supported when using Codex with a ChatGPT account. (code=invalid_request_error)", + ); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const messages: AgentMessage[] = [{ role: "user", content: "aaa", timestamp: 1 } as AgentMessage]; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + notifyFailure: error => failures.push(error), + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + await Bun.sleep(0); + + expect(promptInputs).toHaveLength(3); + expect(failures).toHaveLength(1); + expect(runtime.halted).toBe(true); + + // New deltas must be ignored while halted — no further prompts. + messages.push({ role: "user", content: "bbb", timestamp: 2 } as AgentMessage); + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + expect(promptInputs).toHaveLength(3); + + // The catch-up gate must not park the primary agent on a runtime that + // will never drain again: resolve immediately regardless of maxMs. + await runtime.waitForCatchup(60_000, 0); + + // Explicit reset (config rebuild, /new) re-enables the runtime. + runtime.reset(); + expect(runtime.halted).toBe(false); + }); + + it("halts after three transient drop cycles without an intervening success, but not across successes", async () => { + const promptInputs: string[] = []; + let shouldFail = true; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + if (shouldFail) throw new Error("socket hang up"); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const messages: AgentMessage[] = [{ role: "user", content: "t1", timestamp: 1 } as AgentMessage]; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + notifyFailure: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + const runTurn = async (content: string) => { + messages.push({ role: "user", content, timestamp: messages.length + 1 } as AgentMessage); + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + await Bun.sleep(0); + }; + + // Two failing drop cycles, then a success: the cycle counter resets. + await runTurn("f1"); + await runTurn("f2"); + expect(runtime.halted).toBe(false); + shouldFail = false; + await runTurn("ok"); + expect(runtime.halted).toBe(false); + + // Three CONSECUTIVE drop cycles with no success latch the runtime off. + shouldFail = true; + await runTurn("f3"); + await runTurn("f4"); + expect(runtime.halted).toBe(false); + await runTurn("f5"); + expect(runtime.halted).toBe(true); + const promptsAtHalt = promptInputs.length; + await runTurn("ignored"); + 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 when a post-reset replay rendered a multi-MB transcript. These + // tests pin the correctness contracts for large deltas: complete + // delivery, tool call/result pairing, ordering across interleaved + // turns, and full replay after a mid-render reset. + describe("large-transcript responsiveness", () => { + const bigMessage = (i: number, chars = 5_000): AgentMessage => { + const text = `msg-${i} ${"x".repeat(chars)}`; + return ( + i % 2 + ? { role: "assistant", content: [{ type: "text", text }], timestamp: i } + : { role: "user", content: text, timestamp: i } + ) as AgentMessage; + }; + + const waitForPrompts = async (prompts: string[], count: number, timeoutMs = 10_000): Promise => { + const deadline = Date.now() + timeoutMs; + while (prompts.length < count && Date.now() < deadline) await Bun.sleep(5); + }; + + it("delivers a multi-MB transcript replay completely", async () => { + const promptInputs: string[] = []; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + // ~2000 × 5KB ≈ 10MB replay — the post-reset/first-enable shape. + const messages = Array.from({ length: 2000 }, (_, i) => bigMessage(i)); + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + runtime.onTurnEnd(messages); + await waitForPrompts(promptInputs, 1); + expect(promptInputs).toHaveLength(1); + // Nothing dropped: first and last transcript messages both rendered. + expect(promptInputs[0]).toContain("msg-0 "); + expect(promptInputs[0]).toContain("msg-1999 "); + runtime.dispose(); + }, 20_000); + + it("pairs a toolCall with its non-adjacent toolResult inside one update", async () => { + const promptInputs: string[] = []; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + // The toolCall sits at index 99 and its result arrives 49 messages + // later (index 148), far past any adjacency window: only the + // whole-delta result index can pair them. + const messages: AgentMessage[] = Array.from({ length: 150 }, (_, i) => bigMessage(i, 64)); + messages[99] = { + role: "assistant", + content: [{ type: "toolCall", id: "call-split", name: "read", arguments: { path: "x" } }], + timestamp: 99, + } as unknown as AgentMessage; + messages[100] = { + role: "custom", + customType: "hook", + content: "interleaved", + timestamp: 100, + } as AgentMessage; + messages[148] = { + role: "toolResult", + toolCallId: "call-split", + content: [{ type: "text", text: "result-body" }], + timestamp: 148, + } as AgentMessage; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + runtime.onTurnEnd(messages); + await waitForPrompts(promptInputs, 1); + expect(promptInputs).toHaveLength(1); + expect(promptInputs[0]).toContain("read("); + // The call+result pair rendered as completed, never as a spurious + // in-flight call. + expect(promptInputs[0]).toContain("⇒ ok"); + expect(promptInputs[0]).not.toContain("⇒ pending"); + runtime.dispose(); + }, 20_000); + + it("delivers a single turn carrying a multi-MB payload", async () => { + const promptInputs: string[] = []; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const messages: AgentMessage[] = [{ role: "user", content: "before", timestamp: 1 } as AgentMessage]; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + runtime.onTurnEnd(messages); + await waitForPrompts(promptInputs, 1); + expect(promptInputs).toHaveLength(1); + + // One turn, one message, multi-MB body (an edit-diff-sized payload) + // must deliver completely. + messages.push({ + role: "assistant", + content: [{ type: "text", text: `huge ${"y".repeat(3_000_000)}` }], + timestamp: 2, + } as AgentMessage); + runtime.onTurnEnd(messages); + await waitForPrompts(promptInputs, 2); + expect(promptInputs).toHaveLength(2); + expect(promptInputs[1]).toContain("huge "); + runtime.dispose(); + }, 20_000); + + it("replays the full transcript after a reset lands between renders", async () => { + const promptInputs: string[] = []; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const messages = Array.from({ length: 400 }, (_, i) => bigMessage(i)); + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + runtime.onTurnEnd(messages); + runtime.reset(); + runtime.onTurnEnd(messages); + await waitForPrompts(promptInputs, 1); + // The aborted pre-reset render must not have advanced the cursor: + // the post-reset replay carries the whole transcript. + const replay = promptInputs.find(input => input.includes("msg-0 ") && input.includes("msg-399 ")); + expect(replay).toBeDefined(); + runtime.dispose(); + }, 20_000); + + it("delivers interleaved turns in order without loss", async () => { + const promptInputs: string[] = []; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const messages = Array.from({ length: 300 }, (_, i) => bigMessage(i)); + const host: AdvisorRuntimeHost = { + snapshotMessages: () => messages, + enqueueAdvice: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + runtime.onTurnEnd(messages); + // Second turn arrives immediately behind the first. + messages.push({ role: "user", content: "late-arrival tail", timestamp: 300 } as AgentMessage); + runtime.onTurnEnd(messages); + const deadline = Date.now() + 10_000; + while (Date.now() < deadline && !promptInputs.join("\n").includes("late-arrival tail")) await Bun.sleep(5); + const combined = promptInputs.join("\n"); + // Every message exactly once, ordering preserved. + expect(combined).toContain("msg-0 "); + expect(combined).toContain("msg-299 "); + expect(combined.indexOf("msg-299 ")).toBeGreaterThan(combined.indexOf("msg-0 ")); + expect(combined.indexOf("late-arrival tail")).toBeGreaterThan(combined.indexOf("msg-299 ")); + expect(combined.match(/msg-150 /g)).toHaveLength(1); + expect(combined.match(/late-arrival tail/g)).toHaveLength(1); + runtime.dispose(); + }, 20_000); + }); + it("treats a clean prompt resolution with state.error as a failed turn (real Agent contract)", async () => { // `Agent.#runLoop` catches provider/stream failures internally — it resolves // `prompt()` cleanly and stores the message on `state.error` (e.g. the @@ -2432,7 +2815,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); @@ -2479,7 +2862,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([ @@ -2535,7 +2918,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); @@ -2605,7 +2988,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(1); expect(rollbackCalls).toEqual([0]); @@ -2743,7 +3126,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); @@ -2752,7 +3135,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]); @@ -2793,7 +3176,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"); @@ -2858,6 +3241,445 @@ describe("advisor", () => { }); }); + describe("AdvisorRuntime quota classification", () => { + it("pauses on quota/rate-limit errors and notifies the host without retrying", async () => { + const promptInputs: string[] = []; + let quotaNotified = false; + let failureNotified = false; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + throw new Error("resource_exhausted"); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + notifyFailure: () => { + failureNotified = true; + }, + notifyQuotaExhausted: () => { + quotaNotified = true; + }, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + const messages: AgentMessage[] = [{ role: "user", content: "first", timestamp: 1 } as AgentMessage]; + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + + // Quota path: single prompt attempt, no retries, no generic failure. + expect(promptInputs).toHaveLength(1); + expect(runtime.quotaExhausted).toBe(true); + expect(quotaNotified).toBe(true); + expect(failureNotified).toBe(false); + + // Subsequent turns are skipped while quota-exhausted. + messages.push({ role: "user", content: "second", timestamp: 2 } as AgentMessage); + runtime.onTurnEnd(messages); + await Bun.sleep(0); + expect(promptInputs).toHaveLength(1); + }); + + it("treats 'overloaded' as a transient server error, not quota exhaustion", async () => { + const promptInputs: string[] = []; + const failures: unknown[] = []; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + throw new Error("overloaded: server is at capacity"); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + notifyFailure: error => failures.push(error), + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + const messages: AgentMessage[] = [{ role: "user", content: "first", timestamp: 1 } as AgentMessage]; + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + await Bun.sleep(0); + + // Overloaded follows the 3-retry → notifyFailure path, not the quota path. + expect(promptInputs).toHaveLength(3); + expect(runtime.quotaExhausted).toBe(false); + expect(failures).toHaveLength(1); + }); + it("retains the failed batch in the pending queue on quota error", async () => { + const promptInputs: string[] = []; + let shouldFail = true; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + if (shouldFail) throw new Error("insufficient_quota: rate limit exceeded"); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + notifyQuotaExhausted: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + const messages: AgentMessage[] = [{ role: "user", content: "quota-turn", timestamp: 1 } as AgentMessage]; + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + + // The batch must remain in the queue (backlog > 0) so it's replayed + // once the quota window resets, instead of being silently dropped. + expect(runtime.quotaExhausted).toBe(true); + expect(runtime.backlog).toBeGreaterThan(0); + expect(promptInputs).toHaveLength(1); + expect(promptInputs[0]).toContain("quota-turn"); + + // After reset() clears the quota pause, the next onTurnEnd drains the + // retained batch — proving it was never lost. + shouldFail = false; + runtime.reset(); + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + expect(promptInputs.at(-1)).toContain("quota-turn"); + }); + + it("resolves waitForCatchup immediately when quota is exhausted", async () => { + const agent: AdvisorAgent = { + prompt: async () => { + throw new Error("insufficient_quota"); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + notifyQuotaExhausted: () => {}, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + const messages: AgentMessage[] = [{ role: "user", content: "turn", timestamp: 1 } as AgentMessage]; + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + + expect(runtime.quotaExhausted).toBe(true); + expect(runtime.backlog).toBeGreaterThan(0); + + // waitForCatchup must resolve instantly — a quota-paused advisor can't + // make progress, so blocking the primary agent for 30s is wrong. + const start = Date.now(); + await runtime.waitForCatchup(30_000, 1); + expect(Date.now() - start).toBeLessThan(1000); + }); + it("retries once when onTurnError signals a switched sibling credential", async () => { + const promptInputs: string[] = []; + let firstCall = true; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + if (firstCall) { + firstCall = false; + throw new Error("insufficient_quota: you have exceeded your rate limit"); + } + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + let quotaNotified = false; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + onTurnError: async () => true, + notifyQuotaExhausted: () => { + quotaNotified = true; + }, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + const messages: AgentMessage[] = [{ role: "user", content: "quota-turn", timestamp: 1 } as AgentMessage]; + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + await Bun.sleep(0); + + // Sibling credential switched: retry succeeds, no quota pause. + expect(promptInputs).toHaveLength(2); + expect(runtime.quotaExhausted).toBe(false); + expect(quotaNotified).toBe(false); + expect(runtime.backlog).toBe(0); + }); + + it("requeues when a switched retry produces no assistant response", async () => { + const promptInputs: string[] = []; + const state = { messages: [] as AgentMessage[] }; + let callCount = 0; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + callCount++; + if (callCount === 1) throw new Error("insufficient_quota"); + if (callCount === 2) { + state.messages.push({ role: "user", content: input, timestamp: Date.now() } as AgentMessage); + } + }, + abort: () => {}, + reset: () => {}, + rollbackTo: count => state.messages.splice(count), + state, + }; + const hookErrors: unknown[] = []; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + onTurnError: async error => { + hookErrors.push(error); + return hookErrors.length === 1; + }, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + runtime.onTurnEnd([{ role: "user", content: "quota-turn", timestamp: 1 } as AgentMessage]); + await settleUntil(() => promptInputs.length >= 3 && runtime.backlog === 0); + + expect(promptInputs).toHaveLength(3); + expect(hookErrors).toHaveLength(2); + expect(runtime.backlog).toBe(0); + }); + + it("falls through to quota pause when onTurnError returns false (no sibling)", async () => { + const promptInputs: string[] = []; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + throw new Error("insufficient_quota: you have exceeded your rate limit"); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + let quotaNotified = false; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + onTurnError: async () => false, + notifyQuotaExhausted: () => { + quotaNotified = true; + }, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + const messages: AgentMessage[] = [{ role: "user", content: "first", timestamp: 1 } as AgentMessage]; + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + + // No sibling: single prompt, then quota pause (no retry). + expect(promptInputs).toHaveLength(1); + expect(runtime.quotaExhausted).toBe(true); + expect(quotaNotified).toBe(true); + }); + it("drops stale quota handling when reset happens during onTurnError", async () => { + const promptInputs: string[] = []; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + if (input.includes("stale-turn")) { + throw new Error("insufficient_quota: you have exceeded your rate limit"); + } + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + let quotaNotified = 0; + let hookInvocations = 0; + const { promise: hookEntered, resolve: allowHook } = Promise.withResolvers(); + const { promise: hookProceed, resolve: proceedHook } = Promise.withResolvers(); + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + onTurnError: async () => { + hookInvocations++; + allowHook(); + await hookProceed; + return false; + }, + notifyQuotaExhausted: () => { + quotaNotified++; + }, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + runtime.onTurnEnd([{ role: "user", content: "stale-turn", timestamp: 1 } as AgentMessage]); + await hookEntered; + runtime.reset(); + runtime.onTurnEnd([{ role: "user", content: "fresh-turn", timestamp: 2 } as AgentMessage]); + proceedHook(); + await runtime.waitForCatchup(1000, 1); + + expect(hookInvocations).toBe(1); + expect(promptInputs).toHaveLength(2); + expect(promptInputs[0]).toContain("stale-turn"); + expect(promptInputs[1]).toContain("fresh-turn"); + expect(runtime.quotaExhausted).toBe(false); + expect(runtime.backlog).toBe(0); + expect(quotaNotified).toBe(0); + }); + it("uses generic failure path when switched retry hits a non-quota error", async () => { + const promptInputs: string[] = []; + let callCount = 0; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + callCount++; + if (callCount === 1) { + throw new Error("insufficient_quota: you have exceeded your rate limit"); + } + if (callCount === 2) { + throw new Error("ECONNRESET: socket hang up"); + } + // callCount >= 3: success + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const hookErrors: unknown[] = []; + let quotaNotified = false; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + onTurnError: async error => { + hookErrors.push(error); + return hookErrors.length === 1 ? true : undefined; + }, + notifyQuotaExhausted: () => { + quotaNotified = true; + }, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + const messages: AgentMessage[] = [{ role: "user", content: "mixed-turn", timestamp: 1 } as AgentMessage]; + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + await Bun.sleep(0); + await Bun.sleep(0); + await Bun.sleep(0); + + // Sibling switched (call 1 quota), retry failed with non-quota + // (call 2), then succeeded (call 3). No quota pause, backlog cleared. + expect(promptInputs).toHaveLength(3); + expect(runtime.quotaExhausted).toBe(false); + expect(quotaNotified).toBe(false); + expect(runtime.backlog).toBe(0); + // Hook sees both errors: the original quota (switched) and the + // retry's non-quota (generic path, no switch). + expect(hookErrors).toHaveLength(2); + }); + + it("marks sibling and pauses when switched retry hits a second quota error", async () => { + const promptInputs: string[] = []; + let firstCall = true; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + if (firstCall) { + firstCall = false; + throw new Error("insufficient_quota: you have exceeded your rate limit"); + } + throw new Error("429 Too Many Requests: quota exceeded"); + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const hookErrors: unknown[] = []; + let quotaNotified = false; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + onTurnError: async error => { + hookErrors.push(error); + return hookErrors.length === 1 ? true : undefined; + }, + notifyQuotaExhausted: () => { + quotaNotified = true; + }, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + const messages: AgentMessage[] = [{ role: "user", content: "double-quota", timestamp: 1 } as AgentMessage]; + runtime.onTurnEnd(messages); + await Bun.sleep(0); + await Bun.sleep(0); + await Bun.sleep(0); + + // Both credentials exhausted: retry prompted twice, then entered quota pause. + expect(promptInputs).toHaveLength(2); + expect(runtime.quotaExhausted).toBe(true); + expect(quotaNotified).toBe(true); + // Hook marks both the original credential (switched=true) and the + // newly exhausted sibling on the second quota error. + expect(hookErrors).toHaveLength(2); + }); + + it("keeps rotating while another credential is immediately available", async () => { + const promptInputs: string[] = []; + const agent: AdvisorAgent = { + prompt: async input => { + promptInputs.push(input); + if (promptInputs.length <= 2) { + throw new Error("429 Too Many Requests: quota exceeded"); + } + }, + abort: () => {}, + reset: () => {}, + state: { messages: [] }, + }; + const hookErrors: unknown[] = []; + let quotaNotified = false; + const host: AdvisorRuntimeHost = { + snapshotMessages: () => [], + enqueueAdvice: () => {}, + onTurnError: async error => { + hookErrors.push(error); + return true; + }, + notifyQuotaExhausted: () => { + quotaNotified = true; + }, + }; + const runtime = new AdvisorRuntime(agent, host, 0); + + const messages: AgentMessage[] = [ + { role: "user", content: "triple-credential", timestamp: 1 } as AgentMessage, + ]; + runtime.onTurnEnd(messages); + await settleUntil(() => promptInputs.length >= 3 && runtime.backlog === 0); + + expect(promptInputs).toHaveLength(3); + expect(hookErrors).toHaveLength(2); + expect(runtime.quotaExhausted).toBe(false); + expect(quotaNotified).toBe(false); + expect(runtime.backlog).toBe(0); + }); + }); + describe("advisor default tools", () => { it("defaults to read/grep/glob, a subset of the full grantable tool pool", () => { expect([...ADVISOR_DEFAULT_TOOL_NAMES]).toEqual(["read", "grep", "glob"]); @@ -3183,5 +4005,22 @@ describe("advisor", () => { expect(text).toContain("default"); expect(text).toContain("anthropic/claude-opus"); }); + it("shows disabled advisors with a dim circle marker and toggles them in the detail editor", async () => { + const uiTheme = await getThemeByName("dark"); + if (!uiTheme) throw new Error("theme unavailable"); + setThemeInstance(uiTheme); + const overlay = make({ + advisors: [ + { name: "Active", model: "x-ai/grok-code-fast:high" }, + { name: "Disabled", model: "openai/gpt-4", enabled: false }, + ], + }); + const text = strip(overlay.render(200)); + // The list shows ● for enabled and ○ for disabled. + expect(text).toContain("● Active"); + expect(text).toContain("○ Disabled"); + // The preview of the highlighted (first) advisor shows its enabled status. + expect(text).toContain("● on"); + }); }); }); diff --git a/packages/coding-agent/src/advisor/__tests__/config.test.ts b/packages/coding-agent/src/advisor/__tests__/config.test.ts index 2875d751f..92d334ca0 100644 --- a/packages/coding-agent/src/advisor/__tests__/config.test.ts +++ b/packages/coding-agent/src/advisor/__tests__/config.test.ts @@ -201,9 +201,13 @@ describe("WATCHDOG.yml file round-trip", () => { }); const doc: WatchdogConfigDoc = { - instructions: 'Shared baseline.\nSecond line with: a colon and "quotes".', + instructions: 'Shared baseline.\n\nSecond line with: a colon and "quotes".', advisors: [ - { name: "Architecture", model: "x-ai/grok-code-fast:high", instructions: "Watch module boundaries." }, + { + name: "Architecture", + model: "x-ai/grok-code-fast:high", + instructions: "Watch module boundaries.\nReport coupling.", + }, { name: "Security", tools: ["read", "grep"] }, ], }; @@ -222,11 +226,25 @@ describe("WATCHDOG.yml file round-trip", () => { // Block style (not the flow `{...}` form), so it stays hand-editable. expect(text).toContain("advisors:"); expect(text).not.toMatch(/^\{/); + expect(text).toContain('instructions: |2-\n Shared baseline.\n \n Second line with: a colon and "quotes".'); + expect(text).toContain(" instructions: |2-\n Watch module boundaries.\n Report coupling."); + expect(text).not.toContain("\\n"); const { advisors, sharedInstructions } = await discoverAdvisorConfigs(tmp, tmp); expect(advisors.map(a => a.name)).toEqual(["Architecture", "Security"]); expect(sharedInstructions).toContain("Shared baseline."); }); + it("preserves significant leading whitespace and trailing newlines in block scalars", async () => { + const file = path.join(tmp, "WATCHDOG.yml"); + const whitespaceDoc: WatchdogConfigDoc = { + instructions: " indented first line\nplain second line\n\n", + advisors: [{ name: "Whitespace", instructions: "\n indented after blank\nplain" }], + }; + + await saveWatchdogConfigFile(file, whitespaceDoc); + expect(await loadWatchdogConfigFile(file)).toEqual(whitespaceDoc); + }); + it("round-trips an explicit empty tools list without collapsing it into the default", async () => { const file = path.join(tmp, "WATCHDOG.yml"); const explicitNoToolsDoc: WatchdogConfigDoc = { @@ -291,3 +309,41 @@ describe("resolveAdvisorConfigEditPath", () => { expect(await resolveAdvisorConfigEditPath("project", dirs(tmp))).toBe(path.join(tmp, "WATCHDOG.yml")); }); }); + +describe("per-advisor enabled field", () => { + it("preserves explicit true, explicit false, and absence through save and discovery", async () => { + const tmp = await fsp.mkdtemp(path.join(os.tmpdir(), "omp-advisor-enabled-")); + try { + const doc: WatchdogConfigDoc = { + advisors: [ + { name: "Explicit On", model: "test/model-a", enabled: true }, + { name: "Explicit Off", model: "test/model-b", enabled: false }, + { name: "Default", model: "test/model-c" }, + ], + }; + const file = path.join(tmp, "WATCHDOG.yml"); + await saveWatchdogConfigFile(file, doc); + + const loaded = await loadWatchdogConfigFile(file); + expect(loaded.advisors.map(advisor => advisor.enabled)).toEqual([true, false, undefined]); + + const { advisors } = await discoverAdvisorConfigs(tmp, tmp); + expect(advisors.map(advisor => advisor.enabled)).toEqual([true, false, undefined]); + } finally { + await fsp.rm(tmp, { recursive: true, force: true }); + } + }); + + it("emits explicit boolean values but omits an absent enabled field", () => { + const text = serializeWatchdogConfig({ + advisors: [ + { name: "Explicit On", enabled: true }, + { name: "Explicit Off", enabled: false }, + { name: "Default" }, + ], + }); + expect(text).toContain("enabled: true"); + expect(text).toContain("enabled: false"); + expect(text.match(/enabled:/g)).toHaveLength(2); + }); +}); diff --git a/packages/coding-agent/src/advisor/config.ts b/packages/coding-agent/src/advisor/config.ts index ae1633ed2..01c92c328 100644 --- a/packages/coding-agent/src/advisor/config.ts +++ b/packages/coding-agent/src/advisor/config.ts @@ -22,8 +22,23 @@ export interface AdvisorConfig { model?: string; tools?: string[]; instructions?: string; + /** Per-advisor on/off toggle (default `true`). When `false`, the advisor + * stays in the roster but its runtime is never built — it shows `○` in + * the status line and `/advisor status` rather than disappearing. */ + enabled?: boolean; } +/** + * Runtime health of a single advisor, surfaced in stats and the status line. + * - `running` — actively processing primary turns + * - `paused` — user-toggled off via per-advisor switch (runtime disposed) + * - `quota_exhausted` — provider returned a quota/rate-limit error; the + * runtime auto-retries after a cooldown so it can resume without user action + * - `error` — repeated transient failures; backlog dropped to prevent stall + * - `no_model` — no model resolved for this advisor's role/explicit model + */ +export type AdvisorRuntimeStatus = "running" | "paused" | "quota_exhausted" | "error" | "no_model"; + /** * The result of walking the `WATCHDOG.yml`/`WATCHDOG.yaml` search path: the * deduped advisor roster plus the concatenated top-level `instructions` baseline @@ -39,6 +54,7 @@ const advisorEntrySchema = type({ "model?": "string", "tools?": "string[]", "instructions?": "string", + "enabled?": "boolean", }); const watchdogYamlSchema = type({ @@ -156,6 +172,7 @@ export async function discoverAdvisorConfigs(cwd: string, agentDir?: string): Pr model: entry.model?.trim() || undefined, tools: filterAdvisorTools(entry.tools, item.path), instructions, + enabled: entry.enabled, }); } } @@ -236,38 +253,73 @@ export async function loadWatchdogConfigFile(filePath: string): Promise ({ - name: a.name, - model: a.model?.trim() || undefined, - tools: a.tools === undefined ? undefined : [...a.tools], - instructions: a.instructions?.trim() ? a.instructions : undefined, - })), - }; + const advisors = (result.advisors ?? []).map(a => { + const advisor: AdvisorConfig = { name: a.name }; + if (a.model?.trim()) advisor.model = a.model; + if (a.tools !== undefined) advisor.tools = [...a.tools]; + if (a.instructions?.trim()) advisor.instructions = a.instructions; + if (a.enabled !== undefined) advisor.enabled = a.enabled; + return advisor; + }); + const doc: WatchdogConfigDoc = { advisors }; + if (result.instructions?.trim()) doc.instructions = result.instructions; + return doc; } /** - * Serialize an editable doc back to block-style `WATCHDOG.yml` text via Bun's - * `YAML.stringify` (the same API the repo uses for other hand-editable config), - * omitting empty fields. Round-trips through {@link loadWatchdogConfigFile}. + * Serialize an editable doc back to canonical, hand-editable `WATCHDOG.yml`. + * Multiline instruction fields use literal block scalars while scalar quoting + * delegates to Bun's YAML encoder. Round-trips through {@link loadWatchdogConfigFile}. * Returns `""` for an empty doc. */ -export function serializeWatchdogConfig(doc: WatchdogConfigDoc): string { - const out: { instructions?: string; advisors?: AdvisorConfig[] } = {}; - if (doc.instructions?.trim()) out.instructions = doc.instructions; - if (doc.advisors.length > 0) { - out.advisors = doc.advisors.map(a => { - const entry: AdvisorConfig = { name: a.name }; - if (a.model?.trim()) entry.model = a.model; - if (a.tools !== undefined) entry.tools = [...a.tools]; - if (a.instructions?.trim()) entry.instructions = a.instructions; - return entry; - }); + +function appendYamlString(lines: string[], indent: string, key: string, value: string): void { + const hasSignificantLeadingWhitespace = value.split("\n").some(line => /^[ \t]/.test(line)); + if (!value.includes("\n") || hasSignificantLeadingWhitespace) { + lines.push(`${indent}${key}: ${YAML.stringify(value)}`); + return; } - if (out.instructions === undefined && out.advisors === undefined) return ""; - const text = YAML.stringify(out, null, 2); - return text.endsWith("\n") ? text : `${text}\n`; + const normalized = value.replaceAll("\r\n", "\n"); + let trailingNewlines = 0; + for (let index = normalized.length - 1; index >= 0 && normalized[index] === "\n"; index--) { + trailingNewlines++; + } + const chomp = trailingNewlines === 0 ? "|2-" : trailingNewlines === 1 ? "|2" : "|2+"; + const body = trailingNewlines === 0 ? normalized : normalized.slice(0, -trailingNewlines); + lines.push(`${indent}${key}: ${chomp}`); + for (const line of body.split("\n")) { + lines.push(`${indent} ${line}`); + } + for (let index = 1; index < trailingNewlines; index++) { + lines.push(`${indent} `); + } +} + +export function serializeWatchdogConfig(doc: WatchdogConfigDoc): string { + const lines: string[] = []; + if (doc.instructions?.trim()) appendYamlString(lines, "", "instructions", doc.instructions); + if (doc.advisors.length > 0) { + lines.push("advisors:"); + for (const advisor of doc.advisors) { + lines.push(` - name: ${YAML.stringify(advisor.name)}`); + if (advisor.model?.trim()) lines.push(` model: ${YAML.stringify(advisor.model)}`); + if (advisor.tools !== undefined) { + if (advisor.tools.length === 0) { + lines.push(" tools: []"); + } else { + lines.push(" tools:"); + for (const tool of advisor.tools) { + lines.push(` - ${YAML.stringify(tool)}`); + } + } + } + if (advisor.instructions?.trim()) { + appendYamlString(lines, " ", "instructions", advisor.instructions); + } + if (advisor.enabled !== undefined) lines.push(` enabled: ${advisor.enabled}`); + } + } + return lines.length === 0 ? "" : `${lines.join("\n")}\n`; } /** diff --git a/packages/coding-agent/src/advisor/runtime.ts b/packages/coding-agent/src/advisor/runtime.ts index 92a48abdc..547bc0248 100644 --- a/packages/coding-agent/src/advisor/runtime.ts +++ b/packages/coding-agent/src/advisor/runtime.ts @@ -66,6 +66,22 @@ export interface AdvisorRuntimeHost { onTurnSuccess?(): Promise | void; /** Surface a non-recovering advisor failure to the host UI without adding model-visible context. */ notifyFailure?(error: unknown): void; + /** Signal that the advisor paused on a quota/rate-limit after host-level + * recovery (credential switch, fallback chain) declined. Cleared only by + * an explicit reset (`/new`, config rebuild, session restart). */ + notifyQuotaExhausted?(): void; +} + +/** + * A request rejection that no amount of retrying can fix for this advisor + * configuration: the provider refuses the model/request shape outright (e.g. + * "The 'gpt-5.3-codex-spark' model is not supported when using Codex with a + * ChatGPT account", code=invalid_request_error). Distinct from quota errors, + * which pause via the dedicated quota path and auto-resume on reset. + */ +function isPermanentAdvisorError(error: unknown): boolean { + const message = error instanceof Error ? error.message : String(error); + return /invalid_request_error|model[_ ]not[_ ]found|is not supported when|does not exist/i.test(message); } const ADVISOR_QUARANTINE_PREFIX = "Advisor response quarantined"; @@ -183,6 +199,14 @@ export function buildAdvisorQuarantineSourceText(currentInput: string, messages: */ const MAX_COALESCE_ROUNDS = 3; +const ADVISOR_RENDER_OPTIONS = { + includeThinking: true, + includeToolIntent: true, + watchedRoles: true, + expandPrimaryContext: true, + expandEditDiffs: true, +} as const; + interface PendingDelta { text: string; rawMessages: AgentMessage[]; @@ -235,6 +259,20 @@ export class AdvisorRuntime { #backlog = 0; #consecutiveFailures = 0; #failureNotified = false; + /** Completed 3-failure backlog-drop cycles since the last success/reset. */ + #droppedBacklogs = 0; + /** + * Hard stop after repeated drop cycles or a permanent request rejection + * (e.g. "model not supported"): without it the advisor re-attempts on every + * new delta forever, and in a shared daemon that unbounded churn burns CPU + * and starves every hosted session's event loop. Cleared only by an + * 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 @@ -243,6 +281,13 @@ export class AdvisorRuntime { * being retried/requeued into the post-reset conversation. */ #epoch = 0; disposed = false; + /** Quota/rate-limit pause state. When `true`, the advisor stops processing + * turns and drops new deltas until an explicit {@link reset} clears it + * (triggered by `/new`, config rebuild, or session restart). There is no + * timer-based auto-resume: provider quota windows (5h/7d) are far longer + * than any reasonable timer, and premature retries waste calls and + * re-trigger the same error. */ + #quotaExhausted = false; constructor( private readonly agent: AdvisorAgent, @@ -253,6 +298,16 @@ export class AdvisorRuntime { get backlog(): number { return this.#backlog; } + get quotaExhausted(): boolean { + return this.#quotaExhausted; + } + get failureNotified(): boolean { + return this.#failureNotified; + } + /** True after the runtime hard-stopped on repeated or permanent failures. */ + get halted(): boolean { + return this.#halted; + } /** * True when `#pending` is non-empty while the drain loop is busy — i.e., newer @@ -277,11 +332,32 @@ export class AdvisorRuntime { * the delta and forwarded to the reprime path so it is never silently dropped. */ onTurnEnd(messages?: AgentMessage[], opts?: { willContinue?: boolean }): void { - if (this.disposed) return; + if (this.disposed || this.#quotaExhausted || this.#halted) return; const all = messages ?? this.host.snapshotMessages(); this.#latestMessages = all; const wip = opts?.willContinue ?? false; - const rendered = this.#renderDelta(all, wip); + let rendered: Omit | null = null; + // #renderDelta advances the cursor/prefix/dedup state before formatting + // can throw; snapshot them so a formatter bug loses NOTHING — the next + // turn re-renders this delta (a prefix change mid-render self-heals via + // the fingerprint scan, at worst costing one full replay). + const cursorBefore = this.#lastCount; + const prefixBefore = this.#deliveredPrefix.slice(); + const seenBefore = [...this.#seenContext]; + try { + rendered = 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 until a turn succeeds. + this.#lastCount = cursorBefore; + this.#deliveredPrefix = prefixBefore; + 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 (rendered) { this.#pending.push({ ...rendered, turns: 1 }); this.#backlog++; @@ -291,7 +367,18 @@ export class AdvisorRuntime { } waitForCatchup(maxMs: number, threshold: number, signal?: AbortSignal): Promise { - if (this.disposed || signal?.aborted || this.#backlog < threshold) return Promise.resolve(); + 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; const finish = (): void => { @@ -353,6 +440,25 @@ export class AdvisorRuntime { } } + /** + * Account one completed 3-failure backlog-drop cycle. Repeated cycles (or a + * single permanent request rejection, e.g. "model not supported") hard-stop + * the runtime: in a shared daemon, unbounded advisor churn re-builds heavy + * context on every new delta and starves every hosted session's event loop. + * Only an explicit {@link reset} (config rebuild, /new, restart) resumes. + */ + #noteDroppedBacklog(error: unknown): void { + this.#droppedBacklogs++; + if (this.#droppedBacklogs < 3 && !isPermanentAdvisorError(error)) return; + this.#halted = true; + this.#pending = []; + this.#wakeAllWaiters(); + logger.warn("advisor halted after repeated failures; use /advisor or reload config to re-enable", { + droppedBacklogs: this.#droppedBacklogs, + err: String(error), + }); + } + /** * Re-prime the advisor after a history rewrite (compaction, session * switch/resume, branch). Clears the advisor's own (non-persisted) context @@ -362,6 +468,10 @@ export class AdvisorRuntime { */ reset(): void { this.#epoch++; + this.#quotaExhausted = false; + this.#halted = false; + this.#failing = false; + this.#droppedBacklogs = 0; this.#resetAdvisorContext(true, true); } @@ -380,6 +490,8 @@ export class AdvisorRuntime { this.#pending = []; this.#backlog = 0; this.#consecutiveFailures = 0; + this.#failing = false; + this.#droppedBacklogs = 0; this.#failureNotified = false; this.#clearSeenContext(); this.#wakeAllWaiters(); @@ -392,13 +504,7 @@ export class AdvisorRuntime { if (delta.length === 0) return null; const obfuscator = this.host.obfuscator; const formattedDelta = obfuscator?.hasSecrets() ? obfuscateAdvisorDelta(obfuscator, delta) : delta; - const md = formatSessionHistoryMarkdown(formattedDelta, { - includeThinking: true, - includeToolIntent: true, - watchedRoles: true, - expandPrimaryContext: true, - expandEditDiffs: true, - }); + const md = formatSessionHistoryMarkdown(formattedDelta, ADVISOR_RENDER_OPTIONS); if (!md.trim()) return null; const heading = wip ? "### Session update [in progress — more steps follow]" : "### Session update"; return `${heading}\n\n${md}`; @@ -683,8 +789,10 @@ 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; if (this.host.onTurnSuccess) { try { await this.host.onTurnSuccess(); @@ -696,6 +804,12 @@ export class AdvisorRuntime { // 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; + // 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(); const failedMessages = this.agent.state.messages.slice(messageSnapshot); const terminalFailure = this.#terminalAssistantFailure(messageSnapshot); const terminalFailureId = @@ -742,6 +856,32 @@ export class AdvisorRuntime { }); continue; } + if (AIError.isUsageLimit(err)) { + // Host recovery (credential switch / fallback chain) declined: + // pause on the quota latch instead of burning retries — provider + // quota windows (5h/7d) outlast any retry budget. The batch is + // requeued and the backlog stays visible so reset() replays it. + logger.warn("advisor quota exhausted", { err: String(err) }); + this.#quotaExhausted = true; + this.#consecutiveFailures = 0; + this.#failureNotified = false; + this.#clearSeenContext(); + this.#pending.unshift({ + text: batch, + rawMessages, + renderRevision: this.#renderRevision, + turns: finalTurns, + wip, + overflowRecovery: recoveringOverflow || undefined, + }); + this.#wakeAllWaiters(); + try { + this.host.notifyQuotaExhausted?.(); + } catch (notifyErr) { + logger.warn("advisor quota notification failed", { err: String(notifyErr) }); + } + break; + } if (!terminalFailureRetriable) { logger.warn("advisor terminal failure is non-retriable; dropping bounded batch"); this.#notifyFailureOnce(err); @@ -749,6 +889,7 @@ export class AdvisorRuntime { // The dropped batch may carry primary-context we never delivered; drop // the seen-state too so queued raw deltas re-expand before delivery. this.#clearSeenContext(); + this.#noteDroppedBacklog(err); success = true; } else if (contextOverflow) { this.#clearAdvisorContextAtCurrentCursor(); @@ -782,6 +923,7 @@ export class AdvisorRuntime { // The dropped batch may carry primary-context we never delivered; drop // the seen-state too so queued raw deltas re-expand before delivery. this.#clearSeenContext(); + this.#noteDroppedBacklog(err); success = true; } else { this.#pending.unshift({ diff --git a/packages/coding-agent/src/modes/components/advisor-config.ts b/packages/coding-agent/src/modes/components/advisor-config.ts index 6543a7aae..2d0673f81 100644 --- a/packages/coding-agent/src/modes/components/advisor-config.ts +++ b/packages/coding-agent/src/modes/components/advisor-config.ts @@ -16,7 +16,7 @@ * `save` callback. */ import type { ThinkingLevel } from "@oh-my-pi/pi-agent-core"; -import type { Model } from "@oh-my-pi/pi-ai"; +import type { Model, UsageReport } from "@oh-my-pi/pi-ai"; import { getSupportedEfforts } from "@oh-my-pi/pi-catalog/model-thinking"; import { type Component, @@ -38,6 +38,9 @@ import { import type { ModelRegistry } from "../../config/model-registry"; import { formatModelSelectorValue } from "../../config/model-resolver"; import type { Settings } from "../../config/settings"; +import type { PerAdvisorStat } from "../../session/agent-session"; +import type { OAuthAccountIdentity } from "../../session/auth-storage"; +import { formatCompactQuota } from "../controllers/command-controller"; import { getSelectListTheme, theme } from "../theme/theme"; import { HookEditorComponent } from "./hook-editor"; import { buildBrowserItems, ModelBrowser, sortModelItems } from "./model-browser"; @@ -63,6 +66,11 @@ export interface AdvisorConfigCallbacks { requestRender: () => void; /** Surface a transient status/warning line to the user. */ notify: (message: string) => void; + /** Live advisor usage stats; lets the preview show tokens/cost per advisor. */ + getAdvisorStats?: () => PerAdvisorStat[]; + getUsageReports?: () => Promise; + /** Resolve the active OAuth identity for quota filtering (per-advisor account stickiness). */ + resolveActiveAccount?: (provider: string, sessionId?: string) => OAuthAccountIdentity | undefined; } export interface AdvisorConfigDeps { @@ -126,6 +134,8 @@ export class AdvisorConfigOverlayComponent implements Component { #cb: AdvisorConfigCallbacks; #scope: AdvisorConfigScope; #doc: WatchdogConfigDoc; + /** Cached usage reports (quota/window/reset) prefetched on overlay open. */ + #cachedReports: UsageReport[] | null = null; #dirty = false; #screen: Screen = "list"; @@ -157,6 +167,16 @@ export class AdvisorConfigOverlayComponent implements Component { this.#doc = doc; this.#ensureRosterVisible(); this.#showList(); + // Prefetch usage reports for quota display; non-fatal if unavailable. + if (callbacks.getUsageReports) { + void callbacks + .getUsageReports() + .then(r => { + this.#cachedReports = r; + this.#cb.requestRender(); + }) + .catch(() => {}); + } } // ───────────────────────────── render ───────────────────────────── @@ -275,6 +295,7 @@ export class AdvisorConfigOverlayComponent implements Component { const lines = [ theme.bold(advisor.name || "(unnamed)"), "", + `${theme.fg("dim", "Enabled:")} ${advisor.enabled === false ? "○ off" : "● on"}`, `${theme.fg("dim", "Model:")} ${model}`, `${theme.fg("dim", "Tools:")} ${tools}`, "", @@ -282,6 +303,34 @@ export class AdvisorConfigOverlayComponent implements Component { ]; const instr = advisor.instructions?.trim(); lines.push(...(instr ? wrap(instr, bodyWidth) : [theme.fg("muted", "(none)")])); + // Show live usage stats when available from the session. + const liveStat = this.#cb.getAdvisorStats?.()?.find(s => s.name === (advisor.name || "default")); + if (liveStat && (liveStat.status === "running" || liveStat.status === "quota_exhausted")) { + lines.push("", theme.fg("dim", "Usage:")); + const spendParts: string[] = [ + `${liveStat.tokens.input.toLocaleString()} in`, + `${liveStat.tokens.output.toLocaleString()} out`, + ]; + if (liveStat.tokens.cacheRead > 0) spendParts.push(`${liveStat.tokens.cacheRead.toLocaleString()} cache`); + lines.push(theme.fg("dim", ` Tokens: ${spendParts.join(", ")}`)); + if (liveStat.cost > 0) lines.push(theme.fg("dim", ` Cost: $${liveStat.cost.toFixed(4)}`)); + if (liveStat.contextWindow > 0) { + const pct = Math.round((liveStat.contextTokens / liveStat.contextWindow) * 100); + lines.push( + theme.fg( + "dim", + ` Context: ${liveStat.contextTokens.toLocaleString()}/${liveStat.contextWindow.toLocaleString()} (${pct}%)`, + ), + ); + } + } + const quotaProvider = + (advisor.model?.includes("/") ? advisor.model.split("/")[0] : null) ?? liveStat?.model?.provider; + if (this.#cachedReports && quotaProvider) { + const activeAccount = this.#cb.resolveActiveAccount?.(quotaProvider, liveStat?.sessionId); + const quota = formatCompactQuota(quotaProvider, this.#cachedReports, Date.now(), activeAccount); + if (quota) lines.push(theme.fg("dim", ` ${quota}`)); + } return lines.map(line => truncateToWidth(line, bodyWidth)); } @@ -311,7 +360,8 @@ export class AdvisorConfigOverlayComponent implements Component { advisor.name === "default" && !advisor.model?.trim() && advisor.tools === undefined && - !advisor.instructions?.trim() + !advisor.instructions?.trim() && + advisor.enabled !== false ); } @@ -325,7 +375,7 @@ export class AdvisorConfigOverlayComponent implements Component { this.#ensureRosterVisible(); const items: SelectItem[] = this.#doc.advisors.map((advisor, index) => ({ value: `advisor:${index}`, - label: advisor.name || "(unnamed)", + label: `${advisor.enabled === false ? "○" : "●"} ${advisor.name || "(unnamed)"}`, description: this.#advisorSummary(advisor), })); items.push({ value: "add", label: "+ Add advisor" }); @@ -395,6 +445,11 @@ export class AdvisorConfigOverlayComponent implements Component { const toolsDescription = formatAdvisorTools(advisor.tools, "no tools"); const items: SelectItem[] = [ { value: "name", label: "Name", description: advisor.name }, + { + value: "toggleEnabled", + label: "Enabled", + description: advisor.enabled === false ? "○ off" : "● on", + }, { value: "model", label: "Model", description: modelDescription }, ]; if (advisor.model?.trim()) { @@ -414,6 +469,13 @@ export class AdvisorConfigOverlayComponent implements Component { #onDetailSelect(index: number, field: string): void { switch (field) { + case "toggleEnabled": { + const a = this.#doc.advisors[index]; + a.enabled = a.enabled === false ? undefined : false; + this.#dirty = true; + this.#showDetail(index); + return; + } case "name": this.#showNameEditor(index); return; diff --git a/packages/coding-agent/src/modes/components/hook-editor.ts b/packages/coding-agent/src/modes/components/hook-editor.ts index c56bc9d12..b7a6564d7 100644 --- a/packages/coding-agent/src/modes/components/hook-editor.ts +++ b/packages/coding-agent/src/modes/components/hook-editor.ts @@ -20,6 +20,13 @@ import { DynamicBorder } from "./dynamic-border"; export interface HookEditorOptions { /** When true, use prompt-style keybindings with the legacy ask prompt chrome. */ promptStyle?: boolean; + /** + * Max rows the inner Editor may occupy. When omitted, the editor is + * bounded to the current terminal height minus the component's chrome + * (≈10 rows) so long content scrolls instead of pushing the submit + * hint out of view. + */ + maxHeight?: number; } /** Interactive multiline dialog used by hooks and the ask tool's Other response. */ @@ -64,6 +71,11 @@ export class HookEditorComponent extends Container implements Focusable { this.#editor.setPromptGutter("> "); this.#editor.disableSubmit = true; } + // Bound the editor so long content scrolls instead of pushing the + // submit hint off-screen. Caller may override via options.maxHeight. + const termRows = this.#tui.terminal?.rows ?? process.stdout.rows ?? 40; + this.#editor.setMaxHeight(options?.maxHeight ?? Math.max(3, termRows - 12)); + this.#editor.setScrollbarVisible(true); if (prefill) { this.#editor.setText(prefill); } diff --git a/packages/coding-agent/src/modes/components/status-line/component.test.ts b/packages/coding-agent/src/modes/components/status-line/component.test.ts index 0f5efdeb9..2a63aeb2d 100644 --- a/packages/coding-agent/src/modes/components/status-line/component.test.ts +++ b/packages/coding-agent/src/modes/components/status-line/component.test.ts @@ -36,6 +36,7 @@ function makeSessionWithLastMessage(lastMessage: unknown, prewalkArmed: boolean getPrewalkState: () => (prewalkArmed ? { target: { id: "cheap-model", provider: "openai" } } : undefined), getAsyncJobSnapshot: () => undefined, isAdvisorActive: () => false, + getAdvisorStatusOverview: () => ({ configured: false, advisors: [] }), isFastModeActive: () => false, configuredThinkingLevel: () => undefined, modelRegistry: { diff --git a/packages/coding-agent/src/modes/components/status-line/segments.ts b/packages/coding-agent/src/modes/components/status-line/segments.ts index 079103426..f4700a414 100644 --- a/packages/coding-agent/src/modes/components/status-line/segments.ts +++ b/packages/coding-agent/src/modes/components/status-line/segments.ts @@ -129,9 +129,9 @@ const modelSegment: StatusLineSegment = { // Fast-mode icon and thinking-level suffix trail the model name and are // colored together with it as `statusLineModel`. The advisor "++" badge - // sits between the name and that tail in `accent`, so it reads as a - // distinct marker. theme.fg resets only the fg, so the spans are - // concatenated (not nested) to keep each color intact. + // sits between the name and that tail, so it reads as a distinct marker. + // theme.fg resets only the fg, so the spans are concatenated (not + // nested) to keep each color intact. let tail = ""; if (ctx.session.isFastModeActive() && theme.icon.fast) { tail += ` ${theme.icon.fast}`; @@ -141,10 +141,25 @@ const modelSegment: StatusLineSegment = { } // `statusLineModel` is aliased to `accent` in many themes, so the badge - // uses `success` to stay visibly distinct from the model name color. + // uses status colors to stay visibly distinct from the model name color. let content = theme.fg("statusLineModel", withIcon(modelIcon, modelName)); - if (ctx.session.isAdvisorActive()) { - content += theme.fg("success", "++"); + // Advisor "++" badge, colored by the worst status in the roster: + // success = all running, warning = quota-exhausted, error = failed, + // dim = everything paused/no-model. Per-advisor detail lives in + // `/advisor status`. + // Optional chaining: lightweight session doubles (test mocks) that don't + // implement getAdvisorStatusOverview skip the badge instead of crashing. + const advisorStats = ctx.session.getAdvisorStatusOverview?.(); + if (advisorStats?.configured && advisorStats.advisors.length > 0) { + const statuses = advisorStats.advisors.map(a => a.status); + const badgeColor = statuses.includes("error") + ? "error" + : statuses.includes("quota_exhausted") + ? "warning" + : statuses.includes("running") + ? "success" + : "dim"; + content += theme.fg(badgeColor, "++"); } if (tail) { content += theme.fg("statusLineModel", tail); diff --git a/packages/coding-agent/src/modes/controllers/command-controller.ts b/packages/coding-agent/src/modes/controllers/command-controller.ts index 6f58c51ee..6d21be4de 100644 --- a/packages/coding-agent/src/modes/controllers/command-controller.ts +++ b/packages/coding-agent/src/modes/controllers/command-controller.ts @@ -349,46 +349,119 @@ export class CommandController { this.ctx.present([new Spacer(1), new Text(info, 1, 0)]); } + static readonly #advisorStatusGlyph: Record = { + running: "●", + paused: "○", + no_model: "○", + quota_exhausted: "✕", + error: "✕", + }; + + static readonly #advisorStatusLabel: Record = { + running: "running", + paused: "off", + no_model: "no model", + quota_exhausted: "quota exhausted", + error: "error", + }; + async handleAdvisorStatusCommand(): Promise { const stats = this.ctx.session.getAdvisorStats(); - if (!stats.active) { - this.ctx.present([ - new Spacer(1), - new Text( - stats.configured - ? "Advisor setting is enabled, but no model is assigned to the 'advisor' role." - : "Advisor is disabled.", - 1, - 0, - ), - ]); + if (!stats.configured) { + this.ctx.present([new Spacer(1), new Text("Advisor is disabled.", 1, 0)]); return; } - if (stats.advisors.length > 1) { + // Fetch live quota data (cached 5 min by the auth-gateway) so we can show + // real usage windows/reset timers per advisor provider. Non-fatal when absent. + const usageProvider = this.ctx.session as { fetchUsageReports?: () => Promise }; + let usageReports: UsageReport[] | null = null; + if (usageProvider.fetchUsageReports) { + try { + usageReports = await usageProvider.fetchUsageReports(); + } catch { + // Network/auth failure is non-fatal — just skip the quota line. + } + } + // Resolve the active OAuth identity for each advisor's provider so quota + // filtering matches the credential actually in use (not sibling accounts). + const resolveActiveAdvisorAccount = (provider: string, sessionId?: string): OAuthAccountIdentity | undefined => + this.ctx.session.modelRegistry.authStorage.getOAuthAccountIdentity( + provider, + sessionId ?? this.ctx.session.sessionId, + ); + const nowMs = Date.now(); + // Roster view: show every configured advisor with its status, even when + // none are live (all paused/no-model). The old code returned a generic + // message that hid the per-advisor state the user needs to act on. + if (stats.advisors.length > 1 || (stats.configured && !stats.active)) { let info = `${theme.bold("Advisor Status")} (${stats.advisors.length} advisors)\n`; for (const a of stats.advisors) { - const ctx = - a.contextWindow > 0 - ? `${a.contextTokens.toLocaleString()} / ${a.contextWindow.toLocaleString()} (${Math.round((a.contextTokens / a.contextWindow) * 100)}%)` - : `${a.contextTokens.toLocaleString()}`; - info += `\n${theme.bold(a.name)}\n`; - info += `${theme.fg("dim", "Model:")} ${a.model.provider}/${a.model.id}\n`; - info += `${theme.fg("dim", "Context:")} ${ctx}\n`; - info += `${theme.fg("dim", "Messages:")} ${a.messages.total.toLocaleString()}\n`; - info += `${theme.fg("dim", "Spend:")} ${a.tokens.input.toLocaleString()} in / ${a.tokens.output.toLocaleString()} out`; - if (a.cost > 0) info += `, $${a.cost.toFixed(4)}`; - info += "\n"; + const glyph = CommandController.#advisorStatusGlyph[a.status] ?? "?"; + const label = CommandController.#advisorStatusLabel[a.status] ?? a.status; + const color = + a.status === "running" + ? "success" + : a.status === "quota_exhausted" || a.status === "error" + ? "error" + : "dim"; + info += `\n${theme.fg(color, glyph)} ${theme.bold(a.name)} ${theme.fg("dim", `[${label}]`)}\n`; + if (a.model) { + info += `${theme.fg("dim", "Model:")} ${a.model.provider}/${a.model.id}\n`; + } + if (a.model && usageReports) { + const quota = formatCompactQuota( + a.model.provider, + usageReports, + nowMs, + resolveActiveAdvisorAccount(a.model.provider, a.sessionId), + ); + if (quota) info += `${theme.fg("dim", quota)}\n`; + } + if (a.status === "running" || a.status === "quota_exhausted") { + const ctx = + a.contextWindow > 0 + ? `${a.contextTokens.toLocaleString()} / ${a.contextWindow.toLocaleString()} (${Math.round((a.contextTokens / a.contextWindow) * 100)}%)` + : `${a.contextTokens.toLocaleString()}`; + info += `${theme.fg("dim", "Context:")} ${ctx}\n`; + info += `${theme.fg("dim", "Messages:")} ${a.messages.total.toLocaleString()}\n`; + info += `${theme.fg("dim", "Spend:")} ${a.tokens.input.toLocaleString()} in / ${a.tokens.output.toLocaleString()} out`; + if (a.cost > 0) info += `, $${a.cost.toFixed(4)}`; + info += "\n"; + } + } + if (stats.active) { + info += `\n${theme.bold("Totals")}\n`; + info += `${theme.fg("dim", "Tokens:")} ${stats.tokens.total.toLocaleString()}\n`; + if (stats.cost > 0) info += `${theme.fg("dim", "Cost:")} $${stats.cost.toFixed(4)}\n`; } - info += `\n${theme.bold("Totals")}\n`; - info += `${theme.fg("dim", "Tokens:")} ${stats.tokens.total.toLocaleString()}\n`; - if (stats.cost > 0) info += `${theme.fg("dim", "Cost:")} $${stats.cost.toFixed(4)}\n`; this.ctx.present([new Spacer(1), new Text(info, 1, 0)]); return; } - const model = stats.model!; + // Single active advisor — detailed view. + const model = stats.model; let info = `${theme.bold("Advisor Status")}\n\n`; - info += `${theme.bold("Provider")}\n`; - info += `${theme.fg("dim", "Model:")} ${model.provider}/${model.id}\n`; + if (stats.advisors.length === 1) { + const a = stats.advisors[0]; + const glyph = CommandController.#advisorStatusGlyph[a.status] ?? "?"; + const label = CommandController.#advisorStatusLabel[a.status] ?? a.status; + info += `${theme.fg(a.status === "running" ? "success" : "error", glyph)} ${a.name} ${theme.fg("dim", `[${label}]`)}\n\n`; + } + if (model) { + info += `${theme.bold("Provider")}\n`; + info += `${theme.fg("dim", "Model:")} ${model.provider}/${model.id}\n`; + } + if (model && usageReports) { + const quota = formatCompactQuota( + model.provider, + usageReports, + nowMs, + resolveActiveAdvisorAccount(model.provider, stats.advisors[0]?.sessionId), + ); + if (quota) { + info += `\n${theme.bold("Quota")}\n`; + info += `${theme.fg("dim", quota)}\n`; + } + } info += `\n${theme.bold("Messages")}\n`; info += `${theme.fg("dim", "User:")} ${stats.messages.user.toLocaleString()}\n`; info += `${theme.fg("dim", "Assistant:")} ${stats.messages.assistant.toLocaleString()}\n`; @@ -406,14 +479,7 @@ export class CommandController { if (stats.tokens.cacheRead > 0) { info += `${theme.fg("dim", "Cache Read:")} ${stats.tokens.cacheRead.toLocaleString()}\n`; } - if (stats.tokens.cacheWrite > 0) { - info += `${theme.fg("dim", "Cache Write:")} ${stats.tokens.cacheWrite.toLocaleString()}\n`; - } - info += `${theme.fg("dim", "Total:")} ${stats.tokens.total.toLocaleString()}\n`; - if (stats.cost > 0) { - info += `\n${theme.bold("Cost")}\n`; - info += `${theme.fg("dim", "Total:")} $${stats.cost.toFixed(4)}\n`; - } + if (stats.cost > 0) info += `${theme.fg("dim", "Cost:")} $${stats.cost.toFixed(4)}\n`; this.ctx.present([new Spacer(1), new Text(info, 1, 0)]); } @@ -1542,6 +1608,57 @@ function resolveResetRange(limits: UsageLimit[], nowMs: number): string | null { } return `resets in ${formatDuration(minReset)}`; } +/** + * Compact one-line quota summary for a single advisor's provider. + * Returns `null` when the provider has no usage data. + * When `activeAccount` is provided, only limits matching that credential + * are shown (mirrors `renderUsageReports`'s account-stickiness filtering). + * Example output: `Quota: 7d window · 67% used · resets in 3.2d` + */ +export function formatCompactQuota( + provider: string, + reports: UsageReport[], + nowMs: number, + activeAccount?: OAuthAccountIdentity, +): string | null { + const providerReports = reports.filter(r => r.provider === provider); + if (providerReports.length === 0) return null; + // Group limits by window id so we show BOTH the 5-hour and 7-day windows + // (or any other distinct windows the provider exposes). Within each window, + // pick the highest used fraction across accounts — that's the most pressing. + const byWindow = new Map(); + for (const report of providerReports) { + for (const limit of report.limits) { + // Skip limits that belong to a different credential than the one + // the advisor is actually using, so we don't alarm the user with + // an exhausted account that isn't theirs. + if (activeAccount && !limitMatchesActiveAccount(report, limit, activeAccount)) continue; + const fraction = resolveUsedFraction(limit); + if (fraction === undefined) continue; + const key = limit.window?.id ?? limit.scope.windowId ?? "—"; + const existing = byWindow.get(key); + if (!existing || fraction > existing.fraction) byWindow.set(key, { limit, fraction }); + } + } + if (byWindow.size === 0) return null; + // Sort windows by urgency (highest fraction first) so the most pressing + // quota is always the first thing the user sees. + const entries = [...byWindow.values()].sort((a, b) => b.fraction - a.fraction); + const lines: string[] = []; + for (const { limit, fraction } of entries) { + const pct = Math.round(fraction * 100); + const windowLabel = limit.window?.label ?? limit.scope.windowId ?? "—"; + // Include the limit label (account/tier) when it carries identity beyond + // the window name, so the user can tell which credential's quota is shown. + const identity = limit.label.trim(); + const header = identity && identity !== windowLabel ? `${windowLabel} (${identity})` : windowLabel; + const parts = [`${header}: ${pct}% used`]; + const reset = resolveResetRange([limit], nowMs); + if (reset) parts.push(reset); + lines.push(parts.join(" · ")); + } + return `Quota: ${lines.join(" │ ")}`; +} function resolveStatusIcon(status: UsageLimit["status"], uiTheme: typeof theme): string { if (status === "exhausted") return uiTheme.fg("error", uiTheme.status.error); diff --git a/packages/coding-agent/src/modes/controllers/selector-controller.ts b/packages/coding-agent/src/modes/controllers/selector-controller.ts index d1f95d1bf..a431433e6 100644 --- a/packages/coding-agent/src/modes/controllers/selector-controller.ts +++ b/packages/coding-agent/src/modes/controllers/selector-controller.ts @@ -298,6 +298,13 @@ export class SelectorController { close: done, requestRender: () => this.ctx.ui.requestRender(), notify: message => this.ctx.showStatus(message), + getAdvisorStats: () => this.ctx.session.getAdvisorStats().advisors, + getUsageReports: async () => this.ctx.session.fetchUsageReports?.() ?? null, + resolveActiveAccount: (provider, sessionId) => + this.ctx.session.modelRegistry.authStorage.getOAuthAccountIdentity( + provider, + sessionId ?? this.ctx.session.sessionId, + ), }); overlayHandle = this.ctx.ui.showOverlay(overlay, { anchor: "bottom-center", diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index b0ce15b09..945c7caa6 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -155,6 +155,7 @@ import { type AdvisorNote, AdvisorOutputQuarantinedError, AdvisorRuntime, + type AdvisorRuntimeStatus, type AdvisorSeverity, AdvisorTranscriptRecorder, advisorTranscriptFilename, @@ -1188,15 +1189,20 @@ export interface AdvisorStats { advisors: PerAdvisorStat[]; } -/** One advisor's slice of {@link AdvisorStats}, surfaced for the multi-advisor status panel. */ +/** One advisor's slice of {@link AdvisorStats}. Active advisors carry full + * token/cost data; disabled/no-model/quota-exhausted advisors appear with + * just `name` + `status` so the status line can render a dot for every + * configured advisor. */ export interface PerAdvisorStat { name: string; - model: Model; + status: AdvisorRuntimeStatus; + model?: Model; contextWindow: number; contextTokens: number; tokens: AdvisorStats["tokens"]; cost: number; messages: AdvisorStats["messages"]; + sessionId?: string; } /** @@ -1906,6 +1912,11 @@ export class AgentSession { #advisors: ActiveAdvisor[] = []; /** Configured advisor roster from WATCHDOG.yml; undefined/empty → single legacy advisor. */ #advisorConfigs?: AdvisorConfig[]; + /** Per-advisor runtime status (slug → {name, status}). Tracks disabled/quota/states + * for the configured roster even when the advisor has no live runtime. The name + * is stored alongside the status so {@link getAdvisorStats} doesn't need to + * recompute slugs or resolve config names. */ + #advisorStatuses: Map = new Map(); /** Provider-facing UUIDv7 identities keyed by primary provider session and advisor slug. */ #advisorProviderSessionIds = new Map(); /** Aggregate of the most recent stop's recorder closes; awaited by dispose() and @@ -2728,7 +2739,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") { @@ -2937,6 +2959,12 @@ export class AgentSession { slug = candidate; usedSlugs.add(slug); } + // Per-advisor toggle: skip disabled advisors but keep them in the + // status map so they show `○` rather than disappearing. + if (config.enabled === false) { + this.#advisorStatuses.set(slug, { name: config.name, status: "paused" }); + continue; + } // Resolve the advisor's model: an explicit `model` override wins; else the // `advisor` role chain. A model that fails to resolve skips just this advisor. @@ -2947,6 +2975,7 @@ export class AgentSession { model = resolved.model; thinkingLevel = concreteThinkingLevel(resolved.thinkingLevel); if (!model) { + this.#advisorStatuses.set(slug, { name: config.name, status: "no_model" }); if (emitWarnings) { this.emitNotice("warning", `Advisor "${config.name}": no model matched "${config.model}"`, "advisor"); } @@ -2955,6 +2984,7 @@ export class AgentSession { } else { const sel = resolveAdvisorRoleSelection(this.settings, this.#modelRegistry.getAvailable()); if (!sel) { + this.#advisorStatuses.set(slug, { name: config.name, status: "no_model" }); if (emitWarnings) { logger.debug("advisor enabled but no model assigned to the 'advisor' role; advisor inactive", { advisor: config.name, @@ -2978,6 +3008,11 @@ export class AgentSession { const requestedLevel = thinkingLevel ?? ThinkingLevel.Medium; const resolvedLevel = resolveThinkingLevelForModel(model, requestedLevel); const advisorThinkingLevel: ThinkingLevel = resolvedLevel ?? ThinkingLevel.Inherit; + // Record the status entry now (in roster order) so the Map's insertion + // order matches the configured roster even when earlier advisors were + // skipped as paused/no_model. The build loop overwrites this to "running" + // without changing insertion order. + this.#advisorStatuses.set(slug, { name: config.name, status: "running" }); descriptors.push({ config, name: config.name, @@ -3013,6 +3048,11 @@ export class AgentSession { if (!this.#advisorEnabled) return false; if (this.#agentKind !== "main" && !this.settings.get("advisor.subagents")) return false; + // Rebuild the status map from scratch so removed/renamed advisors don't + // leave stale entries. #resolveAdvisorRuntimeDescriptors populates every + // entry (`paused`/`no_model`/`running`) in roster order; the build loop + // below confirms `running` for successfully built advisors. + this.#advisorStatuses.clear(); const descriptors = this.#resolveAdvisorRuntimeDescriptors(true); // Advisor service tier (`tier.advisor`): "none" (default) runs the advisor @@ -3214,6 +3254,7 @@ export class AgentSession { }); }, notifyFailure: error => { + this.#advisorStatuses.set(slug, { name: advisorName, status: "error" }); const message = error instanceof Error ? error.message : String(error); this.emitNotice( "warning", @@ -3221,6 +3262,10 @@ export class AgentSession { "advisor", ); }, + notifyQuotaExhausted: () => { + this.#advisorStatuses.set(slug, { name: advisorName, status: "quota_exhausted" }); + this.emitNotice("warning", `Advisor "${advisorName}" quota exhausted — pausing until reset.`, "advisor"); + }, }); const advisorRef: ActiveAdvisor = { @@ -3240,6 +3285,7 @@ export class AgentSession { }; this.#attachAdvisorRecorderFeed(advisorRef); if (seedToCurrent) runtime.seedTo(this.agent.state.messages.length); + this.#advisorStatuses.set(slug, { name: advisorName, status: "running" }); this.#advisors.push(advisorRef); } @@ -3466,7 +3512,26 @@ export class AgentSession { const failedMessage = failedMessages.findLast( (message): message is AssistantMessage => message.role === "assistant", ); - if (failedMessage?.stopReason !== "error") return false; + if (failedMessage?.stopReason !== "error") { + // Stream setup can reject before any assistant turn is recorded (e.g. + // an HTTP 429 thrown from prompt()); classify the raw error so a + // structural usage limit still marks the exhausted credential. + const message = error instanceof Error ? error.message : String(error); + if (!AIError.isUsageLimit(error) && !isUsageLimitOutcome(extractHttpStatusFromError(error), message)) { + return false; + } + const currentModel = advisor.agent.state.model; + const outcome = await this.#modelRegistry.authStorage.markUsageLimitReached( + currentModel.provider, + advisor.providerSessionId, + { + retryAfterMs: extractRetryHint(undefined, message), + baseUrl: currentModel.baseUrl, + modelId: currentModel.id, + }, + ); + return outcome.switched; + } if (failedMessage.content.some(block => block.type === "toolCall")) return false; const currentModel = advisor.agent.state.model; @@ -17509,29 +17574,73 @@ export class AgentSession { return this.#advisors[0]?.agent; } + /** + * Lightweight advisor status for the status line: returns just the configured + * flag and per-advisor name/status without computing token/cost breakdowns. + * Avoids re-tokenizing the advisor transcript on every render frame. + */ + getAdvisorStatusOverview(): { configured: boolean; advisors: { name: string; status: AdvisorRuntimeStatus }[] } { + // Override stale map entries with live runtime status: failureNotified/quotaExhausted + // clear on reset() but #advisorStatuses lags until the next build. + const liveStatusBySlug = new Map(); + for (const a of this.#advisors) { + liveStatusBySlug.set( + a.slug, + a.runtime.quotaExhausted ? "quota_exhausted" : a.runtime.failureNotified ? "error" : "running", + ); + } + const advisors = [...this.#advisorStatuses.entries()].map(([slug, { name, status }]) => ({ + name, + status: liveStatusBySlug.get(slug) ?? status, + })); + return { configured: this.#advisorEnabled, advisors }; + } /** * Return structured advisor stats for the status command and TUI panel. */ getAdvisorStats(): AdvisorStats { const configured = this.#advisorEnabled; - const advisors = this.#advisors.map(a => this.#computeAdvisorStat(a)); - if (advisors.length === 0) { + const liveAdvisors = this.#advisors.map(a => this.#computeAdvisorStat(a)); + // Build the complete roster from #advisorStatuses, which already has the + // correct de-duped slugs as keys. Live advisors (from #advisors) carry full + // token/cost data; disabled/no-model/quota-exhausted advisors appear as + // skeleton entries with just name + status so the status line renders a dot. + const liveStatBySlug = new Map(this.#advisors.map((a, i) => [a.slug, liveAdvisors[i]])); + const roster: PerAdvisorStat[] = []; + for (const [slug, entry] of this.#advisorStatuses) { + const live = liveStatBySlug.get(slug); + if (live) { + roster.push(live); + } else { + roster.push({ + name: entry.name, + status: entry.status, + contextWindow: 0, + contextTokens: 0, + tokens: { input: 0, output: 0, reasoning: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + cost: 0, + messages: { user: 0, assistant: 0, total: 0 }, + }); + } + } + const active = liveAdvisors.length > 0; + if (liveAdvisors.length === 0) { return { configured, - active: false, + active, contextWindow: 0, contextTokens: 0, tokens: { input: 0, output: 0, reasoning: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, cost: 0, messages: { user: 0, assistant: 0, total: 0 }, - advisors: [], + advisors: roster, }; } const tokens = { input: 0, output: 0, reasoning: 0, cacheRead: 0, cacheWrite: 0, total: 0 }; const messages = { user: 0, assistant: 0, total: 0 }; let cost = 0; let contextTokens = 0; - for (const a of advisors) { + for (const a of liveAdvisors) { tokens.input += a.tokens.input; tokens.output += a.tokens.output; tokens.reasoning += a.tokens.reasoning; @@ -17548,14 +17657,14 @@ export class AgentSession { // first advisor's so the legacy status line stays byte-identical. return { configured, - active: true, - model: advisors[0].model, - contextWindow: advisors[0].contextWindow, + active, + model: liveAdvisors[0].model, + contextWindow: liveAdvisors[0].contextWindow, contextTokens, tokens, cost, messages, - advisors, + advisors: roster, }; } @@ -17589,12 +17698,18 @@ export class AgentSession { } return { name: advisor.name, + status: advisor.runtime.quotaExhausted + ? "quota_exhausted" + : advisor.runtime.failureNotified + ? "error" + : "running", model, contextWindow: model.contextWindow ?? 0, contextTokens, tokens: { input, output, reasoning, cacheRead, cacheWrite, total: totalTokens }, cost, messages: { user, assistant, total: messages.length }, + sessionId: advisor.agent.sessionId, }; } @@ -17603,13 +17718,18 @@ export class AgentSession { */ formatAdvisorStatus(): string { const stats = this.getAdvisorStats(); - if (!stats.active) { + if (!stats.active && stats.advisors.length === 0) { return stats.configured ? "Advisor setting is enabled, but no model is assigned to the 'advisor' role." : "Advisor is disabled."; } if (stats.advisors.length <= 1) { const s = stats.advisors[0]; + if (s && s.status === "no_model") { + return stats.configured + ? "Advisor setting is enabled, but no model is assigned to the 'advisor' role." + : "Advisor is disabled."; + } const contextLine = s.contextWindow > 0 ? `Context: ${s.contextTokens.toLocaleString()} / ${s.contextWindow.toLocaleString()} tokens (${Math.round((s.contextTokens / s.contextWindow) * 100)}%)` @@ -17618,6 +17738,7 @@ export class AgentSession { if (s.tokens.cacheRead > 0) spendParts.push(`${s.tokens.cacheRead.toLocaleString()} cache read`); if (s.tokens.cacheWrite > 0) spendParts.push(`${s.tokens.cacheWrite.toLocaleString()} cache write`); const spendLine = `Spend: ${spendParts.join(", ")}, $${s.cost.toFixed(4)}`; + if (!s.model || s.status !== "running") return `Advisor "${s.name}" is ${s.status.replace("_", " ")}.`; return `Advisor is enabled (${s.model.provider}/${s.model.id}). ${contextLine}. ${spendLine}.`; } const lines = [`Advisors enabled (${stats.advisors.length}):`]; @@ -17626,7 +17747,9 @@ export class AgentSession { s.contextWindow > 0 ? `${s.contextTokens.toLocaleString()} / ${s.contextWindow.toLocaleString()} (${Math.round((s.contextTokens / s.contextWindow) * 100)}%)` : `${s.contextTokens.toLocaleString()}`; - lines.push(` • ${s.name} (${s.model.provider}/${s.model.id}) — context ${ctx} tokens, $${s.cost.toFixed(4)}`); + lines.push( + ` • ${s.name}${s.model && s.status === "running" ? ` (${s.model.provider}/${s.model.id})` : ` [${s.status}]`} — context ${ctx} tokens, $${s.cost.toFixed(4)}`, + ); } lines.push( `Totals: ${stats.tokens.input.toLocaleString()} input, ${stats.tokens.output.toLocaleString()} output, $${stats.cost.toFixed(4)}.`, diff --git a/packages/coding-agent/src/session/session-history-format.ts b/packages/coding-agent/src/session/session-history-format.ts index de4d07041..a88145d3e 100644 --- a/packages/coding-agent/src/session/session-history-format.ts +++ b/packages/coding-agent/src/session/session-history-format.ts @@ -46,6 +46,15 @@ export interface HistoryFormatOptions { * this so it sees what changed without re-reading the file. */ expandEditDiffs?: boolean; + /** + * Chunked rendering support: a caller formatting one logical transcript in + * several calls (the advisor's chunked delta render) passes a result index + * built over the WHOLE delta plus one shared consumed-id set, so a toolCall + * finds its toolResult across chunk boundaries and the result is never + * re-rendered as an orphan in a later chunk. + */ + toolResultIndex?: ReadonlyMap; + consumedToolCallIds?: Set; } /** Max length of the primary-arg summary inside `→ tool(...)` lines. */ @@ -273,13 +282,19 @@ export function formatSessionHistoryMarkdown(messages: unknown[], opts?: History } // Index tool results by call id so each toolCall collapses to one line. - const resultsByCallId = new Map(); - for (const msg of typed) { - if (msg.role === "toolResult") { - resultsByCallId.set(msg.toolCallId, msg); + // Chunked callers supply a whole-delta index + shared consumed set so + // call/result pairs resolve across chunk boundaries. + let resultsByCallId = opts?.toolResultIndex; + if (!resultsByCallId) { + const local = new Map(); + for (const msg of typed) { + if (msg.role === "toolResult") { + local.set(msg.toolCallId, msg); + } } + resultsByCallId = local; } - const consumed = new Set(); + const consumed = opts?.consumedToolCallIds ?? new Set(); // In watched mode, consecutive same-role messages collapse under one label // (the watched agent emits one assistant message per tool call, so otherwise // every call repeats `**agent**:`). Cleared whenever a diff --git a/packages/coding-agent/test/advisor-toggle.test.ts b/packages/coding-agent/test/advisor-toggle.test.ts index b22067375..5287c6693 100644 --- a/packages/coding-agent/test/advisor-toggle.test.ts +++ b/packages/coding-agent/test/advisor-toggle.test.ts @@ -1,8 +1,10 @@ -import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it } from "bun:test"; +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "bun:test"; import * as fs from "node:fs"; import * as path from "node:path"; import { Agent, type AgentMessage } from "@oh-my-pi/pi-agent-core"; import type { Model } from "@oh-my-pi/pi-ai"; +import * as AIError from "@oh-my-pi/pi-ai/error"; +import { createMockModel } from "@oh-my-pi/pi-ai/providers/mock"; import { getBundledModel } from "@oh-my-pi/pi-catalog/models"; import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry"; import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings"; @@ -284,4 +286,58 @@ describe("AgentSession advisor toggle", () => { expect(sessionB.isAdvisorEnabled()).toBe(true); expect(sessionB.isAdvisorActive()).toBe(true); }); + + it("exposes provider sessionId on live advisor stats", () => { + session.settings.setModelRole("advisor", `${model.provider}/${model.id}`); + session.toggleAdvisorEnabled(); + + const stats = session.getAdvisorStats(); + expect(stats.advisors).toHaveLength(1); + const sid = stats.advisors[0].sessionId!; + // Full UUIDv7 — must not contain the display-label "-advisor" suffix + expect(sid).toMatch(/^[0-9a-f]{8}-[0-9a-f]{4}-7[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/i); + expect(sid).not.toContain("-advisor"); + }); + it("marks structurally classified advisor usage limits", async () => { + const mock = createMockModel({ responses: [{ content: ["primary complete"] }] }); + const primaryAgent = new Agent({ + initialState: { + model, + systemPrompt: ["Test"], + tools: [], + messages: [], + }, + streamFn: mock.stream, + }); + const settings = Settings.isolated({ "compaction.enabled": false }); + settings.setModelRole("advisor", `${model.provider}/${model.id}`); + const quotaSession = new AgentSession({ + agent: primaryAgent, + sessionManager: SessionManager.inMemory(), + settings, + modelRegistry, + advisorTools: [], + }); + + try { + expect(quotaSession.setAdvisorEnabled(true)).toBe(true); + const advisorAgent = quotaSession.getAdvisorAgent(); + if (!advisorAgent) throw new Error("Expected advisor agent to exist"); + vi.spyOn(advisorAgent, "prompt").mockRejectedValue( + new AIError.ProviderHttpError("Generic provider failure", 429, { code: "insufficient_quota" }), + ); + const markUsageLimitReached = vi + .spyOn(authStorage, "markUsageLimitReached") + .mockResolvedValue({ switched: false }); + + await quotaSession.prompt("Trigger advisor"); + await quotaSession.waitForIdle(); + + expect(markUsageLimitReached).toHaveBeenCalledTimes(1); + expect(markUsageLimitReached.mock.calls[0]?.[0]).toBe(model.provider); + } finally { + await quotaSession.dispose(); + vi.restoreAllMocks(); + } + }); }); diff --git a/packages/coding-agent/test/agent-session-retry-fallback.test.ts b/packages/coding-agent/test/agent-session-retry-fallback.test.ts index 68e327858..e2863ce93 100644 --- a/packages/coding-agent/test/agent-session-retry-fallback.test.ts +++ b/packages/coding-agent/test/agent-session-retry-fallback.test.ts @@ -237,6 +237,7 @@ describe("AgentSession retry fallback", () => { const requestedAdvisorModels: string[] = []; const fallbackAppliedEvents: Array> = []; const fallbackSucceededEvents: Array> = []; + const fallbackSucceeded = Promise.withResolvers(); const advisorFailures: string[] = []; const advisorPrimarySelector = `${advisorPrimary.provider}/${advisorPrimary.id}`; const advisorFallbackSelector = `${advisorFallback.provider}/${advisorFallback.id}`; @@ -287,7 +288,10 @@ describe("AgentSession retry fallback", () => { }); session.subscribe(event => { if (event.type === "retry_fallback_applied") fallbackAppliedEvents.push(event); - if (event.type === "retry_fallback_succeeded") fallbackSucceededEvents.push(event); + if (event.type === "retry_fallback_succeeded") { + fallbackSucceededEvents.push(event); + fallbackSucceeded.resolve(); + } if (event.type === "notice" && event.source === "advisor" && event.message.includes("unavailable")) { advisorFailures.push(event.message); } @@ -296,6 +300,10 @@ describe("AgentSession retry fallback", () => { expect(session.setAdvisorEnabled(true)).toBe(true); await session.prompt("Complete one primary turn"); await session.waitForIdle(); + // The catch-up gate releases immediately while the advisor is mid-failure + // (a failing advisor must never park the primary), so waitForIdle can + // return before the fallback retry lands — await the success event. + await fallbackSucceeded.promise; expect(requestedAdvisorModels).toEqual([advisorPrimarySelector, advisorFallbackSelector]); expect(session.getAdvisorAgent()?.state.model).toMatchObject({ diff --git a/packages/coding-agent/test/issue-816-repro.test.ts b/packages/coding-agent/test/issue-816-repro.test.ts index 65ccd12d2..c38811608 100644 --- a/packages/coding-agent/test/issue-816-repro.test.ts +++ b/packages/coding-agent/test/issue-816-repro.test.ts @@ -128,26 +128,16 @@ describe("issue #816 — plan mode pendingModelSwitch leak", () => { const replacementPlanModel = activePlanModel.provider === haiku.provider && activePlanModel.id === haiku.id ? opus : haiku; - // The role-change listener resolves the plan role through real async - // storage hops (project-scoped roles), so await the apply itself rather - // than assuming it lands within one microtask. - const applied = Promise.withResolvers(); - const setModelSpy = vi.spyOn(session, "setModelTemporary").mockImplementation(async () => { - applied.resolve(); - }); + const setModelSpy = vi.spyOn(session, "setModelTemporary").mockResolvedValue(undefined); session.settings.setModelRole("plan", `${replacementPlanModel.provider}/${replacementPlanModel.id}`); - await applied.promise; + await Promise.resolve(); expect(setModelSpy).toHaveBeenCalledWith(replacementPlanModel, undefined); }); it("keeps plan state coherent when restoring the previous model fails", async () => { - // Pick a plan model that differs from the active session model so plan - // entry actually switches models and arms the previous-model restore. - const haiku = modelRegistry.find("anthropic", "claude-haiku-4-5"); - const opus = modelRegistry.find("anthropic", "claude-opus-4-5"); - if (!haiku || !opus) throw new Error("Expected claude models in registry"); - const planModel = session.model?.provider === haiku.provider && session.model.id === haiku.id ? opus : haiku; + const planModel = modelRegistry.find("anthropic", "claude-haiku-4-5"); + if (!planModel) throw new Error("Expected claude-haiku-4-5 in registry"); vi.spyOn(session, "resolveRoleModelWithThinking").mockReturnValue({ model: planModel, diff --git a/packages/coding-agent/test/status-line-model.test.ts b/packages/coding-agent/test/status-line-model.test.ts index 43af37276..3c8b18131 100644 --- a/packages/coding-agent/test/status-line-model.test.ts +++ b/packages/coding-agent/test/status-line-model.test.ts @@ -16,6 +16,10 @@ function createModelContext(advisorActive: boolean): SegmentContext { isAutoThinking: false, autoResolvedThinkingLevel: () => undefined, isAdvisorActive: () => advisorActive, + getAdvisorStatusOverview: () => ({ + configured: advisorActive, + advisors: advisorActive ? [{ name: "default", status: "running" }] : [], + }), } as unknown as SegmentContext["session"], width: 120, compactThinkingLevel: false, @@ -53,14 +57,32 @@ function createModelContext(advisorActive: boolean): SegmentContext { } describe("status line model segment advisor badge", () => { - it("appends a success-colored ++ badge when the advisor is active", () => { + it("appends a success-colored ++ badge when all advisors run", () => { const rendered = renderSegment("model", createModelContext(true)); expect(rendered.content).toContain("Test Model"); - // The badge carries the success color, kept distinct from the statusLineModel - // name color (which several themes alias to `accent`). expect(rendered.content).toContain(theme.fg("success", "++")); }); + it("colors the badge by the worst roster status", () => { + const ctx = createModelContext(true); + ctx.session.getAdvisorStatusOverview = () => ({ + configured: true, + advisors: [ + { name: "a", status: "running" }, + { name: "b", status: "quota_exhausted" }, + ], + }); + expect(renderSegment("model", ctx).content).toContain(theme.fg("warning", "++")); + ctx.session.getAdvisorStatusOverview = () => ({ + configured: true, + advisors: [ + { name: "a", status: "error" }, + { name: "b", status: "quota_exhausted" }, + ], + }); + expect(renderSegment("model", ctx).content).toContain(theme.fg("error", "++")); + }); + it("omits the badge when the advisor is inactive", () => { const rendered = renderSegment("model", createModelContext(false)); expect(rendered.content).toContain("Test Model"); @@ -72,6 +94,7 @@ describe("status line model segment compact thinking level", () => { function createThinkingContext(compactThinkingLevel: boolean): SegmentContext { return { ...createModelContext(false), + compactThinkingLevel, session: { state: { model: { id: "test-model", name: "Test Model", thinking: true }, @@ -81,8 +104,8 @@ describe("status line model segment compact thinking level", () => { isAutoThinking: false, autoResolvedThinkingLevel: () => undefined, isAdvisorActive: () => false, + getAdvisorStatusOverview: () => ({ configured: false, advisors: [] }), } as unknown as SegmentContext["session"], - compactThinkingLevel, }; } diff --git a/packages/coding-agent/test/status-line-overflow.test.ts b/packages/coding-agent/test/status-line-overflow.test.ts index a4a860e00..3f5f204c7 100644 --- a/packages/coding-agent/test/status-line-overflow.test.ts +++ b/packages/coding-agent/test/status-line-overflow.test.ts @@ -93,6 +93,7 @@ function createStatusLineSession(sessionName: string, modelName?: string) { isAutoThinking: false, autoResolvedThinkingLevel: () => undefined, isAdvisorActive: () => false, + getAdvisorStatusOverview: () => ({ configured: false, advisors: [] }), isFastModeActive: () => false, getAsyncJobSnapshot: () => ({ running: [] }), getCurrentModel: () => undefined, diff --git a/packages/coding-agent/test/status-line-settings-cache.test.ts b/packages/coding-agent/test/status-line-settings-cache.test.ts index 7ff92b50f..2caddbba6 100644 --- a/packages/coding-agent/test/status-line-settings-cache.test.ts +++ b/packages/coding-agent/test/status-line-settings-cache.test.ts @@ -46,7 +46,7 @@ function makeSession(sessionName = "Cache Session") { autoResolvedThinkingLevel: () => undefined, isFastModeActive: () => false, isAdvisorActive: () => false, - getGoalModeState: () => null, + getAdvisorStatusOverview: () => ({ configured: false, advisors: [] }), getAsyncJobSnapshot: () => ({ running: [] }), settings: { get: () => false }, modelRegistry: { isUsingOAuth: () => false }, diff --git a/packages/tui/CHANGELOG.md b/packages/tui/CHANGELOG.md index 305a25e92..ffbb49282 100644 --- a/packages/tui/CHANGELOG.md +++ b/packages/tui/CHANGELOG.md @@ -13,6 +13,7 @@ - Fixed Markdown rendering incorrectly turning local file paths containing `www.` or protocol sequences into HTTP links by requiring a valid GFM left boundary for autolinks. - Fixed terminal resize behavior by restoring alternate-screen rendering during drag frames, preventing wrapped fragments from polluting native scrollback while preserving the overlay-exit flicker fix. +- Added optional right-border scrollbar to the `Editor` component (`setScrollbarVisible`): shows a thumb glyph on the right border when content overflows `maxHeight`, enabling scrollable multi-line editors (e.g. advisor instructions) without losing the submit hint off-screen. ## [17.0.1] - 2026-07-16 ### Added diff --git a/packages/tui/src/components/editor.ts b/packages/tui/src/components/editor.ts index 386abe5ed..5736a3990 100644 --- a/packages/tui/src/components/editor.ts +++ b/packages/tui/src/components/editor.ts @@ -404,6 +404,10 @@ export class Editor implements Component, Focusable { #paddingXOverride: number | undefined; #maxHeight?: number; #scrollOffset: number = 0; + /** When true, the right border shows a scrollbar track/thumb when content + * overflows {@link #maxHeight}. Enabled by {@link HookEditorComponent} and + * other multi-line consumers; single-line consumers are unaffected. */ + #scrollbarVisible = false; // Emacs-style kill ring #killRing = new KillRing(); @@ -555,6 +559,11 @@ export class Editor implements Component, Focusable { // Don't reset scrollOffset — #updateScrollOffset will clamp it on next render } + /** Enable/disable the right-border scrollbar. Only shown when content overflows. */ + setScrollbarVisible(visible: boolean): void { + this.#scrollbarVisible = visible; + } + setPaddingX(paddingX: number): void { this.#paddingXOverride = Math.max(0, paddingX); } @@ -831,6 +840,22 @@ export class Editor implements Component, Focusable { const visibleLayoutLines = layoutLines.slice(this.#scrollOffset, this.#scrollOffset + visibleContentHeight); const result: string[] = []; + // Scrollbar: shown only when content overflows and the caller opted in. + const needsScrollbar = this.#scrollbarVisible && layoutLines.length > visibleContentHeight; + let scrollbarThumb: { start: number; end: number } | null = null; + if (needsScrollbar && visibleContentHeight > 0) { + const thumbSize = Math.max( + 1, + Math.min( + Math.floor((visibleContentHeight * visibleContentHeight) / layoutLines.length), + visibleContentHeight, + ), + ); + const travel = visibleContentHeight - thumbSize; + const maxOffset = Math.max(0, layoutLines.length - visibleContentHeight); + const start = maxOffset === 0 ? 0 : Math.round((this.#scrollOffset / maxOffset) * travel); + scrollbarThumb = { start, end: start + thumbSize }; + } if (borderVisible) { // Render top border: ╭─ [status content] ────────────────╮ @@ -1054,7 +1079,11 @@ export class Editor implements Component, Focusable { result.push(`${bottomLeft}${displayText}${linePad}${bottomRightAdjusted}`); } else { const leftBorder = this.borderColor(`${box.vertical}${padding(paddingX)}`); - const rightBorder = this.borderColor(`${padding(Math.max(0, rightChromeCells - 1))}${box.vertical}`); + // When scrollbar is active, replace the right border vertical with a + // thumb glyph (█) on lines inside the thumb range, keeping the track (│) elsewhere. + const inThumb = scrollbarThumb && visibleIndex >= scrollbarThumb.start && visibleIndex < scrollbarThumb.end; + const rightGlyph = inThumb ? "█" : box.vertical; + const rightBorder = this.borderColor(`${padding(Math.max(0, rightChromeCells - 1))}${rightGlyph}`); result.push(leftBorder + displayText + linePad + rightBorder); } }