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.
This commit is contained in:
@@ -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
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -26,6 +26,7 @@ import {
|
||||
finishChatSpan,
|
||||
finishExecuteToolSpan,
|
||||
finishInvokeAgentSpan,
|
||||
fireOnRunEnd,
|
||||
recordSkippedTool,
|
||||
resolveTelemetry,
|
||||
runInActiveSpan,
|
||||
@@ -201,6 +202,9 @@ function buildAgentEndEvent(
|
||||
): Extract<AgentEvent, { type: "agent_end" }> {
|
||||
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 };
|
||||
}
|
||||
|
||||
|
||||
@@ -148,6 +148,25 @@ export class AgentRunCollector {
|
||||
readonly #invokedTools = new Set<string>();
|
||||
readonly #modelsUsed = new Set<string>();
|
||||
readonly #providersUsed = new Set<string>();
|
||||
#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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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<string, AttributeValue | undefined>;
|
||||
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<string, AttributeValue>) {
|
||||
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<AssistantMessage["usage"]> = {},
|
||||
) {
|
||||
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<string, unknown>;
|
||||
}[];
|
||||
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<AgentEvent, { type: "agent_end" }> => 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<AgentEvent, { type: "agent_end" }> => 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<void>((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);
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user