From baddc768552defa4cead3d7fc683f3fb80cdb5b3 Mon Sep 17 00:00:00 2001 From: David Zhao Date: Mon, 14 Sep 2026 23:05:35 -0700 Subject: [PATCH 1/5] feat(telemetry): interruption detail, handoff span, fallback events, text input Port of livekit/agents#7137. Interruptions: `agent_turn` carries `lk.interruption.source`, set by the caller that knows the cause (`audio_activity` for a VAD/STT barge-in and the realtime server's own speech detection, `user_turn` for a committed turn or a final transcript ending a pause, `programmatic` for session.interrupt(), tools and teardown); the first interruption's cause stands. The pipeline and `say` paths stamp `lk.playout.position`, the seconds actually played when the user cut in. The say and realtime paths now also stamp `lk.interrupted`, which only the pipeline reply did before. Agent handoff: `updateAgent()` opens an `update_agent` span under `agent_session` with `lk.previous_agent_label` / `lk.agent_label`; the old agent's `drain_agent_activity` (with `on_exit`) and the new agent's `start_agent_activity` / `resume_agent_activity` nest under it. The initial start stays under `session_start`. Fallback adapters: LLM, TTS and STT `model` / `provider` follow the instance that serves next (first available, else the primary), so `llm_node`, `tts_node` and `start_agent_activity` name a real model. The LLM and TTS attempt span carries `lk.fallback.label` / `lk.fallback.index` and the serving instance's request model and provider; the adapter's request span and the caller's node span get `gen_ai.response.model` / provider of the instance that answered, per request. The STT adapter's last-served tracking is replaced by the same next-instance rule. Text input: the keyterm-detection LLM pass runs in its own `keyterm_detection` span under the `agent_turn` whose reply added the user message, else under `agent_session`, with counts only (`lk.keyterms.count/added/removed`) plus model and provider. Adaptations: agents carry `id` where Python has `label`; JS `SpeechHandle` and `AgentActivity.interrupt` take the source as a second positional parameter / option instead of a keyword. The LLM and TTS stream base classes expose their request span to subclasses (`llmRequestSpan`, `ttsRequestSpan`), the JS counterpart of Python's `_llm_request_span`. A cancelled preemptive attempt names why it was dropped. SpeechHandle._cancel takes the cause like interrupt does: an attempt superseded by more of the user's turn (a later preemptive trigger, or the transcript changing at commit) is user_turn, one dropped by a barge-in audio_activity through interrupt(); a cancel with no cause (teardown, a pause) still reads as programmatic. A cloud export showed such attempts, cancelled when the user kept talking, labelled programmatic. Co-Authored-By: Claude Fable 5.1 --- .changeset/coverage-spans.md | 5 + agents/src/llm/fallback_adapter.test.ts | 33 + agents/src/llm/fallback_adapter.ts | 75 +- agents/src/llm/llm.ts | 15 +- agents/src/stt/fallback_adapter.test.ts | 34 +- agents/src/stt/fallback_adapter.ts | 40 +- agents/src/telemetry/trace_types.test.ts | 9 + agents/src/telemetry/trace_types.ts | 24 + agents/src/tts/fallback_adapter.test.ts | 31 + agents/src/tts/fallback_adapter.ts | 67 ++ agents/src/tts/tts.ts | 10 + agents/src/voice/agent_activity.ts | 128 +++- agents/src/voice/agent_session.ts | 39 +- .../src/voice/agent_session_handoff.test.ts | 21 +- agents/src/voice/coverage_spans.test.ts | 717 ++++++++++++++++++ agents/src/voice/index.ts | 1 + agents/src/voice/keyterm_detection.ts | 114 ++- .../src/voice/keyterm_detection_span.test.ts | 170 +++++ agents/src/voice/speech_handle.ts | 26 +- 19 files changed, 1445 insertions(+), 114 deletions(-) create mode 100644 .changeset/coverage-spans.md create mode 100644 agents/src/voice/coverage_spans.test.ts create mode 100644 agents/src/voice/keyterm_detection_span.test.ts diff --git a/.changeset/coverage-spans.md b/.changeset/coverage-spans.md new file mode 100644 index 0000000000..9894783485 --- /dev/null +++ b/.changeset/coverage-spans.md @@ -0,0 +1,5 @@ +--- +'@livekit/agents': patch +--- + +Telemetry coverage: `lk.interruption.source` and `lk.playout.position` on `agent_turn`, an `update_agent` span grouping agent handoffs, fallback adapters reporting the instance that serves (`lk.fallback.label`/`index`, response model on the request and node spans), and a `keyterm_detection` span for the keyterm-detection LLM pass. diff --git a/agents/src/llm/fallback_adapter.test.ts b/agents/src/llm/fallback_adapter.test.ts index 7ab549ff4f..d583b08531 100644 --- a/agents/src/llm/fallback_adapter.test.ts +++ b/agents/src/llm/fallback_adapter.test.ts @@ -264,4 +264,37 @@ describe('FallbackAdapter', () => { }), ); }); + + it('reports the model and provider of the instance that serves next', () => { + class IdentifiedLLM extends MockLLM { + constructor( + label: string, + private readonly _model: string, + private readonly _provider: string, + ) { + super(label); + } + override get model(): string { + return this._model; + } + override get provider(): string { + return this._provider; + } + } + const primary = new IdentifiedLLM('primary', 'primary-model', 'primary'); + const fallback = new IdentifiedLLM('fallback', 'fallback-model', 'fallback'); + const adapter = new FallbackAdapter({ llms: [primary, fallback] }); + // model and provider follow the instance that serves next, so spans and metrics name the + // model that will answer rather than the adapter; the label stays the adapter's own + expect(adapter.model).toBe('primary-model'); + expect(adapter.provider).toBe('primary'); + expect(adapter.label()).toContain('FallbackAdapter'); + adapter._status[0]!.available = false; + expect(adapter.model).toBe('fallback-model'); + expect(adapter.provider).toBe('fallback'); + // once the primary recovers (its recovery task flips it back to available) the next request + // goes to it first, so that is what model and provider report + adapter._status[0]!.available = true; + expect(adapter.model).toBe('primary-model'); + }); }); diff --git a/agents/src/llm/fallback_adapter.ts b/agents/src/llm/fallback_adapter.ts index 19ca1fa8d3..f91719205e 100644 --- a/agents/src/llm/fallback_adapter.ts +++ b/agents/src/llm/fallback_adapter.ts @@ -2,8 +2,10 @@ // // SPDX-License-Identifier: Apache-2.0 import type { Throws } from '@livekit/throws-transformer/throws'; +import { type Attributes, type Span, trace } from '@opentelemetry/api'; import { APIConnectionError, APIError } from '../_exceptions.js'; import { log } from '../log.js'; +import * as traceTypes from '../telemetry/trace_types.js'; import { type APIConnectOptions, DEFAULT_API_CONNECT_OPTIONS } from '../types.js'; import type { ChatContext } from './chat_context.js'; import type { ChatChunk } from './llm.js'; @@ -103,8 +105,29 @@ export class FallbackAdapter extends LLM { } } - get model(): string { - return 'FallbackAdapter'; + /** + * The instance the next request goes to first: the first one marked available, or the primary + * once all are down (they are then all retried, primary first). A failed instance's recovery + * task flips it back to available, so a recovered primary is reported again before it has + * served. + */ + private nextInstance(): LLM { + const index = this._status.findIndex((status) => status.available); + return this.llms[index === -1 ? 0 : index]!; + } + + /** + * The model of the instance that serves next (see `nextInstance`). Spans and metrics read + * this, so a failover shows the model expected to answer rather than the adapter; the instance + * that actually served is stamped per request by the stream. + */ + override get model(): string { + return this.nextInstance().model; + } + + /** The provider of the instance that serves next (see {@link model}). */ + override get provider(): string { + return this.nextInstance().provider; } label(): string { @@ -147,6 +170,21 @@ export class FallbackAdapter extends LLM { * LLMStream implementation for FallbackAdapter. * Handles fallback logic between multiple LLM providers. */ +function providerAttr(llm: LLM): Attributes { + const normalized = traceTypes.genAIProviderName(llm.provider); + return normalized ? { [traceTypes.ATTR_GEN_AI_PROVIDER_NAME]: normalized } : {}; +} + +/** The instance that served: its label, position, model and provider. */ +function fallbackAttrs(llm: LLM, index: number): Attributes { + return { + [traceTypes.ATTR_FALLBACK_LABEL]: llm.label(), + [traceTypes.ATTR_FALLBACK_INDEX]: index, + [traceTypes.ATTR_GEN_AI_REQUEST_MODEL]: llm.model, + ...providerAttr(llm), + }; +} + class FallbackLLMStream extends LLMStream { private adapter: FallbackAdapter; private parallelToolCalls?: boolean; @@ -154,6 +192,10 @@ class FallbackLLMStream extends LLMStream { private extraKwargs?: Record; private _currentStream?: LLMStream; private _log = log(); + // the span this request was made under (llm_node): told which instance served + private callerSpan: Span | undefined = trace.getActiveSpan(); + /** The instance whose output reached the caller (see recordServed). */ + private servedLlm?: LLM; constructor( adapter: FallbackAdapter, @@ -184,6 +226,32 @@ class FallbackLLMStream extends LLMStream { return this._currentStream?.chatCtx ?? super.chatCtx; } + /** + * The request span's response side names the instance that served, not the one the adapter + * would pick next: after a partial response the failed instance is already marked unavailable, + * so the adapter's own model would name an instance the caller never heard from. + */ + protected override get responseModel(): string { + return this.servedLlm?.model ?? this.adapter.model; + } + + /** + * The instance that served: on the current (attempt) span, and as the response side of the + * adapter's request span and the caller's (llm_node). Request-side attributes named the + * instance expected to serve; the response side names the one that did, read from `llm` + * rather than the adapter since concurrent requests may be served by different instances. + */ + private recordServed(llm: LLM, index: number): void { + this.servedLlm = llm; + trace.getActiveSpan()?.setAttributes(fallbackAttrs(llm, index)); + const responseAttrs: Attributes = { + [traceTypes.ATTR_GEN_AI_RESPONSE_MODEL]: llm.model, + ...providerAttr(llm), + }; + this.llmRequestSpan?.setAttributes(responseAttrs); + this.callerSpan?.setAttributes(responseAttrs); + } + /** * Try to generate with a single LLM. * Returns an async generator that yields chunks. @@ -352,6 +420,7 @@ class FallbackLLMStream extends LLMStream { { llm: llm.label(), totalChunks: chunkCount, textLength: textSent.length }, 'FallbackAdapter: Provider succeeded', ); + this.recordServed(llm, i); return; } catch (error) { // Mark as unavailable if it was available before @@ -374,6 +443,8 @@ class FallbackLLMStream extends LLMStream { { llm: llm.label(), ...extra }, 'failed after sending chunk, skip retrying. Set `retryOnChunkSent` to `true` to enable.', ); + // the caller received this instance's chunks: it served, partially + this.recordServed(llm, i); throw error; } diff --git a/agents/src/llm/llm.ts b/agents/src/llm/llm.ts index eb965a6819..a8260ed615 100644 --- a/agents/src/llm/llm.ts +++ b/agents/src/llm/llm.ts @@ -247,6 +247,19 @@ export abstract class LLMStream implements AsyncIterableIterator { }); } + /** The `llm_request` span of this stream, once the main task has opened it. */ + /** + * The model named on the response side of the request span. The LLM's own model by default; + * a fallback adapter's stream reports the instance that actually served. + */ + protected get responseModel(): string { + return this.#llm.model; + } + + protected get llmRequestSpan(): Span | undefined { + return this.#llmRequestSpan; + } + /** The GenAI inference span's request side, per the OTel GenAI conventions. */ private recordGenAIRequest(span: Span) { genAI.setRequestAttributes(span, { @@ -418,7 +431,7 @@ export abstract class LLMStream implements AsyncIterableIterator { }); genAI.setResponseAttributes(this.#llmRequestSpan, { responseId: requestId || undefined, - model: this.#llm.model, + model: this.responseModel, finishReasons: [finishReason], timeToFirstChunk: metrics.ttftMs >= 0 ? metrics.ttftMs / 1000 : undefined, }); diff --git a/agents/src/stt/fallback_adapter.test.ts b/agents/src/stt/fallback_adapter.test.ts index cc2fdd6c2e..e2c6845bcd 100644 --- a/agents/src/stt/fallback_adapter.test.ts +++ b/agents/src/stt/fallback_adapter.test.ts @@ -574,11 +574,39 @@ describe('FallbackAdapter dynamic model/provider getters', () => { } } - it('returns wrapper defaults before any STT is active', () => { + it('reports the primary before any traffic', () => { + // model and provider follow the instance that serves next, so spans and metrics name the + // model that will answer rather than the adapter; the label stays the adapter's own const a = new IdentifiedFakeSTT({ label: 'a', model: 'a-model', provider: 'a-provider' }); const adapter = new FallbackAdapter({ sttInstances: [a] }); - expect(adapter.model).toBe('FallbackAdapter'); - expect(adapter.provider).toBe('livekit'); + expect(adapter.model).toBe('a-model'); + expect(adapter.provider).toBe('a-provider'); + expect(adapter.label).toContain('FallbackAdapter'); + }); + + it('follows availability: a recovered primary is reported again before it serves', () => { + const primary = new IdentifiedFakeSTT({ + label: 'primary', + model: 'primary-model', + provider: 'primary-provider', + }); + const fallback = new IdentifiedFakeSTT({ + label: 'fallback', + model: 'fallback-model', + provider: 'fallback-provider', + }); + const adapter = new FallbackAdapter({ sttInstances: [primary, fallback] }); + adapter.status[0]!.available = false; + expect(adapter.model).toBe('fallback-model'); + expect(adapter.provider).toBe('fallback-provider'); + // once the primary recovers (its recovery task flips it back to available) the next request + // goes to it first, so that is what model and provider report + adapter.status[0]!.available = true; + expect(adapter.model).toBe('primary-model'); + // all down: they are all retried, primary first + adapter.status[0]!.available = false; + adapter.status[1]!.available = false; + expect(adapter.model).toBe('primary-model'); }); it('reflects the active child after a successful recognize()', async () => { diff --git a/agents/src/stt/fallback_adapter.ts b/agents/src/stt/fallback_adapter.ts index 49bd2fcb4b..edade1edc5 100644 --- a/agents/src/stt/fallback_adapter.ts +++ b/agents/src/stt/fallback_adapter.ts @@ -102,14 +102,6 @@ export class FallbackAdapter extends STT { private _status: STTStatus[] = []; private _logger = log(); private _metricsForwarders = new Map void>(); - // Last child that produced output or returned a recognize result. Surfaced - // via the dynamic label/model/provider getters so OTel attributes like - // `gen_ai.request.model` on `user_turn` (refreshed on every STT event by - // {@link audio_recognition.refreshUserTurnSttAttributes}) reflect the - // actual provider used, not the wrapper. Not safe to share an adapter - // across concurrent sessions/streams — concurrent writers would corrupt - // attribution. - private _activeStt: STT | undefined; label = 'stt.FallbackAdapter'; @@ -169,22 +161,28 @@ export class FallbackAdapter extends STT { // that actually transcribed, not the static wrapper. `audio_recognition. // refreshUserTurnSttAttributes` re-reads these on every STT event, so a // mid-turn fallover surfaces the new child immediately. - override get model(): string { - return this._activeStt?.model ?? 'FallbackAdapter'; - } - - override get provider(): string { - return this._activeStt?.provider ?? 'livekit'; + /** + * The instance the next request goes to first: the first one marked available, or the primary + * once all are down (they are then all retried, primary first). A failed instance's recovery + * task flips it back to available, so a recovered primary is reported again before it has + * served. + */ + private nextInstance(): STT { + const index = this._status.findIndex((status) => status.available); + return this.sttInstances[index === -1 ? 0 : index]!; } /** - * Record the child that most recently produced output. Called by - * `_recognize()` on a successful return and by the streaming path on - * every event yielded by the elected child. - * @internal + * The model of the instance that serves next (see `nextInstance`). Spans and metrics read + * this, so a failover shows the model expected to answer rather than the adapter. */ - _setActiveStt(stt: STT): void { - this._activeStt = stt; + override get model(): string { + return this.nextInstance().model; + } + + /** The provider of the instance that serves next (see {@link model}). */ + override get provider(): string { + return this.nextInstance().provider; } /** @@ -291,7 +289,6 @@ export class FallbackAdapter extends STT { if (status.available || allFailed) { try { const result = await stt.recognize(frame, abortSignal); - this._setActiveStt(stt); return result; } catch (e) { this._logger.warn( @@ -600,7 +597,6 @@ class FallbackSpeechStream extends SpeechStream { if (this.abortSignal.aborted || this.queue.closed) { return; } - this.fallbackAdapter._setActiveStt(sttInstance); this.queue.put(ev); } } finally { diff --git a/agents/src/telemetry/trace_types.test.ts b/agents/src/telemetry/trace_types.test.ts index 5573535f77..7e7b8b73ff 100644 --- a/agents/src/telemetry/trace_types.test.ts +++ b/agents/src/telemetry/trace_types.test.ts @@ -174,6 +174,9 @@ const SAFE_KEYS = new Set([ 'lk.job.launch_latency', 'lk.job.entrypoint_latency', 'lk.job.dispatch_latency', + 'lk.keyterms.count', + 'lk.keyterms.added', + 'lk.keyterms.removed', 'lk.room.auto_subscribe', 'lk.room.e2ee', 'lk.room.remote_participant_count', @@ -253,6 +256,12 @@ const SAFE_KEYS = new Set([ 'lk.amd.delay', // Adaptive interruption 'lk.is_interruption', + // interruptions, handoff, fallback (enums, labels, sizes) + 'lk.interruption.source', + 'lk.playout.position', + 'lk.previous_agent_label', + 'lk.fallback.label', + 'lk.fallback.index', 'lk.interruption.probability', 'lk.interruption.total_duration', 'lk.interruption.prediction_duration', diff --git a/agents/src/telemetry/trace_types.ts b/agents/src/telemetry/trace_types.ts index 4469a3dc18..594056cfd1 100644 --- a/agents/src/telemetry/trace_types.ts +++ b/agents/src/telemetry/trace_types.ts @@ -70,6 +70,13 @@ export const ATTR_JOB_ENTRYPOINT_LATENCY = 'lk.job.entrypoint_latency'; /** Seconds from the availability request to the entrypoint running: the whole chain. */ export const ATTR_JOB_DISPATCH_LATENCY = 'lk.job.dispatch_latency'; +// keyterm detection (keyterm_detection span): counts only, the terms themselves are the +// customer's vocabulary and travel as lk.pii.keyterms in the session report +/** Keyterms in effect after the pass (static + confirmed). */ +export const ATTR_KEYTERMS_COUNT = 'lk.keyterms.count'; +export const ATTR_KEYTERMS_ADDED = 'lk.keyterms.added'; +export const ATTR_KEYTERMS_REMOVED = 'lk.keyterms.removed'; + // room connect / room io export const ATTR_ROOM_AUTO_SUBSCRIBE = 'lk.room.auto_subscribe'; export const ATTR_ROOM_E2EE = 'lk.room.e2ee'; @@ -191,6 +198,23 @@ export const ATTR_AMD_SPEECH_DURATION = 'lk.amd.speech_duration'; export const ATTR_AMD_DELAY = 'lk.amd.delay'; export const ATTR_AMD_TRANSCRIPT = 'lk.pii.amd.transcript'; +// Interruptions (agent_turn) +/** + * What interrupted the speech: `audio_activity` (barge-in), `user_turn` (a committed turn + * preempting the reply), or `programmatic` (session.interrupt(), a tool, teardown). + */ +export const ATTR_INTERRUPTION_SOURCE = 'lk.interruption.source'; +/** Seconds of audio that had actually played when the speech was interrupted. */ +export const ATTR_PLAYOUT_POSITION = 'lk.playout.position'; + +// Agent handoff (update_agent span) +export const ATTR_PREVIOUS_AGENT_LABEL = 'lk.previous_agent_label'; + +// Fallback adapters (the attempt span) +/** Label of the provider that served the request. */ +export const ATTR_FALLBACK_LABEL = 'lk.fallback.label'; +export const ATTR_FALLBACK_INDEX = 'lk.fallback.index'; + // Adaptive Interruption attributes export const ATTR_IS_INTERRUPTION = 'lk.is_interruption'; export const ATTR_INTERRUPTION_PROBABILITY = 'lk.interruption.probability'; diff --git a/agents/src/tts/fallback_adapter.test.ts b/agents/src/tts/fallback_adapter.test.ts index 94c7ad91ac..574017f763 100644 --- a/agents/src/tts/fallback_adapter.test.ts +++ b/agents/src/tts/fallback_adapter.test.ts @@ -822,4 +822,35 @@ describe('TTS FallbackAdapter', () => { expect(primary.listenerCount('metrics_collected')).toBe(0); }); }); + + it('reports the model and provider of the instance that serves next', () => { + class IdentifiedTTS extends MockTTS { + constructor( + label: string, + private readonly _model: string, + private readonly _provider: string, + ) { + super(label); + } + override get model(): string { + return this._model; + } + override get provider(): string { + return this._provider; + } + } + const primary = new IdentifiedTTS('primary', 'primary-model', 'primary'); + const fallback = new IdentifiedTTS('fallback', 'fallback-model', 'fallback'); + const adapter = new FallbackAdapter({ ttsInstances: [primary, fallback] }); + expect(adapter.model).toBe('primary-model'); + expect(adapter.provider).toBe('primary'); + expect(adapter.label).toContain('FallbackAdapter'); + adapter.status[0]!.available = false; + expect(adapter.model).toBe('fallback-model'); + expect(adapter.provider).toBe('fallback'); + // once the primary recovers (its recovery task flips it back to available) the next request + // goes to it first, so that is what model and provider report + adapter.status[0]!.available = true; + expect(adapter.model).toBe('primary-model'); + }); }); diff --git a/agents/src/tts/fallback_adapter.ts b/agents/src/tts/fallback_adapter.ts index bd693f52af..07952c2654 100644 --- a/agents/src/tts/fallback_adapter.ts +++ b/agents/src/tts/fallback_adapter.ts @@ -3,9 +3,11 @@ // SPDX-License-Identifier: Apache-2.0 import { AudioResampler } from '@livekit/rtc-node'; import { type Throws, ThrowsPromise } from '@livekit/throws-transformer/throws'; +import { type Attributes, type Span, trace } from '@opentelemetry/api'; import { APIConnectionError, APIError } from '../_exceptions.js'; import { log } from '../log.js'; import type { TTSMetrics } from '../metrics/base.js'; +import * as traceTypes from '../telemetry/trace_types.js'; import { basic } from '../tokenize/index.js'; import { type APIConnectOptions, DEFAULT_API_CONNECT_OPTIONS } from '../types.js'; import { Task, cancelAndWait } from '../utils.js'; @@ -94,6 +96,31 @@ export class FallbackAdapter extends TTS { label: string = `tts.FallbackAdapter`; + /** + * The instance the next request goes to first: the first one marked available, or the primary + * once all are down (they are then all retried, primary first). A failed instance's recovery + * task flips it back to available, so a recovered primary is reported again before it has + * served. + */ + private nextInstance(): TTS { + const index = this._status.findIndex((status) => status.available); + return this.ttsInstances[index === -1 ? 0 : index]!; + } + + /** + * The model of the instance that serves next (see `nextInstance`). Spans and metrics read + * this, so a failover shows the model expected to answer rather than the adapter; the instance + * that actually served is stamped per request by the stream. + */ + override get model(): string { + return this.nextInstance().model; + } + + /** The provider of the instance that serves next (see {@link model}). */ + override get provider(): string { + return this.nextInstance().provider; + } + constructor(opts: FallbackAdapterOptions) { if (!opts.ttsInstances || opts.ttsInstances.length < 1) { throw new Error('at least one TTS instance must be provided.'); @@ -322,10 +349,44 @@ export class FallbackAdapter extends TTS { } } +/** The instance that served: its label, position, model and provider. */ +function fallbackAttrs(tts: TTS, index: number): Attributes { + const attrs: Attributes = { + [traceTypes.ATTR_FALLBACK_LABEL]: tts.label, + [traceTypes.ATTR_FALLBACK_INDEX]: index, + [traceTypes.ATTR_GEN_AI_REQUEST_MODEL]: tts.model, + }; + const normalized = traceTypes.genAIProviderName(tts.provider); + if (normalized !== undefined) { + attrs[traceTypes.ATTR_GEN_AI_PROVIDER_NAME] = normalized; + } + return attrs; +} + +/** + * The instance that served: on the current (attempt) span, and as the response side of `spans` + * (the adapter's request span and the caller's, tts_node). From `tts`, not the adapter: + * concurrent requests may be served by different instances. + */ +function recordFallbackServed(tts: TTS, index: number, ...spans: (Span | undefined)[]): void { + const attrs = fallbackAttrs(tts, index); + trace.getActiveSpan()?.setAttributes(attrs); + const responseAttrs: Attributes = { [traceTypes.ATTR_GEN_AI_RESPONSE_MODEL]: tts.model }; + const provider = attrs[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]; + if (provider !== undefined) { + responseAttrs[traceTypes.ATTR_GEN_AI_PROVIDER_NAME] = provider; + } + for (const span of spans) { + span?.setAttributes(responseAttrs); + } +} + class FallbackChunkedStream extends ChunkedStream { private adapter: FallbackAdapter; private connOptions: APIConnectOptions; private _logger = log(); + // the span this request was made under (tts_node); see recordFallbackServed + private callerSpan: Span | undefined = trace.getActiveSpan(); label: string = 'tts.FallbackChunkedStream'; @@ -427,6 +488,7 @@ class FallbackChunkedStream extends ChunkedStream { } this._logger.debug({ tts: tts.label }, 'TTS synthesis succeeded'); + recordFallbackServed(tts, i, this.ttsRequestSpan, this.callerSpan); return; } catch (error) { if (sawRawAudio) { @@ -466,6 +528,8 @@ class FallbackSynthesizeStream extends SynthesizeStream { )[] = []; private audioPushed = false; private _logger = log(); + // the span this request was made under (tts_node); see recordFallbackServed + private callerSpan: Span | undefined = trace.getActiveSpan(); label: string = 'tts.FallbackSynthesizeStream'; @@ -656,6 +720,7 @@ class FallbackSynthesizeStream extends SynthesizeStream { this.queue.put(SynthesizeStream.END_OF_STREAM); this._logger.debug({ tts: originalTts.label }, 'TTS stream succeeded'); + recordFallbackServed(originalTts, i, this.ttsRequestSpan, this.callerSpan); await readInputLLMStream.catch(() => {}); return; } catch (error) { @@ -664,6 +729,8 @@ class FallbackSynthesizeStream extends SynthesizeStream { { tts: originalTts.label, error: providerError ?? error }, 'TTS failed after audio pushed, cannot fallback mid-utterance', ); + // the caller heard this instance's audio: it served, partially + recordFallbackServed(originalTts, i, this.ttsRequestSpan, this.callerSpan); throw error; } diff --git a/agents/src/tts/tts.ts b/agents/src/tts/tts.ts index 62530b0969..6126330897 100644 --- a/agents/src/tts/tts.ts +++ b/agents/src/tts/tts.ts @@ -413,6 +413,11 @@ export abstract class SynthesizeStream }); } + /** The `tts_request` span of this stream, once the main task has opened it. */ + protected get ttsRequestSpan(): Span | undefined { + return this.#ttsRequestSpan; + } + private _mainTaskImpl = async (span: Span) => { this.#ttsRequestSpan = span; span.setAttributes({ @@ -828,6 +833,11 @@ export abstract class ChunkedStream implements AsyncIterableIterator { this.#ttsRequestSpan = span; span.setAttributes({ diff --git a/agents/src/voice/agent_activity.ts b/agents/src/voice/agent_activity.ts index fa3ad6dc99..66099c511f 100644 --- a/agents/src/voice/agent_activity.ts +++ b/agents/src/voice/agent_activity.ts @@ -162,7 +162,12 @@ import { } from './generation.js'; import type { PlaybackFinishedEvent, TimedString } from './io.js'; import { releaseTts, retainTts } from './model_refs.js'; -import { type InputDetails, REPLY_TASK_CANCEL_TIMEOUT, SpeechHandle } from './speech_handle.js'; +import { + type InputDetails, + type InterruptionSource, + REPLY_TASK_CANCEL_TIMEOUT, + SpeechHandle, +} from './speech_handle.js'; import { ToolExecutor, cancelTaskTool, @@ -373,6 +378,20 @@ function recordQueueWait(speechHandle: SpeechHandle): void { ?.setAttribute(traceTypes.ATTR_SPEECH_QUEUE_WAIT, queueWait / 1000); } +/** + * Name what interrupted the speech on its agent_turn span. No-op while the speech plays on. + * Module-level: tests drive the reply tasks with stand-in activities. + */ +function recordInterruption(speechHandle: SpeechHandle): void { + if (!speechHandle.interrupted || speechHandle._agentTurnContext === undefined) { + return; + } + trace.getSpan(speechHandle._agentTurnContext)?.setAttributes({ + [traceTypes.ATTR_SPEECH_INTERRUPTED]: true, + [traceTypes.ATTR_INTERRUPTION_SOURCE]: speechHandle._interruptSource ?? 'programmatic', + }); +} + export class AgentActivity implements RecognitionHooks { agent: Agent; agentSession: AgentSession; @@ -641,13 +660,18 @@ export class AgentActivity implements RecognitionHooks { } } - async resume(options?: { reuseResources?: ReusableResources }): Promise { + async resume(options?: { + reuseResources?: ReusableResources; + /** Parent for `resume_agent_activity`: a handoff span; else the session. */ + traceContext?: Context; + }): Promise { const unlock = await this.lock.lock(); try { await this._startSession({ spanName: 'resume_agent_activity', runOnEnter: false, reuseResources: options?.reuseResources, + traceContext: options?.traceContext, }); } finally { unlock(); @@ -862,8 +886,10 @@ export class AgentActivity implements RecognitionHooks { endpointing: createEndpointing(this.endpointingOpts), userTurnLimit: this.agentSession.sessionOptions.turnHandling.userTurnLimit, rootSpanContext: this.agentSession.rootSpanContext, - sttModel: this.stt?.label, - sttProvider: this.getSttProvider(), + // the model and provider, as python passes them (a fallback adapter reports the instance + // expected to serve next); the label is a class name, not a model + sttModel: this.stt?.model, + sttProvider: this.stt?.provider, sttAlignedTranscript: Boolean(this.stt?.capabilities.alignedTranscript), getLinkedParticipant: () => this.agentSession._roomIO?.linkedParticipant, shouldDiscardAudioForStt: () => this.shouldDiscardInputAudio(), @@ -1026,17 +1052,6 @@ export class AgentActivity implements RecognitionHooks { return this.agent._stt !== undefined ? this.agent._stt ?? undefined : this.agentSession.stt; } - private getSttProvider(): string | undefined { - const label = this.stt?.label; - if (!label) { - return undefined; - } - - // Heuristic: most labels look like "-" - const [provider] = label.split('-', 1); - return provider || label; - } - get llm(): LLM | RealtimeModel | undefined { return this.agent._llm !== undefined ? this.agent._llm ?? undefined : this.agentSession.llm; } @@ -1815,7 +1830,8 @@ export class AgentActivity implements RecognitionHooks { // this.interrupt() is going to raise when allow_interruptions is False, // llm.InputSpeechStartedEvent is only fired by the server when the turn_detection is enabled. try { - this.interrupt(); + // the server's own speech detection: a barge-in, like the VAD path + this.interrupt({ source: 'audio_activity' }); } catch (error) { this.logger.error( 'RealtimeAPI input_speech_started, but current speech is not interruptable, this should never happen!', @@ -2128,7 +2144,7 @@ export class AgentActivity implements RecognitionHooks { 'speech interrupted by audio activity', ); this.realtimeSession?.interrupt(); - this._currentSpeech.interrupt(); + this._currentSpeech.interrupt(false, 'audio_activity'); } } else if (this._currentSpeech === undefined || !this._currentSpeech.interrupted) { this.pendingInterruption = undefined; @@ -2266,7 +2282,8 @@ export class AgentActivity implements RecognitionHooks { return; } - this.cancelPreemptiveGeneration(); + // more of the user's turn arrived: the attempt answered a transcript that is now stale + this.cancelPreemptiveGeneration('user_turn'); if ( info.startedSpeakingAt !== undefined && @@ -2392,24 +2409,30 @@ export class AgentActivity implements RecognitionHooks { } } - private cancelPreemptiveGeneration(): void { + private cancelPreemptiveGeneration(source?: InterruptionSource): void { if (this._preemptiveGeneration !== undefined) { - this._preemptiveGeneration.speechHandle._cancel(); + this._preemptiveGeneration.speechHandle._cancel(source); this._preemptiveGeneration = undefined; } } - private _interruptBackgroundSpeeches(force: boolean): SpeechHandle[] { + private _interruptBackgroundSpeeches( + force: boolean, + source: InterruptionSource = 'programmatic', + ): SpeechHandle[] { const interrupted: SpeechHandle[] = []; for (const speech of this._backgroundSpeeches) { if (force || speech.allowInterruptions) { - interrupted.push(speech.interrupt(force)); + interrupted.push(speech.interrupt(force, source)); } } return interrupted; } - private interruptQueuedSpeeches(force: boolean): void { + private interruptQueuedSpeeches( + force: boolean, + source: InterruptionSource = 'programmatic', + ): void { // Heap iteration pops in playout order. Walk a clone so retained speeches stay queued and // interrupted ones remain for mainTask to drain. for (const [, , speech] of this.speechQueue.clone()) { @@ -2423,7 +2446,7 @@ export class AgentActivity implements RecognitionHooks { break; } - speech.interrupt(force); + speech.interrupt(force, source); } } @@ -2972,25 +2995,25 @@ export class AgentActivity implements RecognitionHooks { * Interrupt the current speech generation and any queued speeches. * * A queued speech that disallows interruptions keeps playing, along with the ones behind it, - * unless `force` is set. + * unless `force` is set. `source` names the cause on the speeches' `agent_turn` spans. * * @returns A future that completes when the interruption is fully processed. * @throws Error if the speech currently playing disallows interruptions and `force` is false. */ - interrupt(options: { force?: boolean } = {}): Future { - const { force = false } = options; - this.cancelPreemptiveGeneration(); + interrupt(options: { force?: boolean; source?: InterruptionSource } = {}): Future { + const { force = false, source = 'programmatic' } = options; + this.cancelPreemptiveGeneration(source); const future = new Future(); const currentSpeech = this._currentSpeech; - this._interruptBackgroundSpeeches(force); + this._interruptBackgroundSpeeches(force, source); - currentSpeech?.interrupt(force); + currentSpeech?.interrupt(force, source); this.realtimeSession?.interrupt(); - this.interruptQueuedSpeeches(force); + this.interruptQueuedSpeeches(force, source); if (force) { // Force-interrupt (used during shutdown): cancel all speech tasks so they @@ -3126,7 +3149,8 @@ export class AgentActivity implements RecognitionHooks { return; } - this.interruptQueuedSpeeches(false); + // the committed turn is why the queued replies go too: same cause as the current one below + this.interruptQueuedSpeeches(false, 'user_turn'); if (currentSpeech) { await this.cancelSpeechPause(); @@ -3138,7 +3162,7 @@ export class AgentActivity implements RecognitionHooks { 'speech interrupted, new user turn detected', ); - activeSpeech.interrupt(); + activeSpeech.interrupt(false, 'user_turn'); this.realtimeSession?.interrupt(); } @@ -3267,7 +3291,7 @@ export class AgentActivity implements RecognitionHooks { this.logger.warn( 'preemptive generation invalidated after `onUserTurnCompleted` because the transcript, chat context, tools, or tool choice changed', ); - preemptive.speechHandle._cancel(); + preemptive.speechHandle._cancel('user_turn'); } this._preemptiveGeneration = undefined; @@ -3327,6 +3351,7 @@ export class AgentActivity implements RecognitionHooks { recordQueueWait(speechHandle); if (speechHandle.interrupted) { + recordInterruption(speechHandle); return; } @@ -3427,8 +3452,11 @@ export class AgentActivity implements RecognitionHooks { try { await speechHandle.waitIfNotInterrupted(tasks.map((task) => task.result)); + let playbackEv: PlaybackFinishedEvent | undefined; if (audioOutput) { - await speechHandle.waitIfNotInterrupted([audioOutput.waitForPlayout()]); + const playout = audioOutput.waitForPlayout(); + await speechHandle.waitIfNotInterrupted([playout]); + if (!speechHandle.interrupted) playbackEv = await playout; } if (speechHandle.interrupted) { @@ -3436,7 +3464,20 @@ export class AgentActivity implements RecognitionHooks { await cancelAndWait(tasks, REPLY_TASK_CANCEL_TIMEOUT); if (audioOutput) { audioOutput.clearBuffer(); - await audioOutput.waitForPlayout(); + playbackEv = await audioOutput.waitForPlayout(); + } + } + + recordInterruption(speechHandle); + if (speechHandle.interrupted && audioOutput && audioOut && playbackEv) { + // waitForPlayout returns the previous segment's event when this speech never reached + // the output: only a segment of its own counts as played + const playedOwnFrame = + audioOutput.capturedPlayoutSegments > audioOut.capturedSegmentsBefore; + if (playedOwnFrame) { + trace + .getSpan(speechHandle._agentTurnContext ?? otelContext.active()) + ?.setAttribute(traceTypes.ATTR_PLAYOUT_POSITION, playbackEv.playbackPosition); } } @@ -3711,6 +3752,7 @@ export class AgentActivity implements RecognitionHooks { } if (speechHandle.interrupted) { + recordInterruption(speechHandle); replyAbortController.abort(); await cancelAndWait(tasks, REPLY_TASK_CANCEL_TIMEOUT); return; @@ -3967,6 +4009,13 @@ export class AgentActivity implements RecognitionHooks { } span.setAttribute(traceTypes.ATTR_SPEECH_INTERRUPTED, speechHandle.interrupted); + recordInterruption(speechHandle); + if (speechHandle.interrupted && segmentOutputs.length) { + span.setAttribute( + traceTypes.ATTR_PLAYOUT_POSITION, + segmentOutputs.reduce((total, out) => total + out.playbackPositionInS, 0), + ); + } let hasSpeechMessage = false; if (speechHandle.interrupted) { @@ -4300,6 +4349,7 @@ export class AgentActivity implements RecognitionHooks { recordQueueWait(speechHandle); if (speechHandle.interrupted) { + recordInterruption(speechHandle); return; } @@ -4627,6 +4677,7 @@ export class AgentActivity implements RecognitionHooks { await speechHandle.waitIfNotInterrupted(tasks.map((task) => task.result)); if (speechHandle.interrupted) { + recordInterruption(speechHandle); this.logger.debug( { speech_id: speechHandle.id }, 'Aborting all realtime generation tasks due to interruption', @@ -4686,6 +4737,8 @@ export class AgentActivity implements RecognitionHooks { } finally { this._backgroundSpeeches.delete(speechHandle); } + // the tools may have run past a barge-in: the turn names what cut it short + recordInterruption(speechHandle); if (toolOutput.output.length > 0) { if (this.updateAgentState(stateLease, 'thinking') && !endedAgentSpeechBeforeTool) { @@ -5711,7 +5764,8 @@ export class AgentActivity implements RecognitionHooks { !this.pausedSpeech.handle.interrupted && this.pausedSpeech.handle.allowInterruptions ) { - this.pausedSpeech.handle.interrupt(); + // a final transcript or a committed turn ended the pause + this.pausedSpeech.handle.interrupt(false, 'user_turn'); // ensure the generation is done — but only if a generation // was actually started. Must be raced against interrupt: an interrupted // paused speech may never mark its generation done, and an un-raced diff --git a/agents/src/voice/agent_session.ts b/agents/src/voice/agent_session.ts index 33b0d2ff41..0f45c2c20a 100644 --- a/agents/src/voice/agent_session.ts +++ b/agents/src/voice/agent_session.ts @@ -1526,6 +1526,8 @@ export class AgentSession< const unlock = await this.activityLock.lock(); let onEnterTask: Task | undefined; let reusableResources: ReusableResources | undefined; + let handoffSpan: Span | undefined; + let handoffCtx: Context | undefined; try { if (this.closing && newActivity === 'start') { @@ -1553,16 +1555,36 @@ export class AgentSession< this.nextActivity = agent._agentActivity; } + // one span for the handoff: the old agent's drain/pause and on_exit, then the new one's + // start/resume and on_enter nest under it. Passed explicitly to the calls that spawn + // long-lived tasks, made current only around the ones that do not + if (prevActivityObj && this.nextActivity) { + handoffSpan = tracer.startSpan({ + name: 'update_agent', + context: this.rootSpanContext, + attributes: { + [traceTypes.ATTR_PREVIOUS_AGENT_LABEL]: prevActivityObj.agent.id, + [traceTypes.ATTR_AGENT_LABEL]: this.nextActivity.agent.id, + }, + }); + handoffCtx = trace.setSpan(this.rootSpanContext ?? otelContext.active(), handoffSpan); + } + if (prevActivityObj && prevActivityObj !== this.nextActivity) { if (previousActivity === 'pause') { - reusableResources = await prevActivityObj.pause({ - blockedTasks, - newActivity: this.nextActivity, - }); + const pause = () => + prevActivityObj.pause({ + blockedTasks, + newActivity: this.nextActivity, + }); + reusableResources = handoffCtx + ? await otelContext.with(handoffCtx, pause) + : await pause(); } else { prevActivityObj.blockNewTurns(); reusableResources = await prevActivityObj.drain({ newActivity: this.nextActivity, + traceContext: handoffCtx, }); await prevActivityObj.close(); } @@ -1609,10 +1631,11 @@ export class AgentSession< if (newActivity === 'start') { await activity.start({ reuseResources: reusableResources, - traceContext: options.traceContext, + // the initial start is not a handoff: it lives under session_start + traceContext: handoffCtx ?? options.traceContext, }); } else { - await activity.resume({ reuseResources: reusableResources }); + await activity.resume({ reuseResources: reusableResources, traceContext: handoffCtx }); } reusableResources = undefined; @@ -1622,6 +1645,9 @@ export class AgentSession< activity.attachAudioInput(this._input.audio.stream); } } catch (error) { + if (handoffSpan && error instanceof Error) { + recordException(handoffSpan, error); + } // JS safeguard: session cleanup owns the detached resources until the next activity // starts successfully, preventing leaks when handoff fails mid-transition. if (reusableResources) { @@ -1629,6 +1655,7 @@ export class AgentSession< } throw error; } finally { + handoffSpan?.end(); unlock(); } diff --git a/agents/src/voice/agent_session_handoff.test.ts b/agents/src/voice/agent_session_handoff.test.ts index a0e4e2ec85..8215294312 100644 --- a/agents/src/voice/agent_session_handoff.test.ts +++ b/agents/src/voice/agent_session_handoff.test.ts @@ -83,8 +83,15 @@ describe('AgentSession reusable resources handoff', () => { waitOnEnter: false, }); - expect(previousActivity.drain).toHaveBeenCalledWith({ newActivity: nextActivity }); - expect(nextActivity.resume).toHaveBeenCalledWith({ reuseResources: resources }); + // the handoff runs under an update_agent span: drain and resume are parented to it + expect(previousActivity.drain).toHaveBeenCalledWith({ + newActivity: nextActivity, + traceContext: expect.anything(), + }); + expect(nextActivity.resume).toHaveBeenCalledWith({ + reuseResources: resources, + traceContext: expect.anything(), + }); }); it('cleans up reusable resources if the next activity fails to start', async () => { @@ -161,7 +168,10 @@ describe('AgentSession reusable resources handoff', () => { }), ).rejects.toThrow('attach failed'); - expect(nextActivity.resume).toHaveBeenCalledWith({ reuseResources: resources }); + expect(nextActivity.resume).toHaveBeenCalledWith({ + reuseResources: resources, + traceContext: expect.anything(), + }); // pipeline was already transferred, so cleanup should NOT have been called expect(closeFn).not.toHaveBeenCalled(); }); @@ -190,7 +200,10 @@ describe('AgentSession reusable resources handoff', () => { expect(activity.drain).not.toHaveBeenCalled(); expect(activity.pause).not.toHaveBeenCalled(); - expect(activity.resume).toHaveBeenCalledWith({ reuseResources: undefined }); + expect(activity.resume).toHaveBeenCalledWith({ + reuseResources: undefined, + traceContext: expect.anything(), + }); }); it('emits ConversationItemAdded with an AgentHandoffItem on handoff', async () => { diff --git a/agents/src/voice/coverage_spans.test.ts b/agents/src/voice/coverage_spans.test.ts new file mode 100644 index 0000000000..13aca82ddd --- /dev/null +++ b/agents/src/voice/coverage_spans.test.ts @@ -0,0 +1,717 @@ +// SPDX-FileCopyrightText: 2026 LiveKit, Inc. +// +// SPDX-License-Identifier: Apache-2.0 + +/** + * Coverage additions on existing spans: interruption detail on `agent_turn`, the `update_agent` + * handoff span, and fallback-adapter attribution on the request span. + */ +import { AudioFrame } from '@livekit/rtc-node'; +import { 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, type ReadableStreamDefaultController } from 'node:stream/web'; +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { APIConnectionError } from '../_exceptions.js'; +import { ChatContext } from '../llm/chat_context.js'; +import { FallbackAdapter } from '../llm/fallback_adapter.js'; +import { LLM, LLMStream } from '../llm/llm.js'; +import type { ToolChoice, ToolContextLike } 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 { type APIConnectOptions, DEFAULT_API_CONNECT_OPTIONS } from '../types.js'; +import { delay } from '../utils.js'; +import { VAD, type VADEvent, VADEventType, VADStream } from '../vad.js'; +import { Agent } from './agent.js'; +import { AgentSession } from './agent_session.js'; +import { AudioInput, 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 only(exporter: InMemorySpanExporter, name: string): ReadableSpan { + const found = spansNamed(exporter, name); + expect(found, `expected exactly one ${name} span, got ${found.length}`).toHaveLength(1); + return found[0]!; +} + +function childrenOf(exporter: InMemorySpanExporter, name: string, parent: ReadableSpan) { + return spansNamed(exporter, name).filter( + (span) => span.parentSpanContext?.spanId === parent.spanContext().spanId, + ); +} + +function vadEvent(type: VADEventType, speechDuration = 0): VADEvent { + return { + type, + samplesIndex: 0, + timestamp: Date.now(), + speechDuration, + silenceDuration: 0, + frames: [], + probability: 1, + inferenceDuration: 0, + speaking: type !== VADEventType.END_OF_SPEECH, + rawAccumulatedSilence: 0, + rawAccumulatedSpeech: speechDuration, + }; +} + +class ScriptedVADStream extends VADStream { + constructor(vad: VAD) { + super(vad); + void this.#drain(); + } + + async #drain(): Promise { + try { + while (!this.closed) { + const { done } = await this.inputReader.read(); + if (done) break; + } + } catch { + /* stream detached/closed */ + } + } + + emitEvent(ev: VADEvent): void { + this.sendVADEvent(ev); + } +} + +class ScriptedVAD extends VAD { + label = 'scripted-vad'; + readonly streams: ScriptedVADStream[] = []; + + constructor() { + super({ updateInterval: 32 }); + } + + stream(): ScriptedVADStream { + const stream = new ScriptedVADStream(this); + this.streams.push(stream); + return stream; + } + + startOfSpeech(): void { + this.streams.at(-1)!.emitEvent(vadEvent(VADEventType.START_OF_SPEECH)); + } + + /** The VAD has seen `speechDuration` ms of speech: what the barge-in check reads. */ + inferenceDone(speechDuration: number): void { + this.streams.at(-1)!.emitEvent(vadEvent(VADEventType.INFERENCE_DONE, speechDuration)); + } + + endOfSpeech(): void { + this.streams.at(-1)!.emitEvent(vadEvent(VADEventType.END_OF_SPEECH)); + } +} + +class ScriptedAudioInput extends AudioInput { + #controller!: ReadableStreamDefaultController; + + constructor() { + super(); + this.multiStream.addInputStream( + new ReadableStream({ + start: (controller) => { + this.#controller = controller; + }, + }), + ); + } + + push(durationMs: number, sampleRate = 16_000): void { + const samples = Math.floor((sampleRate * durationMs) / 1000); + this.#controller.enqueue(new AudioFrame(new Int16Array(samples), sampleRate, 1, samples)); + } +} + +/** + * Plays frames as they arrive and reports how far it got: the position of an interrupted segment + * is the wall-clock time since its first frame, as a real output would report. + */ +class PacedOutput extends AudioOutput { + #segmentStartedAt?: number; + #captured = 0; + + constructor() { + super(24_000); + } + + override async captureFrame(frame: AudioFrame): Promise { + const segmentCount = this.capturedPlayoutSegments; + await super.captureFrame(frame); + this.#captured += frame.samplesPerChannel / frame.sampleRate; + if (this.capturedPlayoutSegments > segmentCount) { + this.#segmentStartedAt = Date.now(); + this.#captured = frame.samplesPerChannel / frame.sampleRate; + this.onPlaybackStarted(Date.now()); + } + } + + override flush(): void { + super.flush(); + if (this.pendingPlayoutSegments > 0) { + this.onPlaybackFinished({ playbackPosition: this.#captured, interrupted: false }); + } + } + + override clearBuffer(): void { + if (this.pendingPlayoutSegments > 0) { + const played = + this.#segmentStartedAt !== undefined ? (Date.now() - this.#segmentStartedAt) / 1000 : 0; + this.onPlaybackFinished({ playbackPosition: played, interrupted: true }); + } + } +} + +/** An agent whose TTS streams `frames` 20 ms frames, one every `paceMs`: a long playout. */ +class PacedAgent extends Agent { + constructor( + private readonly frames: number, + private readonly paceMs: number, + ) { + super({ instructions: 'test' }); + } + + override async ttsNode(): Promise> { + const { frames, paceMs } = this; + return new ReadableStream({ + async start(controller) { + for (let i = 0; i < frames; i++) { + controller.enqueue(new AudioFrame(new Int16Array(480), 24_000, 1, 480)); + await delay(paceMs); + } + controller.close(); + }, + }); + } +} + +class FirstAgent extends PacedAgent { + constructor() { + super(1, 0); + } +} + +class SecondAgent extends PacedAgent { + constructor() { + super(1, 0); + } +} + +async function waitFor(predicate: () => boolean, timeoutMs: number, what: string) { + const deadline = Date.now() + timeoutMs; + while (!predicate()) { + if (Date.now() > deadline) throw new Error(`timed out waiting for ${what}`); + await delay(10); + } +} + +describe.sequential('coverage spans', () => { + 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 () => { + setTracerProvider(originalProvider); + await provider.shutdown(); + trace.disable(); + otelContext.disable(); + }); + + // -- interruption detail -- + + it('records the barge-in source and the playout position', async () => { + const vad = new ScriptedVAD(); + const stt = new FakeSTT({ + capabilities: { streaming: true, interimResults: true }, + fakeUserSpeeches: [ + { startTime: 0, endTime: 200, transcript: 'Tell me a story.', sttDelay: 100 }, + // the barge-in: its transcript lands well after the VAD cut the reply short + { startTime: 1_500, endTime: 1_700, transcript: 'Stop!', sttDelay: 100 }, + ], + }); + const llm = new FakeLLM([ + { input: 'Tell me a story.', content: 'Here is a long story for you ... the end.' }, + { input: 'Stop!', content: 'Ok.' }, + ]); + const session = new AgentSession({ + vad, + stt, + llm, + aecWarmupDuration: 0, + turnHandling: { + turnDetection: 'vad', + endpointing: { minDelay: 100, maxDelay: 100 }, + // no pause-and-resume: a barge-in interrupts the reply outright + interruption: { resumeFalseInterruption: false, minDuration: 500 }, + }, + }); + const audioInput = new ScriptedAudioInput(); + session.input.audio = audioInput; + session.output.audio = new PacedOutput(); + // 60 frames at 20 ms every 25 ms: ~1.5 s of playout for the story + await session.start({ agent: new PacedAgent(60, 25) }); + try { + audioInput.push(20); + await delay(20); + vad.startOfSpeech(); + await delay(200); + vad.endOfSpeech(); + await waitFor(() => session.agentState === 'speaking', 10_000, 'the story to start playing'); + await delay(300); + // the user talks over the story: VAD start, then enough speech for the barge-in check + vad.startOfSpeech(); + vad.inferenceDone(600); + await delay(200); + vad.endOfSpeech(); + await waitFor( + () => + spansNamed(exporter, 'agent_turn').filter( + (span) => traceTypes.ATTR_E2E_LATENCY in span.attributes, + ).length >= 2, + 10_000, + 'the reply to the barge-in to play', + ); + } finally { + await session.close(); + } + + const turns = spansNamed(exporter, 'agent_turn'); + const interrupted = turns.filter( + (span) => span.attributes[traceTypes.ATTR_SPEECH_INTERRUPTED] === true, + ); + expect(interrupted).toHaveLength(1); + const turn = interrupted[0]!; + expect(turn.attributes[traceTypes.ATTR_INTERRUPTION_SOURCE]).toBe('audio_activity'); + const position = turn.attributes[traceTypes.ATTR_PLAYOUT_POSITION]; + expect(position).toBeTypeOf('number'); + // ~0.3 s of the ~1.5 s story had played when the user cut in + expect(position as number).toBeGreaterThan(0.1); + expect(position as number).toBeLessThan(1.5); + + // the reply to "Stop!" was not interrupted and names no source + for (const other of turns) { + if (other === turn) continue; + expect(other.attributes[traceTypes.ATTR_INTERRUPTION_SOURCE]).toBeUndefined(); + } + }); + + it('a committed user turn interrupts the queued replies for the same reason', async () => { + // the turn interrupts the reply that is playing and the ones queued behind it: all of + // them name user_turn, not the programmatic default the queue sweep used to fall back to + const vad = new ScriptedVAD(); + const stt = new FakeSTT({ + capabilities: { streaming: true, interimResults: true }, + fakeUserSpeeches: [ + { startTime: 0, endTime: 200, transcript: 'Tell me a story.', sttDelay: 100 }, + { startTime: 1_500, endTime: 1_700, transcript: 'Stop!', sttDelay: 100 }, + ], + }); + const llm = new FakeLLM([ + { input: 'Tell me a story.', content: 'Here is a long story for you ... the end.' }, + { input: 'Stop!', content: 'Ok.' }, + ]); + const session = new AgentSession({ + vad, + stt, + llm, + // AEC warmup holds off barge-ins on audio activity (VAD, or a transcript arriving while + // the agent speaks) without making the speech uninterruptible: the committed turn is + // what interrupts, the case this test is about + aecWarmupDuration: 30_000, + turnHandling: { + turnDetection: 'vad', + endpointing: { minDelay: 100, maxDelay: 100 }, + interruption: { resumeFalseInterruption: false }, + }, + }); + const audioInput = new ScriptedAudioInput(); + session.input.audio = audioInput; + session.output.audio = new PacedOutput(); + // 120 frames at 20 ms every 25 ms: ~3 s of playout, still going when the turn commits + await session.start({ agent: new PacedAgent(120, 25) }); + let queued: SpeechHandle | undefined; + let interrupted: ReadableSpan[] = []; + try { + audioInput.push(20); + await delay(20); + vad.startOfSpeech(); + await delay(200); + vad.endOfSpeech(); + await waitFor(() => session.agentState === 'speaking', 10_000, 'the story to start playing'); + // say with audio: the paced agent has no TTS + queued = session.say('And one more thing.', { + audio: new ReadableStream({ + async start(controller) { + for (let i = 0; i < 20; i++) { + controller.enqueue(new AudioFrame(new Int16Array(480), 24_000, 1, 480)); + await delay(25); + } + controller.close(); + }, + }), + }); + await delay(200); + vad.startOfSpeech(); + await delay(200); + vad.endOfSpeech(); + await waitFor( + () => + spansNamed(exporter, 'agent_turn').some( + (span) => + span.attributes[traceTypes.ATTR_SPEECH_INTERRUPTED] === true && + traceTypes.ATTR_INTERRUPTION_SOURCE in span.attributes, + ) && queued!.interrupted, + 10_000, + 'the user turn to interrupt the story and the queued say', + ); + // before close(), which interrupts whatever is still playing programmatically + interrupted = spansNamed(exporter, 'agent_turn').filter( + (span) => span.attributes[traceTypes.ATTR_SPEECH_INTERRUPTED] === true, + ); + } finally { + await session.close(); + } + + expect(queued!.interrupted).toBe(true); + expect(queued!._interruptSource).toBe('user_turn'); + expect(interrupted.length).toBeGreaterThanOrEqual(1); + for (const turn of interrupted) { + expect(turn.attributes[traceTypes.ATTR_INTERRUPTION_SOURCE]).toBe('user_turn'); + } + }, 15_000); + + it('the first interruption cause stands', () => { + const handle = SpeechHandle.create(); + handle.interrupt(false, 'audio_activity'); + handle.interrupt(false, 'user_turn'); // already interrupted: the first cause stands + expect(handle._interruptSource).toBe('audio_activity'); + expect(SpeechHandle.create().interrupt()._interruptSource).toBe('programmatic'); + }); + + it('a cancelled preemptive attempt names why it was dropped', () => { + // superseded by more of the user's turn, or dropped by a barge-in: the cause travels with + // the cancel like it does with interrupt, and the first one stands + const superseded = SpeechHandle.create()._cancel('user_turn'); + expect(superseded._interruptSource).toBe('user_turn'); + const bargedIn = SpeechHandle.create()._cancel('audio_activity'); + bargedIn.interrupt(false, 'programmatic'); + expect(bargedIn._interruptSource).toBe('audio_activity'); + // a cancel with no cause (teardown, pause) stays programmatic in the trace + expect(SpeechHandle.create()._cancel()._interruptSource).toBeUndefined(); + }); + + // -- user_turn names the STT -- + + class NamedSTT extends FakeSTT { + override get model(): string { + return 'fake-model'; + } + + override get provider(): string { + return 'fake-provider'; + } + } + + it('names the STT model and provider on user_turn, as python does', async () => { + // the label is a class name (`stt.FallbackAdapter`, `fake-stt`): the span carries the + // model and provider the STT reports, which for a fallback adapter is the instance expected + // to serve next + const TRANSCRIPT = 'Hello'; + const vad = new ScriptedVAD(); + const stt = new NamedSTT({ + capabilities: { streaming: true, interimResults: true }, + fakeUserSpeeches: [{ startTime: 0, endTime: 200, transcript: TRANSCRIPT, sttDelay: 100 }], + }); + const llm = new FakeLLM([{ input: TRANSCRIPT, content: 'Hi there' }]); + const session = new AgentSession({ + vad, + stt, + llm, + turnHandling: { turnDetection: 'vad', endpointing: { minDelay: 100, maxDelay: 100 } }, + }); + const audioInput = new ScriptedAudioInput(); + session.input.audio = audioInput; + session.output.audio = new PacedOutput(); + await session.start({ agent: new Agent({ instructions: 'test' }) }); + try { + audioInput.push(20); + await delay(20); + vad.startOfSpeech(); + await delay(200); + vad.endOfSpeech(); + await waitFor( + () => + spansNamed(exporter, 'user_turn').some( + (span) => traceTypes.ATTR_USER_TRANSCRIPT in span.attributes, + ), + 10_000, + 'the user turn to end', + ); + } finally { + await session.close(); + } + + const turn = spansNamed(exporter, 'user_turn').find( + (span) => span.attributes[traceTypes.ATTR_USER_TRANSCRIPT] === TRANSCRIPT, + )!; + expect(turn.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).toBe('fake-model'); + expect(turn.attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]).toBe('fake-provider'); + expect(turn.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).not.toBe(stt.label); + }); + + // -- agent handoff -- + + it('groups the handoff under an update_agent span', async () => { + const TRANSCRIPT = 'Hello'; + const vad = new ScriptedVAD(); + const stt = new FakeSTT({ + capabilities: { streaming: true, interimResults: true }, + fakeUserSpeeches: [{ startTime: 0, endTime: 200, transcript: TRANSCRIPT, sttDelay: 100 }], + }); + const llm = new FakeLLM([{ input: TRANSCRIPT, content: 'Hi there' }]); + const session = new AgentSession({ + vad, + stt, + llm, + turnHandling: { turnDetection: 'vad', endpointing: { minDelay: 100, maxDelay: 100 } }, + }); + const audioInput = new ScriptedAudioInput(); + session.input.audio = audioInput; + session.output.audio = new PacedOutput(); + const first = new FirstAgent(); + const second = new SecondAgent(); + await session.start({ agent: first }); + try { + audioInput.push(20); + await delay(20); + vad.startOfSpeech(); + await delay(200); + vad.endOfSpeech(); + await waitFor( + () => + spansNamed(exporter, 'agent_turn').some( + (span) => traceTypes.ATTR_E2E_LATENCY in span.attributes, + ), + 10_000, + 'the reply to play', + ); + + session.updateAgent(second); + await waitFor( + () => + spansNamed(exporter, 'start_agent_activity').some( + (span) => span.attributes[traceTypes.ATTR_AGENT_LABEL] === second.id, + ), + 10_000, + 'the second agent to start', + ); + } finally { + await session.close(); + } + + const root = only(exporter, 'agent_session'); + const handoff = only(exporter, 'update_agent'); + expect(handoff.parentSpanContext?.spanId).toBe(root.spanContext().spanId); + expect(handoff.attributes[traceTypes.ATTR_PREVIOUS_AGENT_LABEL]).toBe(first.id); + expect(handoff.attributes[traceTypes.ATTR_AGENT_LABEL]).toBe(second.id); + + // the old agent's drain (with on_exit inside it), then the new agent's start, under the handoff + const [drain] = childrenOf(exporter, 'drain_agent_activity', handoff); + expect(drain).toBeDefined(); + expect(childrenOf(exporter, 'on_exit', drain!)).toHaveLength(1); + const [start] = childrenOf(exporter, 'start_agent_activity', handoff); + expect(start?.attributes[traceTypes.ATTR_AGENT_LABEL]).toBe(second.id); + + // the initial start is not a handoff: it lives under session_start, not update_agent + const sessionStart = only(exporter, 'session_start'); + expect(childrenOf(exporter, 'start_agent_activity', sessionStart)).toHaveLength(1); + }); + + // -- fallback adapter attribution -- + + class FailingLLMStream extends LLMStream { + protected async run(): Promise { + throw new APIConnectionError({ message: 'primary down' }); + } + } + + class FailingLLM extends LLM { + label(): string { + return 'broken-llm'; + } + + override get model(): string { + return 'broken-model'; + } + + override get provider(): string { + return 'broken'; + } + + chat(opts: { + chatCtx: ChatContext; + toolCtx?: ToolContextLike; + connOptions?: APIConnectOptions; + parallelToolCalls?: boolean; + toolChoice?: ToolChoice; + extraKwargs?: Record; + }): LLMStream { + return new FailingLLMStream(this, { + chatCtx: opts.chatCtx, + toolCtx: opts.toolCtx, + connOptions: opts.connOptions ?? DEFAULT_API_CONNECT_OPTIONS, + }); + } + } + + class ServingLLM extends FakeLLM { + override get model(): string { + return 'serving-model'; + } + + override get provider(): string { + return 'openai'; + } + } + + class PartialLLMStream extends LLMStream { + protected async run(): Promise { + this.queue.put({ id: 'partial', delta: { role: 'assistant', content: 'Hel' } }); + throw new APIConnectionError({ message: 'primary dropped mid-response' }); + } + } + + class PartialLLM extends LLM { + label(): string { + return 'partial-llm'; + } + + override get model(): string { + return 'partial-model'; + } + + override get provider(): string { + return 'partial'; + } + + chat(opts: { + chatCtx: ChatContext; + toolCtx?: ToolContextLike; + connOptions?: APIConnectOptions; + parallelToolCalls?: boolean; + toolChoice?: ToolChoice; + extraKwargs?: Record; + }): LLMStream { + return new PartialLLMStream(this, { + chatCtx: opts.chatCtx, + toolCtx: opts.toolCtx, + connOptions: opts.connOptions ?? DEFAULT_API_CONNECT_OPTIONS, + }); + } + } + + it('names the instance whose partial response the caller received before it failed', async () => { + // with retryOnChunkSent off the failure is rethrown after chunks reached the caller: the + // spans still say who produced them, as the TTS fallback does for audio already heard + const primary = new PartialLLM(); + const secondary = new ServingLLM([{ input: 'hi', content: 'hello' }]); + const adapter = new FallbackAdapter({ llms: [primary, secondary], attemptTimeout: 1 }); + const chatCtx = ChatContext.empty(); + chatCtx.addMessage({ role: 'user', content: 'hi' }); + let text = ''; + // the stream ends and the failure is reported on the adapter's error event + const failure = new Promise((resolve) => adapter.on('error', (e) => resolve(e.error))); + await tracer.startActiveSpan( + async () => { + const stream = adapter.chat({ chatCtx }); + for await (const chunk of stream) { + text += chunk.delta?.content ?? ''; + } + }, + { name: 'caller' }, + ); + expect(text).toBe('Hel'); + await expect(failure).resolves.toBeInstanceOf(APIConnectionError); + + const run = spansNamed(exporter, 'llm_request_run').find( + (span) => traceTypes.ATTR_FALLBACK_LABEL in span.attributes, + ); + expect(run?.attributes[traceTypes.ATTR_FALLBACK_LABEL]).toBe(primary.label()); + // the adapter's request span (the instance's own llm_request nests under the run) + const request = spansNamed(exporter, 'llm_request').find( + (span) => span.spanContext().spanId === run?.parentSpanContext?.spanId, + )!; + expect(request.attributes[traceTypes.ATTR_GEN_AI_RESPONSE_MODEL]).toBe(primary.model); + expect(request.attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]).toBe('partial'); + const caller = only(exporter, 'caller'); + expect(caller.attributes[traceTypes.ATTR_GEN_AI_RESPONSE_MODEL]).toBe(primary.model); + }); + + it('records the serving provider of an LLM fallback on the request and caller spans', async () => { + const primary = new FailingLLM(); + const secondary = new ServingLLM([{ input: 'hi', content: 'hello' }]); + const adapter = new FallbackAdapter({ llms: [primary, secondary], attemptTimeout: 1 }); + const chatCtx = ChatContext.empty(); + chatCtx.addMessage({ role: 'user', content: 'hi' }); + let text = ''; + // the span the request is made under (llm_node in the pipeline, open until the stream is + // consumed) is told who served + await tracer.startActiveSpan( + async () => { + const stream = adapter.chat({ chatCtx }); + for await (const chunk of stream) { + text += chunk.delta?.content ?? ''; + } + }, + { name: 'caller' }, + ); + expect(text).toBe('hello'); + + // the adapter's request span nests the attempt span that ran the fallback loop + const runs = spansNamed(exporter, 'llm_request_run').filter( + (span) => traceTypes.ATTR_FALLBACK_LABEL in span.attributes, + ); + expect(runs).toHaveLength(1); + const run = runs[0]!; + expect(run.attributes[traceTypes.ATTR_FALLBACK_LABEL]).toBe(secondary.label()); + expect(run.attributes[traceTypes.ATTR_FALLBACK_INDEX]).toBe(1); + // the run names the one that served; the request span kept the one expected to (the primary, + // still available when the request started) and gained the response side + expect(run.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).toBe(secondary.model); + const request = spansNamed(exporter, 'llm_request').find( + (span) => span.spanContext().spanId === run.parentSpanContext?.spanId, + ); + expect(request).toBeDefined(); + expect(request!.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).toBe(primary.model); + expect(request!.attributes[traceTypes.ATTR_GEN_AI_RESPONSE_MODEL]).toBe(secondary.model); + expect(request!.attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]).toBe('openai'); + // the caller's span gets the same response side, per request + const caller = only(exporter, 'caller'); + expect(caller.attributes[traceTypes.ATTR_GEN_AI_RESPONSE_MODEL]).toBe(secondary.model); + // and the adapter itself now reports who serves next + expect(adapter.model).toBe(secondary.model); + expect(adapter.provider).toBe(secondary.provider); + }); +}); diff --git a/agents/src/voice/index.ts b/agents/src/voice/index.ts index 94d2a04831..c084124a91 100644 --- a/agents/src/voice/index.ts +++ b/agents/src/voice/index.ts @@ -71,6 +71,7 @@ export { type InputDetails, type ResolvedSpeechHandle, } from './speech_handle.js'; +export type { InterruptionSource } from './speech_handle.js'; export * from './turn_config/endpointing.js'; export * from './turn_config/user_turn_limit.js'; export * as testing from './testing/index.js'; diff --git a/agents/src/voice/keyterm_detection.ts b/agents/src/voice/keyterm_detection.ts index 4d244e61b6..68f5f477c4 100644 --- a/agents/src/voice/keyterm_detection.ts +++ b/agents/src/voice/keyterm_detection.ts @@ -2,6 +2,7 @@ // // SPDX-License-Identifier: Apache-2.0 import type { TypedEventEmitter as TypedEmitter } from '@livekit/typed-emitter'; +import { type Attributes, type Context, context as otelContext, trace } from '@opentelemetry/api'; import { EventEmitter } from 'node:events'; import { z } from 'zod'; import { LLM as InferenceLLM } from '../inference/llm.js'; @@ -11,6 +12,7 @@ import { tool } from '../llm/tool_context.js'; import { log } from '../log.js'; import type { LLMMetrics } from '../metrics/base.js'; import type { STT } from '../stt/stt.js'; +import { traceTypes, tracer } from '../telemetry/index.js'; import { Task, delay } from '../utils.js'; import { AgentSessionEventTypes, type ConversationItemAddedEvent } from './events.js'; @@ -206,6 +208,8 @@ const recordKeyterms = tool({ */ export interface KeytermDetectorSession { readonly history: ChatContext; + /** The session's root trace context: the pass nests here when no reply is current. */ + readonly rootSpanContext?: Context; on( event: AgentSessionEventTypes.ConversationItemAdded, listener: (ev: ConversationItemAddedEvent) => void, @@ -372,11 +376,16 @@ export class KeytermDetector extends (EventEmitter as new () => TypedEmitter { try { - await this.runOnce(snapshot, controller.signal); + await this.runOnce(snapshot, controller.signal, parent); } catch (error) { this.#logger.child({ error }).error('keyterm detection pass failed'); throw error; @@ -395,31 +404,81 @@ export class KeytermDetector extends (EventEmitter as new () => TypedEmitter { - if (!(this.llm instanceof LLM)) { + /** + * One detection pass, as its own `keyterm_detection` span (otherwise the LLM call reads as a + * second inference step of the reply). `parent` is the span to nest under; the ambient + * context when omitted. + * + * @internal exposed for tests + */ + async runOnce(chatCtx: ChatContext, abortSignal?: AbortSignal, parent?: Context): Promise { + const llm = this.llm; + if (!(llm instanceof LLM)) { return; } - // show static terms as applied too, or the LLM keeps re-proposing them - const current: [string, boolean][] = [ - ...this.staticTerms.map((t): [string, boolean] => [t, true]), - ...this._detectedTerms.map((t): [string, boolean] => [t, true]), - ...[...this._pendingTerms.keys()].map((t): [string, boolean] => [t, false]), - ]; - const [pending, confirm, remove] = await detectKeyterms(this.llm, chatCtx, { - currentKeyterms: current, - instructions: this.instructions, - timeout: this.detectionTimeout, - abortSignal, - }); + const attributes: Attributes = { [traceTypes.ATTR_GEN_AI_REQUEST_MODEL]: llm.model }; + const provider = traceTypes.genAIProviderName(llm.provider); + if (provider !== undefined) { + attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME] = provider; + } + const before = this.keyterms; + const applied = await tracer.startActiveSpan( + async (span) => { + // show static terms as applied too, or the LLM keeps re-proposing them + const current: [string, boolean][] = [ + ...this.staticTerms.map((t): [string, boolean] => [t, true]), + ...this._detectedTerms.map((t): [string, boolean] => [t, true]), + ...[...this._pendingTerms.keys()].map((t): [string, boolean] => [t, false]), + ]; + const [pending, confirm, remove] = await detectKeyterms(llm, chatCtx, { + currentKeyterms: current, + instructions: this.instructions, + timeout: this.detectionTimeout, + abortSignal, + }); + + // cancelled mid-flight (e.g. activity shutdown): don't touch keyterm state + if (abortSignal?.aborted) { + return false; + } - // cancelled mid-flight (e.g. activity shutdown): don't touch keyterm state - if (abortSignal?.aborted) { + this.applyPass(pending, confirm, remove); + const after = this.keyterms; + const beforeSet = new Set(before); + const afterSet = new Set(after); + span.setAttributes({ + [traceTypes.ATTR_KEYTERMS_COUNT]: after.length, + [traceTypes.ATTR_KEYTERMS_ADDED]: after.filter((t) => !beforeSet.has(t)).length, + [traceTypes.ATTR_KEYTERMS_REMOVED]: before.filter((t) => !afterSet.has(t)).length, + }); + return true; + }, + { name: 'keyterm_detection', context: parent, attributes }, + ); + if (!applied) { return; } - const before = this.keyterms; + // update the STT if the keyterms changed + const newKeyterms = this.keyterms; + if ( + this.stt !== undefined && + (newKeyterms.length !== before.length || newKeyterms.some((t, i) => t !== before[i])) + ) { + this.stt._updateSessionKeyterms(newKeyterms); + const beforeSet = new Set(before); + const newSet = new Set(newKeyterms); + this.#logger + .child({ + added: newKeyterms.filter((t) => !beforeSet.has(t)), + removed: before.filter((t) => !newSet.has(t)), + }) + .debug('keyterms changed'); + } + } + + private applyPass(pending: string[], confirm: string[], remove: string[]): void { this.tick += 1; // update the keyterm state @@ -462,23 +521,6 @@ export class KeytermDetector extends (EventEmitter as new () => TypedEmitter t !== before[i])) - ) { - this.stt._updateSessionKeyterms(newKeyterms); - const beforeSet = new Set(before); - const newSet = new Set(newKeyterms); - this.#logger - .child({ - added: newKeyterms.filter((t) => !beforeSet.has(t)), - removed: before.filter((t) => !newSet.has(t)), - }) - .debug('keyterms changed'); - } } } diff --git a/agents/src/voice/keyterm_detection_span.test.ts b/agents/src/voice/keyterm_detection_span.test.ts new file mode 100644 index 0000000000..84bcd09e08 --- /dev/null +++ b/agents/src/voice/keyterm_detection_span.test.ts @@ -0,0 +1,170 @@ +// SPDX-FileCopyrightText: 2026 LiveKit, Inc. +// +// SPDX-License-Identifier: Apache-2.0 + +/** + * The keyterm-detection LLM call is STT context for later turns, not part of any reply: it gets + * its own `keyterm_detection` span, under the `agent_turn` whose reply fired the conversation + * event, else under `agent_session`. + */ +import { ROOT_CONTEXT, trace } from '@opentelemetry/api'; +import { InMemorySpanExporter, SimpleSpanProcessor } from '@opentelemetry/sdk-trace-base'; +import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node'; +import { EventEmitter } from 'node:events'; +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { ChatContext, FunctionCall } from '../llm/chat_context.js'; +import { LLM, type LLMStream } from '../llm/llm.js'; +import { initializeLogger } from '../log.js'; +import type { SpeechEvent, SpeechStream } from '../stt/stt.js'; +import { STT } from '../stt/stt.js'; +import { setTracerProvider, traceTypes, tracer } from '../telemetry/index.js'; +import type { AudioBuffer } from '../utils.js'; +import { createConversationItemAddedEvent } from './events.js'; +import { KeytermDetector } from './keyterm_detection.js'; + +initializeLogger({ pretty: false, level: 'silent' }); + +class FakeStream { + constructor( + private pending: string[], + private confirm: string[], + private remove: string[], + ) {} + + async collect() { + const args = JSON.stringify({ + pending: this.pending, + confirm: this.confirm, + remove: this.remove, + }); + return { + text: '', + toolCalls: [new FunctionCall({ callId: '1', name: 'record_keyterms', args })], + usage: undefined, + extra: {}, + }; + } + + close(): void {} +} + +/** Confirms `Acme` and `LiveKit` on every pass. */ +class RecordingLLM extends LLM { + label(): string { + return 'recording-llm'; + } + + override get model(): string { + return 'recording-model'; + } + + override get provider(): string { + return 'openai'; + } + + chat(): LLMStream { + return new FakeStream([], ['Acme', 'LiveKit'], []) as unknown as LLMStream; + } +} + +class KeytermSTT extends STT { + label = 'keyterm.STT'; + + constructor() { + super({ streaming: true, interimResults: false, keyterms: true }); + } + + protected _recognize(_: AudioBuffer): Promise { + throw new Error('not implemented'); + } + + stream(): SpeechStream { + throw new Error('not implemented'); + } + + override _updateSessionKeyterms(): void {} +} + +class FakeSession extends EventEmitter { + history = ChatContext.empty(); + rootSpanContext = trace.setSpan(ROOT_CONTEXT, tracer.startSpan({ name: 'agent_session' })); + + addUser(text: string): void { + const msg = this.history.addMessage({ role: 'user', content: text }); + this.emit('conversation_item_added', createConversationItemAddedEvent(msg)); + } +} + +describe.sequential('keyterm_detection 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 () => { + setTracerProvider(originalProvider); + await provider.shutdown(); + trace.disable(); + }); + + function detector(): { detector: KeytermDetector; session: FakeSession; llm: RecordingLLM } { + const llm = new RecordingLLM(); + const d = new KeytermDetector({ + staticKeyterms: ['Zed'], + options: { enabled: true, llm, turnInterval: 1 }, + }); + const session = new FakeSession(); + d.start(session, new KeytermSTT()); + return { detector: d, session, llm }; + } + + function keytermSpan() { + const spans = exporter.getFinishedSpans().filter((span) => span.name === 'keyterm_detection'); + expect(spans).toHaveLength(1); + return spans[0]!; + } + + it('nests under the agent_turn whose reply added the user message', async () => { + const { detector: d, session, llm } = detector(); + // the conversation event fires from the reply that answers the user message: the pass is + // that agent's work on the turn + const turn = await tracer.startActiveSpan( + async (turn) => { + session.addUser('order for Acme'); + expect(d._detectTask).toBeDefined(); + await d._detectTask!.result; + return turn; + }, + { name: 'agent_turn' }, + ); + + const span = keytermSpan(); + expect(span.parentSpanContext?.spanId).toBe(turn.spanContext().spanId); + expect(span.attributes[traceTypes.ATTR_KEYTERMS_COUNT]).toBe(3); // Zed + Acme + LiveKit + expect(span.attributes[traceTypes.ATTR_KEYTERMS_ADDED]).toBe(2); + expect(span.attributes[traceTypes.ATTR_KEYTERMS_REMOVED]).toBe(0); + expect(span.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).toBe(llm.model); + expect(span.attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]).toBe('openai'); + expect(d.keyterms).toEqual(['Zed', 'Acme', 'LiveKit']); + // the terms themselves stay out of the trace + expect(JSON.stringify(span.attributes)).not.toContain('Acme'); + }); + + it('falls back to the session root without a reply', async () => { + const { detector: d, session } = detector(); + // no span current: user code edited the history, or the reply was skipped + session.addUser('order for Acme'); + await d._detectTask!.result; + + const span = keytermSpan(); + const root = trace.getSpan(session.rootSpanContext)!; + expect(span.parentSpanContext?.spanId).toBe(root.spanContext().spanId); + }); +}); diff --git a/agents/src/voice/speech_handle.ts b/agents/src/voice/speech_handle.ts index 89d5139c36..b1293b411b 100644 --- a/agents/src/voice/speech_handle.ts +++ b/agents/src/voice/speech_handle.ts @@ -89,6 +89,13 @@ export type ResolvedSpeechHandle = Omit; * * @public */ +/** + * Why a speech was interrupted, for the `agent_turn` trace: the user started talking over it + * (`audio_activity`), a committed user turn preempted it (`user_turn`), or code did + * (`programmatic`: `session.interrupt()`, a tool, teardown). + */ +export type InterruptionSource = 'audio_activity' | 'user_turn' | 'programmatic'; + export class SpeechHandleCircularWaitError extends Error { constructor(functionCallName: string) { super(dedent` @@ -149,6 +156,8 @@ export class SpeechHandle { _scheduledAt?: number; /** @internal - when generation was first authorized, for the queue-wait attribute */ _authorizedAt?: number; + /** @internal - the first interrupt's cause, for the agent_turn trace */ + _interruptSource?: InterruptionSource; /** @internal - used by AgentTask/RunResult final output plumbing */ _maybeRunFinalOutput?: unknown; @@ -289,11 +298,15 @@ export class SpeechHandle { /** * Interrupt the current speech generation. * + * @param force - Interrupt even if this speech disallows interruptions. + * @param source - Why, for the `agent_turn` trace (see {@link InterruptionSource}). The first + * interruption's cause is the one recorded. + * * @throws Error If this speech handle is still running and does not allow interruptions. * * @returns The same speech handle that was interrupted. */ - interrupt(force: boolean = false): SpeechHandle { + interrupt(force: boolean = false, source: InterruptionSource = 'programmatic'): SpeechHandle { if (this.interrupted || this.done()) { // Already cancelled or finished: nothing to interrupt, and protection is moot. return this; @@ -303,6 +316,7 @@ export class SpeechHandle { throw new Error('This generation handle does not allow interruptions'); } + this._interruptSource = source; // first interrupt only: later calls return above this._cancel(); return this; } @@ -394,13 +408,19 @@ export class SpeechHandle { this.doneCallbacks.delete(callback); } - /** @internal */ - _cancel(): SpeechHandle { + /** + * @internal + * @param source - Why, for the `agent_turn` trace: a preemptive attempt superseded by more of + * the user's turn is `user_turn`, one dropped by a barge-in `audio_activity`; a cancel with + * no cause (teardown, a pause) reads as programmatic. The first cause named stands. + */ + _cancel(source?: InterruptionSource): SpeechHandle { if (this.done()) { return this; } if (!this.interruptFut.done) { + if (source !== undefined) this._interruptSource = source; this.interruptFut.resolve(); this.startInterruptTimeout(); } From 190c1772dbc9b3424e4ac103d37c871a0f2da12b Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 25 Sep 2026 16:06:52 -0700 Subject: [PATCH 2/5] review: STT identity fallback, interruption verdict in finally, fallback attribution; port livekit/agents#7373 and #7374 - user_turn names an STT that leaves the base getters at `unknown` by its label (prefix as provider) and normalizes the provider to the GenAI registry spelling - the pipeline reply task stamps its interruption verdict once it is over, so an interruption that lands while the tools run is recorded - an LLM fallback stream's usage metrics name the instance that served, not the one the adapter would pick next after a recovery - #7373: the fallback adapter's request span carries no operation name; the nested provider request is the `chat` - #7374: llm_node records its configured model and provider when the nested request is created, so a failover's serving provider is kept rather than overwritten when the node completes Co-Authored-By: Claude Fable 5.1 --- agents/src/llm/fallback_adapter.test.ts | 43 ++++++ agents/src/llm/fallback_adapter.ts | 9 ++ agents/src/llm/llm.ts | 20 ++- agents/src/telemetry/gen_ai.ts | 27 +++- agents/src/voice/agent_activity.ts | 66 ++++++---- agents/src/voice/coverage_spans.test.ts | 168 ++++++++++++++++++++++-- agents/src/voice/generation.ts | 29 ++-- 7 files changed, 305 insertions(+), 57 deletions(-) diff --git a/agents/src/llm/fallback_adapter.test.ts b/agents/src/llm/fallback_adapter.test.ts index d583b08531..8f426d9c10 100644 --- a/agents/src/llm/fallback_adapter.test.ts +++ b/agents/src/llm/fallback_adapter.test.ts @@ -265,6 +265,49 @@ describe('FallbackAdapter', () => { ); }); + it('names the instance that served in the usage metrics, not the one preferred next', async () => { + class NamedMockLLM extends MockLLM { + constructor( + label: string, + private readonly _model: string, + private readonly _provider: string, + ) { + super(label); + } + override get model(): string { + return this._model; + } + override get provider(): string { + return this._provider; + } + } + const primary = new NamedMockLLM('primary', 'primary-model', 'primary'); + primary.shouldFail = true; + const secondary = new NamedMockLLM('secondary', 'secondary-model', 'secondary'); + const adapter = new FallbackAdapter({ llms: [primary, secondary], attemptTimeout: 1 }); + const metrics: Array<{ + label: string; + metadata?: { modelName?: string; modelProvider?: string }; + }> = []; + adapter.on('metrics_collected', (m) => metrics.push(m)); + + const stream = adapter.chat({ chatCtx: { items: [] } as unknown as ChatContext }); + for await (const _chunk of stream) { + // drain + } + // the primary is preferred again as soon as its recovery probe succeeds; the metrics of + // the request the secondary served must still say so + adapter._status[0]!.available = true; + await delay(20); + + // the instances' own metrics are forwarded as they are; the adapter's own stream reports + // the instance that served its request + const adapterMetrics = metrics.filter((m) => m.label === adapter.label()); + expect(adapterMetrics).toHaveLength(1); + expect(adapterMetrics[0]!.metadata?.modelName).toBe('secondary-model'); + expect(adapterMetrics[0]!.metadata?.modelProvider).toBe('secondary'); + }); + it('reports the model and provider of the instance that serves next', () => { class IdentifiedLLM extends MockLLM { constructor( diff --git a/agents/src/llm/fallback_adapter.ts b/agents/src/llm/fallback_adapter.ts index f91719205e..6703ca0a88 100644 --- a/agents/src/llm/fallback_adapter.ts +++ b/agents/src/llm/fallback_adapter.ts @@ -235,6 +235,15 @@ class FallbackLLMStream extends LLMStream { return this.servedLlm?.model ?? this.adapter.model; } + protected override get responseProvider(): string { + return this.servedLlm?.provider ?? this.adapter.provider; + } + + /** The provider request nested under this stream owns the `chat` operation. */ + protected override get genAIOperationName(): string | undefined { + return undefined; + } + /** * The instance that served: on the current (attempt) span, and as the response side of the * adapter's request span and the caller's (llm_node). Request-side attributes named the diff --git a/agents/src/llm/llm.ts b/agents/src/llm/llm.ts index a8260ed615..9d8f91a134 100644 --- a/agents/src/llm/llm.ts +++ b/agents/src/llm/llm.ts @@ -256,6 +256,20 @@ export abstract class LLMStream implements AsyncIterableIterator { return this.#llm.model; } + /** The provider named on the response side and in the usage metrics (see {@link responseModel}). */ + protected get responseProvider(): string { + return this.#llm.provider; + } + + /** + * The convention's operation for the request span: `chat` for a provider request. An adapter + * that delegates to another stream has none, or a backend counting inference spans would see + * two calls for one. + */ + protected get genAIOperationName(): string | undefined { + return traceTypes.GenAIOperationName.CHAT; + } + protected get llmRequestSpan(): Span | undefined { return this.#llmRequestSpan; } @@ -263,7 +277,7 @@ export abstract class LLMStream implements AsyncIterableIterator { /** The GenAI inference span's request side, per the OTel GenAI conventions. */ private recordGenAIRequest(span: Span) { genAI.setRequestAttributes(span, { - operation: traceTypes.GenAIOperationName.CHAT, + operation: this.genAIOperationName, provider: this.#llm.provider, model: this.#llm.model, stream: true, @@ -414,8 +428,8 @@ export abstract class LLMStream implements AsyncIterableIterator { return (usage?.completionTokens || 0) / (durationMs / 1000); })(), metadata: { - modelProvider: this.#llm.provider, - modelName: this.#llm.model, + modelProvider: this.responseProvider, + modelName: this.responseModel, }, }; diff --git a/agents/src/telemetry/gen_ai.ts b/agents/src/telemetry/gen_ai.ts index d54791cd69..0a78f40b76 100644 --- a/agents/src/telemetry/gen_ai.ts +++ b/agents/src/telemetry/gen_ai.ts @@ -76,11 +76,21 @@ const INFERENCE_RECORDED = Symbol('lkInferenceRecorded'); export interface InferenceMarker { recorded: boolean; + /** Runs once, when the first `llm_request` span is created inside the tracked call. */ + onRecorded?: () => void; } -/** Runs `fn` with a marker that fills in if an `llm_request` span is created inside it. */ -export function withInferenceTracking(fn: (marker: InferenceMarker) => T): T { - const marker: InferenceMarker = { recorded: false }; +/** + * Runs `fn` with a marker that fills in if an `llm_request` span is created inside it. + * `onRecorded` runs at that moment: what the node knew before the request (its configured + * model and provider) is recorded before a fallback can overwrite the provider with the + * instance that actually served, and never afterwards. + */ +export function withInferenceTracking( + fn: (marker: InferenceMarker) => T, + options: { onRecorded?: () => void } = {}, +): T { + const marker: InferenceMarker = { recorded: false, onRecorded: options.onRecorded }; return otelContext.with(otelContext.active().setValue(INFERENCE_RECORDED, marker), () => fn(marker), ); @@ -89,7 +99,9 @@ export function withInferenceTracking(fn: (marker: InferenceMarker) => T): T /** Called where an `llm_request` span is created, so the enclosing node stands down. */ export function markInferenceSpanRecorded(): void { const marker = otelContext.active().getValue(INFERENCE_RECORDED) as InferenceMarker | undefined; - if (marker) marker.recorded = true; + if (!marker || marker.recorded) return; + marker.recorded = true; + marker.onRecorded?.(); } function textPart(content: string): MessagePart { @@ -341,7 +353,8 @@ export function setContentAttributes( export function setRequestAttributes( span: Span, params: { - operation: string; + /** Absent for a delegating span: the nested provider request owns the operation. */ + operation: string | undefined; provider?: string; model?: string; stream?: boolean; @@ -350,7 +363,9 @@ export function setRequestAttributes( ): void { if (!span.isRecording()) return; - const attrs: Attributes = { [traceTypes.ATTR_GEN_AI_OPERATION_NAME]: params.operation }; + const attrs: Attributes = {}; + if (params.operation !== undefined) + attrs[traceTypes.ATTR_GEN_AI_OPERATION_NAME] = params.operation; const provider = traceTypes.genAIProviderName(params.provider); if (provider) attrs[traceTypes.ATTR_GEN_AI_PROVIDER_NAME] = provider; if (params.model) attrs[traceTypes.ATTR_GEN_AI_REQUEST_MODEL] = params.model; diff --git a/agents/src/voice/agent_activity.ts b/agents/src/voice/agent_activity.ts index 66099c511f..d21b330f4d 100644 --- a/agents/src/voice/agent_activity.ts +++ b/agents/src/voice/agent_activity.ts @@ -367,6 +367,20 @@ function recordUserTurnStages(span: Span, userMetrics: MetricsReport): void { } } +/** + * The STT identity for the user_turn span. Plugins that leave the base getters at their + * `unknown` default are named by their label instead (its `-` prefix for the + * provider), and the provider is normalized to the GenAI registry spelling like every other + * `gen_ai.provider.name` the framework writes. + */ +function sttIdentity(stt: STT | undefined): { model?: string; provider?: string } { + if (!stt) return {}; + const label = stt.label; + const model = stt.model !== 'unknown' ? stt.model : label; + const rawProvider = stt.provider !== 'unknown' ? stt.provider : label.split('-', 1)[0]; + return { model, provider: traceTypes.genAIProviderName(rawProvider) ?? rawProvider }; +} + /** Stamp how long the speech sat in the queue on its agent_turn span, in seconds. */ function recordQueueWait(speechHandle: SpeechHandle): void { const queueWait = speechHandle._queueWait(); @@ -871,6 +885,8 @@ export class AgentActivity implements RecognitionHooks { this.llm instanceof RealtimeModel && this.llm.capabilities.turnDetection === true; const recognitionVad = this.usingDefaultVad && realtimeUsesServerVad ? undefined : this.vad; + // as python passes them (a fallback adapter reports the instance expected to serve next) + const sttIdent = sttIdentity(this.stt); this.audioRecognition = new AudioRecognition({ recognitionHooks: this, // Disable stt node if stt is not provided @@ -886,10 +902,8 @@ export class AgentActivity implements RecognitionHooks { endpointing: createEndpointing(this.endpointingOpts), userTurnLimit: this.agentSession.sessionOptions.turnHandling.userTurnLimit, rootSpanContext: this.agentSession.rootSpanContext, - // the model and provider, as python passes them (a fallback adapter reports the instance - // expected to serve next); the label is a class name, not a model - sttModel: this.stt?.model, - sttProvider: this.stt?.provider, + sttModel: sttIdent.model, + sttProvider: sttIdent.provider, sttAlignedTranscript: Boolean(this.stt?.capabilities.alignedTranscript), getLinkedParticipant: () => this.agentSession._roomIO?.linkedParticipant, shouldDiscardAudioForStt: () => this.shouldDiscardInputAudio(), @@ -1399,8 +1413,7 @@ export class AgentActivity implements RecognitionHooks { await this.audioRecognition.updateStt( this.stt ? (...args) => this.agent.sttNode(...args) : undefined, { - model: resolvedStt?.model, - provider: resolvedStt?.provider, + ...sttIdentity(resolvedStt), alignedTranscript: Boolean(resolvedStt?.capabilities.alignedTranscript), resetContext: true, }, @@ -1499,8 +1512,7 @@ export class AgentActivity implements RecognitionHooks { await this.audioRecognition?.updateStt( previous.resolvedStt ? (...args) => this.agent.sttNode(...args) : undefined, { - model: previous.resolvedStt?.model, - provider: previous.resolvedStt?.provider, + ...sttIdentity(previous.resolvedStt), alignedTranscript: Boolean(previous.resolvedStt?.capabilities.alignedTranscript), resetContext: true, }, @@ -4220,20 +4232,30 @@ export class AgentActivity implements RecognitionHooks { _previousUserMetrics?: MetricsReport, ): Promise => tracer.startActiveSpan( - async (span) => ( - this.recordAgentTurn(span), - this._pipelineReplyTaskImpl({ - stateLease, - chatCtx, - toolCtx, - modelSettings, - replyAbortController, - instructions, - newMessage, - span, - _previousUserMetrics, - }) - ), + async (span) => { + this.recordAgentTurn(span); + try { + await this._pipelineReplyTaskImpl({ + stateLease, + chatCtx, + toolCtx, + modelSettings, + replyAbortController, + instructions, + newMessage, + span, + _previousUserMetrics, + }); + } 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, + ); + recordInterruption(stateLease.speechHandle); + } + }, { name: 'agent_turn', context: this.agentSession.rootSpanContext, diff --git a/agents/src/voice/coverage_spans.test.ts b/agents/src/voice/coverage_spans.test.ts index 13aca82ddd..24cf23e004 100644 --- a/agents/src/voice/coverage_spans.test.ts +++ b/agents/src/voice/coverage_spans.test.ts @@ -19,17 +19,18 @@ import { afterEach, beforeEach, describe, expect, it } from 'vitest'; import { APIConnectionError } from '../_exceptions.js'; import { ChatContext } from '../llm/chat_context.js'; import { FallbackAdapter } from '../llm/fallback_adapter.js'; -import { LLM, LLMStream } from '../llm/llm.js'; -import type { ToolChoice, ToolContextLike } from '../llm/tool_context.js'; +import { type ChatChunk, LLM, LLMStream } from '../llm/llm.js'; +import { type ToolChoice, ToolContext, type ToolContextLike, 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 { type APIConnectOptions, DEFAULT_API_CONNECT_OPTIONS } from '../types.js'; -import { delay } from '../utils.js'; +import { Future, delay } from '../utils.js'; import { VAD, type VADEvent, VADEventType, VADStream } from '../vad.js'; import { Agent } from './agent.js'; import { AgentSession } from './agent_session.js'; -import { AudioInput, AudioOutput } from './io.js'; +import { performLLMInference } from './generation.js'; +import { AudioInput, AudioOutput, type LLMNode } from './io.js'; import { SpeechHandle } from './speech_handle.js'; import { FakeLLM } from './testing/fake_llm.js'; @@ -317,6 +318,69 @@ describe.sequential('coverage spans', () => { } }); + it('records an interruption that lands while the tools are running', async () => { + // the reply task returns early when its speech is interrupted while it waits for the + // tools, past the point where it stamped the verdict: the turn must still say so + const TRANSCRIPT = 'Look it up.'; + const toolStarted = new Future(); + const release = new Future(); + const agent = new Agent({ + instructions: 'test', + tools: { + lookup: tool({ + description: 'a slow lookup', + execute: async () => { + toolStarted.resolve(); + await release.await; + return 'found'; + }, + }), + }, + }); + const vad = new ScriptedVAD(); + const stt = new FakeSTT({ + capabilities: { streaming: true, interimResults: true }, + fakeUserSpeeches: [{ startTime: 0, endTime: 200, transcript: TRANSCRIPT, sttDelay: 100 }], + }); + const llm = new FakeLLM([{ input: TRANSCRIPT, toolCalls: [{ name: 'lookup', args: {} }] }]); + const session = new AgentSession({ + vad, + stt, + llm, + turnHandling: { turnDetection: 'vad', endpointing: { minDelay: 100, maxDelay: 100 } }, + }); + const audioInput = new ScriptedAudioInput(); + session.input.audio = audioInput; + session.output.audio = new PacedOutput(); + await session.start({ agent }); + try { + audioInput.push(20); + await delay(20); + vad.startOfSpeech(); + await delay(200); + vad.endOfSpeech(); + await toolStarted.await; + session.interrupt(); + release.resolve(); + await waitFor( + () => + spansNamed(exporter, 'agent_turn').some( + (span) => span.attributes[traceTypes.ATTR_SPEECH_INTERRUPTED] !== undefined, + ), + 10_000, + 'the interrupted turn to end', + ); + } finally { + await session.close(); + } + + const turn = spansNamed(exporter, 'agent_turn').find( + (span) => span.attributes[traceTypes.ATTR_SPEECH_INTERRUPTED] !== undefined, + )!; + expect(turn.attributes[traceTypes.ATTR_SPEECH_INTERRUPTED]).toBe(true); + expect(turn.attributes[traceTypes.ATTR_INTERRUPTION_SOURCE]).toBeTypeOf('string'); + }); + it('a committed user turn interrupts the queued replies for the same reason', async () => { // the turn interrupts the reply that is playing and the ones queued behind it: all of // them name user_turn, not the programmatic default the queue sweep used to fall back to @@ -434,16 +498,10 @@ describe.sequential('coverage spans', () => { } } - it('names the STT model and provider on user_turn, as python does', async () => { - // the label is a class name (`stt.FallbackAdapter`, `fake-stt`): the span carries the - // model and provider the STT reports, which for a fallback adapter is the instance expected - // to serve next + /** One user turn through `stt`, returning its user_turn span. */ + async function userTurnWith(stt: FakeSTT): Promise { const TRANSCRIPT = 'Hello'; const vad = new ScriptedVAD(); - const stt = new NamedSTT({ - capabilities: { streaming: true, interimResults: true }, - fakeUserSpeeches: [{ startTime: 0, endTime: 200, transcript: TRANSCRIPT, sttDelay: 100 }], - }); const llm = new FakeLLM([{ input: TRANSCRIPT, content: 'Hi there' }]); const session = new AgentSession({ vad, @@ -473,14 +531,49 @@ describe.sequential('coverage spans', () => { await session.close(); } - const turn = spansNamed(exporter, 'user_turn').find( + return spansNamed(exporter, 'user_turn').find( (span) => span.attributes[traceTypes.ATTR_USER_TRANSCRIPT] === TRANSCRIPT, )!; + } + + const speeches = { + capabilities: { streaming: true, interimResults: true }, + fakeUserSpeeches: [{ startTime: 0, endTime: 200, transcript: 'Hello', sttDelay: 100 }], + }; + + it('names the STT model and provider on user_turn, as python does', async () => { + // the label is a class name (`stt.FallbackAdapter`, `fake-stt`): the span carries the + // model and provider the STT reports, which for a fallback adapter is the instance expected + // to serve next + const stt = new NamedSTT(speeches); + const turn = await userTurnWith(stt); expect(turn.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).toBe('fake-model'); expect(turn.attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]).toBe('fake-provider'); expect(turn.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).not.toBe(stt.label); }); + it('falls back to the label for an STT that does not report its identity', async () => { + // a plugin that leaves the base getters at `unknown` is still named: by its label, whose + // `-` prefix stands in for the provider + const stt = new FakeSTT({ ...speeches, label: 'sarvam-saarika' }); + const turn = await userTurnWith(stt); + expect(turn.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).toBe('sarvam-saarika'); + expect(turn.attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]).toBe('sarvam'); + }); + + it('normalizes the STT provider to the GenAI registry spelling', async () => { + class DisplayNamedSTT extends FakeSTT { + override get model(): string { + return 'whisper-1'; + } + override get provider(): string { + return 'OpenAI'; + } + } + const turn = await userTurnWith(new DisplayNamedSTT(speeches)); + expect(turn.attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]).toBe('openai'); + }); + // -- agent handoff -- it('groups the handoff under an update_agent span', async () => { @@ -707,6 +800,13 @@ describe.sequential('coverage spans', () => { expect(request!.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).toBe(primary.model); expect(request!.attributes[traceTypes.ATTR_GEN_AI_RESPONSE_MODEL]).toBe(secondary.model); expect(request!.attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]).toBe('openai'); + // the adapter delegates: only the instance's own request under the run is the `chat`, so + // a backend counting inference spans sees one call + expect(request!.attributes[traceTypes.ATTR_GEN_AI_OPERATION_NAME]).toBeUndefined(); + const served = spansNamed(exporter, 'llm_request').find( + (span) => span.parentSpanContext?.spanId === run.spanContext().spanId, + ); + expect(served?.attributes[traceTypes.ATTR_GEN_AI_OPERATION_NAME]).toBe('chat'); // the caller's span gets the same response side, per request const caller = only(exporter, 'caller'); expect(caller.attributes[traceTypes.ATTR_GEN_AI_RESPONSE_MODEL]).toBe(secondary.model); @@ -714,4 +814,46 @@ describe.sequential('coverage spans', () => { expect(adapter.model).toBe(secondary.model); expect(adapter.provider).toBe(secondary.provider); }); + + it('keeps the configured model and names the serving provider on llm_node after a failover', async () => { + // the node span records its configured identity when the request is made, so the serving + // provider a failover stamps afterwards is not overwritten when the node completes + const primary = new FailingLLM(); + const secondary = new ServingLLM([{ input: 'hi', content: 'hello' }]); + const adapter = new FallbackAdapter({ llms: [primary, secondary], attemptTimeout: 1 }); + const chatCtx = ChatContext.empty(); + chatCtx.addMessage({ role: 'user', content: 'hi' }); + const node: LLMNode = async (ctx, tools) => { + const stream = adapter.chat({ chatCtx: ctx, toolCtx: tools }); + return new ReadableStream({ + async start(controller) { + for await (const chunk of stream) controller.enqueue(chunk); + controller.close(); + }, + }); + }; + const [task, data] = performLLMInference( + node, + chatCtx, + ToolContext.empty(), + {}, + new AbortController(), + adapter.model, + adapter.provider, + ); + const reader = data.textStream.getReader(); + let text = ''; + while (true) { + const { done, value } = await reader.read(); + if (done) break; + if (typeof value === 'string') text += value; + } + await task.result; + expect(text).toBe('hello'); + + const nodeSpan = only(exporter, 'llm_node'); + expect(nodeSpan.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).toBe(primary.model); + expect(nodeSpan.attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]).toBe('openai'); + expect(nodeSpan.attributes[traceTypes.ATTR_GEN_AI_RESPONSE_MODEL]).toBe(secondary.model); + }); }); diff --git a/agents/src/voice/generation.ts b/agents/src/voice/generation.ts index a40b75818a..3b6d7760ad 100644 --- a/agents/src/voice/generation.ts +++ b/agents/src/voice/generation.ts @@ -631,6 +631,18 @@ export function performLLMInference( const toolCallWriter = toolCallStream.writable.getWriter(); const data = new _LLMGenerationData(textStream.readable, toolCallStream.readable); + // the configured model and provider describe the inference only once it is known that this + // LLM served it: they go on the node span the moment the nested `llm_request` span is created + // (see withInferenceTracking), so a fallback that failed over can then name the serving + // provider on top of them rather than be overwritten by them when the node completes + const recordConfiguredModel = (span: Span) => { + if (model) span.setAttribute(traceTypes.ATTR_GEN_AI_REQUEST_MODEL, model); + const normalizedProvider = traceTypes.genAIProviderName(provider); + if (normalizedProvider) { + span.setAttribute(traceTypes.ATTR_GEN_AI_PROVIDER_NAME, normalizedProvider); + } + }; + const _performLLMInferenceImpl = async ( signal: AbortSignal, span: Span, @@ -644,15 +656,6 @@ export function performLLMInference( ); span.setAttribute(traceTypes.ATTR_FUNCTION_TOOLS, JSON.stringify(sortedToolNames(toolCtx))); - // the configured model and provider describe the inference only once it is known that - // this LLM served it; that is decided below, when the nested span is (or is not) there - const recordConfiguredModel = () => { - if (model) span.setAttribute(traceTypes.ATTR_GEN_AI_REQUEST_MODEL, model); - const normalizedProvider = traceTypes.genAIProviderName(provider); - if (normalizedProvider) { - span.setAttribute(traceTypes.ATTR_GEN_AI_PROVIDER_NAME, normalizedProvider); - } - }; let nodeError: Error | string | undefined; // the GenAI inference attributes belong to the nested `llm_request` span, which is the @@ -773,9 +776,7 @@ export function performLLMInference( } // a custom node may have generated this itself, with no nested `llm_request` span to // carry the convention's attributes; when there was one, they are already recorded - if (inference.recorded) { - recordConfiguredModel(); - } else { + if (!inference.recorded) { // a third-party engine served this, so the configured model and provider are left // off rather than crediting it with a call it never made genAI.setRequestAttributes(span, { @@ -816,7 +817,9 @@ export function performLLMInference( const inferenceTask = async (signal: AbortSignal) => tracer.startActiveSpan( async (span) => - genAI.withInferenceTracking((marker) => _performLLMInferenceImpl(signal, span, marker)), + genAI.withInferenceTracking((marker) => _performLLMInferenceImpl(signal, span, marker), { + onRecorded: () => recordConfiguredModel(span), + }), { name: 'llm_node', context: currentContext }, ); From 5320e2550b8d726f913655463d2996623cb2c1ca Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 25 Sep 2026 20:20:32 -0700 Subject: [PATCH 3/5] review: served-instance attribution for TTS metrics and partial chunked audio; live STT identity on user_turn - a TTS fallback stream's usage metrics name the instance that served, not the one the adapter would pick next once a failed instance recovered - chunked synthesis that fails after audio reached the caller still names the instance that produced it on the request and caller spans, as the streaming path does - user_turn reads the STT identity when it stamps the turn, at its start and again at its end, so a fallback that failed over names the instance that transcribed rather than the one snapshotted at activity start Co-Authored-By: Claude Fable 5.1 --- agents/src/tts/fallback_adapter.test.ts | 105 ++++++++++++++++++++++++ agents/src/tts/fallback_adapter.ts | 26 ++++++ agents/src/tts/tts.ts | 36 +++++++- agents/src/voice/agent_activity.ts | 3 + agents/src/voice/audio_recognition.ts | 27 ++++-- agents/src/voice/coverage_spans.test.ts | 52 ++++++++++++ 6 files changed, 239 insertions(+), 10 deletions(-) diff --git a/agents/src/tts/fallback_adapter.test.ts b/agents/src/tts/fallback_adapter.test.ts index 574017f763..2100d96dd7 100644 --- a/agents/src/tts/fallback_adapter.test.ts +++ b/agents/src/tts/fallback_adapter.test.ts @@ -2,13 +2,18 @@ // // SPDX-License-Identifier: Apache-2.0 import { AudioFrame } from '@livekit/rtc-node'; +import { context as otelContext, trace } from '@opentelemetry/api'; +import { InMemorySpanExporter, SimpleSpanProcessor } from '@opentelemetry/sdk-trace-base'; +import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node'; import { ReadableStream } from 'node:stream/web'; import { afterEach, beforeAll, beforeEach, describe, expect, it, vi } from 'vitest'; import { APIConnectionError, APIError, APIStatusError } from '../_exceptions.js'; import { initializeLogger, log } from '../log.js'; +import { setTracerProvider, traceTypes, tracer } from '../telemetry/index.js'; import { basic } from '../tokenize/index.js'; import type { APIConnectOptions } from '../types.js'; import { USERDATA_TTS_STARTED_TIME } from '../types.js'; +import { delay } from '../utils.js'; import { FallbackAdapter } from './fallback_adapter.js'; import { StreamAdapter } from './stream_adapter.js'; import { ChunkedStream, SynthesizeStream, TTS, type TTSError } from './tts.js'; @@ -823,6 +828,106 @@ describe('TTS FallbackAdapter', () => { }); }); + it('names the instance that served in the usage metrics, not the one preferred next', async () => { + class IdentifiedTTS extends MockTTS { + constructor( + label: string, + private readonly _model: string, + private readonly _provider: string, + ) { + super(label); + } + override get model(): string { + return this._model; + } + override get provider(): string { + return this._provider; + } + } + const primary = new IdentifiedTTS('primary', 'primary-model', 'primary'); + primary.shouldFail = true; + const secondary = new IdentifiedTTS('secondary', 'secondary-model', 'secondary'); + const adapter = new FallbackAdapter({ + ttsInstances: [primary, secondary], + maxRetryPerTTS: 0, + recoveryDelayMs: 60_000, + }); + const metrics: Array<{ + label: string; + metadata?: { modelName?: string; modelProvider?: string }; + }> = []; + adapter.on('metrics_collected', (m) => metrics.push(m)); + + const stream = adapter.stream(); + stream.updateInputStream( + new ReadableStream({ + start(controller) { + controller.enqueue('hello world'); + controller.close(); + }, + }), + ); + for await (const event of stream) { + if (event === SynthesizeStream.END_OF_STREAM) break; + } + // the primary is preferred again as soon as its probe succeeds; the request the secondary + // served must still say so + adapter.status[0]!.available = true; + await delay(20); + + const own = metrics.filter((m) => m.label === adapter.label); + expect(own.length).toBeGreaterThan(0); + for (const m of own) { + expect(m.metadata?.modelName).toBe('secondary-model'); + expect(m.metadata?.modelProvider).toBe('secondary'); + } + stream.close(); + await adapter.close(); + }); + + it('attributes partial chunked audio to the instance the caller heard', async () => { + // chunked synthesis that fails after emitting audio cannot fall back (the audio was heard); + // the request and caller spans still name the instance that produced it + const exporter = new InMemorySpanExporter(); + const provider = new NodeTracerProvider({ + spanProcessors: [new SimpleSpanProcessor(exporter)], + }); + provider.register(); + const previous = tracer.getProvider(); + setTracerProvider(provider); + try { + const primary = new MockTTS('primary'); + primary.shouldFail = true; + primary.failAfterAudio = true; + const adapter = new FallbackAdapter({ + ttsInstances: [primary, new MockTTS('secondary')], + maxRetryPerTTS: 0, + recoveryDelayMs: 60_000, + }); + await tracer.startActiveSpan( + async () => { + const stream = adapter.synthesize('hello world'); + try { + for await (const _frame of stream) { + // drain + } + } catch { + // the partial failure is reported to the caller + } + }, + { name: 'caller' }, + ); + await adapter.close(); + const caller = exporter.getFinishedSpans().find((span) => span.name === 'caller'); + expect(caller?.attributes[traceTypes.ATTR_GEN_AI_RESPONSE_MODEL]).toBe(primary.model); + } finally { + setTracerProvider(previous); + await provider.shutdown(); + trace.disable(); + otelContext.disable(); + } + }); + it('reports the model and provider of the instance that serves next', () => { class IdentifiedTTS extends MockTTS { constructor( diff --git a/agents/src/tts/fallback_adapter.ts b/agents/src/tts/fallback_adapter.ts index 07952c2654..4cc55b6b5b 100644 --- a/agents/src/tts/fallback_adapter.ts +++ b/agents/src/tts/fallback_adapter.ts @@ -387,9 +387,19 @@ class FallbackChunkedStream extends ChunkedStream { private _logger = log(); // the span this request was made under (tts_node); see recordFallbackServed private callerSpan: Span | undefined = trace.getActiveSpan(); + /** The instance whose audio reached the caller (see recordFallbackServed). */ + private servedTts?: TTS; label: string = 'tts.FallbackChunkedStream'; + protected override get responseModel(): string { + return this.servedTts?.model ?? this.adapter.model; + } + + protected override get responseProvider(): string { + return this.servedTts?.provider ?? this.adapter.provider; + } + constructor( adapter: FallbackAdapter, text: string, @@ -488,6 +498,7 @@ class FallbackChunkedStream extends ChunkedStream { } this._logger.debug({ tts: tts.label }, 'TTS synthesis succeeded'); + this.servedTts = tts; recordFallbackServed(tts, i, this.ttsRequestSpan, this.callerSpan); return; } catch (error) { @@ -496,6 +507,9 @@ class FallbackChunkedStream extends ChunkedStream { { tts: tts.label, error: providerError ?? error }, 'TTS failed after audio pushed, cannot fallback mid-utterance', ); + // the caller heard this instance's audio: it served, partially + this.servedTts = tts; + recordFallbackServed(tts, i, this.ttsRequestSpan, this.callerSpan); throw error; } @@ -530,9 +544,19 @@ class FallbackSynthesizeStream extends SynthesizeStream { private _logger = log(); // the span this request was made under (tts_node); see recordFallbackServed private callerSpan: Span | undefined = trace.getActiveSpan(); + /** The instance whose audio reached the caller (see recordFallbackServed). */ + private servedTts?: TTS; label: string = 'tts.FallbackSynthesizeStream'; + protected override get responseModel(): string { + return this.servedTts?.model ?? this.adapter.model; + } + + protected override get responseProvider(): string { + return this.servedTts?.provider ?? this.adapter.provider; + } + constructor(adapter: FallbackAdapter, connOptions: APIConnectOptions) { super(adapter, connOptions); this.adapter = adapter; @@ -720,6 +744,7 @@ class FallbackSynthesizeStream extends SynthesizeStream { this.queue.put(SynthesizeStream.END_OF_STREAM); this._logger.debug({ tts: originalTts.label }, 'TTS stream succeeded'); + this.servedTts = originalTts; recordFallbackServed(originalTts, i, this.ttsRequestSpan, this.callerSpan); await readInputLLMStream.catch(() => {}); return; @@ -730,6 +755,7 @@ class FallbackSynthesizeStream extends SynthesizeStream { 'TTS failed after audio pushed, cannot fallback mid-utterance', ); // the caller heard this instance's audio: it served, partially + this.servedTts = originalTts; recordFallbackServed(originalTts, i, this.ttsRequestSpan, this.callerSpan); throw error; } diff --git a/agents/src/tts/tts.ts b/agents/src/tts/tts.ts index 6126330897..e839928905 100644 --- a/agents/src/tts/tts.ts +++ b/agents/src/tts/tts.ts @@ -414,6 +414,20 @@ export abstract class SynthesizeStream } /** The `tts_request` span of this stream, once the main task has opened it. */ + /** + * The model named in the usage metrics: the TTS's own by default; a fallback adapter's stream + * reports the instance that actually served, which the adapter's own getters cannot tell once + * a failed instance recovered while the request was still running. + */ + protected get responseModel(): string { + return this.#tts.model; + } + + /** The provider named in the usage metrics; see `responseModel`. */ + protected get responseProvider(): string { + return this.#tts.provider; + } + protected get ttsRequestSpan(): Span | undefined { return this.#ttsRequestSpan; } @@ -627,8 +641,8 @@ export abstract class SynthesizeStream outputTokens: this.#outputTokens, streamed: true, metadata: { - modelProvider: this.#tts.provider, - modelName: this.#tts.model, + modelProvider: this.responseProvider, + modelName: this.responseModel, }, }; if (this.#ttsRequestSpan) { @@ -834,6 +848,20 @@ export abstract class ChunkedStream implements AsyncIterableIterator sttIdentity(this.stt), sttAlignedTranscript: Boolean(this.stt?.capabilities.alignedTranscript), getLinkedParticipant: () => this.agentSession._roomIO?.linkedParticipant, shouldDiscardAudioForStt: () => this.shouldDiscardInputAudio(), @@ -1414,6 +1415,7 @@ export class AgentActivity implements RecognitionHooks { this.stt ? (...args) => this.agent.sttNode(...args) : undefined, { ...sttIdentity(resolvedStt), + identity: () => sttIdentity(resolvedStt), alignedTranscript: Boolean(resolvedStt?.capabilities.alignedTranscript), resetContext: true, }, @@ -1513,6 +1515,7 @@ export class AgentActivity implements RecognitionHooks { previous.resolvedStt ? (...args) => this.agent.sttNode(...args) : undefined, { ...sttIdentity(previous.resolvedStt), + identity: () => sttIdentity(previous.resolvedStt), alignedTranscript: Boolean(previous.resolvedStt?.capabilities.alignedTranscript), resetContext: true, }, diff --git a/agents/src/voice/audio_recognition.ts b/agents/src/voice/audio_recognition.ts index 4ac5e84347..d899575b23 100644 --- a/agents/src/voice/audio_recognition.ts +++ b/agents/src/voice/audio_recognition.ts @@ -310,6 +310,11 @@ export interface AudioRecognitionOptions { sttModel?: string; /** STT provider name for tracing */ sttProvider?: string; + /** + * The STT identity read when a user turn is stamped, over the static names above: a fallback + * adapter that failed over mid-session then names the instance serving, not the one at start. + */ + sttIdentity?: () => { model?: string; provider?: string }; /** Whether the active STT provides transcript alignment timestamps. */ sttAlignedTranscript?: boolean; /** Getter for linked participant for span attribution */ @@ -368,6 +373,7 @@ export class AudioRecognition { private rootSpanContext?: Context; private sttModel?: string; private sttProvider?: string; + private sttIdentity?: () => { model?: string; provider?: string }; private sttAlignedTranscript: boolean; private getLinkedParticipant?: () => ParticipantLike | undefined; @@ -484,6 +490,7 @@ export class AudioRecognition { this.rootSpanContext = opts.rootSpanContext; this.sttModel = opts.sttModel; this.sttProvider = opts.sttProvider; + this.sttIdentity = opts.sttIdentity; this.sttAlignedTranscript = opts.sttAlignedTranscript ?? false; this.getLinkedParticipant = opts.getLinkedParticipant; this.transcriptionTimeout = opts.transcriptionTimeout ?? undefined; @@ -754,6 +761,7 @@ export class AudioRecognition { options: { model?: string; provider?: string; + identity?: () => { model?: string; provider?: string }; alignedTranscript?: boolean; resetContext?: boolean; } = {}, @@ -767,6 +775,9 @@ export class AudioRecognition { if (Object.hasOwn(options, 'provider')) { this.sttProvider = options.provider; } + if (Object.hasOwn(options, 'identity')) { + this.sttIdentity = options.identity; + } if (Object.hasOwn(options, 'alignedTranscript')) { this.sttAlignedTranscript = options.alignedTranscript ?? false; } @@ -1154,16 +1165,18 @@ export class AudioRecognition { setParticipantSpanAttributes(this.userTurnSpan, participant); } - if (this.sttModel) { - this.userTurnSpan.setAttribute(traceTypes.ATTR_GEN_AI_REQUEST_MODEL, this.sttModel); - } - if (this.sttProvider) { - this.userTurnSpan.setAttribute(traceTypes.ATTR_GEN_AI_PROVIDER_NAME, this.sttProvider); - } + this.stampSttIdentity(this.userTurnSpan); return this.userTurnSpan; } + /** Name the STT on a user turn; read live, so a failover during the turn shows through. */ + private stampSttIdentity(span: Span): void { + const { model = this.sttModel, provider = this.sttProvider } = this.sttIdentity?.() ?? {}; + if (model) span.setAttribute(traceTypes.ATTR_GEN_AI_REQUEST_MODEL, model); + if (provider) span.setAttribute(traceTypes.ATTR_GEN_AI_PROVIDER_NAME, provider); + } + private userTurnContext(span: Span): Context { const base = this.rootSpanContext ?? ROOT_CONTEXT; return trace.setSpan(base, span); @@ -2643,6 +2656,8 @@ export class AudioRecognition { // opened after the decided one is a later bounce's to decide if (decided === undefined || this.eouWaitSpan === decided.wait) this.endEouWaitSpan('dropped'); if (this.userTurnSpan && info) { + // the instance that transcribed the turn, after any failover during it + this.stampSttIdentity(this.userTurnSpan); this.userTurnSpan.setAttributes({ [traceTypes.ATTR_USER_TRANSCRIPT]: info.transcript, [traceTypes.ATTR_TRANSCRIPT_CONFIDENCE]: info.confidence, diff --git a/agents/src/voice/coverage_spans.test.ts b/agents/src/voice/coverage_spans.test.ts index 24cf23e004..a7552f2eb9 100644 --- a/agents/src/voice/coverage_spans.test.ts +++ b/agents/src/voice/coverage_spans.test.ts @@ -22,6 +22,8 @@ import { FallbackAdapter } from '../llm/fallback_adapter.js'; import { type ChatChunk, LLM, LLMStream } from '../llm/llm.js'; import { type ToolChoice, ToolContext, type ToolContextLike, tool } from '../llm/tool_context.js'; import { initializeLogger } from '../log.js'; +import { FallbackAdapter as STTFallbackAdapter } from '../stt/fallback_adapter.js'; +import { STT, type SpeechEvent, SpeechStream } from '../stt/stt.js'; import { FakeSTT } from '../stt/testing/fake_stt.js'; import { setTracerProvider, traceTypes, tracer } from '../telemetry/index.js'; import { type APIConnectOptions, DEFAULT_API_CONNECT_OPTIONS } from '../types.js'; @@ -561,6 +563,56 @@ describe.sequential('coverage spans', () => { expect(turn.attributes[traceTypes.ATTR_GEN_AI_PROVIDER_NAME]).toBe('sarvam'); }); + it('names the STT instance that transcribed the turn after a failover', async () => { + // the identity is read when the turn is stamped, not snapshotted at start: a primary that + // failed mid-session must not be credited with the fallback's transcripts + class FailingSTTStream extends SpeechStream { + label = 'failing-stt-stream'; + protected async run(): Promise { + throw new APIConnectionError({ message: 'primary down' }); + } + } + class FailingSTT extends STT { + label = 'primary'; + constructor() { + super({ streaming: true, interimResults: true }); + } + override get model(): string { + return 'primary-model'; + } + override get provider(): string { + return 'fake-provider'; + } + protected async _recognize(): Promise { + throw new Error('not used'); + } + stream(options?: { connOptions?: APIConnectOptions }): SpeechStream { + return new FailingSTTStream(this, undefined, options?.connOptions); + } + } + class ServingSTT extends FakeSTT { + override get model(): string { + return 'secondary-model'; + } + override get provider(): string { + return 'fake-provider'; + } + } + // the fallback stream hands frames to whichever child is current, so the secondary speaks + // on its own rather than from audio the failed primary already consumed + const secondary = new ServingSTT({ + label: 'secondary', + capabilities: { streaming: true, interimResults: true }, + fakeTranscript: 'Hello', + }); + const adapter = new STTFallbackAdapter({ + sttInstances: [new FailingSTT(), secondary], + maxRetryPerSTT: 0, + }); + const turn = await userTurnWith(adapter as unknown as FakeSTT); + expect(turn.attributes[traceTypes.ATTR_GEN_AI_REQUEST_MODEL]).toBe('secondary-model'); + }); + it('normalizes the STT provider to the GenAI registry spelling', async () => { class DisplayNamedSTT extends FakeSTT { override get model(): string { From b934f3233be768c170586e448d4bf22f7439c97b Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 25 Sep 2026 22:21:12 -0700 Subject: [PATCH 4/5] review: the STT fallback adapter credits a transcript to the child serving it model/provider name the instance serving the open stream (or the one that answered the last recognize()), and only fall back to the next-in-line instance between streams. A recovery probe finding the primary back no longer relabels a turn the fallback is transcribing. Co-Authored-By: Claude Fable 5.1 --- agents/src/stt/fallback_adapter.test.ts | 44 +++++++++++++++++++++---- agents/src/stt/fallback_adapter.ts | 26 +++++++++------ 2 files changed, 54 insertions(+), 16 deletions(-) diff --git a/agents/src/stt/fallback_adapter.test.ts b/agents/src/stt/fallback_adapter.test.ts index e2c6845bcd..a07f6a473b 100644 --- a/agents/src/stt/fallback_adapter.test.ts +++ b/agents/src/stt/fallback_adapter.test.ts @@ -550,12 +550,10 @@ describe('FallbackSpeechStream (streaming path)', () => { }); describe('FallbackAdapter dynamic model/provider getters', () => { - // The OTel `gen_ai.request.model` / `gen_ai.provider.name` attributes on - // the `user_turn` span are refreshed on every STT event by - // `audio_recognition.refreshUserTurnSttAttributes`. Without dynamic - // getters that reflect the active child, those attributes are frozen at - // the static wrapper labels (`FallbackAdapter` / `livekit`) regardless - // of which provider actually transcribed, so a mid-turn fallover is + // The OTel `gen_ai.request.model` / `gen_ai.provider.name` attributes on the `user_turn` + // span read these when the turn opens and when it ends. Without dynamic getters that + // reflect the child serving, those attributes would name the wrapper (`FallbackAdapter` / + // `livekit`) regardless of which provider actually transcribed, so a fallover would be // invisible in traces. class IdentifiedFakeSTT extends FakeSTT { @@ -647,6 +645,40 @@ describe('FallbackAdapter dynamic model/provider getters', () => { expect(adapter.provider).toBe('fallback-provider'); }); + it('keeps naming the child serving the stream after the primary recovers', async () => { + const primary = new IdentifiedFakeSTT({ + label: 'primary', + model: 'primary-model', + provider: 'primary-provider', + }); + const fallback = new IdentifiedFakeSTT({ + label: 'fallback', + model: 'fallback-model', + provider: 'fallback-provider', + fakeTranscript: 'hello world', + }); + const adapter = new FallbackAdapter({ sttInstances: [primary, fallback] }); + adapter.status[0]!.available = false; + + const stream = adapter.stream(); + expect(adapter.model).toBe('fallback-model'); + stream.pushFrame(emptyAudioFrame()); + stream.endInput(); + let events = 0; + for await (const ev of stream) { + if (ev.type !== SpeechEventType.FINAL_TRANSCRIPT) continue; + events += 1; + // a recovery probe finds the primary back while the fallback still serves this stream: + // the transcript in flight is the fallback's + adapter.status[0]!.available = true; + expect(adapter.model).toBe('fallback-model'); + expect(adapter.provider).toBe('fallback-provider'); + } + expect(events).toBeGreaterThan(0); + // the stream over, the next request goes to the recovered primary + expect(adapter.model).toBe('primary-model'); + }); + it('reflects the active child once streaming events flow', async () => { const primary = new IdentifiedFakeSTT({ label: 'primary', diff --git a/agents/src/stt/fallback_adapter.ts b/agents/src/stt/fallback_adapter.ts index edade1edc5..97a4be9080 100644 --- a/agents/src/stt/fallback_adapter.ts +++ b/agents/src/stt/fallback_adapter.ts @@ -156,11 +156,13 @@ export class FallbackAdapter extends STT { this.setupEventForwarding(); } - // Reflect the active child's model/provider so OTel `gen_ai.request.model` - // and `gen_ai.provider.name` on `user_turn` spans identify the provider - // that actually transcribed, not the static wrapper. `audio_recognition. - // refreshUserTurnSttAttributes` re-reads these on every STT event, so a - // mid-turn fallover surfaces the new child immediately. + /** + * @internal The instance serving the open stream, or the one that answered the last + * recognize(): what a transcript in flight is credited to. A stream stays on the child it + * elected until that child finishes or fails, whatever a recovery probe finds meanwhile. + */ + _servedStt?: STT; + /** * The instance the next request goes to first: the first one marked available, or the primary * once all are down (they are then all retried, primary first). A failed instance's recovery @@ -173,16 +175,17 @@ export class FallbackAdapter extends STT { } /** - * The model of the instance that serves next (see `nextInstance`). Spans and metrics read - * this, so a failover shows the model expected to answer rather than the adapter. + * The model of the instance serving (see `_servedStt`), else of the one that serves next (see + * `nextInstance`). Spans and metrics read this, so a failover shows the model that answered, + * or is expected to, rather than the adapter. */ override get model(): string { - return this.nextInstance().model; + return (this._servedStt ?? this.nextInstance()).model; } - /** The provider of the instance that serves next (see {@link model}). */ + /** The provider of the instance serving, else of the one that serves next (see {@link model}). */ override get provider(): string { - return this.nextInstance().provider; + return (this._servedStt ?? this.nextInstance()).provider; } /** @@ -289,6 +292,7 @@ export class FallbackAdapter extends STT { if (status.available || allFailed) { try { const result = await stt.recognize(frame, abortSignal); + this._servedStt = stt; return result; } catch (e) { this._logger.warn( @@ -577,6 +581,7 @@ class FallbackSpeechStream extends SpeechStream { // Keep child timestamps anchored to the parent stream's current retry attempt. child.startTimeOffset = this.startTimeOffset + (Date.now() - startTime) / 1000; mainRef.current = child; + this.fallbackAdapter._servedStt = sttInstance; try { if (this.abortSignal.aborted) return; // If the forwarder has already drained and exited (input EOF), it @@ -643,6 +648,7 @@ class FallbackSpeechStream extends SpeechStream { const tasks = [forwarderTask, ...this.recoveringStreams.values()]; closeStreams(); await cancelAndWait(tasks, 1000); + this.fallbackAdapter._servedStt = undefined; } } } From b9709d231b52072528f8003f4bcf717cafa2870d Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 25 Sep 2026 22:51:57 -0700 Subject: [PATCH 5/5] review: the served STT slot belongs to the stream that elected it - only the electing stream clears it, so a stream ending cannot erase a newer stream's attribution; the one-slot limit is documented - recognize() no longer pins the slot: between requests the getters follow availability again, as a completed recognize() has no turn in flight Co-Authored-By: Claude Fable 5.1 --- agents/src/stt/fallback_adapter.test.ts | 35 +++++++++++++++++++++++++ agents/src/stt/fallback_adapter.ts | 26 +++++++++--------- 2 files changed, 49 insertions(+), 12 deletions(-) diff --git a/agents/src/stt/fallback_adapter.test.ts b/agents/src/stt/fallback_adapter.test.ts index a07f6a473b..c35d89640a 100644 --- a/agents/src/stt/fallback_adapter.test.ts +++ b/agents/src/stt/fallback_adapter.test.ts @@ -679,6 +679,41 @@ describe('FallbackAdapter dynamic model/provider getters', () => { expect(adapter.model).toBe('primary-model'); }); + it('a stream ending does not clear what a newer stream elected', async () => { + const primary = new IdentifiedFakeSTT({ + label: 'primary', + model: 'primary-model', + provider: 'primary-provider', + fakeTranscript: 'hello', + }); + const fallback = new IdentifiedFakeSTT({ + label: 'fallback', + model: 'fallback-model', + provider: 'fallback-provider', + fakeTranscript: 'hello', + }); + const adapter = new FallbackAdapter({ sttInstances: [primary, fallback] }); + adapter.status[0]!.available = false; + const older = adapter.stream(); + older.pushFrame(emptyAudioFrame()); + expect(adapter.model).toBe('fallback-model'); + // the primary recovers and a second stream elects it while the first still runs + adapter.status[0]!.available = true; + const newer = adapter.stream(); + newer.pushFrame(emptyAudioFrame()); + expect(adapter.model).toBe('primary-model'); + older.endInput(); + for await (const _ of older) { + /* drain */ + } + expect(adapter.model).toBe('primary-model'); + newer.endInput(); + for await (const _ of newer) { + /* drain */ + } + expect(adapter.model).toBe('primary-model'); + }); + it('reflects the active child once streaming events flow', async () => { const primary = new IdentifiedFakeSTT({ label: 'primary', diff --git a/agents/src/stt/fallback_adapter.ts b/agents/src/stt/fallback_adapter.ts index 97a4be9080..5b38c9641b 100644 --- a/agents/src/stt/fallback_adapter.ts +++ b/agents/src/stt/fallback_adapter.ts @@ -157,11 +157,13 @@ export class FallbackAdapter extends STT { } /** - * @internal The instance serving the open stream, or the one that answered the last - * recognize(): what a transcript in flight is credited to. A stream stays on the child it - * elected until that child finishes or fails, whatever a recovery probe finds meanwhile. + * @internal The instance serving the open stream, and the stream that elected it: what a + * transcript in flight is credited to. A stream stays on the child it elected until that + * child finishes or fails, whatever a recovery probe finds meanwhile. One slot: an adapter + * streaming for two sessions at once reports the child elected last (a session has one + * recognition stream, and an adapter is normally one session's). */ - _servedStt?: STT; + _served?: { stt: STT; stream: object }; // the stream, by identity only /** * The instance the next request goes to first: the first one marked available, or the primary @@ -175,17 +177,17 @@ export class FallbackAdapter extends STT { } /** - * The model of the instance serving (see `_servedStt`), else of the one that serves next (see - * `nextInstance`). Spans and metrics read this, so a failover shows the model that answered, - * or is expected to, rather than the adapter. + * The model of the instance serving the open stream (see `_served`), else of the one that + * serves next (see `nextInstance`). Spans and metrics read this, so a failover shows the + * model transcribing, or expected to, rather than the adapter. */ override get model(): string { - return (this._servedStt ?? this.nextInstance()).model; + return (this._served?.stt ?? this.nextInstance()).model; } /** The provider of the instance serving, else of the one that serves next (see {@link model}). */ override get provider(): string { - return (this._servedStt ?? this.nextInstance()).provider; + return (this._served?.stt ?? this.nextInstance()).provider; } /** @@ -292,7 +294,6 @@ export class FallbackAdapter extends STT { if (status.available || allFailed) { try { const result = await stt.recognize(frame, abortSignal); - this._servedStt = stt; return result; } catch (e) { this._logger.warn( @@ -581,7 +582,7 @@ class FallbackSpeechStream extends SpeechStream { // Keep child timestamps anchored to the parent stream's current retry attempt. child.startTimeOffset = this.startTimeOffset + (Date.now() - startTime) / 1000; mainRef.current = child; - this.fallbackAdapter._servedStt = sttInstance; + this.fallbackAdapter._served = { stt: sttInstance, stream: this }; try { if (this.abortSignal.aborted) return; // If the forwarder has already drained and exited (input EOF), it @@ -648,7 +649,8 @@ class FallbackSpeechStream extends SpeechStream { const tasks = [forwarderTask, ...this.recoveringStreams.values()]; closeStreams(); await cancelAndWait(tasks, 1000); - this.fallbackAdapter._servedStt = undefined; + // only this stream's own election; another stream may have elected since + if (this.fallbackAdapter._served?.stream === this) this.fallbackAdapter._served = undefined; } } }