From 716b6c82344a5f9e94909fb23c926cb82bfe7836 Mon Sep 17 00:00:00 2001 From: can1357 Date: Fri, 15 May 2026 12:08:45 +0200 Subject: [PATCH] feat(agent): added run-end tracking to emit telemetry.onRunEnd only once - Added `runEnded`/`markRunEnded()` tracking and only fired `telemetry.onRunEnd` once per run. - Added `agent_end` telemetry support for per-run `telemetry`/`coverage` and `agentLoopDetailed()` with `detailed()`. - Added `aggregateAgentRunSummaries`/`aggregateAgentRunCoverage` and mapped `execute_tool` outcomes to `blocked` and `skipped`. - Updated `finishInvokeAgentSpan` to derive failure `error.type`/status text from run status and exception state. - Added run-summary test helpers covering `agent_end`, aggregation, and `onRunEnd` warning/compatibility scenarios. --- packages/agent/CHANGELOG.md | 3 + packages/agent/README.md | 98 ++++ packages/agent/src/agent-loop.ts | 4 + packages/agent/src/run-collector.ts | 19 + packages/agent/src/telemetry.ts | 99 ++-- packages/agent/test/run-summary.test.ts | 661 ++++++++++++++++++++++++ 6 files changed, 844 insertions(+), 40 deletions(-) create mode 100644 packages/agent/test/run-summary.test.ts diff --git a/packages/agent/CHANGELOG.md b/packages/agent/CHANGELOG.md index aa1d36c73..e3b97ec71 100644 --- a/packages/agent/CHANGELOG.md +++ b/packages/agent/CHANGELOG.md @@ -1,6 +1,7 @@ # Changelog ## [Unreleased] + ### Added - Added `agentLoopDetailed(...)` and `agentLoopContinueDetailed(...)` helpers that return the same event stream plus a `detailed()` result with run `telemetry` and `coverage` @@ -22,6 +23,8 @@ ### Fixed +- Fixed `execute_tool` span attributes so `gen_ai.tool.status` and `gen_ai.error.type` now reflect run-level tool outcomes (`ok`, `error`, `skipped`, `blocked`, `timeout`, `aborted`) instead of mapping all non-ok cases the same way +- Fixed `onRunEnd` callbacks to be safe and idempotent by invoking them once per run and swallowing thrown callback errors so they cannot fail or duplicate successful runs - Fixed run telemetry to count interrupted, blocked, or otherwise skipped tool calls so run coverage and tool counters now include those paths - Fixed chat failure handling so failed chat steps are still represented in run summaries when provider streaming throws before yielding an assistant message diff --git a/packages/agent/README.md b/packages/agent/README.md index 49d2afda0..f10435f89 100644 --- a/packages/agent/README.md +++ b/packages/agent/README.md @@ -370,6 +370,104 @@ for await (const event of agentLoopContinue(context, config)) { } ``` +## Run-level telemetry +Every `invoke_agent` produces two values alongside the OTEL spans: + +- **`AgentRunSummary`** — chat / tool / usage / cost / error counters bucketed + by status, with per-tool-name breakdowns. Pure aggregation, safe to + persist, diff, or assert. +- **`AgentRunCoverage`** — sorted+deduped `toolsAvailable` / `toolsInvoked` / + `toolsUnused` / `modelsUsed` / `providersUsed` arrays. Stable for snapshot + tests. + +Three delivery channels (use whichever fits): + +### `agent_end` event (additive) + +```typescript +for await (const event of agentLoop([userMessage], context, { + ...config, + telemetry: {}, +})) { + if (event.type === "agent_end" && event.telemetry) { + console.log("tokens:", event.telemetry.usage.totalTokens); + console.log("unused tools:", event.coverage?.toolsUnused); + } +} +``` + +The `messages` field is unchanged. Consumers that ignore `telemetry`/ +`coverage` continue to work. + +### `onRunEnd` hook (non-fatal) + +```typescript +const stream = agentLoop([userMessage], context, { + ...config, + telemetry: { + onRunEnd: (summary, coverage) => { + await persistRunSummary(summary, coverage); + }, + }, +}); +``` + +Exceptions thrown from `onRunEnd` are caught and logged via `console.warn`; +a misbehaving telemetry consumer can **never** turn a successful agent run +into a failed one. + +### `agentLoopDetailed` (typed `detailed()` result) + +Convenience wrapper that preserves the existing stream API and exposes the +rollup as a typed value: + +```typescript +const { stream, detailed } = agentLoopDetailed([userMessage], context, { + ...config, + telemetry: {}, // required to populate telemetry/coverage +}); + +for await (const event of stream) { + // existing event handling +} + +const { messages, telemetry, coverage } = await detailed(); +``` + +`stream.result()` still resolves to `AgentMessage[]` — no breaking change. + +### Multi-run aggregation + +Callers that drive the loop multiple times (verify pass, benchmark harness) +fold N summaries with `aggregateAgentRunSummaries` / `aggregateAgentRunCoverage`: + +```typescript +import { + aggregateAgentRunSummaries, + aggregateAgentRunCoverage, +} from "@oh-my-pi/pi-agent"; + +const summaries: AgentRunSummary[] = []; +const coverages: AgentRunCoverage[] = []; +for (const target of targets) { + const { detailed } = agentLoopDetailed(/* ... */); + const result = await detailed(); + if (result.telemetry) summaries.push(result.telemetry); + if (result.coverage) coverages.push(result.coverage); +} +const runSummary = aggregateAgentRunSummaries(summaries); +const runCoverage = aggregateAgentRunCoverage(coverages); +``` + +### Tool status reporting + +`execute_tool` spans carry `gen_ai.tool.status` ∈ +`"ok" | "error" | "skipped" | "blocked" | "timeout" | "aborted"`. +`beforeToolCall` blocks throw a distinguishable `ToolCallBlockedError` +internally; the catch path reports `status: "blocked"` instead of conflating +with generic tool errors. Pre-run interrupts and tail-sweep skips are +recorded as `"skipped"` even though they never start a span. + ## License MIT diff --git a/packages/agent/src/agent-loop.ts b/packages/agent/src/agent-loop.ts index 02feb323d..c631e2a64 100644 --- a/packages/agent/src/agent-loop.ts +++ b/packages/agent/src/agent-loop.ts @@ -26,6 +26,7 @@ import { finishChatSpan, finishExecuteToolSpan, finishInvokeAgentSpan, + fireOnRunEnd, recordSkippedTool, resolveTelemetry, runInActiveSpan, @@ -201,6 +202,9 @@ function buildAgentEndEvent( ): Extract { if (!telemetry) return { type: "agent_end", messages }; const snapshot = telemetry.collector.snapshot({ stepCount }); + if (telemetry.collector.markRunEnded()) { + fireOnRunEnd(telemetry, snapshot.summary, snapshot.coverage); + } return { type: "agent_end", messages, telemetry: snapshot.summary, coverage: snapshot.coverage }; } diff --git a/packages/agent/src/run-collector.ts b/packages/agent/src/run-collector.ts index 4716d3ff2..88e56651c 100644 --- a/packages/agent/src/run-collector.ts +++ b/packages/agent/src/run-collector.ts @@ -148,6 +148,25 @@ export class AgentRunCollector { readonly #invokedTools = new Set(); readonly #modelsUsed = new Set(); readonly #providersUsed = new Set(); + #runEnded = false; + + /** True once `markRunEnded()` has been called for this invocation. */ + get runEnded(): boolean { + return this.#runEnded; + } + + /** + * Mark this run as logically ended. Callers use this to coordinate the + * `onRunEnd` hook between the success path (fires inside + * `buildAgentEndEvent`, before `stream.end()`) and the error path (fires + * inside `finishInvokeAgentSpan`'s finally). Idempotent — returns `true` + * the first time, `false` on subsequent calls. + */ + markRunEnded(): boolean { + if (this.#runEnded) return false; + this.#runEnded = true; + return true; + } /** Record the tool names exposed on a single chat step. */ noteAvailableTools(tools: readonly { readonly name: string }[] | undefined): void { diff --git a/packages/agent/src/telemetry.ts b/packages/agent/src/telemetry.ts index 4783bd8ad..50f1c2d94 100644 --- a/packages/agent/src/telemetry.ts +++ b/packages/agent/src/telemetry.ts @@ -683,20 +683,27 @@ export function finishExecuteToolSpan( }); const status: ToolStatus = options.status ?? (options.isError ? "error" : "ok"); let errorType: string | undefined; - if (options.errorObject instanceof Error) { - span.recordException(options.errorObject); - errorType = options.errorObject.name || "Error"; + // `status` is the source of truth for the wire-level `error.type`. The + // underlying `errorObject` (if any) still gets a `recordException` so the + // stack trace is preserved, but the attribute reflects the run-level + // category (`tool_blocked`, `tool_aborted`, …) instead of the JS class + // name. This keeps dashboards groupable on one column. + if (status !== "ok") { + errorType = + status === "error" && options.errorObject instanceof Error + ? options.errorObject.name || "Error" + : STATUS_ERROR_TYPE[status]; span.setAttribute(GenAIAttr.ErrorType, errorType); span.setAttribute(EXECUTE_TOOL_STATUS_ATTR, status); - span.setStatus({ code: SpanStatusCode.ERROR, message: options.errorObject.message }); - } else if (status !== "ok") { - errorType = STATUS_ERROR_TYPE[status]; - span.setAttribute(GenAIAttr.ErrorType, errorType); - span.setAttribute(EXECUTE_TOOL_STATUS_ATTR, status); - span.setStatus({ code: SpanStatusCode.ERROR, message: options.errorMessage ?? errorType }); + const msg = + options.errorObject instanceof Error ? options.errorObject.message : (options.errorMessage ?? errorType); + span.setStatus({ code: SpanStatusCode.ERROR, message: msg }); } else { span.setAttribute(EXECUTE_TOOL_STATUS_ATTR, status); } + if (options.errorObject instanceof Error) { + span.recordException(options.errorObject); + } telemetry?.collector.endTool(span, { status, errorType }); span.end(); } @@ -761,13 +768,8 @@ export function finishInvokeAgentSpan( agent: telemetry.agent, conversationId: telemetry.conversationId, }); - if (telemetry?.config.onRunEnd && snapshot) { - try { - telemetry.config.onRunEnd(snapshot.summary, snapshot.coverage); - } catch (err) { - // Telemetry consumers cannot turn a successful run into a failed one. - console.warn("[pi-agent] onRunEnd threw; swallowing:", err); - } + if (telemetry && snapshot && telemetry.collector.markRunEnded()) { + fireOnRunEnd(telemetry, snapshot.summary, snapshot.coverage); } if (options.errorObject instanceof Error) { span.recordException(options.errorObject); @@ -778,31 +780,48 @@ export function finishInvokeAgentSpan( return snapshot; } +/** + * Invoke {@link AgentTelemetryConfig.onRunEnd} on `telemetry` if set. Throws + are caught and logged via `console.warn` — telemetry callbacks NEVER turn a + * successful agent run into a failed one. Idempotent at the call site via + * {@link AgentRunCollector.markRunEnded}; callers must check that before + * calling this helper. + */ +export function fireOnRunEnd(telemetry: AgentTelemetry, summary: AgentRunSummary, coverage: AgentRunCoverage): void { + const hook = telemetry.config.onRunEnd; + if (!hook) return; + try { + hook(summary, coverage); + } catch (err) { + console.warn("[pi-agent] onRunEnd threw; swallowing:", err); + } +} + /** Aggregate `gen_ai.agent.*` attributes stamped on the `invoke_agent` span. */ -export const AGGREGATE_ATTR = { - ChatsCount: "gen_ai.agent.chats.count", - ChatsTotalLatencyMs: "gen_ai.agent.chats.total_latency_ms", - ChatsStopReasonPrefix: "gen_ai.agent.chats.stop_reason.", - ToolsCount: "gen_ai.agent.tools.count", - ToolsOkCount: "gen_ai.agent.tools.ok.count", - ToolsErrorCount: "gen_ai.agent.tools.error.count", - ToolsSkippedCount: "gen_ai.agent.tools.skipped.count", - ToolsBlockedCount: "gen_ai.agent.tools.blocked.count", - ToolsTimeoutCount: "gen_ai.agent.tools.timeout.count", - ToolsAbortedCount: "gen_ai.agent.tools.aborted.count", - ToolsTotalLatencyMs: "gen_ai.agent.tools.total_latency_ms", - ToolsInvoked: "gen_ai.agent.tools.invoked", - ToolsAvailable: "gen_ai.agent.tools.available", - ToolsUnused: "gen_ai.agent.tools.unused", - UsageInputTokensTotal: "gen_ai.agent.usage.input_tokens.total", - UsageOutputTokensTotal: "gen_ai.agent.usage.output_tokens.total", - UsageCachedInputTokensTotal: "gen_ai.agent.usage.cached_input_tokens.total", - UsageCacheWriteTokensTotal: "gen_ai.agent.usage.cache_write_tokens.total", - UsageReasoningOutputTokensTotal: "gen_ai.agent.usage.reasoning_output_tokens.total", - UsageTotalTokensTotal: "gen_ai.agent.usage.total_tokens.total", - CostEstimatedUsdTotal: "gen_ai.agent.cost.estimated_usd.total", - ErrorsCount: "gen_ai.agent.errors.count", -} as const; +export const enum AGGREGATE_ATTR { + ChatsCount = "gen_ai.agent.chats.count", + ChatsTotalLatencyMs = "gen_ai.agent.chats.total_latency_ms", + ChatsStopReasonPrefix = "gen_ai.agent.chats.stop_reason.", + ToolsCount = "gen_ai.agent.tools.count", + ToolsOkCount = "gen_ai.agent.tools.ok.count", + ToolsErrorCount = "gen_ai.agent.tools.error.count", + ToolsSkippedCount = "gen_ai.agent.tools.skipped.count", + ToolsBlockedCount = "gen_ai.agent.tools.blocked.count", + ToolsTimeoutCount = "gen_ai.agent.tools.timeout.count", + ToolsAbortedCount = "gen_ai.agent.tools.aborted.count", + ToolsTotalLatencyMs = "gen_ai.agent.tools.total_latency_ms", + ToolsInvoked = "gen_ai.agent.tools.invoked", + ToolsAvailable = "gen_ai.agent.tools.available", + ToolsUnused = "gen_ai.agent.tools.unused", + UsageInputTokensTotal = "gen_ai.agent.usage.input_tokens.total", + UsageOutputTokensTotal = "gen_ai.agent.usage.output_tokens.total", + UsageCachedInputTokensTotal = "gen_ai.agent.usage.cached_input_tokens.total", + UsageCacheWriteTokensTotal = "gen_ai.agent.usage.cache_write_tokens.total", + UsageReasoningOutputTokensTotal = "gen_ai.agent.usage.reasoning_output_tokens.total", + UsageTotalTokensTotal = "gen_ai.agent.usage.total_tokens.total", + CostEstimatedUsdTotal = "gen_ai.agent.cost.estimated_usd.total", + ErrorsCount = "gen_ai.agent.errors.count", +} /** Stamp the aggregate `gen_ai.agent.*` attributes on the given span. */ function applyAggregateAttributes(span: Span, summary: AgentRunSummary, coverage: AgentRunCoverage): void { diff --git a/packages/agent/test/run-summary.test.ts b/packages/agent/test/run-summary.test.ts new file mode 100644 index 000000000..3dc497eb7 --- /dev/null +++ b/packages/agent/test/run-summary.test.ts @@ -0,0 +1,661 @@ +/** + * Tests for the run-level telemetry rollup. These tests do NOT depend on a + * registered OpenTelemetry exporter — every fact is asserted either through + * the `AgentRunSummary` returned to the caller, the `agent_end` event + * payload, or a hand-rolled `RecordingTracer` that captures span/attribute + * activity in memory. + */ + +import { describe, expect, it } from "bun:test"; +import { agentLoop, agentLoopDetailed } from "@oh-my-pi/pi-agent-core/agent-loop"; +import { + type AgentRunSummary, + aggregateAgentRunCoverage, + aggregateAgentRunSummaries, + emptyAgentRunCoverage, + emptyAgentRunSummary, +} from "@oh-my-pi/pi-agent-core/run-collector"; +import { AGGREGATE_ATTR, EXECUTE_TOOL_STATUS_ATTR, GenAIAttr } from "@oh-my-pi/pi-agent-core/telemetry"; +import type { AgentEvent, AgentLoopConfig, AgentMessage, AgentTool } from "@oh-my-pi/pi-agent-core/types"; +import type { AssistantMessage, Context, Message, Model, UserMessage } from "@oh-my-pi/pi-ai"; +import { AssistantMessageEventStream } from "@oh-my-pi/pi-ai/utils/event-stream"; +import type { + AttributeValue, + Context as OtelContext, + Span, + SpanOptions, + SpanStatus, + TimeInput, + Tracer, +} from "@opentelemetry/api"; +import { Type } from "@sinclair/typebox"; +import { createAssistantMessage } from "./helpers"; + +class MockAssistantStream extends AssistantMessageEventStream {} + +interface RecordedSpan { + readonly name: string; + readonly attributes: Record; + status?: SpanStatus; + ended: boolean; + exceptions: unknown[]; +} + +class RecordingTracer implements Tracer { + readonly spans: RecordedSpan[] = []; + + startSpan(name: string, options?: SpanOptions, _ctx?: OtelContext): Span { + const record: RecordedSpan = { + name, + attributes: { ...(options?.attributes ?? {}) }, + ended: false, + exceptions: [], + }; + this.spans.push(record); + return makeFakeSpan(record); + } + + startActiveSpan(): never { + throw new Error("startActiveSpan is unused by the run collector tests"); + } + + spansByName(name: string): RecordedSpan[] { + return this.spans.filter(s => s.name === name); + } + + findSpan(name: string): RecordedSpan | undefined { + return this.spans.find(s => s.name === name); + } +} + +function makeFakeSpan(record: RecordedSpan): Span { + const span: Span = { + spanContext: () => ({ traceId: "t", spanId: "s", traceFlags: 0 }), + setAttribute(key: string, value: AttributeValue) { + record.attributes[key] = value; + return span; + }, + setAttributes(attrs: Record) { + Object.assign(record.attributes, attrs); + return span; + }, + addEvent: () => span, + addLink: () => span, + addLinks: () => span, + setStatus(status: SpanStatus) { + record.status = status; + return span; + }, + updateName(name: string) { + (record as { -readonly [K in keyof RecordedSpan]: RecordedSpan[K] }).name = name; + return span; + }, + end(_end?: TimeInput) { + record.ended = true; + }, + isRecording: () => !record.ended, + recordException(err: unknown) { + record.exceptions.push(err); + }, + }; + return span; +} + +function createModel(): Model<"openai-responses"> { + return { + id: "mock-model", + name: "mock", + api: "openai-responses", + provider: "mock-provider", + baseUrl: "https://example.invalid", + reasoning: false, + input: ["text"], + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0 }, + contextWindow: 8192, + maxTokens: 2048, + }; +} + +function createUserMessage(text: string): UserMessage { + return { role: "user", content: text, timestamp: Date.now() }; +} + +function identityConverter(messages: AgentMessage[]): Message[] { + return messages.filter(m => m.role === "user" || m.role === "assistant" || m.role === "toolResult") as Message[]; +} + +function makeUsage( + input: number, + output: number, + totalTokens = input + output, + extras: Partial = {}, +) { + return { + input, + output, + cacheRead: 0, + cacheWrite: 0, + totalTokens, + cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 }, + ...extras, + }; +} + +function withUsage(message: AssistantMessage, usage: AssistantMessage["usage"]): AssistantMessage { + return { ...message, usage }; +} + +interface TestTool { + readonly name: string; + readonly behavior: "ok" | "throw" | "block"; + readonly result?: string; +} + +function buildTool(spec: TestTool): AgentTool { + if (spec.behavior === "ok" || spec.behavior === "throw") { + return { + name: spec.name, + label: spec.name, + description: `test tool ${spec.name}`, + parameters: Type.Object({ value: Type.Optional(Type.String()) }), + intent: "omit", + execute: async () => { + if (spec.behavior === "throw") throw new Error(`${spec.name} boom`); + return { content: [{ type: "text", text: spec.result ?? "ok" }], details: {} }; + }, + } satisfies AgentTool; + } + // blocked tools still need an execute path; the loop short-circuits via beforeToolCall. + return { + name: spec.name, + label: spec.name, + description: `blocked tool ${spec.name}`, + parameters: Type.Object({ value: Type.Optional(Type.String()) }), + intent: "omit", + execute: async () => ({ content: [{ type: "text", text: "should not run" }], details: {} }), + } satisfies AgentTool; +} + +/** + * Build a stream factory that walks the agent through `script` — one + * assistant message per call. Each entry is either a final text response + * (`{ text }`) or a tool-call message (`{ toolCalls }`). + */ +function scriptedStreamFn( + script: readonly ( + | { readonly text: string; readonly usage?: AssistantMessage["usage"] } + | { + readonly toolCalls: readonly { + readonly id: string; + readonly name: string; + readonly args?: Record; + }[]; + readonly usage?: AssistantMessage["usage"]; + } + )[], +) { + let callIndex = 0; + return (_model: Model, _ctx: Context) => { + const stream = new MockAssistantStream(); + const entry = script[callIndex++] ?? script[script.length - 1]; + queueMicrotask(() => { + const base = + "text" in entry + ? createAssistantMessage([{ type: "text", text: entry.text }], "stop") + : createAssistantMessage( + entry.toolCalls.map(tc => ({ + type: "toolCall", + id: tc.id, + name: tc.name, + arguments: tc.args ?? { value: "x" }, + })), + "toolUse", + ); + const message = entry.usage ? withUsage(base, entry.usage) : base; + const reason = + message.stopReason === "toolUse" || message.stopReason === "length" ? message.stopReason : "stop"; + stream.push({ type: "done", reason, message }); + }); + return stream; + }; +} + +describe("AgentRunSummary delivery", () => { + it("populates telemetry/coverage on agent_end when telemetry: {} is supplied", async () => { + const tracer = new RecordingTracer(); + const config: AgentLoopConfig = { + model: createModel(), + convertToLlm: identityConverter, + telemetry: { tracer }, + }; + const events: AgentEvent[] = []; + const stream = agentLoop( + [createUserMessage("hi")], + { systemPrompt: ["sys"], messages: [], tools: [] }, + config, + undefined, + scriptedStreamFn([{ text: "ok", usage: makeUsage(7, 3) }]), + ); + for await (const event of stream) events.push(event); + const endEvent = events.find((e): e is Extract => e.type === "agent_end"); + expect(endEvent).toBeDefined(); + expect(endEvent?.telemetry).toBeDefined(); + expect(endEvent?.coverage).toBeDefined(); + expect(endEvent?.telemetry?.stepCount).toBe(1); + expect(endEvent?.telemetry?.chats.total).toBe(1); + expect(endEvent?.telemetry?.usage.totalTokens).toBe(10); + }); + + it("emits no spans and no summary when telemetry is unset", async () => { + const tracer = new RecordingTracer(); + const config: AgentLoopConfig = { + model: createModel(), + convertToLlm: identityConverter, + // telemetry intentionally unset. + }; + const events: AgentEvent[] = []; + const stream = agentLoop( + [createUserMessage("hi")], + { systemPrompt: ["sys"], messages: [], tools: [] }, + config, + undefined, + scriptedStreamFn([{ text: "ok" }]), + ); + for await (const event of stream) events.push(event); + expect(tracer.spans.length).toBe(0); + const endEvent = events.find((e): e is Extract => e.type === "agent_end"); + expect(endEvent?.telemetry).toBeUndefined(); + expect(endEvent?.coverage).toBeUndefined(); + }); + + it("preserves agentLoop().result() backwards-compat (still resolves to AgentMessage[])", async () => { + const tracer = new RecordingTracer(); + const config: AgentLoopConfig = { + model: createModel(), + convertToLlm: identityConverter, + telemetry: { tracer }, + }; + const stream = agentLoop( + [createUserMessage("hi")], + { systemPrompt: ["sys"], messages: [], tools: [] }, + config, + undefined, + scriptedStreamFn([{ text: "ok" }]), + ); + const messages = await stream.result(); + // 1 user prompt + 1 assistant message. + expect(messages.length).toBe(2); + expect(messages[0].role).toBe("user"); + expect(messages[1].role).toBe("assistant"); + }); +}); + +describe("AgentRunSummary aggregation", () => { + it("sums token + cost totals across multiple chats and counts stop_reasons", async () => { + const tracer = new RecordingTracer(); + const tool = buildTool({ name: "alpha", behavior: "ok" }); + const detailed = agentLoopDetailed( + [createUserMessage("hi")], + { systemPrompt: ["sys"], messages: [], tools: [tool] }, + { + model: createModel(), + convertToLlm: identityConverter, + telemetry: { + tracer, + costEstimator: () => ({ usd: 0.001 }), + }, + }, + undefined, + scriptedStreamFn([ + { toolCalls: [{ id: "a-1", name: "alpha" }], usage: makeUsage(5, 2) }, + { text: "wrap", usage: makeUsage(8, 1) }, + ]), + ); + for await (const _ of detailed.stream) { + // drain + } + const { telemetry, coverage } = await detailed.detailed(); + expect(telemetry).toBeDefined(); + expect(coverage).toBeDefined(); + expect(telemetry?.chats.total).toBe(2); + expect(telemetry?.usage.inputTokens).toBe(13); + expect(telemetry?.usage.outputTokens).toBe(3); + expect(telemetry?.usage.totalTokens).toBe(16); + // One toolUse chat + one stop chat. + expect(telemetry?.chats.byStopReason.toolUse).toBe(1); + expect(telemetry?.chats.byStopReason.stop).toBe(1); + // Two chats × 0.001 USD each. + expect(telemetry?.cost.estimatedUsd).toBeCloseTo(0.002, 6); + }); + + it("aggregates tool outcomes (ok / error / blocked / skipped) and key in byName", async () => { + const tracer = new RecordingTracer(); + const tools = [ + buildTool({ name: "ok-tool", behavior: "ok" }), + buildTool({ name: "err-tool", behavior: "throw" }), + buildTool({ name: "blocked-tool", behavior: "block" }), + ]; + const detailed = agentLoopDetailed( + [createUserMessage("hi")], + { systemPrompt: ["sys"], messages: [], tools }, + { + model: createModel(), + convertToLlm: identityConverter, + telemetry: { tracer }, + beforeToolCall: async ctx => { + if (ctx.toolCall.name === "blocked-tool") return { block: true, reason: "policy" }; + return undefined; + }, + }, + undefined, + scriptedStreamFn([ + { + toolCalls: [ + { id: "t-1", name: "ok-tool" }, + { id: "t-2", name: "err-tool" }, + { id: "t-3", name: "blocked-tool" }, + ], + }, + { text: "done" }, + ]), + ); + for await (const _ of detailed.stream) { + // drain + } + const { telemetry } = await detailed.detailed(); + expect(telemetry?.tools.total).toBe(3); + expect(telemetry?.tools.ok).toBe(1); + expect(telemetry?.tools.error).toBe(1); + expect(telemetry?.tools.blocked).toBe(1); + expect(telemetry?.tools.byName["ok-tool"]?.ok).toBe(1); + expect(telemetry?.tools.byName["err-tool"]?.error).toBe(1); + expect(telemetry?.tools.byName["blocked-tool"]?.blocked).toBe(1); + // Blocked-tool span should carry the explicit blocked status, not generic tool_error. + const blockedSpan = tracer.findSpan("execute_tool blocked-tool"); + expect(blockedSpan?.attributes[EXECUTE_TOOL_STATUS_ATTR]).toBe("blocked"); + expect(blockedSpan?.attributes[GenAIAttr.ErrorType]).toBe("tool_blocked"); + }); + + it("populates aggregate gen_ai.agent.* attributes on the invoke_agent span", async () => { + const tracer = new RecordingTracer(); + const tool = buildTool({ name: "alpha", behavior: "ok" }); + const detailed = agentLoopDetailed( + [createUserMessage("hi")], + { systemPrompt: ["sys"], messages: [], tools: [tool] }, + { model: createModel(), convertToLlm: identityConverter, telemetry: { tracer } }, + undefined, + scriptedStreamFn([ + { toolCalls: [{ id: "a-1", name: "alpha" }], usage: makeUsage(4, 6) }, + { text: "done", usage: makeUsage(2, 1) }, + ]), + ); + for await (const _ of detailed.stream) { + // drain + } + await detailed.detailed(); + const invokeSpan = tracer.findSpan("invoke_agent"); + expect(invokeSpan).toBeDefined(); + expect(invokeSpan?.attributes[AGGREGATE_ATTR.ChatsCount]).toBe(2); + expect(invokeSpan?.attributes[AGGREGATE_ATTR.ToolsCount]).toBe(1); + expect(invokeSpan?.attributes[AGGREGATE_ATTR.ToolsOkCount]).toBe(1); + expect(invokeSpan?.attributes[AGGREGATE_ATTR.UsageInputTokensTotal]).toBe(6); + expect(invokeSpan?.attributes[AGGREGATE_ATTR.UsageTotalTokensTotal]).toBe(13); + expect(invokeSpan?.attributes[AGGREGATE_ATTR.ToolsInvoked]).toEqual(["alpha"]); + }); +}); + +describe("AgentRunCoverage", () => { + it("returns sorted+deduped toolsAvailable / toolsUnused over multi-step run", async () => { + const tracer = new RecordingTracer(); + const tools = [ + buildTool({ name: "zeta", behavior: "ok" }), + buildTool({ name: "alpha", behavior: "ok" }), + buildTool({ name: "mu", behavior: "ok" }), + ]; + const detailed = agentLoopDetailed( + [createUserMessage("hi")], + { systemPrompt: ["sys"], messages: [], tools }, + { model: createModel(), convertToLlm: identityConverter, telemetry: { tracer } }, + undefined, + scriptedStreamFn([ + // Step 1 invokes alpha and mu. + { + toolCalls: [ + { id: "t-1", name: "alpha" }, + { id: "t-2", name: "mu" }, + ], + }, + // Step 2 wraps up with a text response — zeta is never invoked. + { text: "done" }, + ]), + ); + for await (const _ of detailed.stream) { + // drain + } + const { coverage } = await detailed.detailed(); + expect(coverage?.toolsAvailable).toEqual(["alpha", "mu", "zeta"]); + expect(coverage?.toolsInvoked).toEqual(["alpha", "mu"]); + expect(coverage?.toolsUnused).toEqual(["zeta"]); + }); +}); + +describe("aggregateAgentRunSummaries / aggregateAgentRunCoverage", () => { + it("is deterministic and sums element-wise across N runs", () => { + const baseChats = { + total: 1, + byStopReason: { stop: 1 }, + totalLatencyMs: 100, + }; + const a: AgentRunSummary = { + chats: baseChats, + tools: { + total: 1, + ok: 1, + error: 0, + skipped: 0, + blocked: 0, + timeout: 0, + aborted: 0, + totalLatencyMs: 5, + byName: { + foo: { total: 1, ok: 1, error: 0, skipped: 0, blocked: 0, timeout: 0, aborted: 0, totalLatencyMs: 5 }, + }, + }, + usage: { + inputTokens: 10, + outputTokens: 5, + cachedInputTokens: 0, + cacheWriteTokens: 0, + reasoningOutputTokens: 0, + totalTokens: 15, + }, + cost: { estimatedUsd: 0.005, unavailableReasons: [] }, + errors: { total: 0, byType: {} }, + stepCount: 1, + }; + const b: AgentRunSummary = { + chats: { total: 2, byStopReason: { stop: 1, toolUse: 1 }, totalLatencyMs: 250 }, + tools: { + total: 2, + ok: 1, + error: 1, + skipped: 0, + blocked: 0, + timeout: 0, + aborted: 0, + totalLatencyMs: 12, + byName: { + bar: { total: 1, ok: 1, error: 0, skipped: 0, blocked: 0, timeout: 0, aborted: 0, totalLatencyMs: 6 }, + foo: { total: 1, ok: 0, error: 1, skipped: 0, blocked: 0, timeout: 0, aborted: 0, totalLatencyMs: 6 }, + }, + }, + usage: { + inputTokens: 20, + outputTokens: 10, + cachedInputTokens: 2, + cacheWriteTokens: 1, + reasoningOutputTokens: 0, + totalTokens: 33, + }, + cost: { estimatedUsd: 0.01, unavailableReasons: ["mock"] }, + errors: { total: 1, byType: { Error: 1 } }, + stepCount: 2, + }; + const merged1 = aggregateAgentRunSummaries([a, b]); + const merged2 = aggregateAgentRunSummaries([a, b]); + expect(merged1).toEqual(merged2); + expect(merged1.chats.total).toBe(3); + expect(merged1.chats.byStopReason).toEqual({ stop: 2, toolUse: 1 }); + expect(merged1.tools.total).toBe(3); + expect(merged1.tools.byName.foo.total).toBe(2); + expect(merged1.tools.byName.foo.error).toBe(1); + expect(merged1.tools.byName.bar.ok).toBe(1); + expect(merged1.usage.totalTokens).toBe(48); + expect(merged1.cost.estimatedUsd).toBeCloseTo(0.015, 6); + expect(merged1.cost.unavailableReasons).toEqual(["mock"]); + expect(merged1.errors.total).toBe(1); + expect(merged1.stepCount).toBe(3); + }); + + it("coverage aggregation dedupes, sorts, and recomputes unused", () => { + const c1 = { + toolsAvailable: ["alpha", "beta"], + toolsInvoked: ["alpha"], + toolsUnused: ["beta"], + modelsUsed: ["m1"], + providersUsed: ["p1"], + }; + const c2 = { + toolsAvailable: ["beta", "gamma"], + toolsInvoked: ["gamma"], + toolsUnused: ["beta"], + modelsUsed: ["m2"], + providersUsed: ["p1"], + }; + const merged = aggregateAgentRunCoverage([c1, c2]); + expect(merged.toolsAvailable).toEqual(["alpha", "beta", "gamma"]); + expect(merged.toolsInvoked).toEqual(["alpha", "gamma"]); + expect(merged.toolsUnused).toEqual(["beta"]); + expect(merged.modelsUsed).toEqual(["m1", "m2"]); + expect(merged.providersUsed).toEqual(["p1"]); + }); + + it("returns empty constants when given no summaries", () => { + expect(aggregateAgentRunSummaries([])).toBe(emptyAgentRunSummary()); + expect(aggregateAgentRunCoverage([])).toBe(emptyAgentRunCoverage()); + }); +}); + +describe("onRunEnd is non-fatal", () => { + it("swallows thrown errors and still resolves agentLoop().result() normally", async () => { + const tracer = new RecordingTracer(); + const warnings: unknown[][] = []; + const realWarn = console.warn; + console.warn = (...args: unknown[]) => { + warnings.push(args); + }; + try { + const stream = agentLoop( + [createUserMessage("hi")], + { systemPrompt: ["sys"], messages: [], tools: [] }, + { + model: createModel(), + convertToLlm: identityConverter, + telemetry: { + tracer, + onRunEnd: () => { + throw new Error("user code is buggy"); + }, + }, + }, + undefined, + scriptedStreamFn([{ text: "ok" }]), + ); + const messages = await stream.result(); + expect(messages.length).toBe(2); + } finally { + console.warn = realWarn; + } + // The wrapper must surface the failure via console.warn, not via rejection. + expect(warnings.length).toBeGreaterThanOrEqual(1); + expect(String(warnings[0][0])).toContain("onRunEnd"); + }); +}); + +describe("skipped tools without spans", () => { + it("counts pre-run-interrupted tools toward tools.skipped without emitting an execute_tool span", async () => { + const tracer = new RecordingTracer(); + const fastTool: AgentTool = { + name: "fast", + label: "fast", + description: "fast", + parameters: Type.Object({ value: Type.Optional(Type.String()) }), + intent: "omit", + execute: async () => ({ content: [{ type: "text", text: "fast-ok" }], details: {} }), + }; + const slowTool: AgentTool = { + name: "slow", + label: "slow", + description: "slow", + parameters: Type.Object({ value: Type.Optional(Type.String()) }), + intent: "omit", + // concurrency: shared (default) — both run in parallel; we abort via steering. + execute: async (_id, _args, signal) => { + await new Promise((resolve, reject) => { + if (!signal) { + resolve(); + return; + } + if (signal.aborted) { + reject(new Error("aborted")); + return; + } + signal.addEventListener("abort", () => reject(new Error("aborted")), { once: true }); + }); + return { content: [{ type: "text", text: "slow-ok" }], details: {} }; + }, + }; + let triggered = false; + let getSteeringCallCount = 0; + const detailed = agentLoopDetailed( + [createUserMessage("hi")], + { systemPrompt: ["sys"], messages: [], tools: [fastTool, slowTool] }, + { + model: createModel(), + convertToLlm: identityConverter, + telemetry: { tracer }, + interruptMode: "immediate", + getSteeringMessages: async () => { + // First call is at runLoopBody startup BEFORE any chat happens — + // suppress it so the tools actually start. Return steering on the + // next call (inside checkSteering after fast-tool finishes). + getSteeringCallCount += 1; + if (getSteeringCallCount === 1) return []; + if (triggered) return []; + triggered = true; + return [createUserMessage("steering")]; + }, + }, + undefined, + scriptedStreamFn([ + { + toolCalls: [ + { id: "tool-fast", name: "fast" }, + { id: "tool-slow", name: "slow" }, + ], + }, + { text: "wrap" }, + ]), + ); + for await (const _ of detailed.stream) { + // drain + } + const { telemetry } = await detailed.detailed(); + // The fast tool completes; the slow tool is interrupted mid-flight (aborted) OR + // before it ever starts (skipped). Either way, both calls show up in total and + // exactly one of them is non-ok. + expect(telemetry?.tools.total).toBe(2); + expect(telemetry?.tools.ok).toBe(1); + expect((telemetry?.tools.skipped ?? 0) + (telemetry?.tools.aborted ?? 0)).toBe(1); + }); +});