Skip to content
5 changes: 5 additions & 0 deletions .changeset/clean-fallback-teardown.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@livekit/agents': patch
---

Keep healthy STT providers available when fallback streams close with transcripts in flight.
10 changes: 6 additions & 4 deletions agents/etc/agents.api.md
Original file line number Diff line number Diff line change
Expand Up @@ -7699,6 +7699,8 @@ abstract class SpeechStream implements AsyncIterableIterator<SpeechEvent> {
// (undocumented)
detachInputStream(): void;
endInput(): void;
// @internal
get _failed(): boolean;
flush(): void;
// (undocumented)
protected static readonly FLUSH_SENTINEL: unique symbol;
Expand Down Expand Up @@ -10149,13 +10151,13 @@ export const zipFunctionCallsAndOutputs: (event: FunctionToolsExecutedEvent) =>
// src/llm/tool_context.ts:746:3 - (ae-unresolved-link) The @link reference could not be resolved: The reference is ambiguous because "ToolFlag" has more than one declaration; you need to add a TSDoc member reference selector
// src/metrics/base.ts:198:3 - (ae-forgotten-export) The symbol "RealtimeModelMetricsInputTokenDetails" needs to be exported by the entry point index.d.ts
// src/metrics/base.ts:202:3 - (ae-forgotten-export) The symbol "RealtimeModelMetricsOutputTokenDetails" needs to be exported by the entry point index.d.ts
// src/stt/stt.ts:364:3 - (ae-unresolved-link) The @link reference could not be resolved: The package "@livekit/agents" does not have an export "STT"
// src/stt/stt.ts:365:3 - (ae-unresolved-link) The @link reference could not be resolved: The package "@livekit/agents" does not have an export "STT"
// src/utils.ts:550:3 - (ae-unresolved-link) The @link reference could not be resolved: The package "@livekit/agents" does not have an export "cancelled"
// src/voice/agent_session.ts:387:3 - (ae-unresolved-link) The @link reference could not be resolved: This type of declaration is not supported yet by the resolver
// src/voice/agent_session.ts:1025:5 - (ae-forgotten-export) The symbol "RecordingOptions" needs to be exported by the entry point index.d.ts
// src/voice/agent_session.ts:1695:5 - (ae-forgotten-export) The symbol "STTError" needs to be exported by the entry point index.d.ts
// src/voice/agent_session.ts:1695:5 - (ae-forgotten-export) The symbol "TTSError" needs to be exported by the entry point index.d.ts
// src/voice/agent_session.ts:1695:5 - (ae-forgotten-export) The symbol "LLMError" needs to be exported by the entry point index.d.ts
// src/voice/agent_session.ts:1696:5 - (ae-forgotten-export) The symbol "STTError" needs to be exported by the entry point index.d.ts
// src/voice/agent_session.ts:1696:5 - (ae-forgotten-export) The symbol "TTSError" needs to be exported by the entry point index.d.ts
// src/voice/agent_session.ts:1696:5 - (ae-forgotten-export) The symbol "LLMError" needs to be exported by the entry point index.d.ts
// src/voice/amd.ts:315:3 - (ae-unresolved-link) The @link reference could not be resolved: The reference is ambiguous because "waitForTrackPublication" has more than one declaration; you need to add a TSDoc member reference selector
// src/voice/amd.ts:315:3 - (ae-unresolved-link) The @link reference could not be resolved: The package "@livekit/agents" does not have an export "gateListening"
// src/voice/amd.ts:323:3 - (ae-unresolved-link) The @link reference could not be resolved: The package "@livekit/agents" does not have an export "aclose"
Expand Down
76 changes: 76 additions & 0 deletions agents/src/stt/fallback_adapter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -314,6 +314,82 @@ describe('FallbackSpeechStream (streaming path)', () => {
expect(adapter.status[1]?.available).toBe(true);
});

it('keeps the provider available when a late transcript arrives after stream close', async () => {
const primary = new FakeSTT({
label: 'primary',
fakeTranscript: 'late transcript',
fakeTimeoutMs: 50,
});
const adapter = new FallbackAdapter({
sttInstances: [primary],
maxRetryPerSTT: 0,
});

const availabilityChanges: Array<{ stt: STT; available: boolean }> = [];
(adapter as unknown as EventEmitter).on(
'stt_availability_changed',
(ev: { stt: STT; available: boolean }) => {
availabilityChanges.push(ev);
},
);

const stream = adapter.stream();
await primary.streamCh.next();
stream.close();
await delay(150);

expect(availabilityChanges).toEqual([]);
expect(adapter.status[0]?.available).toBe(true);

await adapter.close();
});

it('closes recovery probes when a late transcript arrives after stream close', async () => {
const primary = new FakeSTT({
label: 'primary',
fakeException: new APIError('primary down'),
});
const fallback = new FakeSTT({
label: 'fallback',
fakeTranscript: 'late transcript',
fakeTimeoutMs: 50,
});
const adapter = new FallbackAdapter({
sttInstances: [primary, fallback],
maxRetryPerSTT: 0,
});

(adapter as unknown as EventEmitter).on(
'stt_availability_changed',
(ev: { stt: STT; available: boolean }) => {
if (ev.stt === primary && !ev.available) {
// Keep the recovery probe alive until the parent stream tears it down.
primary.updateOptions({ fakeException: null });
}
},
);

const stream = adapter.stream();
stream.endInput();

await primary.streamCh.next(); // failed main stream
const recoveryProbe = (await primary.streamCh.next()).value;
await fallback.streamCh.next();

stream.close();
await delay(150);

let recoveryProbeClosed = false;
try {
recoveryProbe?.endInput();
} catch {
recoveryProbeClosed = true;
}
await adapter.close();

expect(recoveryProbeClosed).toBe(true);
});

it('stream switches to the secondary provider when the primary errors', async () => {
const primary = new FakeSTT({
label: 'primary',
Expand Down
Loading
Loading