Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/coverage-spans.md
Original file line number Diff line number Diff line change
@@ -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.
76 changes: 76 additions & 0 deletions agents/src/llm/fallback_adapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -264,4 +264,80 @@ 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(
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');
});
});
84 changes: 82 additions & 2 deletions agents/src/llm/fallback_adapter.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand Down Expand Up @@ -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;
Comment on lines +124 to +130

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Recovered primary misattributes fallback usage

When primary recovery finishes before metrics finalize, model and provider switch back despite the secondary serving. The secondary's usage is attributed to the primary.

Learn more

Fallback streams emit wrapper-level metrics after output collection. Those metrics read the adapter's dynamic model and provider through monitorMetrics. A failed primary starts recovery immediately, and recovery can mark it available while the secondary request is still running. The getters then select the primary even though servedLlm correctly records the secondary on spans. The same timing can affect wrapper-level TTS metrics because its fallback adapter uses equivalent dynamic getters.

Example: The primary fails, then its probe succeeds while the secondary is generating. The secondary returns 100 tokens. By metric finalization, nextInstance() selects the recovered primary, so those 100 tokens carry the primary model and provider.

Recommended fix: Store the serving child per fallback stream and use that child for wrapper metrics. Extend the base LLM/TTS metric construction with protected model/provider getters, analogous to responseModel, then override them in fallback streams. Keep adapter-level getters for pre-request attribution.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Codex found the same issue:

Step What happens
1 Primary fails. The adapter sends the request to secondary.
2 Secondary starts generating. A background probe checks whether primary has recovered.
3 The probe succeeds. Primary becomes the preferred provider for the next request.
4 Secondary finishes with 100 output tokens.
5 The adapter emits usage metrics. Its model and provider getters now return primary.

}

label(): string {
Expand Down Expand Up @@ -147,13 +170,32 @@ 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;
private toolChoice?: ToolChoice;
private extraKwargs?: Record<string, unknown>;
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,
Expand Down Expand Up @@ -184,6 +226,41 @@ 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;
}

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
* 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.
Expand Down Expand Up @@ -352,6 +429,7 @@ class FallbackLLMStream extends LLMStream {
{ llm: llm.label(), totalChunks: chunkCount, textLength: textSent.length },
'FallbackAdapter: Provider succeeded',
);
this.recordServed(llm, i);
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
return;
} catch (error) {
// Mark as unavailable if it was available before
Expand All @@ -374,6 +452,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;
}

Expand Down
35 changes: 31 additions & 4 deletions agents/src/llm/llm.ts
Original file line number Diff line number Diff line change
Expand Up @@ -247,10 +247,37 @@ export abstract class LLMStream implements AsyncIterableIterator<ChatChunk> {
});
}

/** 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;
}

/** 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;
}

/** 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,
Expand Down Expand Up @@ -401,8 +428,8 @@ export abstract class LLMStream implements AsyncIterableIterator<ChatChunk> {
return (usage?.completionTokens || 0) / (durationMs / 1000);
})(),
metadata: {
modelProvider: this.#llm.provider,
modelName: this.#llm.model,
modelProvider: this.responseProvider,
modelName: this.responseModel,
},
};

Expand All @@ -418,7 +445,7 @@ export abstract class LLMStream implements AsyncIterableIterator<ChatChunk> {
});
genAI.setResponseAttributes(this.#llmRequestSpan, {
responseId: requestId || undefined,
model: this.#llm.model,
model: this.responseModel,
finishReasons: [finishReason],
timeToFirstChunk: metrics.ttftMs >= 0 ? metrics.ttftMs / 1000 : undefined,
});
Expand Down
Loading
Loading