feat(advisor): added advisor agent for passive code review with severity-tagged advice

- Created AdvisorRuntime and AdviseTool to drive a read-only advisor agent that delivers severity-tagged advice (nit, concern, blocker) with interruption policy and transcript delta rendering.
- Added /advisor slash command with on/off/status/dump subcommands to control advisor lifecycle and inspect advisor metrics (model, messages, tokens, cost).
- Added advisor.enabled and advisor.subagents settings to enable passive advisor review on main agent and spawned task/eval subagents.
- Implemented advisor message rendering with severity-color badges (blocker=error, concern=warning, nit=muted) in chat log and status line indicator (++ badge).
- Extended yield-queue and session-history-format to support advisor batching and optional thinking block inclusion.
This commit is contained in:
can1357
2026-06-15 16:32:13 +02:00
parent 0b05911d11
commit 37ecd3e73a
24 changed files with 1438 additions and 22 deletions
+11
View File
@@ -7,8 +7,19 @@
- Renamed the SDK tool format type and resolver from `ToolCallFormat`/`resolveToolCallSyntax` to `DialectFormat`/`resolveDialect`, and the agent option from `toolCallSyntax` to `dialect`.
- Changed `/dump` transcript output to render messages with the selected model's native dialect turn and thinking envelopes instead of markdown role headings.
### Added
- Added `/advisor on`, `/advisor off`, `/advisor status`, and `/advisor dump [raw]` slash-command subcommands to manage the advisor at runtime
- Added `advisor.enabled` and `advisor.subagents` settings to enable the advisor and extend it to spawned task/eval subagents
- Added advisor status badge (`++` in success color) to the status line when an advisor is active
- Added `/dump [raw]` flag to toggle between compact and legacy uncompact transcript output formats
- Added `/advisor on`, `/advisor off`, and `/advisor status` slash-command subcommands to enable or disable the advisor at runtime and view advisor status metrics
- Added a passive advisor: assign a second model to the `advisor` role and enable `advisor.enabled` to have it silently review each primary turn and inject severity-tagged advice notes via the `advise` tool. A `nit` rides the non-interrupting aside queue (batched into one card at the next step boundary), while a `concern` or `blocker` interrupts the running agent through the steering channel — aborting in-flight tools, or resuming the agent when it has already yielded — so high-severity advice is acted on immediately. Advice renders in the primary transcript as a distinct `Advisor` card, and the advisor gets hard-isolated read-only `read`/`search`/`find` access — bound to its own `ToolSession` so its reads never touch the primary's snapshot/seen-lines caches — to investigate the workspace before weighing in. The status line shows a `++` badge (in the success color, kept distinct from the model name) after the model name while an advisor is active, and `/advisor dump` copies the advisor's own transcript to the clipboard. Advisors are created only for the top-level session by default; enable `advisor.subagents` to extend them to spawned task/eval subagents.
### Changed
- Changed `/dump` default output to compact markdown format; use `/dump raw` for the legacy uncompact format
- Changed `/dump` and `/advisor dump` to default to compact transcript output and accept an optional `raw` flag for the legacy uncompact format
- Session dump output now renders message history using the model's native dialect turn envelope instead of markdown role headings
### Fixed
@@ -0,0 +1,302 @@
import { describe, expect, it, vi } from "bun:test";
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { createAdvisorMessageCard } from "../../modes/components/advisor-message";
import { getThemeByName } from "../../modes/theme/theme";
import { formatSessionHistoryMarkdown } from "../../session/session-history-format";
import { YieldQueue } from "../../session/yield-queue";
import {
ADVISOR_READONLY_TOOL_NAMES,
AdviseTool,
type AdvisorAgent,
AdvisorRuntime,
type AdvisorRuntimeHost,
formatAdvisorBatchContent,
isInterruptingSeverity,
} from "..";
describe("advisor", () => {
describe("formatSessionHistoryMarkdown includeThinking", () => {
it("includes thinking text when includeThinking is true", () => {
const thinking = "I should check the edge case first.";
const assistantMsg = {
role: "assistant",
content: [{ type: "thinking", thinking }],
timestamp: Date.now(),
} as AgentMessage;
const md = formatSessionHistoryMarkdown([assistantMsg], { includeThinking: true });
expect(md).toContain(thinking);
expect(md).toContain("_thinking:_");
});
it("elides thinking text by default", () => {
const thinking = "I should check the edge case first.";
const assistantMsg = {
role: "assistant",
content: [{ type: "thinking", thinking }],
timestamp: Date.now(),
} as AgentMessage;
const md = formatSessionHistoryMarkdown([assistantMsg]);
expect(md).not.toContain(thinking);
expect(md).not.toContain("_thinking:_");
});
});
describe("advisor yield-queue dispatcher", () => {
it("batches advice notes into one custom message", async () => {
const injected: AgentMessage[] = [];
const yq = new YieldQueue({
isStreaming: () => false,
injectIdle: async messages => {
injected.push(...messages);
},
scheduleIdleFlush: () => {},
});
yq.register<{ note: string; severity?: "nit" | "concern" | "blocker" }>("advisor", {
build: entries =>
entries.length === 0
? null
: ({
role: "custom",
customType: "advisor",
display: true,
attribution: "agent",
timestamp: Date.now(),
content:
"Advisor (a senior reviewer watching your work — weigh it, don't blindly obey):\n" +
entries.map(e => `- ${e.severity ? `[${e.severity}] ` : ""}${e.note}`).join("\n"),
} as AgentMessage),
});
yq.enqueue("advisor", { note: "first note" });
yq.enqueue("advisor", { note: "second note", severity: "blocker" });
await yq.flush("idle");
expect(injected).toHaveLength(1);
const msg = injected[0] as { role: string; customType?: string; display?: boolean; content: string };
expect(msg.role).toBe("custom");
expect(msg.customType).toBe("advisor");
expect(msg.display).toBe(true);
expect(msg.content).toContain("[blocker] second note");
expect(msg.content).toContain("- first note");
});
it("skipIdleFlush prevents idle scheduling", () => {
let scheduled = 0;
const yq = new YieldQueue({
isStreaming: () => false,
injectIdle: async () => {},
scheduleIdleFlush: () => {
scheduled++;
},
});
yq.register<{ note: string }>("advisor", {
build: entries => (entries.length === 0 ? null : ({ role: "custom", content: "x" } as AgentMessage)),
skipIdleFlush: true,
});
yq.register<{ note: string }>("normal", {
build: entries => (entries.length === 0 ? null : ({ role: "custom", content: "y" } as AgentMessage)),
});
yq.enqueue("advisor", { note: "a" });
expect(scheduled).toBe(0);
yq.enqueue("normal", { note: "b" });
expect(scheduled).toBe(1);
});
});
describe("AdviseTool", () => {
it("forwards advice to the callback and returns details", async () => {
const onAdvice = vi.fn();
const tool = new AdviseTool(onAdvice);
const result = await tool.execute("tc-1", { note: "x", severity: "concern" });
expect(onAdvice).toHaveBeenCalledWith("x", "concern");
expect(result.details).toEqual({ note: "x", severity: "concern" });
expect(result.useless).toBe(true);
});
});
describe("advice delivery policy", () => {
it("interrupts on concern and blocker, queues a plain nit", () => {
expect(isInterruptingSeverity("blocker")).toBe(true);
expect(isInterruptingSeverity("concern")).toBe(true);
expect(isInterruptingSeverity("nit")).toBe(false);
expect(isInterruptingSeverity(undefined)).toBe(false);
});
it("formats a batch with the advisor prefix and severity-tagged bullets", () => {
const content = formatAdvisorBatchContent([
{ note: "first note" },
{ note: "second note", severity: "blocker" },
]);
const lines = content.split("\n");
expect(lines[0]).toContain("senior reviewer");
expect(lines[1]).toBe("- first note");
expect(lines[2]).toBe("- [blocker] second note");
});
});
describe("AdvisorRuntime", () => {
function makeAgent(promptInputs: string[]): AdvisorAgent {
return {
prompt: async input => {
promptInputs.push(input);
},
abort: () => {},
reset: () => {},
state: { messages: [] },
};
}
it("coalesces multiple onTurnEnd calls while a prompt is in-flight", async () => {
const promptInputs: string[] = [];
const { promise: firstPromptPromise, resolve: finishFirstPrompt } = Promise.withResolvers<void>();
const agent: AdvisorAgent = {
prompt: async input => {
promptInputs.push(input);
await firstPromptPromise;
},
abort: () => {},
reset: () => {},
state: { messages: [] },
};
const messages: AgentMessage[] = [{ role: "user", content: "first", timestamp: 1 } as AgentMessage];
const host: AdvisorRuntimeHost = {
snapshotMessages: () => messages,
enqueueAdvice: () => {},
};
const runtime = new AdvisorRuntime(agent, host);
runtime.onTurnEnd();
await Promise.resolve();
expect(promptInputs).toHaveLength(1);
expect(promptInputs[0]).toContain("first");
messages.push({ role: "user", content: "second", timestamp: 2 } as AgentMessage);
runtime.onTurnEnd();
await Promise.resolve();
expect(promptInputs).toHaveLength(1);
finishFirstPrompt();
await Promise.resolve();
await Promise.resolve();
expect(promptInputs).toHaveLength(2);
expect(promptInputs[1]).toContain("second");
});
it("excludes advisor custom messages from the rendered delta", () => {
const promptInputs: string[] = [];
const agent = makeAgent(promptInputs);
const messages: AgentMessage[] = [
{ role: "user", content: "hello", timestamp: 1 } as AgentMessage,
{ role: "custom", customType: "advisor", content: "note", display: true, timestamp: 2 } as AgentMessage,
];
const host: AdvisorRuntimeHost = {
snapshotMessages: () => messages,
enqueueAdvice: () => {},
};
const runtime = new AdvisorRuntime(agent, host);
runtime.onTurnEnd();
expect(promptInputs).toHaveLength(1);
expect(promptInputs[0]).toContain("hello");
expect(promptInputs[0]).not.toContain("note");
});
it("handles compaction shrink without prompting", () => {
const promptInputs: string[] = [];
const agent = makeAgent(promptInputs);
let messages: AgentMessage[] = [
{ role: "user", content: "a", timestamp: 1 } as AgentMessage,
{ role: "user", content: "b", timestamp: 2 } as AgentMessage,
];
const host: AdvisorRuntimeHost = {
snapshotMessages: () => messages,
enqueueAdvice: () => {},
};
const runtime = new AdvisorRuntime(agent, host);
runtime.onTurnEnd();
expect(promptInputs).toHaveLength(1);
messages = [{ role: "user", content: "a", timestamp: 1 } as AgentMessage];
expect(() => runtime.onTurnEnd()).not.toThrow();
expect(promptInputs).toHaveLength(1);
});
it("reset re-primes the advisor with the full current transcript", async () => {
const promptInputs: string[] = [];
const agent = makeAgent(promptInputs);
const messages: AgentMessage[] = [{ role: "user", content: "aaa", timestamp: 1 } as AgentMessage];
const host: AdvisorRuntimeHost = {
snapshotMessages: () => messages,
enqueueAdvice: () => {},
};
const runtime = new AdvisorRuntime(agent, host);
runtime.onTurnEnd();
await Promise.resolve();
expect(promptInputs).toHaveLength(1);
expect(promptInputs[0]).toContain("aaa");
// Simulate a compaction: transcript replaced, then reset.
messages.length = 0;
messages.push({ role: "user", content: "summary-bbb", timestamp: 2 } as AgentMessage);
runtime.reset();
runtime.onTurnEnd();
await Promise.resolve();
// The next turn replays the full post-compaction transcript, not just new tail.
expect(promptInputs).toHaveLength(2);
expect(promptInputs[1]).toContain("summary-bbb");
});
});
describe("read-only tool allowlist", () => {
it("selects only the investigation tools from a mixed toolset", () => {
const toolset = ["read", "edit", "search", "bash", "find", "write", "advise"];
const selected = toolset.filter(name => ADVISOR_READONLY_TOOL_NAMES.has(name));
expect(selected).toEqual(["read", "search", "find"]);
expect(ADVISOR_READONLY_TOOL_NAMES.has("edit")).toBe(false);
expect(ADVISOR_READONLY_TOOL_NAMES.has("bash")).toBe(false);
expect(ADVISOR_READONLY_TOOL_NAMES.has("write")).toBe(false);
});
});
describe("createAdvisorMessageCard", () => {
const strip = (lines: readonly string[]): string => lines.join("\n").replace(/\x1b\[[0-9;]*m/g, "");
it("renders the advisor header, severity badge, and note text", async () => {
const uiTheme = await getThemeByName("dark");
if (!uiTheme) throw new Error("theme unavailable");
const card = createAdvisorMessageCard(
{ notes: [{ note: "deleting the wrong file", severity: "blocker" }, { note: "watch the empty case" }] },
() => true,
uiTheme,
);
const text = strip(card.render(80));
expect(text).toContain("Advisor");
expect(text).toContain("2 notes");
expect(text).toContain("blocker");
expect(text).toContain("deleting the wrong file");
expect(text).toContain("watch the empty case");
});
it("collapses to the first notes with an overflow hint", async () => {
const uiTheme = await getThemeByName("dark");
if (!uiTheme) throw new Error("theme unavailable");
const notes = Array.from({ length: 5 }, (_, i) => ({ note: `note ${i}` }));
const card = createAdvisorMessageCard({ notes }, () => false, uiTheme);
const text = strip(card.render(80));
expect(text).toContain("note 0");
expect(text).toContain("+2 more");
expect(text).not.toContain("note 4");
});
it("wraps long notes across multiple lines based on render width instead of truncating them", async () => {
const uiTheme = await getThemeByName("dark");
if (!uiTheme) throw new Error("theme unavailable");
const note =
"This is a very long advisor note that will definitely exceed the restricted width constraint of thirty characters and should therefore wrap across multiple lines rather than getting truncated.";
const card = createAdvisorMessageCard({ notes: [{ note, severity: "concern" }] }, () => true, uiTheme);
const text = strip(card.render(30));
expect(text).toContain("truncated.");
});
});
});
@@ -0,0 +1,87 @@
import type { AgentTool, AgentToolContext, AgentToolResult, AgentToolUpdateCallback } from "@oh-my-pi/pi-agent-core";
import { z } from "zod/v4";
import adviseDescription from "../prompts/advisor/advise-tool.md" with { type: "text" };
const adviseSchema = z.object({
note: z
.string()
.describe("One concrete piece of advice for the agent you are watching. Terse, specific, actionable."),
severity: z
.enum(["nit", "concern", "blocker"])
.optional()
.describe("How strongly to weigh this. Omit for a plain nit."),
});
export type AdviseParams = z.infer<typeof adviseSchema>;
export type AdvisorSeverity = "nit" | "concern" | "blocker";
export interface AdviseDetails {
note: string;
severity?: AdvisorSeverity;
}
/** One queued advice note. */
export interface AdvisorNote {
note: string;
severity?: AdvisorSeverity;
}
/** Details payload on the batched `advisor` custom message rendered in the transcript. */
export interface AdvisorMessageDetails {
notes: AdvisorNote[];
}
/**
* Prose framing prepended to every batched advisor message. Kept here so the
* non-interrupting YieldQueue dispatcher and the interrupting steer path build
* byte-identical content.
*/
const ADVISOR_BATCH_PREFIX = "Advisor (a senior reviewer watching your work — weigh it, don't blindly obey):";
/** Render one advisor card body from a batch of notes (prefix + one bullet per note). */
export function formatAdvisorBatchContent(notes: readonly AdvisorNote[]): string {
return `${ADVISOR_BATCH_PREFIX}\n${notes.map(n => `- ${n.severity ? `[${n.severity}] ` : ""}${n.note}`).join("\n")}`;
}
/**
* Whether advice at this severity should interrupt the running agent (delivered
* via the steering channel, aborting in-flight tools) rather than ride the
* non-interrupting aside queue that lands at the next step boundary. `concern`
* and `blocker` interrupt; a plain `nit` queues.
*/
export function isInterruptingSeverity(severity: AdvisorSeverity | undefined): boolean {
return severity === "concern" || severity === "blocker";
}
/**
* Side-effect-free investigation tools handed to the advisor agent so it can
* inspect the workspace before weighing in. Names match the primary session's
* tool instances, which the advisor reuses.
*/
export const ADVISOR_READONLY_TOOL_NAMES: ReadonlySet<string> = new Set(["read", "search", "find"]);
export class AdviseTool implements AgentTool<typeof adviseSchema, AdviseDetails> {
readonly name = "advise";
readonly label = "Advise";
readonly description = adviseDescription;
readonly parameters = adviseSchema;
readonly intent = "omit" as const;
constructor(private readonly onAdvice: (note: string, severity?: AdviseDetails["severity"]) => void) {}
async execute(
_toolCallId: string,
args: AdviseParams,
_signal?: AbortSignal,
_onUpdate?: AgentToolUpdateCallback<AdviseDetails>,
_context?: AgentToolContext,
): Promise<AgentToolResult<AdviseDetails>> {
this.onAdvice(args.note, args.severity);
return {
content: [{ type: "text", text: "Recorded." }],
details: { note: args.note, severity: args.severity },
useless: true,
};
}
}
@@ -0,0 +1,2 @@
export * from "./advise-tool";
export * from "./runtime";
@@ -0,0 +1,107 @@
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import { logger } from "@oh-my-pi/pi-utils";
import { formatSessionHistoryMarkdown } from "../session/session-history-format";
/** Minimal slice of `Agent` the runtime drives — satisfied by pi-agent-core `Agent`. */
export interface AdvisorAgent {
prompt(input: string): Promise<void>;
abort(reason?: unknown): void;
reset(): void;
readonly state: { messages: AgentMessage[] };
}
export interface AdvisorRuntimeHost {
/** Live primary transcript (use `agent.state.messages`). */
snapshotMessages(): AgentMessage[];
/** Surface one advice note to the primary (enqueues into the session YieldQueue). */
enqueueAdvice(note: string, severity?: "nit" | "concern" | "blocker"): void;
}
export class AdvisorRuntime {
#lastCount = 0;
#pending: string[] = [];
#busy = false;
#disposed = false;
constructor(
private readonly agent: AdvisorAgent,
private readonly host: AdvisorRuntimeHost,
) {}
onTurnEnd(): void {
if (this.#disposed) return;
const render = this.#renderDelta();
if (render) {
this.#pending.push(render);
void this.#drain();
}
}
dispose(): void {
this.#disposed = true;
this.#pending = [];
try {
this.agent.abort("advisor disposed");
} catch {}
}
/**
* Re-prime the advisor after a history rewrite (compaction, session
* switch/resume, branch). Clears the advisor's own (non-persisted) context
* and rewinds the cursor to 0 so the NEXT turn replays the full current —
* post-compaction — transcript, giving the advisor fresh context instead of
* leaving it blind to everything before the rewrite.
*/
reset(): void {
this.#lastCount = 0;
this.#pending = [];
try {
this.agent.reset();
} catch {}
try {
this.agent.abort("advisor reset");
} catch {}
}
/**
* Seed the cursor to the current transcript length when the advisor is enabled
* mid-session. Prevents the next turn from replaying the entire history to the
* advisor (which would be expensive and likely stale).
*/
seedTo(count: number): void {
this.#lastCount = count;
this.#pending = [];
}
#renderDelta(): string | null {
const all = this.host.snapshotMessages();
if (all.length < this.#lastCount) {
this.#lastCount = all.length;
return null;
}
const delta = all
.slice(this.#lastCount)
.filter(m => !(m.role === "custom" && (m as { customType?: string }).customType === "advisor"));
this.#lastCount = all.length;
if (delta.length === 0) return null;
const md = formatSessionHistoryMarkdown(delta, { includeThinking: true });
return md.trim() ? md : null;
}
async #drain(): Promise<void> {
if (this.#busy) return;
this.#busy = true;
try {
while (!this.#disposed && this.#pending.length) {
const batch = this.#pending.splice(0).join("\n\n---\n\n");
try {
await this.agent.prompt(batch);
} catch (err) {
logger.debug("advisor turn failed", { err: String(err) });
}
}
} finally {
this.#busy = false;
}
}
}
@@ -5,7 +5,17 @@
import { isValidThemeColor, type ThemeColor } from "../modes/theme/theme";
import type { Settings } from "./settings";
export type ModelRole = "default" | "smol" | "slow" | "vision" | "plan" | "designer" | "commit" | "title" | "task";
export type ModelRole =
| "default"
| "smol"
| "slow"
| "vision"
| "plan"
| "designer"
| "commit"
| "title"
| "task"
| "advisor";
export interface ModelRoleInfo {
tag?: string;
@@ -25,6 +35,7 @@ export const MODEL_ROLES: Record<ModelRole, ModelRoleInfo> = {
commit: { tag: "COMMIT", name: "Commit", color: "dim" },
title: { tag: "TITLE", name: "Title", color: "dim", hidden: true },
task: { tag: "TASK", name: "Subtask", color: "muted" },
advisor: { tag: "ADVISOR", name: "Advisor", color: "accent" },
};
export const MODEL_ROLE_IDS: ModelRole[] = [
@@ -37,6 +48,7 @@ export const MODEL_ROLE_IDS: ModelRole[] = [
"commit",
"title",
"task",
"advisor",
];
export type RoleInfo = ModelRoleInfo;
@@ -381,6 +381,27 @@ export const SETTINGS_SCHEMA = {
description: "Keep the display from idle-sleeping while a session is open (caffeinate -d)",
},
},
"advisor.enabled": {
type: "boolean",
default: false,
ui: {
tab: "model",
group: "Advisor",
label: "Enable Advisor",
description:
"Pair a second model (assigned to the 'advisor' role) that passively reviews each turn and injects notes.",
},
},
"advisor.subagents": {
type: "boolean",
default: false,
ui: {
tab: "model",
group: "Advisor",
label: "Advisor for Subagents",
description: "Also enable the advisor on spawned task/eval subagents.",
},
},
shellPath: { type: "string", default: undefined },
extensions: { type: "array", default: EMPTY_STRING_ARRAY },
@@ -0,0 +1,103 @@
import { type Component, visibleWidth } from "@oh-my-pi/pi-tui";
import type { AdvisorMessageDetails, AdvisorSeverity } from "../../advisor";
import {
createCachedComponent,
formatBadge,
replaceTabs,
type ToolUIColor,
wrapTextWithAnsi,
} from "../../tools/render-utils";
import { Ellipsis, renderStatusLine, truncateToWidth } from "../../tui";
import type { Theme } from "../theme/theme";
const COLLAPSED_NOTES = 3;
const NOTE_LINE_WIDTH = 110;
function wrapVarying(text: string, w1: number, w2: number): string[] {
if (text.length === 0) return [];
const firstWrap = wrapTextWithAnsi(text, w1);
if (firstWrap.length <= 1) {
return firstWrap;
}
const firstLine = firstWrap[0];
const idx = text.indexOf(firstLine);
if (idx === -1) {
return wrapTextWithAnsi(text, w2);
}
const remainder = text.slice(idx + firstLine.length).trimStart();
const restWrap = wrapTextWithAnsi(remainder, w2);
return [firstLine, ...restWrap];
}
function severityColor(severity: AdvisorSeverity | undefined): ToolUIColor {
switch (severity) {
case "blocker":
return "error";
case "concern":
return "warning";
default:
return "muted";
}
}
/**
* Display-only transcript card for advisor notes injected into the primary
* session. Mirrors the IRC card's glyph + quote-border conventions so passive
* advice reads as a distinct, non-interrupting aside rather than a user turn.
*/
export function createAdvisorMessageCard(
details: AdvisorMessageDetails | undefined,
getExpanded: () => boolean,
uiTheme: Theme,
): Component {
const notes = details?.notes ?? [];
const blockers = notes.filter(note => note.severity === "blocker").length;
const meta: string[] = [`${notes.length} ${notes.length === 1 ? "note" : "notes"}`];
if (blockers > 0) meta.push(uiTheme.fg("error", `${blockers} blocker${blockers === 1 ? "" : "s"}`));
return createCachedComponent(
getExpanded,
(width, expanded) => {
const glyph = uiTheme.styledSymbol("status.info", "accent");
const lines = [renderStatusLine({ iconOverride: glyph, title: "Advisor", meta }, uiTheme)];
const quote = uiTheme.fg("dim", uiTheme.md.quoteBorder);
const shown = expanded ? notes : notes.slice(0, COLLAPSED_NOTES);
for (const entry of shown) {
const badge = entry.severity
? `${formatBadge(entry.severity, severityColor(entry.severity), uiTheme)} `
: "";
const quotePrefix = ` ${quote} `;
const quoteWidth = visibleWidth(quotePrefix);
const badgeWidth = visibleWidth(badge);
const w1 = Math.max(10, Math.min(NOTE_LINE_WIDTH, width) - quoteWidth - badgeWidth);
const w2 = Math.max(10, Math.min(NOTE_LINE_WIDTH, width) - quoteWidth);
const paragraphs = entry.note.split("\n").filter(p => p.trim());
let bodyLines: string[] = [];
for (let i = 0; i < paragraphs.length; i++) {
const p = paragraphs[i];
if (i === 0) {
bodyLines.push(...wrapVarying(p, w1, w2));
} else {
bodyLines.push(...wrapTextWithAnsi(p, w2));
}
}
if (!expanded && bodyLines.length > 2) {
bodyLines = [bodyLines[0], truncateToWidth(bodyLines.slice(1).join(" "), w2, Ellipsis.Unicode)];
}
bodyLines.forEach((line, index) => {
const prefix = index === 0 ? badge : "";
lines.push(` ${quote} ${prefix}${uiTheme.fg("toolOutput", replaceTabs(line))}`);
});
}
const hidden = notes.length - shown.length;
if (hidden > 0) {
lines.push(` ${quote} ${uiTheme.fg("dim", `… +${hidden} more ${hidden === 1 ? "note" : "notes"}`)}`);
}
return lines.map(line => truncateToWidth(line, width, Ellipsis.Unicode));
},
{ paddingX: 1 },
);
}
@@ -19,6 +19,7 @@ import type { AgentMessage, AgentTool } from "@oh-my-pi/pi-agent-core";
import type { Usage } from "@oh-my-pi/pi-ai";
import { Container, Editor, matchesKey, ScrollView, Text, type TUI } from "@oh-my-pi/pi-tui";
import { formatAge, formatBytes, formatDuration, formatNumber, getProjectDir, logger } from "@oh-my-pi/pi-utils";
import type { AdvisorMessageDetails } from "../../advisor";
import { COLLAB_PROMPT_MESSAGE_TYPE, type CollabPromptDetails } from "../../collab/protocol";
import type { KeyId } from "../../config/keybindings";
import { settings } from "../../config/settings";
@@ -45,6 +46,7 @@ import { canonicalizeMessage } from "../../utils/thinking-display";
import type { ObservableSession, SessionObserverRegistry } from "../session-observer-registry";
import { getEditorTheme, theme } from "../theme/theme";
import { matchesSelectDown, matchesSelectUp } from "../utils/keybinding-matchers";
import { createAdvisorMessageCard } from "./advisor-message";
import { AssistantMessageComponent } from "./assistant-message";
import { createBackgroundTanDispatchBlock } from "./background-tan-message";
import { BashExecutionComponent } from "./bash-execution";
@@ -1241,6 +1243,11 @@ export class AgentHubOverlayComponent extends Container {
this.#chatLog.addChild(card);
return;
}
if (message.customType === "advisor") {
const details = (message as CustomMessage<AdvisorMessageDetails>).details;
this.#chatLog.addChild(createAdvisorMessageCard(details, () => this.#chatExpanded, theme));
return;
}
if (message.customType === BACKGROUND_TAN_DISPATCH_MESSAGE_TYPE) {
this.#chatLog.addChild(createBackgroundTanDispatchBlock(message as CustomMessage<unknown>));
return;
@@ -86,32 +86,45 @@ const modelSegment: StatusLineSegment = {
modelName = modelName.slice(7);
}
let content = withIcon(theme.icon.model, modelName);
// 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.
let tail = "";
if (ctx.session.isFastModeActive() && theme.icon.fast) {
content += ` ${theme.icon.fast}`;
tail += ` ${theme.icon.fast}`;
}
// Add thinking level with dot separator
if (opts.showThinkingLevel !== false && state.model?.thinking) {
if (ctx.session.isAutoThinking) {
// Pending (no turn classified yet / classifying) shows a symbol-theme
// question-box marker; once resolved it shows `<level>`.
const resolved = ctx.session.autoResolvedThinkingLevel();
const resolvedText = resolved ? (theme.thinking[resolved as keyof typeof theme.thinking] ?? resolved) : "";
content += `${theme.sep.dot}${resolved ? resolvedText : `${theme.thinking.autoPending} auto`}`;
tail += `${theme.sep.dot}${resolved ? resolvedText : `${theme.thinking.autoPending} auto`}`;
} else {
const level = state.thinkingLevel ?? ThinkingLevel.Off;
if (level !== ThinkingLevel.Off) {
const thinkingText = theme.thinking[level as keyof typeof theme.thinking];
if (thinkingText) {
content += `${theme.sep.dot}${thinkingText}`;
tail += `${theme.sep.dot}${thinkingText}`;
}
}
}
}
return { content: theme.fg("statusLineModel", content), visible: true };
// `statusLineModel` is aliased to `accent` in many themes, so the badge
// uses `success` to stay visibly distinct from the model name color.
let content = theme.fg("statusLineModel", withIcon(theme.icon.model, modelName));
if (ctx.session.isAdvisorActive()) {
content += theme.fg("success", "++");
}
if (tail) {
content += theme.fg("statusLineModel", tail);
}
return { content, visible: true };
},
};
@@ -84,9 +84,9 @@ export class CommandController {
}
}
handleDumpCommand() {
handleDumpCommand(isRaw = false) {
try {
const formatted = this.ctx.session.formatSessionAsText();
const formatted = this.ctx.session.formatSessionAsText({ compact: !isRaw });
if (!formatted) {
this.ctx.showError("No messages to dump yet.");
return;
@@ -98,6 +98,26 @@ export class CommandController {
}
}
handleAdvisorDumpCommand(isRaw = false) {
try {
const advisorHistory = this.ctx.session.formatAdvisorHistoryAsText({ compact: !isRaw });
if (advisorHistory === null) {
this.ctx.showError("Advisor is not active for this session.");
return;
}
if (!advisorHistory) {
this.ctx.showError("Advisor has no history yet.");
return;
}
copyToClipboard(advisorHistory);
this.ctx.showStatus("Advisor history copied to clipboard");
} catch (error: unknown) {
this.ctx.showError(
`Failed to copy advisor history: ${error instanceof Error ? error.message : "Unknown error"}`,
);
}
}
async handleDebugTranscriptCommand(): Promise<void> {
try {
const width = Math.max(1, this.ctx.ui.terminal.columns);
@@ -305,6 +325,53 @@ export class CommandController {
this.ctx.present([new Spacer(1), new Text(info, 1, 0)]);
}
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,
),
]);
return;
}
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`;
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`;
info += `${theme.fg("dim", "Total:")} ${stats.messages.total.toLocaleString()}\n`;
info += `\n${theme.bold("Context")}\n`;
if (stats.contextWindow > 0) {
const percent = Math.round((stats.contextTokens / stats.contextWindow) * 100);
info += `${theme.fg("dim", "Tokens:")} ${stats.contextTokens.toLocaleString()} / ${stats.contextWindow.toLocaleString()} (${percent}%)\n`;
} else {
info += `${theme.fg("dim", "Tokens:")} ${stats.contextTokens.toLocaleString()}\n`;
}
info += `\n${theme.bold("Spend")}\n`;
info += `${theme.fg("dim", "Input:")} ${stats.tokens.input.toLocaleString()}\n`;
info += `${theme.fg("dim", "Output:")} ${stats.tokens.output.toLocaleString()}\n`;
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`;
}
this.ctx.present([new Spacer(1), new Text(info, 1, 0)]);
}
async handleJobsCommand(): Promise<void> {
const snapshot = this.ctx.session.getAsyncJobSnapshot({ recentLimit: 5 });
if (!snapshot) {
@@ -3335,8 +3335,12 @@ export class InteractiveMode implements InteractiveModeContext {
return this.#commandController.handleExportCommand(text);
}
handleDumpCommand() {
return this.#commandController.handleDumpCommand();
handleDumpCommand(isRaw?: boolean) {
return this.#commandController.handleDumpCommand(isRaw);
}
handleAdvisorDumpCommand(isRaw?: boolean) {
return this.#commandController.handleAdvisorDumpCommand(isRaw);
}
handleDebugTranscriptCommand(): Promise<void> {
@@ -3355,6 +3359,10 @@ export class InteractiveMode implements InteractiveModeContext {
return this.#commandController.handleSessionCommand();
}
handleAdvisorStatusCommand(): Promise<void> {
return this.#commandController.handleAdvisorStatusCommand();
}
handleJobsCommand(): Promise<void> {
return this.#commandController.handleJobsCommand();
}
+3 -1
View File
@@ -270,13 +270,15 @@ export interface InteractiveModeContext {
handleShareCommand(): Promise<void>;
handleTodoCommand(args: string): Promise<void>;
handleSessionCommand(): Promise<void>;
handleAdvisorStatusCommand(): Promise<void>;
handleJobsCommand(): Promise<void>;
handleUsageCommand(reports?: UsageReport[] | null): Promise<void>;
handleChangelogCommand(showFull?: boolean): Promise<void>;
handleHotkeysCommand(): void;
handleToolsCommand(): void;
handleContextCommand(): void;
handleDumpCommand(): void;
handleDumpCommand(isRaw?: boolean): void;
handleAdvisorDumpCommand(isRaw?: boolean): void;
handleDebugTranscriptCommand(): Promise<void>;
handleClearCommand(): Promise<void>;
handleFreshCommand(): Promise<void>;
@@ -1,9 +1,11 @@
import type { AgentMessage } from "@oh-my-pi/pi-agent-core";
import type { AssistantMessage, ImageContent, Message, Usage } from "@oh-my-pi/pi-ai";
import { type Component, Spacer, Text, TruncatedText } from "@oh-my-pi/pi-tui";
import type { AdvisorMessageDetails } from "../../advisor";
import { COLLAB_PROMPT_MESSAGE_TYPE, type CollabPromptDetails } from "../../collab/protocol";
import { settings } from "../../config/settings";
import { getFileSnapshotStore } from "../../edit/file-snapshot-store";
import { createAdvisorMessageCard } from "../../modes/components/advisor-message";
import { AssistantMessageComponent } from "../../modes/components/assistant-message";
import { createBackgroundTanDispatchBlock } from "../../modes/components/background-tan-message";
import { BashExecutionComponent } from "../../modes/components/bash-execution";
@@ -240,6 +242,13 @@ export class UiHelpers {
this.ctx.chatContainer.addChild(card);
return [card];
}
if (message.customType === "advisor") {
const details = (message as CustomMessage<AdvisorMessageDetails>).details;
this.ctx.chatContainer.addChild(
createAdvisorMessageCard(details, () => this.ctx.toolOutputExpanded, theme),
);
break;
}
if (message.customType === BACKGROUND_TAN_DISPATCH_MESSAGE_TYPE) {
this.ctx.chatContainer.addChild(createBackgroundTanDispatchBlock(message as CustomMessage<unknown>));
break;
@@ -0,0 +1 @@
Send one concrete, terse piece of advice to the agent you are watching. Use sparingly; stay silent when nothing matters.
@@ -0,0 +1,9 @@
You are a senior engineer silently watching another agent work. You receive that agent's transcript incrementally, including its private thinking, rendered as concise markdown.
You cannot change anything or run commands. You have read-only access to the workspace through `read`, `search`, and `find` — use them sparingly to verify a suspicion (confirm an API exists, check a callsite, read the function under edit) before weighing in. The only way you can speak to the agent is by calling the `advise` tool.
Call `advise` only for things that materially matter: a wrong approach, a missed edge case or failure mode, a hallucinated fact or API, scope creep beyond the task, going in circles, or a likely bug about to be written. Prefer silence. If the agent is on track, do not call any tool.
At most one `advise` per update. Keep each note to one or two sentences. Address the agent in second person. Never restate what it already knows. Never give meta-instructions about how to use the advisor.
Severity controls delivery: a `nit` is folded in non-interruptingly at the next step boundary, while a `concern` or `blocker` interrupts the agent mid-work to reach it immediately. Reserve `concern`/`blocker` for advice worth stopping the agent for; default to `nit` for anything that can wait. Use `blocker` only when continuing will clearly waste the turn.
+31
View File
@@ -35,6 +35,7 @@ import {
prompt,
Snowflake,
} from "@oh-my-pi/pi-utils";
import { ADVISOR_READONLY_TOOL_NAMES } from "./advisor";
import { type AsyncJob, AsyncJobManager } from "./async";
import { AutoLearnController, buildAutoLearnInstructions } from "./autolearn/controller";
import { loadCapability } from "./capability";
@@ -2537,6 +2538,35 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
}
}
// Hard-isolated read-only toolset for the advisor (built unconditionally so
// it can be toggled at runtime). Fresh ReadTool/SearchTool/FindTool bound to a
// DISTINCT ToolSession so the advisor's investigative reads never touch the
// primary's snapshot, seen-lines, conflict, or summary caches (all keyed on
// session identity). `cwd` stays dynamic; edit/yield capabilities are off.
const advisorToolSession: ToolSession = {
...toolSession,
get cwd() {
return sessionManager.getCwd();
},
hasEditTool: false,
requireYieldTool: false,
conflictHistory: undefined,
fileSnapshotStore: undefined,
getSessionId: () => {
const id = sessionManager.getSessionId?.();
return id ? `${id}-advisor` : null;
},
getAgentId: () => "advisor",
};
const built = await Promise.all(
[...ADVISOR_READONLY_TOOL_NAMES].map(name =>
BUILTIN_TOOLS[name as keyof typeof BUILTIN_TOOLS](advisorToolSession),
),
);
const advisorReadOnlyTools: Tool[] = built
.filter((tool): tool is Tool => tool != null)
.map(wrapToolWithMetaNotice);
session = new AgentSession({
agent,
thinkingLevel: autoThinking ? AUTO_THINKING : effectiveThinkingLevel,
@@ -2592,6 +2622,7 @@ export async function createAgentSession(options: CreateAgentSessionOptions = {}
agentKind,
providerSessionId: options.providerSessionId,
parentEvalSessionId: options.parentEvalSessionId,
advisorReadOnlyTools,
});
hasSession = true;
if (asyncJobManager) {
@@ -22,7 +22,7 @@ import type { InMemorySnapshotStore } from "@oh-my-pi/hashline";
import {
type AfterToolCallContext,
type AfterToolCallResult,
type Agent,
Agent,
AgentBusyError,
type AgentEvent,
type AgentMessage,
@@ -114,6 +114,16 @@ import {
Snowflake,
} from "@oh-my-pi/pi-utils";
import * as snapcompact from "@oh-my-pi/snapcompact";
import {
AdviseTool,
type AdvisorAgent,
type AdvisorMessageDetails,
type AdvisorNote,
AdvisorRuntime,
type AdvisorSeverity,
formatAdvisorBatchContent,
isInterruptingSeverity,
} from "../advisor";
import { type AsyncJob, type AsyncJobDeliveryState, AsyncJobManager } from "../async";
import { classifyDifficulty } from "../auto-thinking/classifier";
import { reset as resetCapabilities } from "../capability";
@@ -129,6 +139,7 @@ import {
parseModelString,
type ResolvedModelRoleValue,
resolveModelRoleValue,
resolveRoleSelection,
} from "../config/model-resolver";
import { MODEL_ROLE_IDS } from "../config/model-roles";
import { expandPromptTemplate, type PromptTemplate } from "../config/prompt-templates";
@@ -189,6 +200,7 @@ import { computeNonMessageTokens } from "../modes/utils/context-usage";
import { containsWorkflow, WORKFLOW_NOTICE } from "../modes/workflow";
import { createPlanReadMatcher } from "../plan-mode/plan-protection";
import type { PlanModeState } from "../plan-mode/state";
import advisorSystemPrompt from "../prompts/advisor/system.md" with { type: "text" };
import autoContinuePrompt from "../prompts/system/auto-continue.md" with { type: "text" };
import eagerTaskPrompt from "../prompts/system/eager-task.md" with { type: "text" };
import eagerTodoPrompt from "../prompts/system/eager-todo.md" with { type: "text" };
@@ -268,6 +280,7 @@ import { getLatestCompactionEntry, getRestorableSessionModels } from "./session-
import { formatSessionDumpText } from "./session-dump-format";
import type { BranchSummaryEntry, CompactionEntry, NewSessionOptions } from "./session-entries";
import { EPHEMERAL_MODEL_CHANGE_ROLE } from "./session-entries";
import { formatSessionHistoryMarkdown } from "./session-history-format";
import type { SessionManager } from "./session-manager";
import type { ShakeMode, ShakeResult } from "./shake-types";
import { ToolChoiceQueue } from "./tool-choice-queue";
@@ -457,6 +470,13 @@ export interface AgentSessionConfig {
* so that credential sticky selection is consistent with the session's streaming calls.
*/
providerSessionId?: string;
/**
* Hard-isolated read-only tools (read/search/find) for the advisor agent,
* pre-built in `createAgentSession` against a distinct `ToolSession` so the
* advisor's reads never share the primary's snapshot/seen-lines/conflict
* caches. Undefined when the advisor is disabled.
*/
advisorReadOnlyTools?: AgentTool[];
}
/** Options for AgentSession.prompt() */
@@ -539,6 +559,28 @@ export interface SessionStats {
cost: number;
}
/** Advisor statistics for /advisor status command. */
export interface AdvisorStats {
configured: boolean;
active: boolean;
model?: Model;
contextWindow: number;
contextTokens: number;
tokens: {
input: number;
output: number;
cacheRead: number;
cacheWrite: number;
total: number;
};
cost: number;
messages: {
user: number;
assistant: number;
total: number;
};
}
export interface FreshSessionResult {
previousSessionId: string;
sessionId: string;
@@ -944,6 +986,11 @@ export class AgentSession {
#planModeState: PlanModeState | undefined;
#goalModeState: GoalModeState | undefined;
#goalRuntime: GoalRuntime;
#advisorRuntime?: AdvisorRuntime;
/** The advisor's own agent, retained so `/dump advisor` can serialize its transcript. Undefined when no advisor is active. */
#advisorAgent?: Agent;
#advisorReadOnlyTools?: AgentTool[];
#advisorYieldQueueUnsubscribe?: () => void;
#goalTurnCounter = 0;
#planReferenceSent = false;
#planReferencePath = "local://PLAN.md";
@@ -1234,6 +1281,7 @@ export class AgentSession {
this.#customCommands = config.customCommands ?? [];
this.#skillsSettings = config.skillsSettings;
this.#modelRegistry = config.modelRegistry;
this.#advisorReadOnlyTools = config.advisorReadOnlyTools;
this.#validateRetryFallbackChains();
this.#toolRegistry = config.toolRegistry ?? new Map();
this.#requestedToolNames = config.requestedToolNames;
@@ -1377,12 +1425,133 @@ export class AgentSession {
},
});
if (this.settings.get("advisor.enabled")) this.#buildAdvisorRuntime();
// Always subscribe to agent events for internal handling
// (session persistence, hooks, auto-compaction, retry logic)
this.#unsubscribeAgent = this.agent.subscribe(this.#handleAgentEvent);
// Re-evaluate append-only context mode when the setting changes at runtime.
this.#unsubscribeAppendOnly = onAppendOnlyModeChanged(_value => this.#syncAppendOnlyContext(this.model));
}
// -------------------------------------------------------------------------
// Advisor runtime lifecycle
// -------------------------------------------------------------------------
#buildAdvisorRuntime(seedToCurrent = false): boolean {
if (this.#isDisposed) return false;
if (this.#advisorRuntime) return true;
if (!this.settings.get("advisor.enabled")) return false;
if (this.#agentKind !== "main" && !this.settings.get("advisor.subagents")) return false;
const advisorSel = resolveRoleSelection(
["advisor"],
this.settings,
this.#modelRegistry.getAvailable(),
this.#modelRegistry,
);
if (!advisorSel) {
logger.debug("advisor enabled but no model assigned to the 'advisor' role; advisor inactive");
return false;
}
// Concern and blocker interrupt the running agent through the steering
// channel (aborting in-flight tools at the next steering boundary); when
// the loop has already yielded, triggerTurn resumes it so the advice is
// acted on immediately rather than waiting for the next user prompt. A
// plain nit rides the non-interrupting YieldQueue aside.
const enqueueAdvice = (note: string, severity?: AdvisorSeverity) => {
if (isInterruptingSeverity(severity)) {
const notes: AdvisorNote[] = [{ note, severity }];
void this.sendCustomMessage(
{
customType: "advisor",
content: formatAdvisorBatchContent(notes),
display: true,
attribution: "agent",
details: { notes } satisfies AdvisorMessageDetails,
},
{ deliverAs: "steer", triggerTurn: true },
).catch(err => logger.debug("advisor steer failed", { err: String(err) }));
return;
}
this.yieldQueue.enqueue("advisor", { note, severity });
};
const adviseTool = new AdviseTool(enqueueAdvice);
const advisorReadOnlyTools = this.#advisorReadOnlyTools ?? [];
const appendOnlyContext = new AppendOnlyContextManager();
const advisorThinkingLevel = advisorSel.thinkingLevel ?? ThinkingLevel.Medium;
const advisorAgent = new Agent({
initialState: {
systemPrompt: [advisorSystemPrompt],
model: advisorSel.model,
thinkingLevel: toReasoningEffort(advisorThinkingLevel),
tools: [adviseTool, ...advisorReadOnlyTools],
},
appendOnlyContext,
sessionId: this.sessionId ? `${this.sessionId}-advisor` : undefined,
getApiKey: async provider => {
const key = await this.#modelRegistry.getApiKeyForProvider(
provider,
this.sessionId ? `${this.sessionId}-advisor` : undefined,
);
if (!key) throw new Error(`No API key for advisor provider "${provider}"`);
return key;
},
intentTracing: false,
});
advisorAgent.setDisableReasoning(shouldDisableReasoning(advisorThinkingLevel));
const advisorAgentFacade: AdvisorAgent = {
prompt: input => advisorAgent.prompt(input),
abort: reason => advisorAgent.abort(reason),
reset: () => {
advisorAgent.reset();
appendOnlyContext.log.clear();
},
state: advisorAgent.state,
};
this.#advisorAgent = advisorAgent;
this.#advisorRuntime = new AdvisorRuntime(advisorAgentFacade, {
snapshotMessages: () => this.agent.state.messages,
enqueueAdvice,
});
if (seedToCurrent) {
this.#advisorRuntime.seedTo(this.agent.state.messages.length);
}
// Batch non-blocking advisor notes into one injected custom message.
this.#advisorYieldQueueUnsubscribe = this.yieldQueue.register<AdvisorNote>("advisor", {
build: entries =>
entries.length === 0
? null
: ({
role: "custom",
customType: "advisor",
display: true,
attribution: "agent",
timestamp: Date.now(),
content: formatAdvisorBatchContent(entries),
details: { notes: entries } satisfies AdvisorMessageDetails,
} satisfies CustomMessage),
skipIdleFlush: true,
});
return true;
}
#stopAdvisorRuntime(): void {
if (this.#advisorRuntime) {
this.#advisorRuntime.dispose();
this.#advisorRuntime = undefined;
}
if (this.#advisorAgent) {
this.#advisorAgent = undefined;
}
this.#advisorYieldQueueUnsubscribe?.();
this.#advisorYieldQueueUnsubscribe = undefined;
}
/** Model registry for API key resolution and model discovery */
get modelRegistry(): ModelRegistry {
@@ -1708,6 +1877,8 @@ export class AgentSession {
await this.#emitSessionEvent(displayEvent);
if (event.type === "turn_end") this.#advisorRuntime?.onTurnEnd();
if (event.type === "turn_start") {
this.#resetStreamingEditState();
// TTSR: Reset buffer on turn start
@@ -3213,6 +3384,7 @@ export class AgentSession {
this.#pendingIrcAsides = [];
this.yieldQueue.clear();
this.agent.setAsideMessageProvider(undefined);
this.#stopAdvisorRuntime();
this.#evalExecutionDisposing = true;
}
@@ -5617,6 +5789,7 @@ export class AgentSession {
this.#todoReminderAwaitingProgress = false;
this.#planReferenceSent = false;
this.#planReferencePath = "local://PLAN.md";
this.#advisorRuntime?.reset();
this.#reconnectToAgent();
// Emit session_switch event with reason "new" to hooks
@@ -6234,6 +6407,7 @@ export class AgentSession {
await this.sessionManager.rewriteEntries();
const sessionContext = this.buildDisplaySessionContext();
this.agent.replaceMessages(sessionContext.messages);
this.#advisorRuntime?.reset();
this.#syncTodoPhasesFromBranch();
this.#closeCodexProviderSessionsForHistoryRewrite();
return result;
@@ -6263,9 +6437,9 @@ export class AgentSession {
return undefined;
}
await this.sessionManager.rewriteEntries();
const sessionContext = this.buildDisplaySessionContext();
this.agent.replaceMessages(sessionContext.messages);
this.#advisorRuntime?.reset();
this.#syncTodoPhasesFromBranch();
this.#closeCodexProviderSessionsForHistoryRewrite();
return result;
@@ -6316,6 +6490,7 @@ export class AgentSession {
await this.sessionManager.rewriteEntries();
const sessionContext = this.buildDisplaySessionContext();
this.agent.replaceMessages(sessionContext.messages);
this.#advisorRuntime?.reset();
this.#closeCodexProviderSessionsForHistoryRewrite();
return { removed };
}
@@ -6366,6 +6541,7 @@ export class AgentSession {
await this.sessionManager.rewriteEntries();
const sessionContext = this.buildDisplaySessionContext();
this.agent.replaceMessages(sessionContext.messages);
this.#advisorRuntime?.reset();
this.#closeCodexProviderSessionsForHistoryRewrite();
return {
@@ -6582,6 +6758,7 @@ export class AgentSession {
const newEntries = this.sessionManager.getEntries();
const sessionContext = this.buildDisplaySessionContext();
this.agent.replaceMessages(sessionContext.messages);
this.#advisorRuntime?.reset();
this.#syncTodoPhasesFromBranch();
this.#closeCodexProviderSessionsForHistoryRewrite();
@@ -6801,6 +6978,7 @@ export class AgentSession {
// Rebuild agent messages from session
const sessionContext = this.buildDisplaySessionContext();
this.agent.replaceMessages(sessionContext.messages);
this.#advisorRuntime?.reset();
this.#syncTodoPhasesFromBranch();
return { document: handoffText, savedPath };
@@ -7238,6 +7416,7 @@ export class AgentSession {
}
const safeCount = Math.max(0, Math.min(checkpointState.checkpointMessageCount, this.agent.state.messages.length));
this.agent.replaceMessages(this.agent.state.messages.slice(0, safeCount));
this.#advisorRuntime?.reset();
try {
this.sessionManager.branchWithSummary(checkpointState.checkpointEntryId, report, {
startedAt: checkpointState.startedAt,
@@ -8423,6 +8602,7 @@ export class AgentSession {
const newEntries = this.sessionManager.getEntries();
const sessionContext = this.buildDisplaySessionContext();
this.agent.replaceMessages(sessionContext.messages);
this.#advisorRuntime?.reset();
this.#syncTodoPhasesFromBranch();
this.#closeCodexProviderSessionsForHistoryRewrite();
@@ -10068,6 +10248,7 @@ export class AgentSession {
this.#applyThinkingLevelToAgent(previousThinkingLevel);
this.agent.serviceTier = previousServiceTier;
this.#syncTodoPhasesFromBranch();
this.#advisorRuntime?.reset();
this.#reconnectToAgent();
if (restoreMcpError) {
throw restoreMcpError;
@@ -10149,6 +10330,7 @@ export class AgentSession {
if (!skipConversationRestore) {
this.agent.replaceMessages(sessionContext.messages);
this.#advisorRuntime?.reset();
this.#closeCodexProviderSessionsForHistoryRewrite();
}
@@ -10315,6 +10497,7 @@ export class AgentSession {
const displayContext = deobfuscateSessionContext(stateContext, this.#obfuscator);
await this.#restoreMCPSelectionsForSessionContext(displayContext);
this.agent.replaceMessages(displayContext.messages);
this.#advisorRuntime?.reset();
this.#syncTodoPhasesFromBranch();
this.#closeCodexProviderSessionsForHistoryRewrite();
@@ -10811,7 +10994,10 @@ export class AgentSession {
* Format the entire session as plain text for clipboard export.
* Includes user messages, assistant text, thinking blocks, tool calls, and tool results.
*/
formatSessionAsText(): string {
formatSessionAsText(options?: { compact?: boolean }): string {
if (options?.compact) {
return formatSessionHistoryMarkdown(this.messages);
}
return formatSessionDumpText({
messages: this.messages,
systemPrompt: this.agent.state.systemPrompt,
@@ -10821,6 +11007,178 @@ export class AgentSession {
});
}
/**
* Enable or disable the advisor for this session. The setting is persisted,
* and the runtime is started or stopped to match.
*
* @returns true when the advisor is actively running after the call.
*/
setAdvisorEnabled(enabled: boolean): boolean {
if (enabled) {
this.settings.set("advisor.enabled", true);
return this.#buildAdvisorRuntime(true);
}
this.settings.set("advisor.enabled", false);
this.#stopAdvisorRuntime();
return false;
}
/**
* Toggle the advisor setting and start/stop the runtime accordingly.
*
* @returns true when the advisor is actively running after the call.
*/
toggleAdvisorEnabled(): boolean {
return this.setAdvisorEnabled(!this.settings.get("advisor.enabled"));
}
/**
* Whether a live advisor agent is attached to this session. True only when
* `advisor.enabled` is set AND a model resolved for the `advisor` role AND
* the advisor applies to this agent kind — i.e. the actual runtime exists,
* not merely the setting. Drives the status-line badge and `/dump advisor`.
*/
isAdvisorActive(): boolean {
return this.#advisorAgent !== undefined;
}
/**
* Return structured advisor stats for the status command and TUI panel.
*/
getAdvisorStats(): AdvisorStats {
const configured = this.settings.get("advisor.enabled") as boolean;
const advisor = this.#advisorAgent;
if (!advisor) {
return {
configured,
active: false,
contextWindow: 0,
contextTokens: 0,
tokens: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
cost: 0,
messages: { user: 0, assistant: 0, total: 0 },
};
}
const model = advisor.state.model;
const messages = advisor.state.messages;
const contextTokens = this.#estimateAdvisorContextTokens(messages);
let input = 0;
let output = 0;
let cacheRead = 0;
let cacheWrite = 0;
let cost = 0;
let user = 0;
let assistant = 0;
for (const message of messages) {
if (message.role === "user") user++;
if (message.role === "assistant") {
assistant++;
const assistantMsg = message as AssistantMessage;
input += assistantMsg.usage.input;
output += assistantMsg.usage.output;
cacheRead += assistantMsg.usage.cacheRead;
cacheWrite += assistantMsg.usage.cacheWrite;
cost += assistantMsg.usage.cost.total;
}
}
return {
configured,
active: true,
model,
contextWindow: model.contextWindow ?? 0,
contextTokens,
tokens: {
input,
output,
cacheRead,
cacheWrite,
total: input + output + cacheRead + cacheWrite,
},
cost,
messages: { user, assistant, total: messages.length },
};
}
/**
* Format a concise advisor status line for ACP/text output.
*/
formatAdvisorStatus(): string {
const stats = this.getAdvisorStats();
if (!stats.active) {
return stats.configured
? "Advisor setting is enabled, but no model is assigned to the 'advisor' role."
: "Advisor is disabled.";
}
const model = stats.model!;
const contextLine =
stats.contextWindow > 0
? `Context: ${stats.contextTokens.toLocaleString()} / ${stats.contextWindow.toLocaleString()} tokens (${Math.round((stats.contextTokens / stats.contextWindow) * 100)}%)`
: `Context: ${stats.contextTokens.toLocaleString()} tokens`;
const spendParts = [
`${stats.tokens.input.toLocaleString()} input`,
`${stats.tokens.output.toLocaleString()} output`,
];
if (stats.tokens.cacheRead > 0) spendParts.push(`${stats.tokens.cacheRead.toLocaleString()} cache read`);
if (stats.tokens.cacheWrite > 0) spendParts.push(`${stats.tokens.cacheWrite.toLocaleString()} cache write`);
const spendLine = `Spend: ${spendParts.join(", ")}, $${stats.cost.toFixed(4)}`;
return `Advisor is enabled (${model.provider}/${model.id}). ${contextLine}. ${spendLine}.`;
}
/**
* Estimate the advisor's current context tokens. When the advisor has a
* recent non-aborted assistant message with usage, use that prompt's token
* count and add a trailing estimate for messages after it. Otherwise estimate
* every message.
*/
#estimateAdvisorContextTokens(messages: AgentMessage[]): number {
let lastUsageIndex: number | null = null;
let lastUsage: AssistantMessage["usage"] | undefined;
for (let i = messages.length - 1; i >= 0; i--) {
const msg = messages[i];
if (msg.role === "assistant") {
const assistantMsg = msg as AssistantMessage;
if (assistantMsg.stopReason !== "aborted" && assistantMsg.stopReason !== "error" && assistantMsg.usage) {
lastUsage = assistantMsg.usage;
lastUsageIndex = i;
break;
}
}
}
if (!lastUsage || lastUsageIndex === null) {
let estimated = 0;
for (const message of messages) {
estimated += estimateTokens(message);
}
return estimated;
}
let trailingTokens = 0;
for (let i = lastUsageIndex + 1; i < messages.length; i++) {
trailingTokens += estimateTokens(messages[i]);
}
return calculatePromptTokens(lastUsage) + trailingTokens;
}
/**
* Format the advisor agent's own transcript (its system prompt, config,
* tools, and the markdown deltas it received plus its thinking/advise/read
* calls) as plain text — the advisor-side equivalent of
* {@link formatSessionAsText}. Returns null when no advisor is active.
*/
formatAdvisorHistoryAsText(options?: { compact?: boolean }): string | null {
const advisor = this.#advisorAgent;
if (!advisor) return null;
if (options?.compact) {
return formatSessionHistoryMarkdown(advisor.state.messages);
}
return formatSessionDumpText({
messages: advisor.state.messages,
systemPrompt: advisor.state.systemPrompt,
model: advisor.state.model,
thinkingLevel: advisor.state.thinkingLevel,
tools: advisor.state.tools,
});
}
// =========================================================================
// Extension System
// =========================================================================
@@ -22,6 +22,8 @@ import type {
export interface HistoryFormatOptions {
/** Optional H1 prepended to the transcript. */
title?: string;
/** Render assistant thinking blocks (default: elided). */
includeThinking?: boolean;
}
/** Max length of the primary-arg summary inside `→ tool(...)` lines. */
@@ -194,8 +196,10 @@ export function formatSessionHistoryMarkdown(messages: unknown[], opts?: History
const result = resultsByCallId.get(block.id);
if (result) consumed.add(block.id);
body.push(toolCallLine(block.name, block.arguments, result));
} else if (opts?.includeThinking && block.type === "thinking" && block.thinking.trim()) {
body.push(`_thinking:_ ${block.thinking}`);
}
// thinking / redactedThinking elided entirely
// redactedThinking elided entirely (no readable text)
}
if (body.length === 0) break;
lines.push("## assistant", "", ...body, "");
@@ -6,6 +6,8 @@ export interface YieldDispatcher<P> {
isStale?(entry: P): boolean;
/** Produce one batched AgentMessage from non-stale entries. Return null to skip. */
build(survivors: P[]): AgentMessage | null;
/** If true, entries for this kind are drained only by {@link drainLazy} and never trigger the idle flush. */
skipIdleFlush?: boolean;
}
export interface YieldQueueOptions {
@@ -20,6 +22,7 @@ type YieldFlushMode = "streaming" | "idle";
interface StoredDispatcher {
isStale?: (entry: unknown) => boolean;
build: (survivors: unknown[]) => AgentMessage | null;
skipIdleFlush?: boolean;
}
function formatError(error: unknown): string {
@@ -40,6 +43,7 @@ export class YieldQueue {
const stored: StoredDispatcher = {
...(dispatcher.isStale ? { isStale: entry => dispatcher.isStale?.(entry as P) ?? false } : {}),
build: survivors => dispatcher.build(survivors as P[]),
...(dispatcher.skipIdleFlush ? { skipIdleFlush: true } : {}),
};
this.#dispatchers.set(kind, stored);
return () => {
@@ -60,7 +64,7 @@ export class YieldQueue {
this.#entries.set(kind, entries);
}
entries.push(entry);
if (!this.#options.isStreaming()) {
if (!this.#options.isStreaming() && !this.#dispatchers.get(kind)!.skipIdleFlush) {
this.#scheduleIdleFlush();
}
}
@@ -419,6 +419,100 @@ const BUILTIN_SLASH_COMMAND_REGISTRY: ReadonlyArray<SlashCommandSpec> = [
runtime.ctx.editor.setText("");
},
},
{
name: "advisor",
description: "Toggle the advisor (a second model that reviews each turn and injects notes)",
acpDescription: "Toggle advisor",
acpInputHint: "[on|off|status|dump [raw]]",
subcommands: [
{ name: "on", description: "Enable the advisor" },
{ name: "off", description: "Disable the advisor" },
{ name: "status", description: "Show advisor status" },
{ name: "dump", description: "Copy the advisor's transcript to clipboard", usage: "[raw]" },
],
allowArgs: true,
handle: async (command, runtime) => {
const { verb, rest } = parseSubcommand(command.args);
if (!verb || verb === "toggle") {
const active = runtime.session.toggleAdvisorEnabled();
const configured = runtime.session.settings.get("advisor.enabled") as boolean;
if (active) {
await runtime.output("Advisor enabled.");
} else if (configured) {
await runtime.output("Advisor setting enabled, but no model is assigned to the 'advisor' role.");
} else {
await runtime.output("Advisor disabled.");
}
return commandConsumed();
}
if (verb === "on") {
const active = runtime.session.setAdvisorEnabled(true);
await runtime.output(
active ? "Advisor enabled." : "Advisor setting enabled, but no model is assigned to the 'advisor' role.",
);
return commandConsumed();
}
if (verb === "off") {
runtime.session.setAdvisorEnabled(false);
await runtime.output("Advisor disabled.");
return commandConsumed();
}
if (verb === "status") {
await runtime.output(runtime.session.formatAdvisorStatus());
return commandConsumed();
}
if (verb === "dump") {
const isRaw = rest.toLowerCase() === "raw";
const text = runtime.session.formatAdvisorHistoryAsText({ compact: !isRaw });
await runtime.output(text ?? "Advisor is not active for this session.");
return commandConsumed();
}
return usage("Usage: /advisor [on|off|status|dump [raw]]", runtime);
},
handleTui: async (command, runtime) => {
const { verb, rest } = parseSubcommand(command.args);
if (!verb || verb === "toggle") {
const active = runtime.ctx.session.toggleAdvisorEnabled();
const configured = runtime.ctx.session.settings.get("advisor.enabled") as boolean;
if (active) {
runtime.ctx.showStatus("Advisor enabled.");
} else if (configured) {
runtime.ctx.showStatus("Advisor setting enabled, but no model is assigned to the 'advisor' role.");
} else {
runtime.ctx.showStatus("Advisor disabled.");
}
runtime.ctx.editor.setText("");
return;
}
if (verb === "on") {
const active = runtime.ctx.session.setAdvisorEnabled(true);
runtime.ctx.showStatus(
active ? "Advisor enabled." : "Advisor setting enabled, but no model is assigned to the 'advisor' role.",
);
runtime.ctx.editor.setText("");
return;
}
if (verb === "off") {
runtime.ctx.session.setAdvisorEnabled(false);
runtime.ctx.showStatus("Advisor disabled.");
runtime.ctx.editor.setText("");
return;
}
if (verb === "status") {
await runtime.ctx.handleAdvisorStatusCommand();
runtime.ctx.editor.setText("");
return;
}
if (verb === "dump") {
const isRaw = rest.toLowerCase() === "raw";
runtime.ctx.handleAdvisorDumpCommand(isRaw);
runtime.ctx.editor.setText("");
return;
}
runtime.ctx.showStatus("Usage: /advisor [on|off|status|dump [raw]]");
runtime.ctx.editor.setText("");
},
},
{
name: "export",
description: "Export session to HTML file",
@@ -450,13 +544,17 @@ const BUILTIN_SLASH_COMMAND_REGISTRY: ReadonlyArray<SlashCommandSpec> = [
name: "dump",
description: "Copy session transcript to clipboard",
acpDescription: "Return full transcript as plain text",
handle: async (_command, runtime) => {
const text = runtime.session.formatSessionAsText();
inlineHint: "[raw]",
allowArgs: true,
handle: async (command, runtime) => {
const isRaw = command.args.trim().toLowerCase() === "raw";
const text = runtime.session.formatSessionAsText({ compact: !isRaw });
await runtime.output(text || "No messages to dump yet.");
return commandConsumed();
},
handleTui: async (_command, runtime) => {
await runtime.ctx.handleDumpCommand();
handleTui: (command, runtime) => {
const isRaw = command.args.trim().toLowerCase() === "raw";
runtime.ctx.handleDumpCommand(isRaw);
runtime.ctx.editor.setText("");
},
},
@@ -0,0 +1,101 @@
import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it } from "bun:test";
import * as path from "node:path";
import { Agent } from "@oh-my-pi/pi-agent-core";
import type { Model } from "@oh-my-pi/pi-ai";
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";
import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { TempDir } from "@oh-my-pi/pi-utils";
describe("AgentSession advisor toggle", () => {
let sharedDir: TempDir;
let authStorage: AuthStorage;
let modelRegistry: ModelRegistry;
let model: Model;
beforeAll(async () => {
sharedDir = TempDir.createSync("@pi-advisor-toggle-shared-");
authStorage = await AuthStorage.create(path.join(sharedDir.path(), "testauth.db"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
modelRegistry = new ModelRegistry(authStorage);
const bundled = getBundledModel("anthropic", "claude-sonnet-4-5");
if (!bundled) throw new Error("Expected built-in anthropic model to exist");
model = bundled;
});
afterAll(async () => {
authStorage.close();
try {
await sharedDir.remove();
} catch {}
});
let tempDir: TempDir;
let session: AgentSession;
let sessionManager: SessionManager;
beforeEach(async () => {
tempDir = TempDir.createSync("@pi-advisor-toggle-");
sessionManager = SessionManager.create(tempDir.path(), tempDir.path());
const agent = new Agent({
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
messages: [],
},
});
const settings = Settings.isolated({ "compaction.enabled": false });
session = new AgentSession({
agent,
sessionManager,
settings,
modelRegistry,
advisorReadOnlyTools: [],
});
});
afterEach(async () => {
await session.dispose();
try {
await tempDir.remove();
} catch {}
});
it("starts with advisor disabled", () => {
expect(session.isAdvisorActive()).toBe(false);
expect(session.settings.get("advisor.enabled")).toBe(false);
expect(session.formatAdvisorStatus()).toBe("Advisor is disabled.");
});
it("toggle enables the advisor and runtime", () => {
session.settings.setModelRole("advisor", "anthropic/claude-sonnet-4-5");
const active = session.toggleAdvisorEnabled();
expect(active).toBe(true);
expect(session.isAdvisorActive()).toBe(true);
expect(session.settings.get("advisor.enabled")).toBe(true);
expect(session.formatAdvisorStatus()).toContain("Advisor is enabled (anthropic/claude-sonnet-4-5)");
});
it("toggle disables the advisor and runtime", () => {
session.settings.setModelRole("advisor", "anthropic/claude-sonnet-4-5");
session.toggleAdvisorEnabled();
const active = session.toggleAdvisorEnabled();
expect(active).toBe(false);
expect(session.isAdvisorActive()).toBe(false);
expect(session.settings.get("advisor.enabled")).toBe(false);
});
it("setAdvisorEnabled reports inactive when no advisor model is assigned", () => {
const active = session.setAdvisorEnabled(true);
expect(active).toBe(false);
expect(session.isAdvisorActive()).toBe(false);
expect(session.settings.get("advisor.enabled")).toBe(true);
expect(session.formatAdvisorStatus()).toBe(
"Advisor setting is enabled, but no model is assigned to the 'advisor' role.",
);
});
});
@@ -0,0 +1,58 @@
import { beforeAll, describe, expect, it } from "bun:test";
import type { SegmentContext } from "@oh-my-pi/pi-coding-agent/modes/components/status-line/segments";
import { renderSegment } from "@oh-my-pi/pi-coding-agent/modes/components/status-line/segments";
import { initTheme, theme } from "@oh-my-pi/pi-coding-agent/modes/theme/theme";
beforeAll(async () => {
await initTheme();
});
function createModelContext(advisorActive: boolean): SegmentContext {
return {
session: {
state: { model: { id: "test-model", name: "Test Model" } },
isFastModeActive: () => false,
isAutoThinking: false,
autoResolvedThinkingLevel: () => undefined,
isAdvisorActive: () => advisorActive,
} as unknown as SegmentContext["session"],
width: 120,
options: {},
planMode: null,
loopMode: null,
goalMode: null,
collab: null,
usageStats: {
input: 0,
output: 0,
cacheRead: 0,
cacheWrite: 0,
premiumRequests: 0,
cost: 0,
tokensPerSecond: null,
},
contextPercent: 0,
contextWindow: 0,
autoCompactEnabled: false,
subagentCount: 0,
sessionStartTime: Date.now(),
git: { branch: null, status: null, pr: null },
usage: null,
};
}
describe("status line model segment advisor badge", () => {
it("appends a success-colored ++ badge when the advisor is active", () => {
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("omits the badge when the advisor is inactive", () => {
const rendered = renderSegment("model", createModelContext(false));
expect(rendered.content).toContain("Test Model");
expect(rendered.content).not.toContain("++");
});
});
@@ -44,6 +44,7 @@ function makeSession(sessionName = "Cache Session") {
isAutoThinking: false,
autoResolvedThinkingLevel: () => undefined,
isFastModeActive: () => false,
isAdvisorActive: () => false,
getGoalModeState: () => null,
getAsyncJobSnapshot: () => ({ running: [] }),
settings: { get: () => false },