diff --git a/docs/advisor-watchdog.md b/docs/advisor-watchdog.md index 0ed4768bf..d64d10eab 100644 --- a/docs/advisor-watchdog.md +++ b/docs/advisor-watchdog.md @@ -9,6 +9,7 @@ The advisor is not a second executor. It cannot edit files, run commands, approv - [`src/advisor/runtime.ts`](../packages/coding-agent/src/advisor/runtime.ts) - [`src/advisor/advise-tool.ts`](../packages/coding-agent/src/advisor/advise-tool.ts) - [`src/advisor/watchdog.ts`](../packages/coding-agent/src/advisor/watchdog.ts) +- [`src/advisor/transcript-recorder.ts`](../packages/coding-agent/src/advisor/transcript-recorder.ts) - [`src/prompts/advisor/system.md`](../packages/coding-agent/src/prompts/advisor/system.md) - [`src/prompts/advisor/advise-tool.md`](../packages/coding-agent/src/prompts/advisor/advise-tool.md) - [`src/session/agent-session.ts`](../packages/coding-agent/src/session/agent-session.ts) @@ -203,4 +204,22 @@ The advisor has its own append-only context. Before each advisor prompt, `AgentS 2. if promotion cannot fit enough context, compact the advisor's own message history 3. if compaction has no candidates or still cannot fit, re-prime from the current bounded primary transcript -The advisor transcript is in-memory for the session. It is retained while the session runs so `/advisor dump` can inspect it, but advisor state is not a replacement for the primary persisted transcript. +The advisor's live context is in-memory and append-only; it is retained while the session runs so `/advisor dump` can inspect it, and is independently promoted/compacted/re-primed (above). It is not a replacement for the primary persisted transcript. + +## Transcript persistence and observability + +The advisor is a passive reviewer with its own model usage, so — like a task subagent — every finalized advisor turn is appended to a JSONL inside the owning session's artifacts dir: + +- main session: `/__advisor.jsonl` +- subagent advisor (`advisor.subagents: true`): `//__advisor.jsonl` + +The path is derived from the session file (not the artifacts dir, which subagents share with their parent), so each advisor writes a distinct file. The reserved `__advisor` stem cannot collide with a task subagent's `.jsonl` (task id allocation reserves it). + +Why a file: + +- **Usage attribution.** `omp stats` scans each session folder recursively, so advisor assistant turns (with their usage/cost) are attributed to the same project/session like any other subagent. Advisor "session update" prompts are persisted as `synthetic`, agent-attributed user messages so they never inflate user-message metrics. +- **Observability.** The Agent Hub discovers `__advisor.jsonl` on open and shows it as a read-only `advisor`-kind transcript under its owning session. + +The file follows session switches: on `/new`, resume/switch, and branch the recorder reopens at the new session's path on the next advisor turn; before a `/drop` deletes the old artifacts dir the recorder feed is detached and drained so a queued write cannot recreate the deleted file. The on-disk log is append-only and independent of the in-memory context — re-primes and compaction never truncate it. + +The advisor is never a peer. The `advisor`-kind registry ref is excluded from every agent-facing surface — the `irc` peer roster and broadcast targets, the subagent peer prompt, and the `history://` index/lookup/completions — and cannot be messaged (`irc send` and collab chat refuse it) or revived/killed from the Agent Hub or collab. It is observability only. diff --git a/packages/coding-agent/CHANGELOG.md b/packages/coding-agent/CHANGELOG.md index c085e0a14..ace7f1139 100644 --- a/packages/coding-agent/CHANGELOG.md +++ b/packages/coding-agent/CHANGELOG.md @@ -1,15 +1,22 @@ # Changelog ## [Unreleased] + ### Added +- Added `__advisor.jsonl` transcript persistence for advisor model usage attribution and visibility in the Agent Hub +- Added defensive reservation against naming a task `__advisor` to prevent filesystem collisions with internal transcripts - Added `mode` property to the `CompactOptions` interface for extension-based compaction control - Updated `/compact` slash command to support `soft`, `remote`, and `snapcompact` subcommands for per-run strategy overrides - Added `/compact` mode subcommands so a manual compaction can override the configured `compaction.strategy`/`remoteEnabled` for that one run: `/compact soft` (summarize locally, skip remote endpoints), `/compact remote` (summarize via the remote endpoint / provider-native compaction; warns and falls back to a local summary when no remote path is available), and `/compact snapcompact` (archive history onto dense bitmap images, no LLM call). Bare `/compact` and `/compact ` keep their existing behavior; `soft`/`remote` still accept trailing focus instructions, `snapcompact` rejects them. The subcommands are advertised to ACP clients and surface in TUI autocomplete. Extensions can pass the same selection via `compact(instructions, { mode })`. +- The advisor now persists its own turns to a subagent-style transcript (`/__advisor.jsonl`) so the advisor model's usage is attributed in `omp stats` and its transcript is observable in the Agent Hub (read-only). The transcript follows session switches (`/new`, resume, branch) and is drained/released before a `/drop` deletes the old artifacts dir. The advisor remains a non-peer: it is hidden from agent-facing rosters (`irc` list, `history://`, subagent peer prompt, broadcast targets), is not messageable (`irc send` / collab chat), and is not revivable/killable from the Hub or collab. ### Changed +- Advisor transcripts are now excluded from agent-facing surfaces like `irc`, `history://`, and peer rosters +- Advisor transcripts are read-only and cannot be messaged, revived, or killed via the Agent Hub or IRC - Refined `/compact` argument parsing to reject focus instructions for modes that do not support them (e.g., `snapcompact`) +- Protocol hosts (RPC/`rpc-ui`/ACP) now host-default the full advisor settings group — `advisor.syncBacklog` and `advisor.immuneTurns` in addition to `advisor.enabled`/`advisor.subagents` — so a host that opts the advisor in gets the default tuning instead of inheriting the user's local advisor preferences. ### Fixed diff --git a/packages/coding-agent/src/advisor/index.ts b/packages/coding-agent/src/advisor/index.ts index 30dd105e5..3a7362afb 100644 --- a/packages/coding-agent/src/advisor/index.ts +++ b/packages/coding-agent/src/advisor/index.ts @@ -1,3 +1,4 @@ export * from "./advise-tool"; export * from "./runtime"; +export * from "./transcript-recorder"; export * from "./watchdog"; diff --git a/packages/coding-agent/src/advisor/transcript-recorder.ts b/packages/coding-agent/src/advisor/transcript-recorder.ts new file mode 100644 index 000000000..26b03d78e --- /dev/null +++ b/packages/coding-agent/src/advisor/transcript-recorder.ts @@ -0,0 +1,136 @@ +import * as path from "node:path"; +import type { AgentMessage } from "@oh-my-pi/pi-agent-core"; +import type { Message, UserMessage } from "@oh-my-pi/pi-ai"; +import { logger } from "@oh-my-pi/pi-utils"; +import { SessionManager } from "../session/session-manager"; + +/** + * Reserved transcript stem for advisor session files. Chosen so it cannot + * collide with a task subagent's `.jsonl` (task ids are reserved against + * this exact stem in {@link AgentOutputManager}). + */ +export const ADVISOR_TRANSCRIPT_STEM = "__advisor"; +export const ADVISOR_TRANSCRIPT_FILENAME = `${ADVISOR_TRANSCRIPT_STEM}.jsonl`; + +const JSONL_SUFFIX = ".jsonl"; + +/** + * Append-only persister for an advisor agent's transcript. + * + * The advisor is a passive reviewer with its own model usage, so — like a task + * subagent — its turns are written to a JSONL inside the owning session's + * artifacts dir (`/__advisor.jsonl`, `//__advisor.jsonl` + * for subagent advisors). That single file gives the advisor model proper usage + * attribution in `omp stats` (the stats parser scans the session dir + * recursively) and a read-only transcript in the Agent Hub, without making the + * advisor a registered, messageable peer. + * + * The target is derived from the *session file* (`getSessionFile()`), never + * `getArtifactsDir()` — subagents adopt the parent's artifact manager, so the + * artifacts dir points at the parent root and every subagent advisor would + * collide. The file path is resolved synchronously when a message finalizes and + * captured for the queued write, so a `/new`, resume, or session switch in + * flight can never misattribute an old advisor turn into the new session's file. + * On such a switch the previous writer is closed and the new file opened on the + * next recorded turn. The recorder never truncates: the advisor's in-memory + * context resets/compacts independently, but every billed turn is appended here. + */ +export class AdvisorTranscriptRecorder { + #manager: SessionManager | undefined; + #file: string | undefined; + /** Serializes the async open/close against synchronous appends so records land in order. */ + #queue: Promise; + + /** + * @param after Optional barrier the queue starts behind — used on the advisor + * on→off→on toggle so a fresh recorder's first `open` waits for the prior + * recorder's `close` and the two never hold the same `__advisor.jsonl` at once. + */ + constructor( + private readonly resolveSessionFile: () => string | undefined, + private readonly resolveCwd: () => string, + after?: Promise, + ) { + this.#queue = after + ? after.then( + () => {}, + () => {}, + ) + : Promise.resolve(); + } + + /** + * Persist one finalized advisor message. Assistant turns carry the usage the + * stats parser reads; tool results round out the Hub transcript; user deltas + * (the advisor's "session update" prompts) are persisted but flagged + * `synthetic`/agent-attributed so they never inflate user-message metrics. + * Non-conversational message kinds are skipped. + */ + record(message: AgentMessage): void { + let persisted: Message; + switch (message.role) { + case "assistant": + case "toolResult": + persisted = message; + break; + case "user": + // Clone so the live advisor message stays untouched; mark synthetic so + // stats' user-message metrics skip these agent-internal review prompts. + persisted = { ...(message as UserMessage), synthetic: true, attribution: "agent" }; + break; + default: + return; + } + const sessionFile = this.resolveSessionFile(); + if (!sessionFile?.endsWith(JSONL_SUFFIX)) return; + const file = path.join(sessionFile.slice(0, -JSONL_SUFFIX.length), ADVISOR_TRANSCRIPT_FILENAME); + const cwd = this.resolveCwd(); + this.#enqueue(async () => { + if (file !== this.#file) { + await this.#closeManager(); + this.#manager = await SessionManager.open(file, undefined, undefined, { + initialCwd: cwd, + suppressBreadcrumb: true, + }); + this.#file = file; + } + this.#manager?.appendMessage(persisted); + }); + } + + /** Flush pending writes (best-effort). */ + flush(): Promise { + return this.#enqueueResult(async () => { + if (this.#manager) await this.#manager.flush(); + }); + } + + /** Flush and close the writer, releasing the session file. */ + close(): Promise { + return this.#enqueueResult(() => this.#closeManager()); + } + + async #closeManager(): Promise { + const manager = this.#manager; + this.#manager = undefined; + this.#file = undefined; + if (!manager) return; + try { + await manager.close(); + } catch (err) { + logger.debug("advisor transcript close failed", { err: String(err) }); + } + } + + #enqueue(work: () => Promise): void { + this.#queue = this.#queue.then(work, work).catch(err => { + logger.debug("advisor transcript record failed", { err: String(err) }); + }); + } + + #enqueueResult(work: () => Promise): Promise { + const next = this.#queue.then(work, work); + this.#queue = next.catch(() => {}); + return next; + } +} diff --git a/packages/coding-agent/src/collab/host.ts b/packages/coding-agent/src/collab/host.ts index 5f62f12ae..5f6aac84e 100644 --- a/packages/coding-agent/src/collab/host.ts +++ b/packages/coding-agent/src/collab/host.ts @@ -17,7 +17,7 @@ import { logger } from "@oh-my-pi/pi-utils"; import type { BusChannel, AgentEvent as WireAgentEvent, SessionEntry as WireSessionEntry } from "@oh-my-pi/pi-wire"; import type { InteractiveModeContext } from "../modes/types"; import { AgentLifecycleManager } from "../registry/agent-lifecycle"; -import { AgentRegistry } from "../registry/agent-registry"; +import { type AgentRef, AgentRegistry } from "../registry/agent-registry"; import type { AgentSessionEvent } from "../session/agent-session"; import { stripImagesFromMessage, USER_INTERRUPT_LABEL } from "../session/messages"; import type { SessionEntry as StoredSessionEntry } from "../session/session-entries"; @@ -445,18 +445,24 @@ export class CollabHost { } #snapshotAgents(): AgentSnapshot[] { - return AgentRegistry.global() - .list() - .map(ref => ({ - id: ref.id, - displayName: ref.displayName, - kind: ref.kind, - parentId: ref.parentId, - status: ref.status, - hasSessionFile: !!ref.sessionFile, - createdAt: ref.createdAt, - lastActivity: ref.lastActivity, - })); + return ( + AgentRegistry.global() + .list() + // Advisor transcripts are local observability only; never mirror them to + // guests (the wire AgentSnapshot kind has no `advisor`, and guests must not + // be able to chat/kill/revive them). + .filter((ref): ref is AgentRef & { kind: "main" | "sub" } => ref.kind !== "advisor") + .map(ref => ({ + id: ref.id, + displayName: ref.displayName, + kind: ref.kind, + parentId: ref.parentId, + status: ref.status, + hasSessionFile: !!ref.sessionFile, + createdAt: ref.createdAt, + lastActivity: ref.lastActivity, + })) + ); } #scheduleAgentsBroadcast(): void { @@ -472,6 +478,12 @@ export class CollabHost { this.#rejectReadOnly("agent control", fromPeer); return; } + // Advisor refs are excluded from snapshots, but reject control by id defensively: + // a stale/malicious client must never chat/kill/revive a read-only advisor transcript. + if (AgentRegistry.global().get(agentId)?.kind === "advisor") { + this.#socket?.send({ t: "error", message: `agent ${agentId}: advisor transcripts are read-only` }, fromPeer); + return; + } const fail = (err: unknown) => { logger.warn("collab agent-cmd failed", { cmd, agentId, error: String(err) }); this.#socket?.send({ t: "error", message: `agent ${agentId}: ${String(err)}` }, fromPeer); diff --git a/packages/coding-agent/src/internal-urls/history-protocol.ts b/packages/coding-agent/src/internal-urls/history-protocol.ts index 89e96714e..e84a64449 100644 --- a/packages/coding-agent/src/internal-urls/history-protocol.ts +++ b/packages/coding-agent/src/internal-urls/history-protocol.ts @@ -40,9 +40,12 @@ export class HistoryProtocolHandler implements ProtocolHandler { async resolve(url: InternalUrl): Promise { const agentId = url.rawHost || url.hostname; const registry = AgentRegistry.global(); + // Advisor transcripts are observability-only — surfaced in the Agent Hub, never + // in the agent-facing roster. Hide them from the index, lookup, and completions. + const visible = registry.list().filter(ref => ref.kind !== "advisor"); if (!agentId) { - const content = this.#renderIndex(registry.list()); + const content = this.#renderIndex(visible); return { url: url.href, content, @@ -52,13 +55,14 @@ export class HistoryProtocolHandler implements ProtocolHandler { } let ref = registry.get(agentId); + if (ref?.kind === "advisor") ref = undefined; if (!ref) { // Case-insensitive fallback: agent ids are human-typed (e.g. AuthLoader). const lower = agentId.toLowerCase(); - ref = registry.list().find(candidate => candidate.id.toLowerCase() === lower); + ref = visible.find(candidate => candidate.id.toLowerCase() === lower); } if (!ref) { - const known = registry.list().map(candidate => candidate.id); + const known = visible.map(candidate => candidate.id); const knownStr = known.length > 0 ? known.join(", ") : "none"; throw new Error(`Unknown agent: ${agentId}\nKnown agents: ${knownStr}\nList all with history://`); } @@ -105,6 +109,7 @@ export class HistoryProtocolHandler implements ProtocolHandler { async complete(): Promise { return AgentRegistry.global() .list() + .filter(ref => ref.kind !== "advisor") .map(ref => ({ value: ref.id, description: `${ref.status} · ${ref.kind}${ref.parentId ? ` · parent ${ref.parentId}` : ""}`, diff --git a/packages/coding-agent/src/irc/bus.ts b/packages/coding-agent/src/irc/bus.ts index ec35fe5cb..a40d59ec5 100644 --- a/packages/coding-agent/src/irc/bus.ts +++ b/packages/coding-agent/src/irc/bus.ts @@ -98,6 +98,14 @@ export class IrcBus { if (!ref || ref.status === "aborted") { return { to: message.to, outcome: "failed", error: `Unknown or terminated agent "${message.to}".` }; } + // Advisor refs are observability-only transcripts, never messageable peers. + if (ref.kind === "advisor") { + return { + to: message.to, + outcome: "failed", + error: `Agent "${message.to}" is a read-only advisor transcript and cannot be messaged.`, + }; + } let revived = false; if (ref.status === "parked") { diff --git a/packages/coding-agent/src/modes/components/agent-hub.ts b/packages/coding-agent/src/modes/components/agent-hub.ts index 1ff98a92d..3ebcb486d 100644 --- a/packages/coding-agent/src/modes/components/agent-hub.ts +++ b/packages/coding-agent/src/modes/components/agent-hub.ts @@ -19,7 +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, Ellipsis, 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 { ADVISOR_TRANSCRIPT_FILENAME, 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"; @@ -120,8 +120,33 @@ function registerPersistedSubagentsFromDir(registry: AgentRegistry, dir: string, } for (const entry of entries) { if (!entry.isFile() || !entry.name.endsWith(".jsonl") || entry.name.includes(".bak")) continue; - const id = entry.name.slice(0, -6); const sessionFile = path.join(dir, entry.name); + // The advisor transcript is observability-only: register it as a non-peer + // `advisor` kind under its owning session so the Hub can show its read-only + // transcript, but it never joins agent-facing rosters and is not revivable. + if (entry.name === ADVISOR_TRANSCRIPT_FILENAME) { + const owner = parentId ?? MAIN_AGENT_ID; + const advisorId = `${owner}/advisor`; + const existing = registry.get(advisorId); + // Never clobber a non-advisor ref that happens to share this id (a freak + // user task literally named `/advisor`): leave it, skip the advisor. + if (existing && existing.kind !== "advisor") continue; + if (existing?.sessionFile !== sessionFile) { + // The id is reused across `/new`; refresh it to the current session's file. + if (existing) registry.unregister(advisorId); + registry.register({ + id: advisorId, + displayName: "advisor", + kind: "advisor", + parentId: owner, + session: null, + sessionFile, + status: "parked", + }); + } + continue; + } + const id = entry.name.slice(0, -6); if (!registry.get(id)) { registry.register({ id, @@ -553,7 +578,9 @@ export class AgentHubOverlayComponent extends Container { #activateAgent(ref: AgentRef): void { this.#notice = undefined; const focusAgent = this.#focusAgent; - if (this.#remote || !focusAgent) { + // Advisor refs are read-only transcripts with no live/ revivable session; + // open the in-hub chat view (file-backed) instead of trying to focus one. + if (ref.kind === "advisor" || this.#remote || !focusAgent) { this.openChat(ref.id); return; } @@ -571,6 +598,11 @@ export class AgentHubOverlayComponent extends Container { #reviveSelected(): void { const ref = this.#rows[this.#selectedRow]; if (!ref) return; + if (ref.kind === "advisor") { + this.#notice = `"${ref.id}" is a read-only advisor transcript — nothing to revive.`; + this.#requestRender(); + return; + } if (ref.status !== "parked") { this.#notice = `Agent "${ref.id}" is ${ref.status} — only parked agents can be revived.`; this.#requestRender(); @@ -595,6 +627,11 @@ export class AgentHubOverlayComponent extends Container { #killSelected(): void { const ref = this.#rows[this.#selectedRow]; if (!ref) return; + if (ref.kind === "advisor") { + this.#notice = `"${ref.id}" is a read-only advisor transcript — cannot be killed.`; + this.#requestRender(); + return; + } this.#notice = undefined; if (this.#remote) { this.#remote.kill(ref.id); @@ -826,6 +863,11 @@ export class AgentHubOverlayComponent extends Container { if (!id || !trimmed) return; this.#editor.setText(""); this.#notice = undefined; + if (this.#registry.get(id)?.kind === "advisor") { + this.#notice = "Advisor transcripts are read-only — the advisor cannot be messaged."; + this.#requestRender(); + return; + } if (this.#remote) { this.#remote.chat(id, trimmed); this.#scheduleChatRefresh(); diff --git a/packages/coding-agent/src/registry/agent-registry.ts b/packages/coding-agent/src/registry/agent-registry.ts index b5d52847e..84849bc9a 100644 --- a/packages/coding-agent/src/registry/agent-registry.ts +++ b/packages/coding-agent/src/registry/agent-registry.ts @@ -22,7 +22,13 @@ export const MAIN_AGENT_ID = "Main"; * - `aborted`: hard-killed, terminal. */ export type AgentStatus = "running" | "idle" | "parked" | "aborted"; -export type AgentKind = "main" | "sub"; +/** + * - `main`/`sub`: the user-facing agent tree (driving agent + task subagents). + * - `advisor`: a passive review transcript persisted like a subagent for usage + * attribution and Agent Hub observability, but never a peer — hidden from + * agent-facing rosters (`irc`, `history://`) and not messageable/revivable. + */ +export type AgentKind = "main" | "sub" | "advisor"; export interface AgentRef { id: string; @@ -157,11 +163,14 @@ export class AgentRegistry { } /** - * Returns every alive agent (running | idle) except the caller. - * Flat namespace: every agent can see every other agent. + * Returns every alive agent (running | idle) except the caller. Advisor refs + * are observability-only transcripts, never peers, so they are excluded. + * Flat namespace: every other agent is visible. */ listVisibleTo(id: string): AgentRef[] { - return this.list().filter(ref => ref.id !== id && (ref.status === "running" || ref.status === "idle")); + return this.list().filter( + ref => ref.id !== id && ref.kind !== "advisor" && (ref.status === "running" || ref.status === "idle"), + ); } onChange(listener: RegistryListener): () => void { diff --git a/packages/coding-agent/src/session/agent-session.ts b/packages/coding-agent/src/session/agent-session.ts index 472210fff..2dadb1935 100644 --- a/packages/coding-agent/src/session/agent-session.ts +++ b/packages/coding-agent/src/session/agent-session.ts @@ -127,6 +127,7 @@ import { type AdvisorNote, AdvisorRuntime, type AdvisorSeverity, + AdvisorTranscriptRecorder, formatAdvisorBatchContent, isAdvisorInterruptImmuneTurnActive, isInterruptingSeverity, @@ -1112,6 +1113,13 @@ export class AgentSession { #advisorReadOnlyTools?: AgentTool[]; #advisorWatchdogPrompt?: string; #advisorYieldQueueUnsubscribe?: () => void; + /** Persists the advisor agent's turns to `/__advisor.jsonl` for stats + * attribution and Agent Hub observability. Undefined when no advisor is active. */ + #advisorTranscriptRecorder?: AdvisorTranscriptRecorder; + /** Unsubscribe for the advisor agent's event stream feeding the recorder. */ + #advisorAgentUnsubscribe?: () => void; + /** Latest advisor-recorder close, awaited by dispose() so the final turn lands on disk. */ + #advisorRecorderClosed: Promise = Promise.resolve(); #goalTurnCounter = 0; #planReferenceSent = false; #planReferencePath = "local://PLAN.md"; @@ -1710,7 +1718,13 @@ export class AgentSession { * so none of them inject into the new conversation. */ #resetAdvisorSessionState(): void { + // Mute the recorder across the re-prime: AdvisorRuntime.reset() aborts the advisor + // loop, and that abort can emit an `aborted` message_end we must not attribute to + // either session's transcript. Detach, reset, then re-attach the live agent's feed. + this.#advisorAgentUnsubscribe?.(); + this.#advisorAgentUnsubscribe = undefined; this.#advisorRuntime?.reset(); + this.#attachAdvisorRecorderFeed(); this.#advisorPrimaryTurnsCompleted = 0; this.#advisorInterruptImmuneTurnStart = undefined; this.#advisorAutoResumeSuppressed = false; @@ -1843,6 +1857,18 @@ export class AgentSession { }; this.#advisorAgent = advisorAgent; + // Persist the advisor's turns to `/__advisor.jsonl` (resolved lazily + // so it follows session switches) so its model usage is attributed in stats + // and its transcript shows in the Agent Hub — without registering it as a peer. + const recorder = new AdvisorTranscriptRecorder( + () => this.sessionManager.getSessionFile(), + () => this.sessionManager.getCwd(), + // On the advisor on→off→on toggle, wait for the prior recorder's close so + // two SessionManagers never hold the same __advisor.jsonl at once. + this.#advisorRecorderClosed, + ); + this.#advisorTranscriptRecorder = recorder; + this.#attachAdvisorRecorderFeed(); this.#advisorRuntime = new AdvisorRuntime(advisorAgentFacade, { snapshotMessages: () => this.agent.state.messages, enqueueAdvice, @@ -1873,10 +1899,21 @@ export class AgentSession { } #stopAdvisorRuntime(): void { + // Detach the recorder feed BEFORE aborting the advisor agent: dispose() aborts + // the loop, and an abort emits a final `message_end` we must not enqueue against + // a closing recorder (it would reopen and resurrect an already-released file). + this.#advisorAgentUnsubscribe?.(); + this.#advisorAgentUnsubscribe = undefined; if (this.#advisorRuntime) { this.#advisorRuntime.dispose(); this.#advisorRuntime = undefined; } + if (this.#advisorTranscriptRecorder) { + // Capture the close so dispose()/`/drop` can await the queued open+append+close — + // the last advisor turn would otherwise be lost on a fast process exit. + this.#advisorRecorderClosed = this.#advisorTranscriptRecorder.close(); + this.#advisorTranscriptRecorder = undefined; + } if (this.#advisorAgent) { this.#advisorAgent = undefined; } @@ -1884,6 +1921,18 @@ export class AgentSession { this.#advisorYieldQueueUnsubscribe = undefined; } + /** Subscribe the advisor agent's finalized messages into the transcript recorder. + * Idempotent-by-replacement: callers detach the prior feed first. Kept separate + * so the re-prime path can mute the feed across an abort-driven reset. */ + #attachAdvisorRecorderFeed(): void { + const agent = this.#advisorAgent; + const recorder = this.#advisorTranscriptRecorder; + if (!agent || !recorder) return; + this.#advisorAgentUnsubscribe = agent.subscribe(event => { + if (event.type === "message_end") recorder.record(event.message); + }); + } + async #promoteAdvisorContextModel(currentModel: Model): Promise { const promotionSettings = this.settings.getGroup("contextPromotion"); if (!promotionSettings.enabled) return false; @@ -4062,6 +4111,9 @@ export class AgentSession { await shutdownTinyTitleClient(); this.#releasePowerAssertion(); await this.sessionManager.close(); + // beginDispose() stopped the advisor and captured its recorder close; await + // it so the final advisor turn is flushed before the process may exit. + await this.#advisorRecorderClosed; this.#closeAllProviderSessions("dispose"); // Disconnect the MCP manager this session OWNS so its stdio servers are // not orphaned at exit. Best-effort: a failure here must never throw out @@ -6599,6 +6651,14 @@ export class AgentSession { this.#closeAllProviderSessions("new session"); this.agent.reset(); if (options?.drop && previousSessionFile) { + // Detach the advisor recorder feed and drain its writer BEFORE deleting the + // old artifacts dir: `await this.abort()` only stops the primary, so a still- + // running advisor turn could otherwise finish, emit `message_end`, and recreate + // `/__advisor.jsonl`. #resetAdvisorSessionState (after newSession) re-primes + // the advisor and re-attaches the feed at the new session's path. + this.#advisorAgentUnsubscribe?.(); + this.#advisorAgentUnsubscribe = undefined; + if (this.#advisorTranscriptRecorder) await this.#advisorTranscriptRecorder.close(); try { await this.sessionManager.dropSession(previousSessionFile); } catch (err) { diff --git a/packages/coding-agent/src/task/executor.ts b/packages/coding-agent/src/task/executor.ts index 551518467..514973193 100644 --- a/packages/coding-agent/src/task/executor.ts +++ b/packages/coding-agent/src/task/executor.ts @@ -196,7 +196,7 @@ function installSubagentRetryFallbackChain(args: { function renderIrcPeerRoster(selfId: string): string { const peers = AgentRegistry.global() .list() - .filter(ref => ref.id !== selfId && ref.status !== "aborted"); + .filter(ref => ref.id !== selfId && ref.status !== "aborted" && ref.kind !== "advisor"); if (peers.length === 0) return "- (no other agents)"; const lines = peers.map( peer => diff --git a/packages/coding-agent/src/task/output-manager.ts b/packages/coding-agent/src/task/output-manager.ts index 74fba7832..a4954d91f 100644 --- a/packages/coding-agent/src/task/output-manager.ts +++ b/packages/coding-agent/src/task/output-manager.ts @@ -11,6 +11,7 @@ * collisions across repeated or nested task invocations. */ import * as fs from "node:fs/promises"; +import { ADVISOR_TRANSCRIPT_STEM } from "../advisor/transcript-recorder"; /** * Manages agent output ID allocation to ensure uniqueness. @@ -29,6 +30,10 @@ export class AgentOutputManager { constructor(getArtifactsDir: () => string | null, options?: { parentPrefix?: string }) { this.#getArtifactsDir = getArtifactsDir; this.#parentPrefix = options?.parentPrefix; + // Reserve the advisor transcript stem: a subagent allocated this id would + // write `.jsonl`, clobbering the advisor's `__advisor.jsonl` in the same + // artifacts dir. Reserving bumps such a request to `__advisor-2`. + this.#taken.add(ADVISOR_TRANSCRIPT_STEM); } /** diff --git a/packages/coding-agent/src/tools/irc.ts b/packages/coding-agent/src/tools/irc.ts index a99abd517..6cc6bed02 100644 --- a/packages/coding-agent/src/tools/irc.ts +++ b/packages/coding-agent/src/tools/irc.ts @@ -182,7 +182,7 @@ export class IrcTool implements AgentTool { const bus = IrcBus.global(); const peers = registry .list() - .filter(ref => ref.id !== senderId && ref.status !== "aborted") + .filter(ref => ref.id !== senderId && ref.status !== "aborted" && ref.kind !== "advisor") .map(ref => ({ id: ref.id, displayName: ref.displayName, diff --git a/packages/coding-agent/test/advisor/advisor-visibility.test.ts b/packages/coding-agent/test/advisor/advisor-visibility.test.ts new file mode 100644 index 000000000..18cb07b4a --- /dev/null +++ b/packages/coding-agent/test/advisor/advisor-visibility.test.ts @@ -0,0 +1,58 @@ +/** + * Contracts: an advisor-kind registry ref is observability-only — present for the + * Agent Hub, hidden from every agent-facing surface, and never messageable. + * + * - `AgentRegistry.listVisibleTo` (irc roster / broadcast targets) excludes advisors. + * - `IrcBus.send` to an advisor ref fails as non-messageable, without reviving it. + */ +import { afterEach, beforeEach, describe, expect, it } from "bun:test"; +import { IrcBus } from "@oh-my-pi/pi-coding-agent/irc/bus"; +import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry"; + +describe("advisor registry visibility", () => { + beforeEach(() => { + AgentRegistry.resetGlobalForTests(); + IrcBus.resetGlobalForTests(); + }); + + afterEach(() => { + IrcBus.resetGlobalForTests(); + AgentRegistry.resetGlobalForTests(); + }); + + it("excludes advisor refs from listVisibleTo", () => { + const registry = AgentRegistry.global(); + registry.register({ id: "Main", displayName: "Main", kind: "main", session: null, status: "running" }); + registry.register({ id: "Worker", displayName: "Worker", kind: "sub", session: null, status: "idle" }); + registry.register({ + id: "Main/advisor", + displayName: "advisor", + kind: "advisor", + session: null, + status: "idle", + }); + + const visible = registry.listVisibleTo("Main").map(ref => ref.id); + expect(visible).toContain("Worker"); + expect(visible).not.toContain("Main/advisor"); + }); + + it("refuses to message an advisor ref", async () => { + const registry = AgentRegistry.global(); + registry.register({ + id: "Main/advisor", + displayName: "advisor", + kind: "advisor", + session: null, + sessionFile: "/tmp/x/__advisor.jsonl", + status: "parked", + }); + const bus = new IrcBus(registry); + + const receipt = await bus.send({ from: "Main", to: "Main/advisor", body: "hi" }); + expect(receipt.outcome).toBe("failed"); + expect(receipt.error).toContain("advisor"); + // It must still be parked — a refused send never revives an advisor transcript. + expect(registry.get("Main/advisor")?.status).toBe("parked"); + }); +}); diff --git a/packages/coding-agent/test/advisor/transcript-recorder.test.ts b/packages/coding-agent/test/advisor/transcript-recorder.test.ts new file mode 100644 index 000000000..23849266b --- /dev/null +++ b/packages/coding-agent/test/advisor/transcript-recorder.test.ts @@ -0,0 +1,163 @@ +/** + * Contracts: AdvisorTranscriptRecorder persists the advisor agent's turns to a + * subagent-style JSONL (`/__advisor.jsonl`) so the advisor model's usage + * is attributed in stats and its transcript shows in the Agent Hub. + * + * - Assistant turns land as `{type:"message", message:{role:"assistant", usage}}` + * entries — exactly the shape the stats parser reads for usage. + * - User deltas are persisted but flagged `synthetic`/agent-attributed so stats' + * user-message metrics skip them. + * - Non-conversational message kinds are not persisted. + * - The target follows the session file: a switch routes later turns to the new + * session's `__advisor.jsonl`, leaving the prior file intact. + */ +import { describe, expect, it } from "bun:test"; +import * as fs from "node:fs/promises"; +import * as os from "node:os"; +import * as path from "node:path"; +import type { AgentMessage } from "@oh-my-pi/pi-agent-core"; +import { + ADVISOR_TRANSCRIPT_FILENAME, + AdvisorTranscriptRecorder, +} from "@oh-my-pi/pi-coding-agent/advisor/transcript-recorder"; + +interface AdvisorEntry { + type?: string; + id?: unknown; + message?: { + role?: string; + model?: string; + usage?: { input?: number }; + synthetic?: boolean; + attribution?: string; + }; +} + +async function withTempDir(fn: (dir: string) => Promise): Promise { + const dir = await fs.mkdtemp(path.join(os.tmpdir(), "advisor-recorder-")); + try { + return await fn(dir); + } finally { + await fs.rm(dir, { recursive: true, force: true }); + } +} + +/** Parse the message entries (skipping the session header) from an advisor JSONL. */ +async function readMessageEntries(file: string): Promise { + const text = await Bun.file(file).text(); + // JSON.parse returns `any`; assigning to the typed array narrows reads below. + const entries: AdvisorEntry[] = text + .trim() + .split("\n") + .map(line => JSON.parse(line)); + return entries.filter(entry => entry.type === "message"); +} + +function assistantMessage(text: string, inputTokens: number): AgentMessage { + const message = { + role: "assistant" as const, + content: [{ type: "text" as const, text }], + api: "anthropic-messages", + provider: "anthropic", + model: "test-advisor-model", + usage: { + input: inputTokens, + output: 3, + cacheRead: 0, + cacheWrite: 0, + totalTokens: inputTokens + 3, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + }, + stopReason: "stop" as const, + timestamp: 1, + }; + return message as unknown as AgentMessage; +} + +function userMessage(text: string): AgentMessage { + const message = { role: "user" as const, content: [{ type: "text" as const, text }], timestamp: 1 }; + return message as unknown as AgentMessage; +} + +function developerMessage(text: string): AgentMessage { + const message = { role: "developer" as const, content: [{ type: "text" as const, text }], timestamp: 1 }; + return message as unknown as AgentMessage; +} + +describe("AdvisorTranscriptRecorder", () => { + it("persists assistant turns with usage to /__advisor.jsonl", async () => { + await withTempDir(async dir => { + const sessionFile = path.join(dir, "sess.jsonl"); + const recorder = new AdvisorTranscriptRecorder( + () => sessionFile, + () => dir, + ); + recorder.record(assistantMessage("reviewing", 42)); + await recorder.close(); + + const messages = await readMessageEntries(path.join(dir, "sess", ADVISOR_TRANSCRIPT_FILENAME)); + expect(messages).toHaveLength(1); + expect(messages[0].message?.role).toBe("assistant"); + expect(messages[0].message?.model).toBe("test-advisor-model"); + expect(messages[0].message?.usage?.input).toBe(42); + // Stats keys on a non-empty entry id; SessionManager must assign one. + expect(typeof messages[0].id).toBe("string"); + expect(String(messages[0].id).length).toBeGreaterThan(0); + }); + }); + + it("marks advisor user deltas synthetic and agent-attributed", async () => { + await withTempDir(async dir => { + const sessionFile = path.join(dir, "sess.jsonl"); + const recorder = new AdvisorTranscriptRecorder( + () => sessionFile, + () => dir, + ); + recorder.record(userMessage("### Session update")); + await recorder.close(); + + const messages = await readMessageEntries(path.join(dir, "sess", ADVISOR_TRANSCRIPT_FILENAME)); + expect(messages).toHaveLength(1); + expect(messages[0].message?.role).toBe("user"); + expect(messages[0].message?.synthetic).toBe(true); + expect(messages[0].message?.attribution).toBe("agent"); + }); + }); + + it("skips non-conversational message kinds", async () => { + await withTempDir(async dir => { + const sessionFile = path.join(dir, "sess.jsonl"); + const recorder = new AdvisorTranscriptRecorder( + () => sessionFile, + () => dir, + ); + recorder.record(developerMessage("noise")); + recorder.record(assistantMessage("kept", 1)); + await recorder.close(); + + const messages = await readMessageEntries(path.join(dir, "sess", ADVISOR_TRANSCRIPT_FILENAME)); + expect(messages.map(m => m.message?.role)).toEqual(["assistant"]); + }); + }); + + it("routes later turns to the new session file after a switch", async () => { + await withTempDir(async dir => { + let sessionFile = path.join(dir, "first.jsonl"); + const recorder = new AdvisorTranscriptRecorder( + () => sessionFile, + () => dir, + ); + recorder.record(assistantMessage("before switch", 1)); + sessionFile = path.join(dir, "second.jsonl"); + recorder.record(assistantMessage("after switch", 2)); + await recorder.close(); + + const first = await readMessageEntries(path.join(dir, "first", ADVISOR_TRANSCRIPT_FILENAME)); + const second = await readMessageEntries(path.join(dir, "second", ADVISOR_TRANSCRIPT_FILENAME)); + expect(first).toHaveLength(1); + expect(first[0].message?.usage?.input).toBe(1); + expect(second).toHaveLength(1); + expect(second[0].message?.usage?.input).toBe(2); + }); + }); +}); diff --git a/packages/coding-agent/test/internal-urls/history-protocol.test.ts b/packages/coding-agent/test/internal-urls/history-protocol.test.ts index 503170f41..348b4f7c7 100644 --- a/packages/coding-agent/test/internal-urls/history-protocol.test.ts +++ b/packages/coding-agent/test/internal-urls/history-protocol.test.ts @@ -13,6 +13,7 @@ import * as fs from "node:fs/promises"; import * as os from "node:os"; import * as path from "node:path"; import { InternalUrlRouter } from "@oh-my-pi/pi-coding-agent/internal-urls"; +import { HistoryProtocolHandler } from "@oh-my-pi/pi-coding-agent/internal-urls/history-protocol"; import { AgentRegistry } from "@oh-my-pi/pi-coding-agent/registry/agent-registry"; import type { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session"; import { CURRENT_SESSION_VERSION } from "@oh-my-pi/pi-coding-agent/session/session-entries"; @@ -186,4 +187,67 @@ describe("history:// protocol", () => { expect(error?.message).toContain("no transcript"); }); + + it("hides advisor transcripts from the index and direct lookup", async () => { + AgentRegistry.global().register({ + id: "HubAgent", + displayName: "task", + kind: "sub", + session: fakeLiveSession([]), + status: "idle", + }); + AgentRegistry.global().register({ + id: "Main/advisor", + displayName: "advisor", + kind: "advisor", + session: fakeLiveSession([{ role: "user", content: "should stay hidden", timestamp: 1 }]), + status: "parked", + }); + AgentRegistry.global().register({ + id: "AdvisorProbe", + displayName: "advisor", + kind: "advisor", + session: fakeLiveSession([{ role: "user", content: "should stay hidden", timestamp: 1 }]), + status: "parked", + }); + + // Index lists the subagent but never the advisor. + const index = await InternalUrlRouter.instance().resolve("history://"); + expect(index.content).toContain("HubAgent"); + expect(index.content).not.toContain("advisor"); + + // Direct lookup of an advisor-kind ref is reported as unknown — the driving + // agent must not be able to read it via history://. + const error = await InternalUrlRouter.instance() + .resolve("history://AdvisorProbe") + .then( + () => null, + err => err as Error, + ); + expect(error).toBeInstanceOf(Error); + expect(error?.message).toContain("Unknown agent"); + }); + + it("omits advisor refs from history:// completions", async () => { + AgentRegistry.global().register({ + id: "HubAgent", + displayName: "task", + kind: "sub", + session: fakeLiveSession([]), + status: "idle", + }); + AgentRegistry.global().register({ + id: "AdvisorProbe", + displayName: "advisor", + kind: "advisor", + session: null, + sessionFile: "/tmp/x/__advisor.jsonl", + status: "parked", + }); + + const completions = await new HistoryProtocolHandler().complete(); + const values = completions.map(c => c.value); + expect(values).toContain("HubAgent"); + expect(values).not.toContain("AdvisorProbe"); + }); }); diff --git a/packages/coding-agent/test/task/output-manager.test.ts b/packages/coding-agent/test/task/output-manager.test.ts index b374105fe..de33b13ad 100644 --- a/packages/coding-agent/test/task/output-manager.test.ts +++ b/packages/coding-agent/test/task/output-manager.test.ts @@ -65,4 +65,13 @@ describe("AgentOutputManager", () => { expect(await mgr.allocate("Bob")).toBe("Anna.Bob-2"); expect(await mgr.allocate("Dave")).toBe("Anna.Dave"); }); + + it("reserves the advisor transcript stem so a task can't clobber __advisor.jsonl", async () => { + const mgr = new AgentOutputManager(() => null); + // A subagent allocated `__advisor` would write `__advisor.jsonl`, colliding with + // the advisor transcript in the same artifacts dir; the stem is pre-reserved. + expect(await mgr.allocate("__advisor")).toBe("__advisor-2"); + // Unrelated names sharing the prefix are unaffected. + expect(await mgr.allocate("__advisor-notes")).toBe("__advisor-notes"); + }); }); diff --git a/packages/coding-agent/test/tools/irc.test.ts b/packages/coding-agent/test/tools/irc.test.ts index 37ed6e562..b246f92a0 100644 --- a/packages/coding-agent/test/tools/irc.test.ts +++ b/packages/coding-agent/test/tools/irc.test.ts @@ -434,6 +434,24 @@ describe("IRC", () => { expect(text).toContain("Parked agents are revived automatically"); }); + it("op=list hides advisor-kind refs from the peer roster", async () => { + const sub = makeFakeSession(); + registry.register({ id: "0-Worker", displayName: "task", kind: "sub", session: sub.session }); + registry.register({ + id: "0-Main/advisor", + displayName: "advisor", + kind: "advisor", + session: null, + status: "parked", + }); + + const tool = new IrcTool(makeToolSession(registry, "0-Main")); + const result = await tool.execute("call-1", { op: "list" }); + const peerIds = result.details?.peers?.map(peer => peer.id) ?? []; + expect(peerIds).toContain("0-Worker"); + expect(peerIds).not.toContain("0-Main/advisor"); + }); + it("op=send returns receipts immediately without waiting for a reply", async () => { const sub = makeFakeSession(); registry.register({ id: "0-Sub", displayName: "task", kind: "sub", session: sub.session }); diff --git a/packages/snapcompact/CHANGELOG.md b/packages/snapcompact/CHANGELOG.md index 2d79cbaa5..ef7e0ef74 100644 --- a/packages/snapcompact/CHANGELOG.md +++ b/packages/snapcompact/CHANGELOG.md @@ -1,6 +1,9 @@ # Changelog ## [Unreleased] +### Changed + +- Updated summary text for consistent descriptions of archived tool output ## [16.0.8] - 2026-06-18 @@ -102,4 +105,4 @@ ### Fixed - Fixed frame rendering at archive chunk boundaries to reopen dim spans when a chunk ends inside a dimmed tool-result segment -- Fixed message serialization to strip user- and assistant-provided dim markers so only renderer-generated dim spans can be applied +- Fixed message serialization to strip user- and assistant-provided dim markers so only renderer-generated dim spans can be applied \ No newline at end of file diff --git a/packages/snapcompact/test/snapcompact.test.ts b/packages/snapcompact/test/snapcompact.test.ts index 6cb6df38a..85b584b17 100644 --- a/packages/snapcompact/test/snapcompact.test.ts +++ b/packages/snapcompact/test/snapcompact.test.ts @@ -661,9 +661,8 @@ describe("compact", () => { expect(result.firstKeptEntryId).toBe("kept-1"); expect(result.tokensBefore).toBe(99000); // Reading instructions reflect the default (anthropic 11on16-bw) shape. - expect(result.summary).toContain("29 characters per row"); + expect(result.summary).toContain("29 characters wide"); expect(result.summary).toContain("dim gray"); - expect(result.summary).toContain("plain black ink"); expect(result.summary).toContain("snapcompact frame"); // File operations are upserted like every other compaction summary: // one grouped tree with per-file access markers. @@ -703,7 +702,7 @@ describe("compact", () => { // Conversation text outside the span stays in black bw ink (frame 1). const first = decodePng(Buffer.from(archive?.frames[0].data ?? "", "base64")); expect(new Set(first.pixels).has(7)).toBe(true); - expect(result.summary).toContain("dim gray ink"); + expect(result.summary).toContain("archived tool output"); }); it("keeps frames free of dim ink when dimToolResults is false", async () => { @@ -716,7 +715,7 @@ describe("compact", () => { const archive = snapcompact.getPreservedArchive(result.preserveData); const decoded = decodePng(Buffer.from(archive?.frames[0].data ?? "", "base64")); expect(new Set(decoded.pixels).has(9)).toBe(false); - expect(result.summary).not.toContain("dim gray ink"); + expect(result.summary).not.toContain("archived tool output"); }); it("keeps history past the frame budget as a text tail instead of dropping it", async () => {