From f06f2aaa5930257511a8a02d589733e3e83a91c7 Mon Sep 17 00:00:00 2001 From: David Zhao Date: Mon, 14 Sep 2026 23:16:07 -0700 Subject: [PATCH 1/4] feat(telemetry): one agent_turn span per speech handle Port of livekit/agents#7143. A reply that calls a tool runs two generations (LLM steps) in two tasks; they were two agent_turn spans linked only by lk.parent_generation_id, so one response rendered as two turns. A speech handle is now exactly one agent_turn. - The speech handle owns the span: the first reply task (pipeline, realtime, or say) opens agent_turn under agent_session; the follow-up generation after a tool call continues the open span. It ends with the speech in SpeechHandle._markDone, recording the speech's error redaction-aware and the gen_ai.invoke_agent.duration histogram for the whole turn (new in JS: otel_metrics.recordInvokeAgentDuration and trace_types.METRIC_GEN_AI_INVOKE_AGENT_DURATION; the Python metric already existed). - Each generation is a `generation` event with lk.generation_id and lk.parent_generation_id (SpeechHandle._generationId / _parentGenerationId, `_` like Python); the span carries the latest generation id and the new lk.generation_count. - A preemptive generation discarded for a successor answering the same user turn hands its open agent_turn over (preemptive_generation_discarded event, lk.speech_id follows the speech that answered), both on a newer attempt and on the real reply after onUserTurnCompleted invalidated it. The queue-wait and interruption helpers tolerate an ended span. - `say` gets an agent_turn too, as in Python's _tts_task; JS had none. Tests: agent_turn_span.test.ts (tool call is one turn with two generation events and every step nested; plain reply is one generation; discarded preemptive hand-off; LLM failure fails the turn; duration metric when sampled out; sampled-out hand-off). The preemptive-guard stand-in handle gained _takeAgentTurn; the PII key test skips METRIC_* names, which are not attribute keys. The handoff runs inside generateReply, before the reply task is created: Task starts its body synchronously and the task opens the turn in its first statements, so a handoff performed after generateReply() returned came too late, leaving the successor's own agent_turn unended and its llm_node, tts_node and function_tool dangling from a span that was never exported (seen in a cloud export as bare nodes and turns whose generation count did not match their events). A regression test drives the real reply task. Co-Authored-By: Claude Fable 5.1 --- .changeset/agent-turn-per-speech.md | 5 + agents/src/telemetry/otel_metrics.ts | 63 +++- agents/src/telemetry/trace_types.test.ts | 4 +- agents/src/telemetry/trace_types.ts | 13 + agents/src/voice/agent_activity.test.ts | 7 +- agents/src/voice/agent_activity.ts | 189 ++++++++-- agents/src/voice/agent_turn_span.test.ts | 437 +++++++++++++++++++++++ agents/src/voice/speech_handle.ts | 145 +++++++- 8 files changed, 812 insertions(+), 51 deletions(-) create mode 100644 .changeset/agent-turn-per-speech.md create mode 100644 agents/src/voice/agent_turn_span.test.ts diff --git a/.changeset/agent-turn-per-speech.md b/.changeset/agent-turn-per-speech.md new file mode 100644 index 0000000000..629ea0f62d --- /dev/null +++ b/.changeset/agent-turn-per-speech.md @@ -0,0 +1,5 @@ +--- +'@livekit/agents': patch +--- + +One `agent_turn` span per speech handle: the follow-up generation after a tool call continues the open span instead of opening a second turn, each generation is a `generation` event with `lk.generation_count` on the span, the turn ends with the speech, and a discarded preemptive generation hands its turn to the reply that answered. diff --git a/agents/src/telemetry/otel_metrics.ts b/agents/src/telemetry/otel_metrics.ts index ecf06d9eaa..5ce82cd322 100644 --- a/agents/src/telemetry/otel_metrics.ts +++ b/agents/src/telemetry/otel_metrics.ts @@ -3,28 +3,41 @@ // SPDX-License-Identifier: Apache-2.0 import { type Attributes, type Histogram, type MeterProvider, metrics } from '@opentelemetry/api'; import { getJobContext } from '../job.js'; +import * as traceTypes from './trace_types.js'; +// Instruments are looked up per call: the global meter provider may be installed after this +// module loads (the cloud pipeline is set up when the job registers), and a histogram created +// on the no-op provider would stay a no-op. let meterProvider: MeterProvider | undefined; -let blockedDuration: Histogram | undefined; +const histograms = new Map(); -function eventLoopBlockedHistogram(): Histogram { +function histogram(name: string, description: string): Histogram { const currentProvider = metrics.getMeterProvider(); - if (currentProvider !== meterProvider || !blockedDuration) { + if (currentProvider !== meterProvider) { meterProvider = currentProvider; - blockedDuration = metrics - .getMeter('livekit-agents') - .createHistogram('lk.agents.event_loop.blocked_duration', { - unit: 's', - description: 'Duration of synchronous blocks detected on an agent event loop', - }); + histograms.clear(); } - return blockedDuration; + let instrument = histograms.get(name); + if (!instrument) { + instrument = metrics.getMeter('livekit-agents').createHistogram(name, { + unit: 's', + description, + }); + histograms.set(name, instrument); + } + return instrument; } -/** Record an event-loop stall in seconds, with the severity and what caused it. */ -export function recordEventLoopBlocked(duration: number, severity: string, cause: string): void { +/** + * Per-measurement job attribution. + * + * The meter provider has process lifetime (the OTel metrics global is set-once), so per-job + * fields cannot live on its resource. Each measurement carries the same per-job attributes that + * are stamped on spans and logs instead. Returns a fresh object; callers may add to it. + */ +function jobAttrs(): Attributes { const ctx = getJobContext(false); - const attributes: Attributes = { severity, cause }; + const attributes: Attributes = {}; if (ctx) { Object.assign(attributes, ctx._otelMetadata()); const roomId = ctx.job.room?.sid; @@ -32,5 +45,27 @@ export function recordEventLoopBlocked(duration: number, severity: string, cause if (ctx.job.id) attributes.job_id = ctx.job.id; if (ctx.job.agentName) attributes['lk.agent_name'] = ctx.job.agentName; } - eventLoopBlockedHistogram().record(duration, attributes); + return attributes; +} + +/** Record an event-loop stall in seconds, with the severity and what caused it. */ +export function recordEventLoopBlocked(duration: number, severity: string, cause: string): void { + const attributes = jobAttrs(); + attributes.severity = severity; + attributes.cause = cause; + histogram( + 'lk.agents.event_loop.blocked_duration', + 'Duration of synchronous blocks detected on an agent event loop', + ).record(duration, attributes); +} + +/** `gen_ai.invoke_agent.duration` for one agent turn, in seconds. */ +export function recordInvokeAgentDuration(duration: number, agentName: string): void { + const attributes = jobAttrs(); + attributes[traceTypes.ATTR_GEN_AI_OPERATION_NAME] = traceTypes.GenAIOperationName.INVOKE_AGENT; + attributes[traceTypes.ATTR_GEN_AI_AGENT_NAME] = agentName; + histogram(traceTypes.METRIC_GEN_AI_INVOKE_AGENT_DURATION, 'Agent invocation duration').record( + duration, + attributes, + ); } diff --git a/agents/src/telemetry/trace_types.test.ts b/agents/src/telemetry/trace_types.test.ts index 7e7b8b73ff..b10c5ae3dc 100644 --- a/agents/src/telemetry/trace_types.test.ts +++ b/agents/src/telemetry/trace_types.test.ts @@ -132,6 +132,7 @@ const SAFE_KEYS = new Set([ 'lk.deployment_id', 'lk.session_options', 'lk.generation_id', + 'lk.generation_count', 'lk.parent_generation_id', 'lk.interrupted', // LLM node metadata @@ -298,7 +299,8 @@ const SAFE_KEYS = new Set([ function declaredKeys(): Record { return Object.fromEntries( Object.entries(traceTypes).filter((entry): entry is [string, string] => { - return typeof entry[1] === 'string'; + // metric names are not attribute keys: they carry no values to classify + return typeof entry[1] === 'string' && !entry[0].startsWith('METRIC_'); }), ); } diff --git a/agents/src/telemetry/trace_types.ts b/agents/src/telemetry/trace_types.ts index 594056cfd1..ca9a4957a2 100644 --- a/agents/src/telemetry/trace_types.ts +++ b/agents/src/telemetry/trace_types.ts @@ -120,8 +120,17 @@ export const ATTR_SHUTDOWN_USER_INITIATED = 'lk.shutdown.user_initiated'; export const ATTR_CALLBACK_NAME = 'lk.callback.name'; // assistant turn +/** + * On `agent_turn`: the latest generation (LLM step) of the speech; each step is also a + * `generation` event carrying its own id. + */ export const ATTR_AGENT_TURN_ID = 'lk.generation_id'; export const ATTR_AGENT_PARENT_TURN_ID = 'lk.parent_generation_id'; +/** + * On `agent_turn`: how many generations (LLM steps) the speech took; more than one means tool + * calls were executed before the final reply. + */ +export const ATTR_GENERATION_COUNT = 'lk.generation_count'; export const ATTR_USER_INPUT = 'lk.pii.user_input'; export const ATTR_INSTRUCTIONS = 'lk.pii.instructions'; export const ATTR_SPEECH_INTERRUPTED = 'lk.interrupted'; @@ -486,3 +495,7 @@ export const ATTR_EXCEPTION_MESSAGE = 'exception.message'; // Platform-specific attributes export const ATTR_LANGFUSE_COMPLETION_START_TIME = 'langfuse.observation.completion_start_time'; + +// metric names (OpenTelemetry GenAI semantic conventions) +/** Histogram, seconds: one agent turn (`invoke_agent`), however many LLM steps it took. */ +export const METRIC_GEN_AI_INVOKE_AGENT_DURATION = 'gen_ai.invoke_agent.duration'; diff --git a/agents/src/voice/agent_activity.test.ts b/agents/src/voice/agent_activity.test.ts index 5093f758ca..144d442ad1 100644 --- a/agents/src/voice/agent_activity.test.ts +++ b/agents/src/voice/agent_activity.test.ts @@ -916,7 +916,12 @@ function buildPreemptiveRunner(opts: Partial = {}) { }; const generateReply = vi.fn( - () => ({ id: 'speech_fake', _cancel: () => {} }) as unknown as SpeechHandle, + () => + ({ + id: 'speech_fake', + _cancel: () => {}, + _takeAgentTurn: () => undefined, + }) as unknown as SpeechHandle, ); const cancelPreemptiveGeneration = vi.fn(); diff --git a/agents/src/voice/agent_activity.ts b/agents/src/voice/agent_activity.ts index 0de0c892be..63e291e81c 100644 --- a/agents/src/voice/agent_activity.ts +++ b/agents/src/voice/agent_activity.ts @@ -387,9 +387,10 @@ function recordQueueWait(speechHandle: SpeechHandle): void { if (queueWait === undefined || speechHandle._agentTurnContext === undefined) { return; // no agent_turn span yet: never fall back to whatever span is current } - trace - .getSpan(speechHandle._agentTurnContext) - ?.setAttribute(traceTypes.ATTR_SPEECH_QUEUE_WAIT, queueWait / 1000); + const span = trace.getSpan(speechHandle._agentTurnContext); + if (span?.isRecording()) { + span.setAttribute(traceTypes.ATTR_SPEECH_QUEUE_WAIT, queueWait / 1000); + } } /** @@ -400,12 +401,96 @@ function recordInterruption(speechHandle: SpeechHandle): void { if (!speechHandle.interrupted || speechHandle._agentTurnContext === undefined) { return; } - trace.getSpan(speechHandle._agentTurnContext)?.setAttributes({ + const span = trace.getSpan(speechHandle._agentTurnContext); + if (!span?.isRecording()) return; // the turn may have ended with the speech + span.setAttributes({ [traceTypes.ATTR_SPEECH_INTERRUPTED]: true, [traceTypes.ATTR_INTERRUPTION_SOURCE]: speechHandle._interruptSource ?? 'programmatic', }); } +/** + * Run `fn` under the speech's `agent_turn` span, made current for one generation. + * + * One speech handle is one agent turn, however many LLM steps it takes: the follow-up + * generation after a tool call runs in a new task but continues the open span instead of + * opening a second turn. Each generation is a `generation` event on the span, whose + * `lk.generation_id` names the latest one and `lk.generation_count` how many there were. The + * span ends with the speech (`SpeechHandle._markDone`), not with the step. + * + * Module-level for the same reason as `recordQueueWait`. + * @internal + */ +export async function withAgentTurn( + speechHandle: SpeechHandle, + options: { rootContext: Context | undefined; agentLabel: string }, + fn: (span: Span) => Promise, +): Promise { + let span = speechHandle._agentTurnSpan; + if (span === undefined) { + span = tracer.startSpan({ + name: 'agent_turn', + context: options.rootContext, + attributes: { [traceTypes.ATTR_SPEECH_ID]: speechHandle.id }, + }); + // an agent turn is the convention's `invoke_agent`: the framework running the agent + // in-process, with the inference and tool spans nested underneath + genAI.setAgentAttributes(span, { + operation: traceTypes.GenAIOperationName.INVOKE_AGENT, + agentName: options.agentLabel, + }); + speechHandle._agentTurnSpan = span; + speechHandle._agentTurnContext = trace.setSpan(options.rootContext ?? ROOT_CONTEXT, span); + speechHandle._agentTurnStartedAt = performance.now(); + speechHandle._agentTurnAgentName = options.agentLabel; + } + + const generationAttrs: Record = { + [traceTypes.ATTR_AGENT_TURN_ID]: speechHandle._generationId, + }; + const parentId = speechHandle._parentGenerationId; + if (parentId) generationAttrs[traceTypes.ATTR_AGENT_PARENT_TURN_ID] = parentId; + span.addEvent('generation', generationAttrs); + // the count is of generation events on the span, whichever speech emitted them: a turn + // continued from a discarded preemptive attempt or a tool call keeps counting + speechHandle._agentTurnGenerations += 1; + speechHandle._emittedGenerationStep = speechHandle._generationStep; + span.setAttributes({ + [traceTypes.ATTR_AGENT_TURN_ID]: speechHandle._generationId, + [traceTypes.ATTR_GENERATION_COUNT]: speechHandle._agentTurnGenerations, + }); + const turnSpan = span; + return otelContext.with(speechHandle._agentTurnContext!, () => fn(turnSpan)); +} + +/** + * A preemptive generation discarded for `successor` (a newer attempt, or the real reply after + * the transcript changed) hands its open `agent_turn` over, so one turn shows the wasted + * generation and the one that answered. Module-level like `recordQueueWait`. + * @internal + */ +export function continueDiscardedTurn( + discarded: SpeechHandle | undefined, + successor: SpeechHandle, +): void { + if (discarded === undefined || discarded === successor) return; + const carry = discarded._takeAgentTurn(); + if (carry !== undefined) successor._continueAgentTurn(carry, discarded); +} + +/** + * A realtime tool reply runs on a new speech handle after the tool calls of `speech`; Python + * runs it on the same handle as its next step. The reply continues the tool call's open + * `agent_turn`, with its generation numbered after (and parented to) the tool call's, so the + * trace shows one turn either way. Module-level like `continueDiscardedTurn`. + * @internal + */ +export function continueToolReplyTurn(speech: SpeechHandle, reply: SpeechHandle): void { + if (speech === reply) return; + const carry = speech._takeAgentTurn(); + if (carry !== undefined) reply._continueAgentTurn(carry, speech, 'tool_reply'); +} + export class AgentActivity implements RecognitionHooks { agent: Agent; agentSession: AgentSession; @@ -2297,7 +2382,10 @@ export class AgentActivity implements RecognitionHooks { return; } - // more of the user's turn arrived: the attempt answered a transcript that is now stale + // a newer attempt supersedes the current one; if one is created below it continues the + // discarded attempt's agent_turn (the cancelled speech only ends once the loop runs). More + // of the user's turn arrived: the attempt answered a transcript that is now stale + const discarded = this._preemptiveGeneration?.speechHandle; this.cancelPreemptiveGeneration('user_turn'); if ( @@ -2332,6 +2420,8 @@ export class AgentActivity implements RecognitionHooks { chatCtx, scheduleSpeech: false, inputDetails: { modality: 'audio' }, + // a newer attempt supersedes the current one: it continues that attempt's agent_turn + continueTurnFrom: discarded, }); this._preemptiveGeneration = { @@ -2484,7 +2574,16 @@ export class AgentActivity implements RecognitionHooks { } if (ownedSpeechHandle) { - return speechHandleStorage.run(ownedSpeechHandle, () => taskFn(ctrl)); + return speechHandleStorage.run(ownedSpeechHandle, () => + taskFn(ctrl).catch((error: unknown) => { + // the first failure of an owned task fails the speech (its agent_turn ends with + // the error and exception() reports it); a cancellation is not a failure + if ((error as Error | undefined)?.name !== 'AbortError') { + ownedSpeechHandle._error ??= error; + } + throw error; + }), + ); } return taskFn(ctrl); }); @@ -2893,6 +2992,12 @@ export class AgentActivity implements RecognitionHooks { allowInterruptions?: boolean; scheduleSpeech?: boolean; inputDetails?: InputDetails; + /** + * A discarded preemptive attempt answering the same user turn: the new speech continues + * its open `agent_turn` (see {@link continueDiscardedTurn}). Handed over before the reply + * task starts, since the task opens the turn in its first synchronous statements. + */ + continueTurnFrom?: SpeechHandle; }): SpeechHandle { const { userMessage, @@ -2902,6 +3007,7 @@ export class AgentActivity implements RecognitionHooks { allowInterruptions: defaultAllowInterruptions, scheduleSpeech = true, inputDetails, + continueTurnFrom, } = options; let instructions: string | Instructions | undefined = defaultInstructions; @@ -2940,6 +3046,9 @@ export class AgentActivity implements RecognitionHooks { allowInterruptions: allowInterruptions ?? this.allowInterruptions, inputDetails, }); + // before the reply task below runs (Task starts its body synchronously): the task must find + // the adopted turn on the handle, or it opens a second one that nothing ever ends + continueDiscardedTurn(continueTurnFrom, handle); this.agentSession.emit( AgentSessionEventTypes.SpeechCreated, @@ -3276,6 +3385,7 @@ export class AgentActivity implements RecognitionHooks { } let speechHandle: SpeechHandle | undefined; + let discardedPreemptive: SpeechHandle | undefined; if (this._preemptiveGeneration !== undefined) { const preemptive = this._preemptiveGeneration; // make sure the onUserTurnCompleted didn't change some request parameters @@ -3306,6 +3416,7 @@ export class AgentActivity implements RecognitionHooks { this.logger.warn( 'preemptive generation invalidated after `onUserTurnCompleted` because the transcript, chat context, tools, or tool choice changed', ); + discardedPreemptive = preemptive.speechHandle; preemptive.speechHandle._cancel('user_turn'); } @@ -3319,6 +3430,8 @@ export class AgentActivity implements RecognitionHooks { userMessage, chatCtx, inputDetails: { modality: 'audio' }, + // the invalidated preemptive attempt answered this same turn: one agent_turn + continueTurnFrom: discardedPreemptive, }); } @@ -3345,6 +3458,29 @@ export class AgentActivity implements RecognitionHooks { modelSettings: ModelSettings, replyAbortController: AbortController, audio?: ReadableStream | null, + ): Promise { + return withAgentTurn( + stateLease.speechHandle, + { rootContext: this.agentSession.rootSpanContext, agentLabel: this.agent.id }, + () => + this.ttsTaskImpl( + stateLease, + text, + addToChatCtx, + modelSettings, + replyAbortController, + audio, + ), + ); + } + + private async ttsTaskImpl( + stateLease: AgentStateLease, + text: string | ReadableStream, + addToChatCtx: boolean, + modelSettings: ModelSettings, + replyAbortController: AbortController, + audio?: ReadableStream | null, ): Promise { const { speechHandle } = stateLease; @@ -3570,8 +3706,6 @@ export class AgentActivity implements RecognitionHooks { _previousUserMetrics?: MetricsReport; }): Promise => { const { speechHandle } = stateLease; - // the turn's context is built from its span, never from whatever is current here - speechHandle._agentTurnContext = trace.setSpan(otelContext.active(), span); span.setAttribute(traceTypes.ATTR_SPEECH_ID, speechHandle.id); if (instructions) { @@ -3662,6 +3796,12 @@ export class AgentActivity implements RecognitionHooks { this.llm?.provider, ); tasks.push(llmTask); + // as python's _on_llm_task_done: a genuine LLM failure (not a cancellation) fails the + // speech, through exception() and the agent_turn span. Nothing else awaits this task's + // rejection: the pipeline reads the node's streams, not its result + void llmTask.result.catch((error: unknown) => { + if ((error as Error | undefined)?.name !== 'AbortError') speechHandle._error ??= error; + }); interface SpeechSegment { textStream: ReadableStream; @@ -4213,17 +4353,6 @@ export class AgentActivity implements RecognitionHooks { } }; - /** - * An agent turn is the convention's `invoke_agent`: the framework running the agent - * in-process, with the inference (`chat`) and tool (`execute_tool`) spans nested underneath. - */ - private recordAgentTurn(span: Span): void { - genAI.setAgentAttributes(span, { - operation: traceTypes.GenAIOperationName.INVOKE_AGENT, - agentName: this.agent.id, - }); - } - private pipelineReplyTask = async ( stateLease: AgentStateLease, chatCtx: ChatContext, @@ -4234,9 +4363,10 @@ export class AgentActivity implements RecognitionHooks { newMessage?: ChatMessage, _previousUserMetrics?: MetricsReport, ): Promise => - tracer.startActiveSpan( + withAgentTurn( + stateLease.speechHandle, + { rootContext: this.agentSession.rootSpanContext, agentLabel: this.agent.id }, async (span) => { - this.recordAgentTurn(span); try { await this._pipelineReplyTaskImpl({ stateLease, @@ -4259,10 +4389,6 @@ export class AgentActivity implements RecognitionHooks { recordInterruption(stateLease.speechHandle); } }, - { - name: 'agent_turn', - context: this.agentSession.rootSpanContext, - }, ); private async realtimeGenerationTask( @@ -4272,9 +4398,10 @@ export class AgentActivity implements RecognitionHooks { replyAbortController: AbortController, addToChatCtx: boolean = true, ): Promise { - return tracer.startActiveSpan( + return withAgentTurn( + stateLease.speechHandle, + { rootContext: this.agentSession.rootSpanContext, agentLabel: this.agent.id }, async (span) => { - this.recordAgentTurn(span); const inferenceSpan = tracer.startSpan({ name: 'realtime_inference' }); try { return await this._realtimeGenerationTaskImpl({ @@ -4290,10 +4417,6 @@ export class AgentActivity implements RecognitionHooks { inferenceSpan.end(); } }, - { - name: 'agent_turn', - context: this.agentSession.rootSpanContext, - }, ); } @@ -4315,8 +4438,6 @@ export class AgentActivity implements RecognitionHooks { inferenceSpan: Span; }): Promise { const { speechHandle } = stateLease; - // the turn's context is built from its span, never from whatever is current here - speechHandle._agentTurnContext = trace.setSpan(otelContext.active(), span); span.setAttribute(traceTypes.ATTR_SPEECH_ID, speechHandle.id); @@ -4901,6 +5022,8 @@ export class AgentActivity implements RecognitionHooks { stepIndex: speechHandle.numSteps + 1, parent: speechHandle, }); + // one agent_turn for the tool call and its reply, as when they share a handle + continueToolReplyTurn(speechHandle, replySpeechHandle); this.agentSession.emit( AgentSessionEventTypes.SpeechCreated, createSpeechCreatedEvent({ diff --git a/agents/src/voice/agent_turn_span.test.ts b/agents/src/voice/agent_turn_span.test.ts new file mode 100644 index 0000000000..08839437c8 --- /dev/null +++ b/agents/src/voice/agent_turn_span.test.ts @@ -0,0 +1,437 @@ +// SPDX-FileCopyrightText: 2026 LiveKit, Inc. +// +// SPDX-License-Identifier: Apache-2.0 + +/** + * One `agent_turn` span per speech handle. + * + * A reply that calls a tool runs two generations (LLM steps) in two tasks; they used to be two + * `agent_turn` spans linked only by `lk.parent_generation_id`. The speech handle now owns a + * single span for its whole life: each generation is an event on it, tool and inference spans + * nest under it, and it ends with the speech. + */ +import { AudioFrame } from '@livekit/rtc-node'; +import { + INVALID_SPAN_CONTEXT, + ROOT_CONTEXT, + SpanStatusCode, + context as otelContext, + trace, +} from '@opentelemetry/api'; +import { + InMemorySpanExporter, + type ReadableSpan, + SimpleSpanProcessor, +} from '@opentelemetry/sdk-trace-base'; +import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node'; +import { ReadableStream } from 'node:stream/web'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { ChatMessage } from '../llm/chat_context.js'; +import { tool } from '../llm/tool_context.js'; +import { initializeLogger } from '../log.js'; +import { FakeSTT } from '../stt/testing/fake_stt.js'; +import { setTracerProvider, traceTypes, tracer } from '../telemetry/index.js'; +import * as otelMetrics from '../telemetry/otel_metrics.js'; +import { Agent } from './agent.js'; +import { continueDiscardedTurn, continueToolReplyTurn, withAgentTurn } from './agent_activity.js'; +import { AgentSession } from './agent_session.js'; +import { AudioOutput } from './io.js'; +import { SpeechHandle } from './speech_handle.js'; +import { FakeLLM } from './testing/fake_llm.js'; + +initializeLogger({ pretty: false, level: 'silent' }); + +function spansNamed(exporter: InMemorySpanExporter, name: string): ReadableSpan[] { + return exporter.getFinishedSpans().filter((span) => span.name === name); +} + +function childrenOf(exporter: InMemorySpanExporter, parent: ReadableSpan, name: string) { + return spansNamed(exporter, name).filter( + (span) => span.parentSpanContext?.spanId === parent.spanContext().spanId, + ); +} + +function ms(time: [number, number]): number { + return time[0] * 1000 + time[1] / 1e6; +} + +/** Reports playout as soon as frames arrive, so a reply "plays" instantly. */ +class ImmediateOutput extends AudioOutput { + constructor() { + super(24_000); + } + + override async captureFrame(frame: AudioFrame): Promise { + const segmentCount = this.capturedPlayoutSegments; + await super.captureFrame(frame); + if (this.capturedPlayoutSegments > segmentCount) { + this.onPlaybackStarted(Date.now()); + } + } + + override flush(): void { + super.flush(); + if (this.pendingPlayoutSegments > 0) { + this.onPlaybackFinished({ playbackPosition: 0.02, interrupted: false }); + } + } + + override clearBuffer(): void { + if (this.pendingPlayoutSegments > 0) { + this.onPlaybackFinished({ playbackPosition: 0, interrupted: true }); + } + } +} + +class WeatherAgent extends Agent { + constructor() { + super({ + instructions: 'You are a helpful assistant.', + tools: { + get_weather: tool({ + description: 'Look up the weather', + execute: async () => 'sunny in Tokyo', + }), + }, + }); + } + + override async ttsNode(): Promise> { + return new ReadableStream({ + start(controller) { + controller.enqueue(new AudioFrame(new Int16Array(480), 24_000, 1, 480)); + controller.close(); + }, + }); + } +} + +describe.sequential('agent_turn span', () => { + let exporter: InMemorySpanExporter; + let provider: NodeTracerProvider; + let originalProvider: ReturnType; + + beforeEach(() => { + originalProvider = tracer.getProvider(); + exporter = new InMemorySpanExporter(); + provider = new NodeTracerProvider({ spanProcessors: [new SimpleSpanProcessor(exporter)] }); + provider.register(); + setTracerProvider(provider); + }); + + afterEach(async () => { + vi.restoreAllMocks(); + setTracerProvider(originalProvider); + await provider.shutdown(); + trace.disable(); + otelContext.disable(); + }); + + async function runReply(llm: FakeLLM, agent: Agent, userInput: string): Promise { + const session = new AgentSession({ llm, stt: new FakeSTT() }); + session.output.audio = new ImmediateOutput(); + await session.start({ agent }); + try { + const speech = session.generateReply({ userInput }); + await speech.waitForPlayout(); + } finally { + await session.close(); + } + } + + it('a tool call is one agent turn', async () => { + const llm = new FakeLLM([ + { + input: "What's the weather in Tokyo?", + content: '', + toolCalls: [{ name: 'get_weather', args: { location: 'Tokyo' } }], + }, + { input: '"sunny in Tokyo"', content: 'It is sunny in Tokyo.' }, + ]); + await runReply(llm, new WeatherAgent(), "What's the weather in Tokyo?"); + + const [root] = spansNamed(exporter, 'agent_session'); + const turns = spansNamed(exporter, 'agent_turn'); + expect( + turns.map((turn) => turn.attributes[traceTypes.ATTR_SPEECH_ID]), + 'one agent_turn per speech', + ).toHaveLength(1); + const [turn] = turns; + expect(turn!.parentSpanContext?.spanId).toBe(root!.spanContext().spanId); + + const attrs = turn!.attributes; + const speechId = attrs[traceTypes.ATTR_SPEECH_ID] as string; + expect(attrs[traceTypes.ATTR_GENERATION_COUNT]).toBe(2); + expect(attrs[traceTypes.ATTR_AGENT_TURN_ID]).toBe(`${speechId}_2`); + const generations = turn!.events.filter((event) => event.name === 'generation'); + expect(generations.map((event) => event.attributes?.[traceTypes.ATTR_AGENT_TURN_ID])).toEqual([ + `${speechId}_1`, + `${speechId}_2`, + ]); + expect(generations[0]!.attributes?.[traceTypes.ATTR_AGENT_PARENT_TURN_ID]).toBeUndefined(); + expect(generations[1]!.attributes?.[traceTypes.ATTR_AGENT_PARENT_TURN_ID]).toBe( + `${speechId}_1`, + ); + + // both generations' inference, the tool between them, and the speech all nest under it + expect(childrenOf(exporter, turn!, 'llm_node')).toHaveLength(2); + const [toolSpan] = childrenOf(exporter, turn!, 'function_tool'); + const [tts] = childrenOf(exporter, turn!, 'tts_node'); + const [speaking] = childrenOf(exporter, turn!, 'agent_speaking'); + expect(toolSpan).toBeDefined(); + expect(tts).toBeDefined(); + expect(speaking).toBeDefined(); + expect(ms(toolSpan!.startTime)).toBeLessThan(ms(tts!.startTime)); + // and the turn covers everything, ending with the speech rather than with the first step + // (2 ms of slack: the SDK anchors each span's clock at creation) + for (const child of [toolSpan!, tts!, speaking!]) { + expect(ms(turn!.startTime)).toBeLessThanOrEqual(ms(child.startTime) + 2); + expect(ms(child.endTime)).toBeLessThanOrEqual(ms(turn!.endTime) + 2); + } + }); + + it('a plain reply is one generation', async () => { + const llm = new FakeLLM([{ input: 'Hello', content: 'Hi there' }]); + await runReply(llm, new WeatherAgent(), 'Hello'); + + const [turn] = spansNamed(exporter, 'agent_turn'); + expect(spansNamed(exporter, 'agent_turn')).toHaveLength(1); + const attrs = turn!.attributes; + expect(attrs[traceTypes.ATTR_GENERATION_COUNT]).toBe(1); + expect(attrs[traceTypes.ATTR_AGENT_TURN_ID]).toBe(`${attrs[traceTypes.ATTR_SPEECH_ID]}_1`); + expect(turn!.events.filter((event) => event.name === 'generation')).toHaveLength(1); + expect(attrs[traceTypes.ATTR_AGENT_PARENT_TURN_ID]).toBeUndefined(); + }); + + it('a discarded preemptive generation hands its turn to the successor', async () => { + // a preemptive attempt discarded for the real reply (or a newer attempt) must not leave a + // second agent_turn behind: the successor continues the span, the discarded speech ends + // without touching it + const root = tracer.startSpan({ name: 'agent_session' }); + const rootCtx = trace.setSpan(ROOT_CONTEXT, root); + const attempt = SpeechHandle.create({ allowInterruptions: true }); + await withAgentTurn(attempt, { rootContext: rootCtx, agentLabel: 'a' }, async () => { + // the attempt's first generation ran here + }); + + const reply = SpeechHandle.create({ allowInterruptions: true }); + continueDiscardedTurn(attempt, reply); + attempt._markDone(); // the cancelled attempt finishes: the span must survive it + expect(spansNamed(exporter, 'agent_turn')).toEqual([]); + + await withAgentTurn(reply, { rootContext: rootCtx, agentLabel: 'a' }, async () => {}); + reply._markDone(); + root.end(); + + const [turn] = spansNamed(exporter, 'agent_turn'); + expect(spansNamed(exporter, 'agent_turn')).toHaveLength(1); + const attrs = turn!.attributes; + expect(attrs[traceTypes.ATTR_SPEECH_ID]).toBe(reply.id); + // the discarded attempt's generation and the reply's: the count is of the turn, not the handle + expect(attrs[traceTypes.ATTR_GENERATION_COUNT]).toBe(2); + expect(turn!.events.map((event) => event.name)).toEqual([ + 'generation', + 'preemptive_generation_discarded', + 'generation', + ]); + const discarded = turn!.events.find( + (event) => event.name === 'preemptive_generation_discarded', + ); + expect(discarded?.attributes?.[traceTypes.ATTR_SPEECH_ID]).toBe(attempt.id); + + // nothing to hand over: a plain successor is untouched + continueDiscardedTurn(undefined, reply); + continueDiscardedTurn(reply, reply); + }); + + it('hands the turn over before the successor task opens its own', async () => { + // Task runs its body synchronously, and the reply task opens agent_turn in its first + // statements: a handoff performed after generateReply() returned came too late, leaving the + // successor's own span unended and its llm_node / tts_node dangling from it. The handoff now + // happens inside generateReply, before the task starts. + const llm = new FakeLLM([{ input: 'Hello', content: 'Hi there' }]); + const session = new AgentSession({ llm, stt: new FakeSTT() }); + session.output.audio = new ImmediateOutput(); + const agent = new WeatherAgent(); + await session.start({ agent }); + try { + const activity = agent._agentActivity!; + const attempt = SpeechHandle.create({ allowInterruptions: true }); + await withAgentTurn( + attempt, + { rootContext: session.rootSpanContext, agentLabel: agent.id }, + async () => {}, + ); + const reply = activity.generateReply({ + userMessage: ChatMessage.create({ role: 'user', content: 'Hello' }), + inputDetails: { modality: 'audio' }, + continueTurnFrom: attempt, + }); + attempt._markDone(); + await reply.waitForPlayout(); + } finally { + await session.close(); + } + + const turns = spansNamed(exporter, 'agent_turn'); + expect(turns).toHaveLength(1); + const [turn] = turns; + expect(turn!.attributes[traceTypes.ATTR_GENERATION_COUNT]).toBe(2); + expect(turn!.events.map((event) => event.name)).toEqual([ + 'generation', + 'preemptive_generation_discarded', + 'generation', + ]); + // the successor's work nests under the adopted turn + expect(childrenOf(exporter, turn!, 'llm_node')).toHaveLength(1); + expect(childrenOf(exporter, turn!, 'tts_node')).toHaveLength(1); + expect(childrenOf(exporter, turn!, 'agent_speaking')).toHaveLength(1); + // and nothing dangles from a span that never ended + const exported = new Set(exporter.getFinishedSpans().map((span) => span.spanContext().spanId)); + const dangling = exporter + .getFinishedSpans() + .filter((span) => span.parentSpanContext && !exported.has(span.parentSpanContext.spanId)); + expect(dangling.map((span) => span.name)).toEqual([]); + }); + + it('keeps counting generations across repeated handoffs', async () => { + // attempt A is replaced by attempt B, which is replaced by the reply: three generations on + // one turn, and the count says so however many hands the turn went through + const root = tracer.startSpan({ name: 'agent_session' }); + const rootCtx = trace.setSpan(ROOT_CONTEXT, root); + const opts = { rootContext: rootCtx, agentLabel: 'a' }; + const a = SpeechHandle.create({ allowInterruptions: true }); + await withAgentTurn(a, opts, async () => {}); + const b = SpeechHandle.create({ allowInterruptions: true }); + continueDiscardedTurn(a, b); + a._markDone(); + await withAgentTurn(b, opts, async () => {}); + const reply = SpeechHandle.create({ allowInterruptions: true }); + continueDiscardedTurn(b, reply); + b._markDone(); + await withAgentTurn(reply, opts, async () => {}); + reply._markDone(); + root.end(); + + const [turn] = spansNamed(exporter, 'agent_turn'); + expect(spansNamed(exporter, 'agent_turn')).toHaveLength(1); + expect(turn!.attributes[traceTypes.ATTR_GENERATION_COUNT]).toBe(3); + expect(turn!.events.filter((event) => event.name === 'generation')).toHaveLength(3); + // each speech still numbers its own generations, as python does + expect(turn!.attributes[traceTypes.ATTR_AGENT_TURN_ID]).toBe(`${reply.id}_1`); + }); + + it('a realtime tool reply continues the tool call turn as its next generation', async () => { + // the framework runs the reply on a new handle; python runs it on the same one as step 2. + // Either way the trace is one turn: the reply's generation numbered after the tool call's + // and parented to it, under the tool call's speech id + const root = tracer.startSpan({ name: 'agent_session' }); + const rootCtx = trace.setSpan(ROOT_CONTEXT, root); + const opts = { rootContext: rootCtx, agentLabel: 'a' }; + const speech = SpeechHandle.create({ allowInterruptions: true }); + let reply: SpeechHandle | undefined; + await withAgentTurn(speech, opts, async () => { + // the tool calls ran; the framework creates the reply inside the tool call's turn + speech._numSteps += 1; // as the realtime path does before scheduling the reply + reply = SpeechHandle.create({ allowInterruptions: true, parent: speech }); + continueToolReplyTurn(speech, reply); + }); + speech._markDone(); // the tool call's own handle ends: the span must survive it + expect(spansNamed(exporter, 'agent_turn')).toEqual([]); + await withAgentTurn(reply!, opts, async () => {}); + reply!._markDone(); + root.end(); + + const turns = spansNamed(exporter, 'agent_turn'); + expect(turns).toHaveLength(1); + const [turn] = turns; + const attrs = turn!.attributes; + expect(attrs[traceTypes.ATTR_SPEECH_ID]).toBe(speech.id); + expect(attrs[traceTypes.ATTR_GENERATION_COUNT]).toBe(2); + expect(attrs[traceTypes.ATTR_AGENT_TURN_ID]).toBe(`${speech.id}_2`); + const generations = turn!.events.filter((event) => event.name === 'generation'); + expect(generations.map((event) => event.attributes?.[traceTypes.ATTR_AGENT_TURN_ID])).toEqual([ + `${speech.id}_1`, + `${speech.id}_2`, + ]); + expect(generations[1]!.attributes?.[traceTypes.ATTR_AGENT_PARENT_TURN_ID]).toBe( + `${speech.id}_1`, + ); + expect(turn!.events.some((event) => event.name === 'preemptive_generation_discarded')).toBe( + false, + ); + // no handoff to make: the same handle + continueToolReplyTurn(reply!, reply!); + }); + + it('a task failure fails the turn and surfaces on the handle', async () => { + // an LLM node that throws rejects the speech task; the turn ends with the error and the + // handle reports it, instead of an unremarkable success + class BrokenAgent extends WeatherAgent { + override async llmNode(): Promise { + throw new Error('provider unavailable'); + } + } + const llm = new FakeLLM([{ input: 'Hello', content: 'Hi there' }]); + const session = new AgentSession({ llm, stt: new FakeSTT() }); + session.output.audio = new ImmediateOutput(); + await session.start({ agent: new BrokenAgent() }); + let speech: SpeechHandle | undefined; + try { + speech = session.generateReply({ userInput: 'Hello' }); + await speech.waitForPlayout(); + } finally { + await session.close(); + } + + expect(speech!.exception()).toBeInstanceOf(Error); + expect((speech!.exception() as Error).message).toBe('provider unavailable'); + const [turn] = spansNamed(exporter, 'agent_turn'); + expect(turn!.status.code).toBe(SpanStatusCode.ERROR); + expect( + turn!.events.find((event) => event.name === 'exception')?.attributes?.['exception.message'], + ).toBe('provider unavailable'); + }); + + it('an LLM failure stored on the handle fails the turn', async () => { + // the pipeline stores the failure on the handle when it marks it done; the turn must end as + // failed whichever step it was on + const handle = SpeechHandle.create({ allowInterruptions: true }); + await withAgentTurn(handle, { rootContext: undefined, agentLabel: 'a' }, async () => {}); + handle._markDone(new Error('llm down')); + + const [turn] = spansNamed(exporter, 'agent_turn'); + expect(turn!.status.code).toBe(SpanStatusCode.ERROR); + expect(turn!.events.filter((event) => event.name === 'exception')).toHaveLength(1); + }); + + it('records the turn duration metric even when the span is sampled out', () => { + const record = vi.spyOn(otelMetrics, 'recordInvokeAgentDuration').mockImplementation(() => {}); + const handle = SpeechHandle.create({ allowInterruptions: true }); + handle._agentTurnSpan = trace.wrapSpanContext(INVALID_SPAN_CONTEXT); + handle._agentTurnStartedAt = performance.now() - 1_000; + handle._agentTurnAgentName = 'a'; + handle._markDone(); + expect(record).toHaveBeenCalledTimes(1); + const [duration, agentName] = record.mock.calls[0]!; + expect(agentName).toBe('a'); + expect(duration).toBeGreaterThanOrEqual(1); + expect(duration).toBeLessThan(5); + }); + + it('hands a sampled-out turn to the successor too', () => { + // the successor adopts a non-recording turn as well, so the duration metric keeps the + // discarded attempt's start time + const attempt = SpeechHandle.create({ allowInterruptions: true }); + attempt._agentTurnSpan = trace.wrapSpanContext(INVALID_SPAN_CONTEXT); + attempt._agentTurnStartedAt = 1; + attempt._agentTurnAgentName = 'a'; + + const reply = SpeechHandle.create({ allowInterruptions: true }); + continueDiscardedTurn(attempt, reply); + expect(attempt._agentTurnSpan).toBeUndefined(); + expect(reply._agentTurnSpan).toBeDefined(); + expect(reply._agentTurnStartedAt).toBe(1); + expect(reply._agentTurnAgentName).toBe('a'); + }); +}); diff --git a/agents/src/voice/speech_handle.ts b/agents/src/voice/speech_handle.ts index b1293b411b..6b15f90925 100644 --- a/agents/src/voice/speech_handle.ts +++ b/agents/src/voice/speech_handle.ts @@ -2,9 +2,12 @@ // // SPDX-License-Identifier: Apache-2.0 import { ThrowsPromise } from '@livekit/throws-transformer/throws'; -import type { Context } from '@opentelemetry/api'; +import { type Context, type Span, context as otelContext, trace } from '@opentelemetry/api'; import type { ChatItem } from '../llm/index.js'; import { log } from '../log.js'; +import { recordInvokeAgentDuration } from '../telemetry/otel_metrics.js'; +import * as traceTypes from '../telemetry/trace_types.js'; +import { recordException } from '../telemetry/utils.js'; import type { Task } from '../utils.js'; import { Event, Future, dedent, shortuuid } from '../utils.js'; import { functionCallStorage } from './agent.js'; @@ -96,6 +99,18 @@ export type ResolvedSpeechHandle = Omit; */ export type InterruptionSource = 'audio_activity' | 'user_turn' | 'programmatic'; +/** An open `agent_turn` handed from a discarded speech to its successor. @internal */ +export interface AgentTurnCarry { + span: Span; + startedAt: number | undefined; + agentName: string | undefined; + /** Generation events already on the span, so `lk.generation_count` keeps counting. */ + generations: number; +} + +/** How a speech came to continue another's `agent_turn` (see `SpeechHandle._continueAgentTurn`). */ +export type AgentTurnContinuation = 'preemptive_discarded' | 'tool_reply'; + export class SpeechHandleCircularWaitError extends Error { constructor(functionCallName: string) { super(dedent` @@ -139,7 +154,8 @@ export class SpeechHandle { private doneFut = new Future(); private generations: Future[] = []; private _chatItems: ChatItem[] = []; - private _error: unknown; + /** @internal The first failure of an owned task, or the one the pipeline stored; see exception(). */ + _error: unknown; private interruptionHolds = 0; private interruptionHoldsRestore: boolean; @@ -148,9 +164,30 @@ export class SpeechHandle { /** @internal */ _numSteps = 1; + /** + * @internal Generation ids continue another speech's numbering when this speech carries on + * its turn (a realtime tool reply, which the framework runs on a new handle): the base id and + * the step the other speech had reached. + */ + _generationBaseId?: string; + /** @internal */ + _generationStepBase = 0; + /** @internal The step of the last generation event this speech emitted on its turn. */ + _emittedGenerationStep?: number; + /** @internal Generation events on the turn this speech owns, carried over on a handoff. */ + _agentTurnGenerations = 0; + /** + * @internal One `agent_turn` span for the whole speech, however many generations (LLM steps) + * it takes; opened by the first reply task, ended with the speech in `_markDone`. + */ + _agentTurnSpan?: Span; /** @internal - OpenTelemetry context for the agent turn span */ _agentTurnContext?: Context; + /** @internal - when the turn opened (performance.now), for the duration metric */ + _agentTurnStartedAt?: number; + /** @internal - the agent the turn was opened for, for the duration metric */ + _agentTurnAgentName?: string; /** @internal - when the speech was scheduled, for the queue-wait attribute */ _scheduledAt?: number; @@ -224,6 +261,23 @@ export class SpeechHandle { return this._id; } + /** @internal The step of the current generation in the turn's numbering (see `_generationBaseId`). */ + get _generationStep(): number { + return this._generationStepBase + this._numSteps; + } + + /** @internal The id of the current generation (LLM step) of this speech. */ + get _generationId(): string { + return `${this._generationBaseId ?? this._id}_${this._generationStep}`; + } + + /** @internal The id of the generation before the current one; undefined on the first. */ + get _parentGenerationId(): string | undefined { + const step = this._generationStep; + if (step <= 1) return undefined; + return `${this._generationBaseId ?? this._id}_${step - 1}`; + } + get scheduled(): boolean { return this.scheduledFut.done; } @@ -528,6 +582,8 @@ export class SpeechHandle { } this.doneFut.resolve(); } + // a pipeline LLM failure is stored on the handle before the tasks finish + this.endAgentTurn(error !== undefined ? error : this._error); // Keep this outside the doneFut guard: if the handle is already done but a // generation future is still active, _waitForGeneration() must be released. @@ -538,6 +594,91 @@ export class SpeechHandle { this.clearInterruptTimeout(); } + /** + * Detach this speech's open `agent_turn` so a successor can continue it. + * + * Used when a preemptive generation is discarded for another speech answering the same user + * turn: the wasted generation stays visible under the one turn instead of becoming a turn of + * its own. After this the speech ends without touching the span. + * @internal + */ + _takeAgentTurn(): AgentTurnCarry | undefined { + const span = this._agentTurnSpan; + if (span === undefined) return undefined; + const carry: AgentTurnCarry = { + span, + startedAt: this._agentTurnStartedAt, + agentName: this._agentTurnAgentName, + generations: this._agentTurnGenerations, + }; + this._agentTurnSpan = undefined; + this._agentTurnContext = undefined; + this._agentTurnStartedAt = undefined; + this._agentTurnAgentName = undefined; + this._agentTurnGenerations = 0; + return carry; + } + + /** + * @internal Adopt the `agent_turn` taken from `from` (see {@link _takeAgentTurn}). + * + * - `preemptive_discarded` (the default): `from` was a preemptive attempt dropped for this + * speech; the span records that and takes this speech's id. Generation ids stay this + * speech's own, as in Python. + * - `tool_reply`: this speech is the realtime tool reply the framework runs on a new handle + * after `from`'s tool calls; Python runs it on the same handle as its next step. The turn + * keeps `from`'s speech id and this speech's generations continue `from`'s numbering, so the + * trace reads as Python's: one turn, the reply's generation parented to the tool call's. + */ + _continueAgentTurn( + carry: AgentTurnCarry, + from: SpeechHandle, + continuation: AgentTurnContinuation = 'preemptive_discarded', + ): void { + // adopted even when sampled out: the duration metric still needs the start time + const { span, startedAt, agentName } = carry; + const own = this._agentTurnSpan; + if (own !== undefined && own !== span) { + // this speech already opened a turn of its own (the handoff came after its task started): + // close it rather than leak an unended span that its children would dangle from + own.addEvent('superseded_by_adopted_turn', { [traceTypes.ATTR_SPEECH_ID]: from.id }); + if (own.isRecording()) own.end(); + } + if (continuation === 'preemptive_discarded') { + span.addEvent('preemptive_generation_discarded', { + [traceTypes.ATTR_SPEECH_ID]: from.id, + }); + span.setAttribute(traceTypes.ATTR_SPEECH_ID, this.id); + } else { + this._generationBaseId = from._generationBaseId ?? from.id; + this._generationStepBase = from._emittedGenerationStep ?? from._generationStep; + } + this._agentTurnSpan = span; + this._agentTurnContext = trace.setSpan(otelContext.active(), span); + this._agentTurnStartedAt = startedAt; + this._agentTurnAgentName = agentName; + this._agentTurnGenerations = carry.generations; + } + + /** Close the speech's `agent_turn` span: the speech is done, whatever step it was on. */ + private endAgentTurn(error: unknown): void { + const span = this._agentTurnSpan; + this._agentTurnSpan = undefined; + if (span === undefined) return; + // the duration metric does not depend on the span being sampled in + if (this._agentTurnStartedAt !== undefined && this._agentTurnAgentName !== undefined) { + recordInvokeAgentDuration( + (performance.now() - this._agentTurnStartedAt) / 1000, + this._agentTurnAgentName, + ); + } + if (!span.isRecording()) return; + if (error instanceof Error) { + recordException(span, error); + } + span.end(); + } + /** @internal */ _markScheduled(): void { if (this._authorizedAt !== undefined) { From 2c883a6d7ef03091fa526dd779c029a621d156d0 Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 25 Sep 2026 16:09:11 -0700 Subject: [PATCH 2/4] review: non-Error failures fail the turn; generation steps file under the turn's speech id - endAgentTurn records a thrown string or object as an exception, so the turn does not read as a success while the handle reports a failure - the reply steps stamp the turn's speech id (the root speech's when a realtime tool reply continues its parent's turn on a new handle) instead of overwriting it with the step's own handle id Co-Authored-By: Claude Fable 5.1 --- agents/src/voice/agent_activity.ts | 8 ++++--- agents/src/voice/agent_turn_span.test.ts | 28 +++++++++++++++++++----- agents/src/voice/speech_handle.ts | 13 +++++++++-- 3 files changed, 39 insertions(+), 10 deletions(-) diff --git a/agents/src/voice/agent_activity.ts b/agents/src/voice/agent_activity.ts index 63e291e81c..fdc77ec5d7 100644 --- a/agents/src/voice/agent_activity.ts +++ b/agents/src/voice/agent_activity.ts @@ -431,7 +431,7 @@ export async function withAgentTurn( span = tracer.startSpan({ name: 'agent_turn', context: options.rootContext, - attributes: { [traceTypes.ATTR_SPEECH_ID]: speechHandle.id }, + attributes: { [traceTypes.ATTR_SPEECH_ID]: speechHandle._turnSpeechId }, }); // an agent turn is the convention's `invoke_agent`: the framework running the agent // in-process, with the inference and tool spans nested underneath @@ -3707,7 +3707,8 @@ export class AgentActivity implements RecognitionHooks { }): Promise => { const { speechHandle } = stateLease; - span.setAttribute(traceTypes.ATTR_SPEECH_ID, speechHandle.id); + // the turn's id, not this step's handle: a tool reply on a new handle continues its parent + span.setAttribute(traceTypes.ATTR_SPEECH_ID, speechHandle._turnSpeechId); if (instructions) { span.setAttribute(traceTypes.ATTR_INSTRUCTIONS, renderInstructions(instructions)); } @@ -4439,7 +4440,8 @@ export class AgentActivity implements RecognitionHooks { }): Promise { const { speechHandle } = stateLease; - span.setAttribute(traceTypes.ATTR_SPEECH_ID, speechHandle.id); + // the turn's id, not this step's handle: a tool reply on a new handle continues its parent + span.setAttribute(traceTypes.ATTR_SPEECH_ID, speechHandle._turnSpeechId); const localParticipant = this.agentSession._roomIO?.localParticipant; if (localParticipant) { diff --git a/agents/src/voice/agent_turn_span.test.ts b/agents/src/voice/agent_turn_span.test.ts index 08839437c8..0b55b14068 100644 --- a/agents/src/voice/agent_turn_span.test.ts +++ b/agents/src/voice/agent_turn_span.test.ts @@ -347,6 +347,9 @@ describe.sequential('agent_turn span', () => { const [turn] = turns; const attrs = turn!.attributes; expect(attrs[traceTypes.ATTR_SPEECH_ID]).toBe(speech.id); + // what the reply's own step stamps on the turn: the tool call's id, not the new handle's + expect(reply!._turnSpeechId).toBe(speech.id); + expect(speech._turnSpeechId).toBe(speech.id); expect(attrs[traceTypes.ATTR_GENERATION_COUNT]).toBe(2); expect(attrs[traceTypes.ATTR_AGENT_TURN_ID]).toBe(`${speech.id}_2`); const generations = turn!.events.filter((event) => event.name === 'generation'); @@ -364,12 +367,16 @@ describe.sequential('agent_turn span', () => { continueToolReplyTurn(reply!, reply!); }); - it('a task failure fails the turn and surfaces on the handle', async () => { + it.each([ + ['an Error', new Error('provider unavailable')], + ['a string', 'provider unavailable'], + ])('a task failure fails the turn and surfaces on the handle (%s)', async (_kind, failure) => { // an LLM node that throws rejects the speech task; the turn ends with the error and the - // handle reports it, instead of an unremarkable success + // handle reports it, instead of an unremarkable success. A thrown value that is not an + // Error is a failure all the same class BrokenAgent extends WeatherAgent { override async llmNode(): Promise { - throw new Error('provider unavailable'); + throw failure; } } const llm = new FakeLLM([{ input: 'Hello', content: 'Hi there' }]); @@ -384,8 +391,7 @@ describe.sequential('agent_turn span', () => { await session.close(); } - expect(speech!.exception()).toBeInstanceOf(Error); - expect((speech!.exception() as Error).message).toBe('provider unavailable'); + expect(speech!.exception()).toBe(failure); const [turn] = spansNamed(exporter, 'agent_turn'); expect(turn!.status.code).toBe(SpanStatusCode.ERROR); expect( @@ -393,6 +399,18 @@ describe.sequential('agent_turn span', () => { ).toBe('provider unavailable'); }); + it('a failure that is not an Error still fails the turn', async () => { + const handle = SpeechHandle.create({ allowInterruptions: true }); + await withAgentTurn(handle, { rootContext: undefined, agentLabel: 'a' }, async () => {}); + handle._markDone('llm down'); + + const [turn] = spansNamed(exporter, 'agent_turn'); + expect(turn!.status.code).toBe(SpanStatusCode.ERROR); + expect( + turn!.events.find((event) => event.name === 'exception')?.attributes?.['exception.message'], + ).toBe('llm down'); + }); + it('an LLM failure stored on the handle fails the turn', async () => { // the pipeline stores the failure on the handle when it marks it done; the turn must end as // failed whichever step it was on diff --git a/agents/src/voice/speech_handle.ts b/agents/src/voice/speech_handle.ts index 6b15f90925..08caf664c4 100644 --- a/agents/src/voice/speech_handle.ts +++ b/agents/src/voice/speech_handle.ts @@ -271,6 +271,14 @@ export class SpeechHandle { return `${this._generationBaseId ?? this._id}_${this._generationStep}`; } + /** + * @internal The speech id the turn is filed under: the root speech's when this handle + * continues its turn (a realtime tool reply runs on a new handle), else its own. + */ + get _turnSpeechId(): string { + return this._generationBaseId ?? this._id; + } + /** @internal The id of the generation before the current one; undefined on the first. */ get _parentGenerationId(): string | undefined { const step = this._generationStep; @@ -673,8 +681,9 @@ export class SpeechHandle { ); } if (!span.isRecording()) return; - if (error instanceof Error) { - recordException(span, error); + if (error !== undefined) { + // a thrown string or object is a failure too: the turn must not read as a success + recordException(span, error instanceof Error ? error : new Error(String(error))); } span.end(); } From 2464edbdfffd41fe9dadaa374428f7a4dc9de16e Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 25 Sep 2026 19:27:45 -0700 Subject: [PATCH 3/4] review: the turn's queue wait is its first generation's With queue waits measured per generation, a tool reply on the same handle would replace the turn's lk.speech_queue_wait with its own; the span keeps the wait before the reply started. Co-Authored-By: Claude Fable 5.1 --- agents/src/voice/agent_activity.ts | 4 ++++ agents/src/voice/speech_handle.ts | 2 ++ 2 files changed, 6 insertions(+) diff --git a/agents/src/voice/agent_activity.ts b/agents/src/voice/agent_activity.ts index fdc77ec5d7..cc6481020a 100644 --- a/agents/src/voice/agent_activity.ts +++ b/agents/src/voice/agent_activity.ts @@ -387,9 +387,13 @@ function recordQueueWait(speechHandle: SpeechHandle): void { if (queueWait === undefined || speechHandle._agentTurnContext === undefined) { return; // no agent_turn span yet: never fall back to whatever span is current } + // the turn's wait is its first generation's: the time before the reply started. A tool + // reply scheduled later on the same handle measures its own wait, which must not replace it + if (speechHandle._queueWaitRecorded) return; const span = trace.getSpan(speechHandle._agentTurnContext); if (span?.isRecording()) { span.setAttribute(traceTypes.ATTR_SPEECH_QUEUE_WAIT, queueWait / 1000); + speechHandle._queueWaitRecorded = true; } } diff --git a/agents/src/voice/speech_handle.ts b/agents/src/voice/speech_handle.ts index 08caf664c4..2179e04824 100644 --- a/agents/src/voice/speech_handle.ts +++ b/agents/src/voice/speech_handle.ts @@ -193,6 +193,8 @@ export class SpeechHandle { _scheduledAt?: number; /** @internal - when generation was first authorized, for the queue-wait attribute */ _authorizedAt?: number; + /** @internal The turn's queue wait is stamped once, for the first generation. */ + _queueWaitRecorded = false; /** @internal - the first interrupt's cause, for the agent_turn trace */ _interruptSource?: InterruptionSource; From 2cc06f97c26c1dddbb52a9dedf560100223d699d Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 25 Sep 2026 22:24:00 -0700 Subject: [PATCH 4/4] review: one turn for a realtime model's own tool reply; ownership guard; chained reply ids - a realtime model that answers its tool calls itself (autoToolReplyGeneration) opens the reply generation after the tool call's speech is done: the turn is now held open for it (AutoToolReplyTurnHold) and adopted by the model's next generation, ended after 5 s or at activity close when none comes - the pipeline task stamps lk.interrupted only while its speech still owns the turn, so a discarded preemptive attempt unwinding cannot write the successor's verdict - a handle that adopted a turn but never emitted a generation passes the numbering on from where it stood, so the next id follows an existing one; chained replies tested to three generations Co-Authored-By: Claude Fable 5.1 --- agents/src/voice/agent_activity.ts | 66 ++++++++++++-- agents/src/voice/agent_turn_span.test.ts | 109 ++++++++++++++++++++++- agents/src/voice/speech_handle.ts | 44 ++++++--- 3 files changed, 197 insertions(+), 22 deletions(-) diff --git a/agents/src/voice/agent_activity.ts b/agents/src/voice/agent_activity.ts index cc6481020a..b59f6867d0 100644 --- a/agents/src/voice/agent_activity.ts +++ b/agents/src/voice/agent_activity.ts @@ -163,10 +163,12 @@ import { import type { PlaybackFinishedEvent, TimedString } from './io.js'; import { releaseTts, retainTts } from './model_refs.js'; import { + type AgentTurnCarry, type InputDetails, type InterruptionSource, REPLY_TASK_CANCEL_TIMEOUT, SpeechHandle, + endCarriedAgentTurn, } from './speech_handle.js'; import { ToolExecutor, @@ -495,6 +497,45 @@ export function continueToolReplyTurn(speech: SpeechHandle, reply: SpeechHandle) if (carry !== undefined) reply._continueAgentTurn(carry, speech, 'tool_reply'); } +/** + * The turn of a realtime tool call whose reply the model generates on its own + * (`autoToolReplyGeneration`): held open past the tool call's speech until the model's next + * generation adopts it (see {@link continueToolReplyTurn}), or ended after `timeout` ms when + * none comes (the user spoke first, the model declined). Module-level like the helpers above. + * @internal + */ +export class AutoToolReplyTurnHold { + private held?: { carry: AgentTurnCarry; from: SpeechHandle; timer: NodeJS.Timeout }; + + /** Take `speech`'s open turn, ending any turn still held. */ + hold(speech: SpeechHandle, timeout = 5000): void { + this.end(); + const carry = speech._takeAgentTurn(); + if (carry === undefined) return; + const timer = setTimeout(() => this.end(), timeout); + this.held = { carry, from: speech, timer }; + } + + /** Continue the held turn on `reply`, the generation that answers the tool calls. */ + adopt(reply: SpeechHandle): boolean { + const held = this.held; + if (held === undefined) return false; + clearTimeout(held.timer); + this.held = undefined; + reply._continueAgentTurn(held.carry, held.from, 'tool_reply'); + return true; + } + + /** End the held turn, if any: no generation is coming for it. */ + end(): void { + const held = this.held; + if (held === undefined) return; + clearTimeout(held.timer); + this.held = undefined; + endCarriedAgentTurn(held.carry); + } +} + export class AgentActivity implements RecognitionHooks { agent: Agent; agentSession: AgentSession; @@ -524,6 +565,7 @@ export class AgentActivity implements RecognitionHooks { // Placeholder used to hold a RunResult open while waiting for a realtime // model to auto-generate a tool reply (autoToolReplyGeneration=true). private pendingAutoToolReplyFut?: Future; + private autoToolReplyTurn = new AutoToolReplyTurnHold(); private lock = new Mutex(); private inlineTaskLock = new Mutex(); private audioStream = new MultiInputStream(); @@ -2017,6 +2059,9 @@ export class AgentActivity implements RecognitionHooks { const handle = SpeechHandle.create({ allowInterruptions: this.allowInterruptions, }); + // the model answering its tool calls: one agent_turn for the call and the reply, adopted + // before the reply task below opens a turn of its own + this.autoToolReplyTurn.adopt(handle); this.agentSession.emit( AgentSessionEventTypes.SpeechCreated, createSpeechCreatedEvent({ @@ -4386,11 +4431,14 @@ export class AgentActivity implements RecognitionHooks { }); } finally { // an interruption while the tools run makes the task return early: the verdict is - // stamped here, whichever way the task left - span.setAttribute( - traceTypes.ATTR_SPEECH_INTERRUPTED, - stateLease.speechHandle.interrupted, - ); + // stamped here, whichever way the task left. Not once the turn was handed to a + // successor (a discarded preemptive attempt unwinding): the verdict is then its + if (stateLease.speechHandle._agentTurnContext !== undefined) { + span.setAttribute( + traceTypes.ATTR_SPEECH_INTERRUPTED, + stateLease.speechHandle.interrupted, + ); + } recordInterruption(stateLease.speechHandle); } }, @@ -5016,8 +5064,11 @@ export class AgentActivity implements RecognitionHooks { } } - // skip realtime reply if not required or auto-generated - if (!shouldGenerateToolReply || realtimeModel.capabilities.autoToolReplyGeneration) { + if (!shouldGenerateToolReply) return; + + // the model generates the reply itself: its next generation continues this turn + if (realtimeModel.capabilities.autoToolReplyGeneration) { + this.autoToolReplyTurn.hold(speechHandle); return; } @@ -5565,6 +5616,7 @@ export class AgentActivity implements RecognitionHooks { this.closed = true; this.cancelPreemptiveGeneration(); + this.autoToolReplyTurn.end(); // Commit the in-flight assistant turn before teardown (#2041): a room // disconnect mid-playout parks the reply task on a playout promise that diff --git a/agents/src/voice/agent_turn_span.test.ts b/agents/src/voice/agent_turn_span.test.ts index 0b55b14068..e18e5aff48 100644 --- a/agents/src/voice/agent_turn_span.test.ts +++ b/agents/src/voice/agent_turn_span.test.ts @@ -33,7 +33,12 @@ import { FakeSTT } from '../stt/testing/fake_stt.js'; import { setTracerProvider, traceTypes, tracer } from '../telemetry/index.js'; import * as otelMetrics from '../telemetry/otel_metrics.js'; import { Agent } from './agent.js'; -import { continueDiscardedTurn, continueToolReplyTurn, withAgentTurn } from './agent_activity.js'; +import { + AutoToolReplyTurnHold, + continueDiscardedTurn, + continueToolReplyTurn, + withAgentTurn, +} from './agent_activity.js'; import { AgentSession } from './agent_session.js'; import { AudioOutput } from './io.js'; import { SpeechHandle } from './speech_handle.js'; @@ -367,6 +372,108 @@ describe.sequential('agent_turn span', () => { continueToolReplyTurn(reply!, reply!); }); + it('numbers chained realtime tool replies in order', async () => { + // a reply that calls another tool is answered by a further reply: generations 1, 2, 3 on + // one turn, each parented to the one before + const root = tracer.startSpan({ name: 'agent_session' }); + const opts = { rootContext: trace.setSpan(ROOT_CONTEXT, root), agentLabel: 'a' }; + const first = SpeechHandle.create({ allowInterruptions: true }); + let speech = first; + for (let i = 0; i < 2; i++) { + const current = speech; + await withAgentTurn(current, opts, async () => { + current._numSteps += 1; + const reply = SpeechHandle.create({ allowInterruptions: true, parent: current }); + continueToolReplyTurn(current, reply); + speech = reply; + }); + current._markDone(); + } + await withAgentTurn(speech, opts, async () => {}); + speech._markDone(); + root.end(); + + const [turn, ...rest] = spansNamed(exporter, 'agent_turn'); + expect(rest).toEqual([]); + expect(turn!.attributes[traceTypes.ATTR_SPEECH_ID]).toBe(first.id); + expect(turn!.attributes[traceTypes.ATTR_GENERATION_COUNT]).toBe(3); + const generations = turn!.events.filter((event) => event.name === 'generation'); + expect( + generations.map((event) => [ + event.attributes?.[traceTypes.ATTR_AGENT_TURN_ID], + event.attributes?.[traceTypes.ATTR_AGENT_PARENT_TURN_ID], + ]), + ).toEqual([ + [`${first.id}_1`, undefined], + [`${first.id}_2`, `${first.id}_1`], + [`${first.id}_3`, `${first.id}_2`], + ]); + // a handle that adopted the turn but never ran (cancelled first) passes the numbering on + // from where it stood, so the next id follows a generation that exists + const idle = SpeechHandle.create({ allowInterruptions: true }); + idle._generationBaseId = first.id; + idle._generationStepBase = 3; + const next = SpeechHandle.create({ allowInterruptions: true }); + next._continueAgentTurn( + { + span: tracer.startSpan({ name: 'agent_turn' }), + startedAt: undefined, + agentName: undefined, + generations: 3, + }, + idle, + 'tool_reply', + ); + expect(next._generationId).toBe(`${first.id}_4`); + expect(next._parentGenerationId).toBe(`${first.id}_3`); + next._markDone(); + }); + + it('a realtime model answering its tool calls itself continues the tool call turn', async () => { + // with autoToolReplyGeneration the model opens the reply generation on its own, after the + // tool call's speech is done: the turn is held open for it and adopted by the next + // generation, so the trace reads as the manual reply's + const root = tracer.startSpan({ name: 'agent_session' }); + const opts = { rootContext: trace.setSpan(ROOT_CONTEXT, root), agentLabel: 'a' }; + const hold = new AutoToolReplyTurnHold(); + const speech = SpeechHandle.create({ allowInterruptions: true }); + await withAgentTurn(speech, opts, async () => { + speech._numSteps += 1; + hold.hold(speech); + }); + speech._markDone(); // the tool call's handle ends: the held turn survives it + expect(spansNamed(exporter, 'agent_turn')).toEqual([]); + + const reply = SpeechHandle.create({ allowInterruptions: true }); + expect(hold.adopt(reply)).toBe(true); + await withAgentTurn(reply, opts, async () => {}); + reply._markDone(); + const [turn, ...rest] = spansNamed(exporter, 'agent_turn'); + expect(rest).toEqual([]); + expect(turn!.attributes[traceTypes.ATTR_SPEECH_ID]).toBe(speech.id); + expect(turn!.attributes[traceTypes.ATTR_GENERATION_COUNT]).toBe(2); + expect(turn!.attributes[traceTypes.ATTR_AGENT_TURN_ID]).toBe(`${speech.id}_2`); + // nothing held any more + expect(hold.adopt(SpeechHandle.create({ allowInterruptions: true }))).toBe(false); + + // no generation comes (the user spoke first): the held turn ends on its own + const unanswered = SpeechHandle.create({ allowInterruptions: true }); + await withAgentTurn(unanswered, opts, async () => hold.hold(unanswered, 20)); + unanswered._markDone(); + expect(spansNamed(exporter, 'agent_turn')).toHaveLength(1); + await new Promise((resolve) => setTimeout(resolve, 40)); + expect( + spansNamed(exporter, 'agent_turn').map((s) => s.attributes[traceTypes.ATTR_SPEECH_ID]), + ).toEqual([speech.id, unanswered.id]); + // and ending the hold explicitly (the activity closing) ends a held turn at once + const closing = SpeechHandle.create({ allowInterruptions: true }); + await withAgentTurn(closing, opts, async () => hold.hold(closing)); + closing._markDone(); + hold.end(); + expect(spansNamed(exporter, 'agent_turn')).toHaveLength(3); + root.end(); + }); + it.each([ ['an Error', new Error('provider unavailable')], ['a string', 'provider unavailable'], diff --git a/agents/src/voice/speech_handle.ts b/agents/src/voice/speech_handle.ts index 2179e04824..8e2184e049 100644 --- a/agents/src/voice/speech_handle.ts +++ b/agents/src/voice/speech_handle.ts @@ -111,6 +111,24 @@ export interface AgentTurnCarry { /** How a speech came to continue another's `agent_turn` (see `SpeechHandle._continueAgentTurn`). */ export type AgentTurnContinuation = 'preemptive_discarded' | 'tool_reply'; +/** + * End a turn taken off its speech (see `SpeechHandle._takeAgentTurn`) that no successor adopted. + * The duration metric is recorded from the turn's start, as when the speech ends its own. + * @internal + */ +export function endCarriedAgentTurn(carry: AgentTurnCarry, error?: unknown): void { + // the duration metric does not depend on the span being sampled in + if (carry.startedAt !== undefined && carry.agentName !== undefined) { + recordInvokeAgentDuration((performance.now() - carry.startedAt) / 1000, carry.agentName); + } + if (!carry.span.isRecording()) return; + if (error !== undefined) { + // a thrown string or object is a failure too: the turn must not read as a success + recordException(carry.span, error instanceof Error ? error : new Error(String(error))); + } + carry.span.end(); +} + export class SpeechHandleCircularWaitError extends Error { constructor(functionCallName: string) { super(dedent` @@ -661,7 +679,9 @@ export class SpeechHandle { span.setAttribute(traceTypes.ATTR_SPEECH_ID, this.id); } else { this._generationBaseId = from._generationBaseId ?? from.id; - this._generationStepBase = from._emittedGenerationStep ?? from._generationStep; + // after the last generation `from` emitted; from where its own numbering started when it + // emitted none (it then never ran), so the next id follows a generation that exists + this._generationStepBase = from._emittedGenerationStep ?? from._generationStepBase; } this._agentTurnSpan = span; this._agentTurnContext = trace.setSpan(otelContext.active(), span); @@ -675,19 +695,15 @@ export class SpeechHandle { const span = this._agentTurnSpan; this._agentTurnSpan = undefined; if (span === undefined) return; - // the duration metric does not depend on the span being sampled in - if (this._agentTurnStartedAt !== undefined && this._agentTurnAgentName !== undefined) { - recordInvokeAgentDuration( - (performance.now() - this._agentTurnStartedAt) / 1000, - this._agentTurnAgentName, - ); - } - if (!span.isRecording()) return; - if (error !== undefined) { - // a thrown string or object is a failure too: the turn must not read as a success - recordException(span, error instanceof Error ? error : new Error(String(error))); - } - span.end(); + endCarriedAgentTurn( + { + span, + startedAt: this._agentTurnStartedAt, + agentName: this._agentTurnAgentName, + generations: this._agentTurnGenerations, + }, + error, + ); } /** @internal */