import type { AgentSettledEvent, AgentStartEvent, MessageEndEvent, MessageStartEvent, MessageUpdateEvent, ToolExecutionStartEvent, TurnEndEvent, TurnStartEvent, Theme, } from "@earendil-works/pi-coding-agent"; import type { IconMode, TelemetryConfig } from "./config.ts"; import { resolveGlyphs } from "./icons.ts"; import { fmtTokens, formatDuration } from "./utils.ts"; const STALL_THRESHOLD_MS = 1000; type TelemetryEvent = | AgentStartEvent | AgentSettledEvent | TurnStartEvent | MessageStartEvent | MessageUpdateEvent | MessageEndEvent | ToolExecutionStartEvent | TurnEndEvent; type AgentMessage = MessageStartEvent["message"]; type AssistantMessage = Extract; interface MessageTiming { lastUpdateMs: number; firstOutputMs: number | null; inStall: boolean; } interface TurnTiming { startMs: number; firstTokenMs: number | null; currentMessage: MessageTiming | null; messages: AssistantMessage[]; generationMs: number; stallMs: number; stallCount: number; } export interface TurnTelemetry { tps: number | null; ttftMs: number; totalMs: number; inputTokens: number; outputTokens: number; cacheReadTokens: number; cacheWriteTokens: number; cacheHitRate: number | undefined; stallMs: number; stallCount: number; rateUsdPerMTokens: number | null; generationMs: number; totalTokens: number; costUsd: number; measurementMs: number | null; } function isAssistantMessage(message: AgentMessage): message is AssistantMessage { return message.role === "assistant"; } function round(value: number, decimals: number): number { const factor = 10 ** decimals; return Math.round(value * factor) / factor; } function getCacheHitRate(inputTokens: number, cacheReadTokens: number, cacheWriteTokens: number): number | undefined { const promptTokens = inputTokens + cacheReadTokens + cacheWriteTokens; if ((cacheReadTokens <= 0 && cacheWriteTokens <= 0) || promptTokens <= 0) return undefined; return (cacheReadTokens / promptTokens) * 100; } export class TurnTelemetryTracker { private readonly now: () => number; private turn: TurnTiming | undefined; private agentStartMs: number | null = null; private agentTurns: TurnTelemetry[] = []; constructor(now: () => number = () => performance.now()) { this.now = now; } handle(event: TelemetryEvent): TurnTelemetry | undefined { switch (event.type) { case "agent_start": if (this.agentStartMs === null) { this.agentStartMs = this.now(); this.agentTurns = []; } return; case "agent_settled": return this.endAgent(); case "turn_start": this.startTurn(); return; case "message_start": this.startMessage(event.message); return; case "message_update": this.updateMessage(event); return; case "message_end": this.endMessage(event.message); return; case "tool_execution_start": return; case "turn_end": return this.endTurnAndCollect(); } } private startTurn(): void { this.turn = { startMs: this.now(), firstTokenMs: null, currentMessage: null, messages: [], generationMs: 0, stallMs: 0, stallCount: 0, }; } private startMessage(message: AgentMessage): void { if (!this.turn || !isAssistantMessage(message)) return; const now = this.now(); this.turn.currentMessage = { lastUpdateMs: now, firstOutputMs: null, inStall: false, }; } private updateMessage(event: MessageUpdateEvent): void { const turn = this.turn; const current = turn?.currentMessage; const streamEvent = event.assistantMessageEvent; if ( streamEvent.type !== "text_delta" && streamEvent.type !== "thinking_delta" && streamEvent.type !== "toolcall_delta" ) return; if (streamEvent.delta.length === 0) return; const message = event.message; if (!turn || !current || !isAssistantMessage(message)) return; const now = this.now(); if (current.firstOutputMs === null) { current.firstOutputMs = now; turn.firstTokenMs ??= now; current.lastUpdateMs = now; return; } const gap = now - current.lastUpdateMs; if (gap >= STALL_THRESHOLD_MS) { if (!current.inStall) turn.stallCount++; current.inStall = true; turn.stallMs += gap; } else { current.inStall = false; } current.lastUpdateMs = now; } private endMessage(message: AgentMessage): void { const turn = this.turn; if (!turn || !isAssistantMessage(message)) return; const current = turn.currentMessage; if (current) { const endMs = this.now(); turn.generationMs = endMs - turn.startMs; if (current.firstOutputMs === null && message.usage.output > 0) { turn.firstTokenMs ??= endMs; } turn.currentMessage = null; } turn.messages.push(message); } private endTurnAndCollect(): TurnTelemetry | undefined { const telemetry = this.endTurn(); if (telemetry && this.agentStartMs !== null) this.agentTurns.push(telemetry); return telemetry; } private endTurn(): TurnTelemetry | undefined { const turn = this.turn; this.turn = undefined; if (!turn || turn.firstTokenMs === null || turn.messages.length === 0) return; const endMs = this.now(); let inputTokens = 0; let outputTokens = 0; let cacheReadTokens = 0; let cacheWriteTokens = 0; let totalTokens = 0; let costUsd = 0; for (const message of turn.messages) { inputTokens += message.usage.input; outputTokens += message.usage.output; cacheReadTokens += message.usage.cacheRead; cacheWriteTokens += message.usage.cacheWrite; totalTokens += message.usage.totalTokens; costUsd += message.usage.cost.total; } if (![inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens, totalTokens, costUsd].every(Number.isFinite)) { throw new Error("Invalid assistant usage in turn telemetry"); } const measurementMs = outputTokens > 0 && turn.generationMs > 0 ? turn.generationMs : null; const tps = measurementMs === null ? null : round(outputTokens / (measurementMs / 1000), 1); const validCost = Number.isFinite(costUsd) && costUsd > 0; const validTokens = Number.isFinite(totalTokens) && totalTokens > 0; return { tps, ttftMs: turn.firstTokenMs - turn.startMs, totalMs: endMs - turn.startMs, inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens, cacheHitRate: getCacheHitRate(inputTokens, cacheReadTokens, cacheWriteTokens), stallMs: turn.stallMs, stallCount: turn.stallCount, rateUsdPerMTokens: validCost && validTokens ? round(costUsd / (totalTokens / 1_000_000), 2) : null, generationMs: turn.generationMs, totalTokens, costUsd: validCost ? costUsd : 0, measurementMs, }; } private endAgent(): TurnTelemetry | undefined { const startMs = this.agentStartMs; const turns = this.agentTurns; this.agentStartMs = null; this.agentTurns = []; if (startMs === null || turns.length === 0) return; const outputTokens = turns.reduce((sum, turn) => sum + turn.outputTokens, 0); const inputTokens = turns.reduce((sum, turn) => sum + turn.inputTokens, 0); const cacheReadTokens = turns.reduce((sum, turn) => sum + turn.cacheReadTokens, 0); const cacheWriteTokens = turns.reduce((sum, turn) => sum + turn.cacheWriteTokens, 0); const totalTokens = turns.reduce((sum, turn) => sum + turn.totalTokens, 0); const costUsd = turns.reduce((sum, turn) => sum + turn.costUsd, 0); const stallMs = turns.reduce((sum, turn) => sum + turn.stallMs, 0); const stallCount = turns.reduce((sum, turn) => sum + turn.stallCount, 0); const generationMs = turns.reduce((sum, turn) => sum + turn.generationMs, 0); const measurementMs = outputTokens > 0 && generationMs > 0 ? generationMs : null; const tps = measurementMs === null ? null : round(outputTokens / (measurementMs / 1000), 1); const validRate = costUsd > 0 && totalTokens > 0; return { tps, ttftMs: turns[0]!.ttftMs, totalMs: this.now() - startMs, inputTokens, outputTokens, cacheReadTokens, cacheWriteTokens, cacheHitRate: getCacheHitRate(inputTokens, cacheReadTokens, cacheWriteTokens), stallMs, stallCount, rateUsdPerMTokens: validRate ? round(costUsd / (totalTokens / 1_000_000), 2) : null, generationMs, totalTokens, costUsd, measurementMs, }; } } function formatTurnDuration(ms: number): string { return ms < 60_000 ? `${(ms / 1000).toFixed(1)}s` : formatDuration(ms); } function formatClockTime(now: Date): string { return [now.getHours(), now.getMinutes(), now.getSeconds()] .map((part) => part.toString().padStart(2, "0")) .join(":"); } type TelemetryTone = "dim" | "muted"; export function formatTurnTelemetry( telemetry: TurnTelemetry, theme: Theme, config: TelemetryConfig, iconMode: IconMode, now = new Date(), tone: TelemetryTone = "dim", ): string { if (!config.timestamp && !config.inputTokens && !config.outputTokens && !config.cacheRate && !config.tps && !config.ttft && !config.duration && !config.cost) return ""; const glyphs = resolveGlyphs(iconMode); const parts: string[] = []; if (config.timestamp) parts.push(theme.fg(tone, `${glyphs.working} ${formatClockTime(now)}`)); if (config.inputTokens) parts.push(theme.fg(tone, `${glyphs.input} ${fmtTokens(telemetry.inputTokens)}`)); if (config.outputTokens) parts.push(theme.fg(tone, `${glyphs.output} ${fmtTokens(telemetry.outputTokens)}`)); if (config.cacheRate) { const cache = telemetry.cacheHitRate === undefined ? "—" : `${telemetry.cacheHitRate.toFixed(1)}%`; parts.push(theme.fg(tone, `${glyphs.cacheHit} ${cache}`)); } if (config.cost) { parts.push(theme.fg(tone, `${glyphs.cost} ${telemetry.costUsd.toFixed(3)}`)); } if (config.duration) { parts.push(theme.fg(tone, `${glyphs.done} ${formatTurnDuration(telemetry.totalMs)}`)); } if (config.tps) { const value = telemetry.tps === null ? "—" : `${telemetry.tps.toFixed(1)} Tok/s`; parts.push(theme.fg(tone, `${glyphs.speed} ${value}`)); } if (config.ttft) { parts.push(theme.fg(tone, `${glyphs.latency} ${formatTurnDuration(telemetry.ttftMs)}`)); } return parts.join(" "); }