Skip to content

deepgram/xai STT plugin: unexpected WebSocket close is never detected, reconnect never fires (bare Task in Promise.all) #2469

Description

@obaverstam

Summary

SpeechStream#runWS in plugins/deepgram/src/stt.ts never notices when the Deepgram WebSocket closes unexpectedly, so the reconnect loop in run() never fires. The stream just stops producing transcripts, silently, forever, with no error and no retry — until the caller tears it down.

The same bug is present in plugins/xai/src/stt.ts.

Root cause

stt.ts's #runWS:

const wsMonitor = Task.from(async (controller) => {
  const closed = new Promise<never>(async (_, reject) => {
    ws.once('close', (code, reason) => {
      if (!closing) {
        this.#logger.error(`WebSocket closed with code ${code}: ${reason}`);
        reject(new Error('WebSocket closed'));
      }
    });
  });
  await Promise.race([closed, waitForAbort(controller.signal)]);
});

...

await Promise.race([
  this.#resetWS.await,
  Promise.all([sendTask(), listenTask.result, wsMonitor]), // <-- bare wsMonitor
]);

(plugins/deepgram/src/stt.ts, currently line 593 on main.)

Task (from @livekit/agents' utils.ts) is a plain class with no .then() — it is not a thenable. Promise.all/Promise.race resolve a non-thenable array member immediately as a plain value, so wsMonitor's actual completion (wsMonitor.result, which correctly rejects when the socket closes) is never observed by this Promise.all. The close IS detected and logged (logger.error("WebSocket closed with code ...") does fire), but the rejection that's supposed to propagate out and trigger run()'s retry/backoff loop never does.

Why this produces a total hang, not just a delayed recovery: with the close-detector's rejection swallowed, the only other way Promise.all([sendTask(), listenTask.result, wsMonitor]) could settle is if sendTask() or listenTask.result settle on their own:

  • listenTask.result only resolves on a NEW ws.on('message', ...) event — none will ever arrive on a closed socket.
  • sendTask()'s ws.send(frame.data.buffer) is called with no callback. ws's own send() (confirmed in ws@8.21.0's source) only throws synchronously when readyState === CONNECTING; on a closed socket it silently routes to an internal sendAfterClose() that does nothing observable without a callback. So repeated sends into the closed socket never throw and never surface the problem either.

With neither path able to settle, #runWS — and therefore run()'s own await this.#runWS(ws) — hangs indefinitely.

The fix is one token

- Promise.all([sendTask(), listenTask.result, wsMonitor]),
+ Promise.all([sendTask(), listenTask.result, wsMonitor.result]),

Same in plugins/xai/src/stt.ts (currently line 364 on main).

This exact pattern is already correct in three sibling plugins — wsMonitor.result (not bare wsMonitor) — which is strong evidence this is a copy/paste slip rather than intentional:

  • plugins/inworld/src/stt.ts
  • plugins/sarvam/src/stt.ts
  • plugins/assemblyai/src/stt.ts

Minimal reproduction

Standalone, no real Deepgram credentials needed — points the plugin's own baseUrl override at a local WebSocket server instead:

import { WebSocketServer } from 'ws'
import { initializeLogger } from '@livekit/agents'
import { AudioFrame } from '@livekit/rtc-node'
import * as deepgram from '@livekit/agents-plugin-deepgram'

initializeLogger({ level: 'error', pretty: false })

const wss = new WebSocketServer({ port: 0 })
let connections = 0
wss.on('connection', ws => {
  connections++
  if (connections === 1) {
    // simulate a mid-stream drop, same as a real provider-side close
    setTimeout(() => ws.close(1000, 'induced'), 150)
  } else {
    console.log('RECONNECTED — a second connection arrived')
    process.exit(0)
  }
})
await new Promise(resolve => wss.once('listening', resolve))
const port = wss.address().port

const stt = new deepgram.STT({
  apiKey: 'not-real-baseUrl-is-localhost',
  baseUrl: `ws://127.0.0.1:${port}`,
})
const stream = stt.stream()

const frame = AudioFrame.create(16000, 1, 1600, {})
frame.data.fill(2000) // non-silent, clears AudioEnergyFilter's RMS threshold
setInterval(() => {
  try { stream.pushFrame(frame) } catch {}
}, 50)

setTimeout(() => {
  console.log('TIMED OUT — no reconnect after the close')
  process.exit(1)
}, 3000)

On the currently-installed/published version (checked against 1.5.1 and 1.8.0's dist/stt.js — the bug is present in both): this always prints TIMED OUT. With wsMonitor changed to wsMonitor.result, it reliably prints RECONNECTED within the same run (the plugin's own backoff, Math.min(retries * 5, 10) seconds with retries starting at 0, means the very first retry attempt has zero delay).

Scope

Checked, not assumed: plugins/stt_v2.ts-style implementations (Deepgram's STTv2/Flux, and the ones in inworld/sarvam/assemblyai) don't use this Task.from(...) + bare-Promise.all pattern for close detection at all, so they're unaffected by this specific bug — only deepgram's classic STT/SpeechStream (stt.ts, not stt_v2.ts) and xai's STT/SpeechStream have it.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions