import type { AgentIdentity, AgentTelemetryConfig, AgentTool, AgentToolContext, AgentToolResult, AgentToolUpdateCallback, } from "@oh-my-pi/pi-agent-core"; import { escapeXmlAttribute, escapeXmlText } from "@oh-my-pi/pi-utils"; import { type } from "arktype"; import adviseDescription from "../prompts/advisor/advise-tool.md" with { type: "text" }; const adviseSchema = type({ note: type("string").describe( "One concrete piece of advice for the agent you are watching. Terse, specific, actionable.", ), "severity?": type("'nit' | 'concern' | 'blocker'").describe("How strongly to weigh this. Omit for a plain nit."), }); export type AdviseParams = typeof adviseSchema.infer; export type AdvisorSeverity = "nit" | "concern" | "blocker"; export interface AdviseDetails { note: string; severity?: AdvisorSeverity; /** Which configured advisor produced this note (omitted for the default advisor). */ advisor?: string; } /** One queued advice note. */ export interface AdvisorNote { note: string; severity?: AdvisorSeverity; /** Which configured advisor produced this note (omitted for the default advisor). */ advisor?: string; } /** Details payload on the batched `advisor` custom message rendered in the transcript. */ export interface AdvisorMessageDetails { notes: AdvisorNote[]; } /** * Behavioral framing for the watched agent — advice, not orders. Carried as a * tag attribute (rather than a prose header) so the rendered agent-facing output * stays a clean `` block. The primary agent's system prompt never * mentions advisories, so this is its only cue for how to treat them. */ const ADVISOR_GUIDANCE = "weigh, don't blindly obey"; /** * Render a batch of advisor notes as the agent-facing message body: one * `` element per note, severity as an attribute. Shared by the * non-interrupting YieldQueue dispatcher and the interrupting steer path so both * build byte-identical content. */ export function formatAdvisorBatchContent(notes: readonly AdvisorNote[]): string { return notes .map(n => { const severity = n.severity ? ` severity="${n.severity}"` : ""; const who = n.advisor ? ` advisor="${escapeXmlAttribute(n.advisor)}"` : ""; return `\n${escapeXmlText(n.note)}\n`; }) .join("\n"); } /** * Whether advice at this severity should interrupt the running agent (delivered * via the steering channel, aborting in-flight tools) rather than ride the * non-interrupting aside queue that lands at the next step boundary. `concern` * and `blocker` interrupt; a plain `nit` queues. */ export function isInterruptingSeverity(severity: AdvisorSeverity | undefined): boolean { return severity === "concern" || severity === "blocker"; } /** * Append a staleness caveat to an advisor note when newer primary turns arrived * after the reviewed transcript window (i.e. `hasFreshBacklog` is true on the * advisor runtime at delivery time). Pure function — no session coupling — so it * can be unit-tested in isolation and called from `AgentSession#routeAdvice`. */ export function annotateForStaleness(note: string, hasFreshBacklog: boolean): string { if (!hasFreshBacklog) return note; return `${note}\n\n_(Note: newer primary turns arrived after this reviewed window — verify this still applies.)_`; } /** How an advisor note is routed to the primary. */ export type AdvisorDeliveryChannel = "aside" | "steer" | "preserve"; /** Half-open turn-count fence for the post-interrupt cooldown. */ export function isAdvisorInterruptImmuneTurnActive(opts: { completedTurns: number; immuneTurnStart: number | undefined; immuneTurns: number; }): boolean { if (opts.immuneTurnStart === undefined || opts.immuneTurns <= 0) return false; return opts.completedTurns < opts.immuneTurnStart + opts.immuneTurns; } /** * Decide how one advisor note reaches the primary agent. * * - A non-interrupting `nit` always rides the non-interrupting aside queue. * - An interrupting `concern`/`blocker` is normally steered into the agent: into * the live turn while one is streaming, or (when idle) a triggered turn so the * advice is acted on immediately. * - If the primary tail is already a terminal text answer and there is no queued * work, late interrupting advice is preserved as a visible card instead of * waking the primary to restate completion. * - After a deliberate user interrupt (`autoResumeSuppressed`) the advisor must * not auto-resume the stopped run. While the agent is idle — or still tearing * the interrupted turn down (`aborting`) — the note is preserved as a visible * card instead of restarting the run. But once a turn is actively streaming * again (a resume the user already drove), steering the note in does NOT * auto-resume anything, so it is delivered live. Parking it during an active * run instead strands it (it never reaches the running agent) and the withheld * notes dump as one burst at the next user prompt — the bug this guards. * - During the post-interrupt immune-turn window, further `concern`/`blocker` * notes are downgraded to asides; preservation still wins. */ export function resolveAdvisorDeliveryChannel(opts: { severity: AdvisorSeverity | undefined; autoResumeSuppressed: boolean; streaming: boolean; aborting: boolean; terminalAnswerNoQueuedWork?: boolean; interruptImmuneTurnActive?: boolean; }): AdvisorDeliveryChannel { if (!isInterruptingSeverity(opts.severity)) return "aside"; if (opts.autoResumeSuppressed && (opts.aborting || !opts.streaming)) return "preserve"; if (opts.terminalAnswerNoQueuedWork && !opts.streaming && !opts.aborting) return "preserve"; if (opts.interruptImmuneTurnActive) return "aside"; return "steer"; } /** * Derive the advisor loop's telemetry from the primary session's config so the * advisor model's GenAI spans and usage/cost hooks (onChatUsage, onCostDelta, * costEstimator) fire under the same pipeline as every other model call — * stamped with the advisor's own agent identity. `conversationId` is cleared so * the advisor loop falls back to its own `-advisor` session id for * `gen_ai.conversation.id` instead of inheriting the primary's conversation. * * Returns undefined when the primary has no telemetry (instrumentation off), so * the advisor `Agent` stays a zero-overhead no-op as well. */ export function deriveAdvisorTelemetry( primaryTelemetry: AgentTelemetryConfig | undefined, identity: AgentIdentity, ): AgentTelemetryConfig | undefined { if (!primaryTelemetry) return undefined; return { ...primaryTelemetry, agent: identity, conversationId: undefined }; } /** * The tools an advisor receives by default when its config omits `tools` — the * read-only investigative set. The full available pool is every built tool the * session has (the advisor is a full agent); a config's `tools` selects from it. */ export const ADVISOR_DEFAULT_TOOL_NAMES: ReadonlySet = new Set(["read", "grep", "glob"]); function advisorNoteDedupeKey(note: string): string { return note.trim().replace(/\s+/g, " "); } /** Rank advisor severities so the dedupe state can detect a real escalation * (nit → concern → blocker) versus a verbatim repeat. `undefined` defers to * `nit` because the schema treats an omitted severity as a plain nit. */ const ADVISOR_SEVERITY_RANK: Record = { nit: 1, concern: 2, blocker: 3 }; function advisorSeverityRank(severity: AdvisorSeverity | undefined): number { return ADVISOR_SEVERITY_RANK[severity ?? "nit"]; } export class AdviseTool implements AgentTool { readonly name = "advise"; readonly label = "Advise"; readonly description = adviseDescription; readonly parameters = adviseSchema; readonly intent = "omit" as const; /** Highest delivered severity rank per normalized note. A new call passes * through only when its rank strictly exceeds the recorded one (a real * escalation: nit → concern → blocker), so an advisor cannot bypass dedupe * by retagging the same text at a lower or equal severity. */ #deliveredNoteSeverities = new Map(); constructor(private readonly onAdvice: (note: string, severity?: AdviseDetails["severity"]) => void) {} /** Clear delivered-note memory when the advisor starts a fresh conversation. */ resetDeliveredNotes(): void { this.#deliveredNoteSeverities.clear(); } async execute( _toolCallId: string, args: AdviseParams, _signal?: AbortSignal, _onUpdate?: AgentToolUpdateCallback, _context?: AgentToolContext, ): Promise> { const key = advisorNoteDedupeKey(args.note); const rank = advisorSeverityRank(args.severity); const previousRank = this.#deliveredNoteSeverities.get(key) ?? 0; if (rank <= previousRank) { return { content: [{ type: "text", text: "Duplicate advice ignored." }], details: { note: args.note, severity: args.severity }, useless: true, }; } this.#deliveredNoteSeverities.set(key, rank); this.onAdvice(args.note, args.severity); return { content: [{ type: "text", text: "Recorded." }], details: { note: args.note, severity: args.severity }, useless: true, }; } }