import * as fs from "node:fs/promises"; 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 { visitEntriesFromFileStream } from "../session/session-loader"; 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"; /** * Transcript filename for an advisor: `__advisor.jsonl` for the legacy/default * advisor (empty slug), `__advisor..jsonl` for a named advisor. The `.` * separator keeps named files out of the output manager's `-` bump namespace. */ export function advisorTranscriptFilename(slug: string): string { return slug ? `${ADVISOR_TRANSCRIPT_STEM}.${slug}${JSONL_SUFFIX}` : ADVISOR_TRANSCRIPT_FILENAME; } /** Whether a filename is any advisor transcript (`__advisor.jsonl` or `__advisor..jsonl`). */ export function isAdvisorTranscriptName(name: string): boolean { return ( name === ADVISOR_TRANSCRIPT_FILENAME || (name.startsWith(`${ADVISOR_TRANSCRIPT_STEM}.`) && name.endsWith(JSONL_SUFFIX)) ); } /** * Sum the advisor spend already persisted next to a primary session transcript, * keyed by advisor slug. * * The ledger a session keeps in memory only covers the current process, so a * resumed session would report zero until the next advisor turn. The recorded * transcripts are the durable copy of exactly the same finalized messages, so * they are read back through the shared loader - no lock, no writer, and no * second parser to keep in step with the session format. * * Only the session's own advisors count: subagent advisors write to * `//__advisor.jsonl`, and their spend belongs to the subagent, * not to this roster. Hence the scan stays at the top level of the directory. */ export async function loadAdvisorTranscriptCosts(sessionFile: string | undefined): Promise> { const costs = new Map(); if (!sessionFile?.endsWith(JSONL_SUFFIX)) return costs; const directory = sessionFile.slice(0, -JSONL_SUFFIX.length); const dirents = await fs.readdir(directory, { withFileTypes: true }).catch(() => []); for (const dirent of dirents) { if (!dirent.isFile() || !isAdvisorTranscriptName(dirent.name)) continue; const slug = dirent.name === ADVISOR_TRANSCRIPT_FILENAME ? "" : dirent.name.slice(`${ADVISOR_TRANSCRIPT_STEM}.`.length, -JSONL_SUFFIX.length); let total = 0; let validHeader: boolean | undefined; try { await visitEntriesFromFileStream(path.join(directory, dirent.name), entry => { const isObject = typeof entry === "object" && entry !== null; if (validHeader === undefined) { validHeader = isObject && entry.type === "session" && typeof entry.id === "string"; return; } // A syntactically valid but non-object entry (e.g. a bare `null` // line) must cost only itself, not crash entry.type access and // discard everything accumulated for this transcript. if (!validHeader || !isObject || entry.type !== "message") return; const message = entry.message; if (!message || typeof message !== "object" || message.role !== "assistant") return; // One malformed usage block must cost that entry only, not the // whole transcript's total. const total_ = message.usage?.cost?.total; if (typeof total_ === "number" && Number.isFinite(total_)) total += total_; }); } catch (err) { logger.debug("advisor transcript cost read failed", { file: dirent.name, err: String(err) }); continue; } if (total > 0) costs.set(slug, total); } return costs; } /** * 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; #filename: string; /** Serializes the async open/close against synchronous appends so records land in order. */ #queue: Promise; /** * @param filename Transcript filename within the session dir. Defaults to * `__advisor.jsonl`; named advisors pass `__advisor..jsonl` via * {@link advisorTranscriptFilename}. * @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 file at once. */ constructor( private readonly resolveSessionFile: () => string | undefined, private readonly resolveCwd: () => string, filename: string = ADVISOR_TRANSCRIPT_FILENAME, after?: Promise, ) { this.#filename = filename; 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), this.#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; } }