feat(coding-agent): supported advisor transcript persistence

- Implemented `AdvisorTranscriptRecorder` to persist advisor sessions to append-only `__advisor.jsonl` files.
- Integrated transcript recording into agent sessions with managed flushing, atomic file switching, and synthetic turn attribution.
- Restricted advisor-kind agents by excluding them from rosters, history protocols, messaging, and interactive agent commands.
- Reserved the `__advisor` filename stem across the output manager and task registry to prevent task ID collisions.
This commit is contained in:
can1357
2026-06-19 03:50:00 +02:00
parent 2b843dc744
commit 29d250fae2
20 changed files with 649 additions and 31 deletions
+20 -1
View File
@@ -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: `<session>/__advisor.jsonl`
- subagent advisor (`advisor.subagents: true`): `<session>/<SubId>/__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 `<id>.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.
+7
View File
@@ -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 <focus text>` 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 (`<session>/__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
@@ -1,3 +1,4 @@
export * from "./advise-tool";
export * from "./runtime";
export * from "./transcript-recorder";
export * from "./watchdog";
@@ -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 `<id>.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 (`<session>/__advisor.jsonl`, `<session>/<SubId>/__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<void>;
/**
* @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<unknown>,
) {
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<void> {
return this.#enqueueResult(async () => {
if (this.#manager) await this.#manager.flush();
});
}
/** Flush and close the writer, releasing the session file. */
close(): Promise<void> {
return this.#enqueueResult(() => this.#closeManager());
}
async #closeManager(): Promise<void> {
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>): void {
this.#queue = this.#queue.then(work, work).catch(err => {
logger.debug("advisor transcript record failed", { err: String(err) });
});
}
#enqueueResult(work: () => Promise<void>): Promise<void> {
const next = this.#queue.then(work, work);
this.#queue = next.catch(() => {});
return next;
}
}
+25 -13
View File
@@ -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);
@@ -40,9 +40,12 @@ export class HistoryProtocolHandler implements ProtocolHandler {
async resolve(url: InternalUrl): Promise<InternalResource> {
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<UrlCompletion[]> {
return AgentRegistry.global()
.list()
.filter(ref => ref.kind !== "advisor")
.map(ref => ({
value: ref.id,
description: `${ref.status} · ${ref.kind}${ref.parentId ? ` · parent ${ref.parentId}` : ""}`,
+8
View File
@@ -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") {
@@ -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 `<owner>/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();
@@ -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 {
@@ -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 `<session>/__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<void> = 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 `<session>/__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<boolean> {
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
// `<old>/__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) {
+1 -1
View File
@@ -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 =>
@@ -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 `<id>.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);
}
/**
+1 -1
View File
@@ -182,7 +182,7 @@ export class IrcTool implements AgentTool<typeof ircSchema, IrcDetails> {
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,
@@ -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");
});
});
@@ -0,0 +1,163 @@
/**
* Contracts: AdvisorTranscriptRecorder persists the advisor agent's turns to a
* subagent-style JSONL (`<session>/__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<T>(fn: (dir: string) => Promise<T>): Promise<T> {
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<AdvisorEntry[]> {
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 <session>/__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);
});
});
});
@@ -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");
});
});
@@ -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");
});
});
@@ -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 });
+4 -1
View File
@@ -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
@@ -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 <files> 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 () => {