/** * Turn reducer — folds parsed agy stream events into provider-facing state. * * agy streams assistant text deltas on `agent_response` steps (the visible * answer — thought text is never exposed by print mode) and the final text * also arrives in the terminal `result` event. Tool steps stream live * (ACTIVE -> DONE/ERROR) and are surfaced as structured {@link AgyActivity} * events the provider renders as native pi tool cards. */ import type { JsonObject } from "@earendil-works/pi-ai"; import type { AgyStepUpdate, AgyUsage, ParsedAgyEvent } from "./events.ts"; import { parseAgyLine } from "./events.ts"; /** * Print mode reports tiny `thinking_tokens` counts even for planner/tool glue * that agy's interactive TUI renders without a Thought row. The observed TUI * boundary is not part of stream-json, so suppress only the clearly incidental * traces; substantive reasoning summaries remain visible. */ const MIN_VISIBLE_THOUGHT_TOKENS = 64; export type AgyActivity = | { type: "tool_start"; stepId?: number; name: string; args: JsonObject } | { type: "tool_done"; stepId?: number; name: string; args: JsonObject; output?: string; durationSeconds?: number; } | { type: "tool_error"; stepId?: number; name: string; args: JsonObject; message: string; } | { /** Synthetic — pushed by the bridge when agy invokes a `pi__*` tool. * Never produced by applyEvent. */ type: "bridge_call"; id: string; name: string; args: JsonObject; } | { /** Synthetic — a persisted native conversation disappeared, so the * runtime started fresh from the bounded Pi branch history. */ type: "conversation_fallback"; } | { /** Synthetic — pushed by the runtime when a stalled agy turn is killed * and retried. Never produced by applyEvent. */ type: "stall"; /** 1-based retry number. */ retry: number; maxRetries: number; stalledMs: number; toolActive: boolean; } | { type: "text"; delta: string; stepId?: number } | { /** agy's own collapsed reasoning line — thought text is never streamed, * only its token count (and the response step's duration) are. */ type: "thought"; tokens: number; durationSeconds?: number; } | { type: "usage"; usage: AgyUsage } | { /** Terminal status, already normalized: anything agy did not explicitly * report as successful arrives as "ERROR". */ type: "result"; status: "OK" | "ERROR"; response: string; error: string | undefined; usage: AgyUsage | undefined; }; export interface AgyTurnOutcome { /** Conversation id for `--conversation` resume; set by init/result. */ conversationId: string | undefined; /** "UNKNOWN" until the result event lands; then "OK" or "ERROR". */ status: "OK" | "ERROR" | "UNKNOWN"; /** Final assistant text (empty until the result event). */ response: string; /** Error message when status is ERROR. */ error: string | undefined; usage: AgyUsage | undefined; /** Structured activity events, appended as the stream unfolds. */ activities: AgyActivity[]; /** True once the result event has been seen. */ finished: boolean; } /** * Flatten `subagent_info.subagents` into the arg map the roster's * pickName/pickDetail and the call renderer already read (Name/Type/Task). * The first spawn leads; extra spawns in the same step fold into the name. */ function subagentStepArgs(step: AgyStepUpdate): JsonObject { const subs = step.subagent_info?.subagents; if (!subs?.length) return {}; const first = subs[0]; const args: JsonObject = {}; if (first.role) args.Name = subs.length > 1 ? `${first.role} +${subs.length - 1}` : first.role; if (first.type_name) args.Type = first.type_name; if (first.initial_prompt) args.Task = first.initial_prompt; if (first.conversation_id) args.ConversationId = first.conversation_id; if (first.log_uri) args.LogUri = first.log_uri; return args; } /** DONE subagent steps carry no output field; summarize the spawn records. */ function subagentStepOutput(step: AgyStepUpdate): string | undefined { const subs = step.subagent_info?.subagents; if (!subs?.length) return undefined; return subs .map((sub) => [ `subagent ${sub.role ?? sub.type_name ?? "spawned"}`, sub.conversation_id ? `conversation ${sub.conversation_id}` : undefined, ] .filter(Boolean) .join(" · "), ) .join("\n"); } export function newTurnOutcome(): AgyTurnOutcome { return { conversationId: undefined, status: "UNKNOWN", response: "", error: undefined, usage: undefined, activities: [], finished: false, }; } export type { AgyUsage }; /** Fold one parsed event into the outcome; returns new activity events (if any). */ export function applyEvent(outcome: AgyTurnOutcome, event: ParsedAgyEvent): AgyActivity[] { const activities: AgyActivity[] = []; switch (event.kind) { case "init": { outcome.conversationId = event.conversationId ?? outcome.conversationId; break; } case "step": { const step = event.step; if (step.conversation_id && !outcome.conversationId) { outcome.conversationId = step.conversation_id; } if (step.step_type === "tool" || step.step_type === "subagent") { // Subagent steps (invoke_subagent/send_message/manage_subagents/…) // arrive as step_type "subagent" with the payload under subagent_info, // not tool_info — normalize both into the activity arg/output shape. const name = step.tool_name ?? step.tool_info?.name ?? "tool"; const args = step.tool_info?.parameters ?? subagentStepArgs(step); if (step.state === "ACTIVE") { activities.push({ type: "tool_start", stepId: step.step_index, name, args, }); } else if (step.state === "DONE") { activities.push({ type: "tool_done", stepId: step.step_index, name, // agy nests tool output under tool_info on DONE steps. output: step.output ?? step.tool_info?.output ?? subagentStepOutput(step), durationSeconds: typeof step.duration_seconds === "number" ? step.duration_seconds : undefined, // Kept so call_mcp_tool completions can be correlated with the // bridged server (ServerName) by the provider. args, }); } else if (step.state === "ERROR") { // agy puts the error detail under tool_info on ERROR steps (e.g. // generate_image 429 rate limits); step.error is often absent. const message = step.error?.message ?? step.tool_info?.error?.message ?? "tool error"; activities.push({ type: "tool_error", stepId: step.step_index, name, args, message: message.replace(/\s+/g, " ").slice(0, 160), }); } } else if (step.step_type === "agent_response") { // agent_response steps stream the visible answer text as continuation // chunks on ACTIVE steps plus a final chunk on DONE. agy never exposes // thought text in print mode; thinking_tokens in usage is the only // reasoning trace. if (typeof step.text_delta === "string" && step.text_delta) { activities.push({ type: "text", delta: step.text_delta, stepId: step.step_index }); } if (step.usage) { // Running per-response usage; the result event carries the totals. outcome.usage = step.usage; activities.push({ type: "usage", usage: step.usage }); const thoughtTokens = step.usage.thinking_tokens; if ( step.state === "DONE" && typeof thoughtTokens === "number" && thoughtTokens >= MIN_VISIBLE_THOUGHT_TOKENS ) { activities.push({ type: "thought", tokens: thoughtTokens, durationSeconds: typeof step.duration_seconds === "number" ? step.duration_seconds : undefined, }); } } } break; } case "result": { const result = event.result; outcome.finished = true; outcome.conversationId = result.conversation_id ?? outcome.conversationId; // Fail closed: only an explicitly successful status completes the turn. // agy >= 1.1.22 reports SUCCESS (older builds: OK); FAILURE, CANCELLED, // and TIMEOUT all exist too, and an unrecognized status must never be // rendered as a normal answer. outcome.status = result.status === "SUCCESS" || result.status === "OK" ? "OK" : "ERROR"; outcome.response = result.response ?? ""; outcome.error = result.error; outcome.usage = result.usage ?? outcome.usage; activities.push({ type: "result", status: outcome.status, response: outcome.response, error: outcome.error, usage: outcome.usage, }); break; } case "unknown": break; } outcome.activities.push(...activities); return activities; } /** Parse a full NDJSON document (test helper) and reduce it. */ export function reduceAgyStream(text: string): AgyTurnOutcome { const outcome = newTurnOutcome(); for (const line of text.split("\n")) { applyEvent(outcome, parseAgyLine(line)); } return outcome; }