Merge: darkphilosophy/feat/advisor-per-agent-toggle

Brings the per-advisor toggle, status-line glyphs, quota display, and the
failing-advisor stall/abort fix (f4c8143) onto main's rewritten advisor
runtime. Conflict reconciliation kept main's architecture (fingerprint
prefix reconciliation, host-level onTurnError recovery + fallback chains,
terminal-failure classification) and ported the branch semantics onto it:

- #failing latch: waitForCatchup resolves immediately while an advisor is
  mid-failure; parked waiters wake the moment a turn fails, before any
  async hook or retry sleep.
- Turn-end render containment: a formatter bug restores the cursor/prefix/
  dedup snapshot and never propagates into the primary's turn-end callback
  (per-advisor try/catch boundary in AgentSession).
- Quota pause: when host recovery declines a usage-limit failure, the
  runtime latches quotaExhausted, requeues the batch, and notifies —
  cleared only by an explicit reset.
- Hard halt after a permanent rejection or three backlog-drop cycles.
- #recoverAdvisorTurn also marks usage limits for structural errors thrown
  before any assistant turn is recorded.
This commit is contained in:
can1357
2026-07-17 07:37:29 +02:00
23 changed files with 1719 additions and 154 deletions
@@ -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}`],
+26 -17
View File
@@ -48,30 +48,39 @@ function stalledBody(bytes: Uint8Array[] = []): ReadableStream<Uint8Array> {
}
function delayedBody(chunks: Array<{ atMs: number; bytes: Uint8Array }>): ReadableStream<Uint8Array> {
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<Uint8Array>({
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();
},
});
}
+6
View File
@@ -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 `<context>` to `<repo-rules>` 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
@@ -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<void> {
const deadline = Date.now() + timeoutMs;
while (!predicate() && Date.now() < deadline) await Bun.sleep(2);
}
describe("advisor", () => {
describe("advisor system prompt", () => {
it("forbids concrete claims about tool arguments hidden from the advisor transcript", () => {
@@ -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<undefined>(() => {}),
};
const runtime = new AdvisorRuntime(agent, host, 60_000);
runtime.onTurnEnd(messages);
const started = performance.now();
// Parked with a huge budget: must release on the failure, not the timer.
await runtime.waitForCatchup(60_000, 1);
expect(performance.now() - started).toBeLessThan(2_000);
// While the advisor is mid-failure (retry pending), new waits are free.
const again = performance.now();
await runtime.waitForCatchup(60_000, 1);
expect(performance.now() - again).toBeLessThan(100);
runtime.dispose();
}, 10_000);
it("survives a poisoned message without throwing into the caller or losing the delta", async () => {
// CRITICAL contract: an advisor render failure (throwing getter,
// formatter bug) must neither propagate into the primary agent's
// turn-end callback nor park it on the catch-up gate — and the
// unrendered delta must survive for the next turn.
const promptInputs: string[] = [];
const agent: AdvisorAgent = {
prompt: async input => {
promptInputs.push(input);
},
abort: () => {},
reset: () => {},
state: { messages: [] },
};
const messages: AgentMessage[] = [{ role: "user", content: "aaa", timestamp: 1 } as AgentMessage];
const host: AdvisorRuntimeHost = {
snapshotMessages: () => messages,
enqueueAdvice: () => {},
};
const runtime = new AdvisorRuntime(agent, host, 0);
runtime.onTurnEnd(messages);
await settleUntil(() => promptInputs.length >= 1);
expect(promptInputs).toHaveLength(1);
// Poison: reading `content` throws — during the size probe or render.
const poisoned = {
role: "user",
get content(): string {
throw new Error("poisoned message");
},
timestamp: 2,
} as AgentMessage;
messages.push(poisoned);
expect(() => runtime.onTurnEnd(messages)).not.toThrow();
// A parked primary must not wait out the catch-up budget.
const started = performance.now();
await runtime.waitForCatchup(60_000, 1);
expect(performance.now() - started).toBeLessThan(2_000);
await settleUntil(() => runtime.backlog === 0);
// Replace the poison with a healthy message: the cursor was restored,
// so the next turn re-renders from the failed position.
messages[1] = { role: "user", content: "bbb-recovered", timestamp: 2 } as AgentMessage;
messages.push({ role: "user", content: "ccc", timestamp: 3 } as AgentMessage);
runtime.onTurnEnd(messages);
await settleUntil(() => promptInputs.length >= 2);
expect(promptInputs).toHaveLength(2);
expect(promptInputs[1]).toContain("bbb-recovered");
expect(promptInputs[1]).toContain("ccc");
runtime.dispose();
}, 10_000);
// The live incident shape: ONE agent + ONE advisor froze the whole
// process 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<void> => {
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<void>();
const { promise: hookProceed, resolve: proceedHook } = Promise.withResolvers<void>();
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");
});
});
});
@@ -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);
});
});
+78 -26
View File
@@ -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<Watchdog
logger.warn("Advisor config: invalid schema for edit", { path: filePath, error: result.summary });
return { advisors: [] };
}
return {
instructions: result.instructions?.trim() ? result.instructions : undefined,
advisors: (result.advisors ?? []).map(a => ({
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`;
}
/**
+152 -10
View File
@@ -66,6 +66,22 @@ export interface AdvisorRuntimeHost {
onTurnSuccess?(): Promise<void> | 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<PendingDelta, "turns" | "overflowRecovery"> | 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<void> {
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<void>();
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({
@@ -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<UsageReport[] | null>;
/** 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;
@@ -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);
}
@@ -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: {
@@ -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);
@@ -349,46 +349,119 @@ export class CommandController {
this.ctx.present([new Spacer(1), new Text(info, 1, 0)]);
}
static readonly #advisorStatusGlyph: Record<string, string> = {
running: "●",
paused: "○",
no_model: "○",
quota_exhausted: "✕",
error: "✕",
};
static readonly #advisorStatusLabel: Record<string, string> = {
running: "running",
paused: "off",
no_model: "no model",
quota_exhausted: "quota exhausted",
error: "error",
};
async handleAdvisorStatusCommand(): Promise<void> {
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<UsageReport[] | null> };
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<string, { limit: UsageLimit; fraction: number }>();
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);
@@ -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",
@@ -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<string, { name: string; status: AdvisorRuntimeStatus }> = new Map();
/** Provider-facing UUIDv7 identities keyed by primary provider session and advisor slug. */
#advisorProviderSessionIds = new Map<string, string>();
/** 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<string, AdvisorRuntimeStatus>();
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)}.`,
@@ -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<string, ToolResultMessage>;
consumedToolCallIds?: Set<string>;
}
/** 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<string, ToolResultMessage>();
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<string, ToolResultMessage>();
for (const msg of typed) {
if (msg.role === "toolResult") {
local.set(msg.toolCallId, msg);
}
}
resultsByCallId = local;
}
const consumed = new Set<string>();
const consumed = opts?.consumedToolCallIds ?? new Set<string>();
// 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
@@ -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();
}
});
});
@@ -237,6 +237,7 @@ describe("AgentSession retry fallback", () => {
const requestedAdvisorModels: string[] = [];
const fallbackAppliedEvents: Array<Extract<AgentSessionEvent, { type: "retry_fallback_applied" }>> = [];
const fallbackSucceededEvents: Array<Extract<AgentSessionEvent, { type: "retry_fallback_succeeded" }>> = [];
const fallbackSucceeded = Promise.withResolvers<void>();
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({
@@ -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<void>();
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,
@@ -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,
};
}
@@ -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,
@@ -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 },
+1
View File
@@ -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
+30 -1
View File
@@ -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);
}
}