•14 min read

Real-Time Voice AI: Streaming Whisper Transcription with WebRTC & WebSockets in Rust

Real-Time Voice AI: Streaming Whisper Transcription with WebRTC & WebSockets in Rust

This guide details the construction of a sub-200ms real-time audio transcription pipeline. We will cover client audio ingestion via WebRTC and WebSockets, audio chunking, Whisper inference using Rust C bindings, Silero VAD, and streaming partial speech-to-text tokens.

Audio Briefing
0:00 / 0:00

Architecture Overview

The system comprises a client-side WebRTC/WebSocket interface and a Rust-based backend. The client captures audio, encodes it, and transmits it to the server. The server decodes, processes, and transcribes the audio using a quantized Whisper model, leveraging Rust's performance characteristics and C bindings for whisper.cpp. Voice Activity Detection (VAD) is integrated to optimize inference.

Client-Side Audio Capture and Transmission

Client audio capture is handled via WebRTC's MediaDevices.getUserMedia() and MediaRecorder API, or directly via AudioContext for more granular control over processing. For real-time, low-latency streaming, WebSockets are preferred for sending raw audio buffers. WebRTC offers peer-to-peer capabilities and built-in media handling, but for a server-centric transcription service, a WebSocket connection for audio data is often simpler to manage and scale.

We'll use AudioContext to capture raw PCM data, resample it to 16kHz, and send it over a WebSocket.

// client/src/audioStreamer.ts
export class AudioStreamer {
  private audioContext: AudioContext | null = null;
  private mediaStream: MediaStream | null = null;
  private audioInput: MediaStreamAudioSourceNode | null = null;
  private processor: AudioWorkletNode | ScriptProcessorNode | null = null;
  private ws: WebSocket | null = null;
  private readonly sampleRate = 16000; // Target sample rate for Whisper
  private readonly bufferSize = 4096; // Audio buffer size for processing

  constructor(private websocketUrl: string) {}

  public async start(): Promise<void> {
    if (this.audioContext) {
      console.warn("AudioStreamer already started.");
      return;
    }

    this.audioContext = new (window.AudioContext || (window as any).webkitAudioContext)({
      sampleRate: this.sampleRate,
    });

    try {
      this.mediaStream = await navigator.mediaDevices.getUserMedia({ audio: true });
      this.audioInput = this.audioContext.createMediaStreamSource(this.mediaStream);

      // Use AudioWorklet for better performance and off-main-thread processing
      // Fallback to ScriptProcessorNode if AudioWorklet is not available
      if (this.audioContext.audioWorklet) {
        await this.audioContext.audioWorklet.addModule('/audio-processor.js');
        this.processor = new AudioWorkletNode(this.audioContext, 'audio-processor');
      } else {
        console.warn("AudioWorklet not supported, falling back to ScriptProcessorNode.");
        this.processor = this.audioContext.createScriptProcessor(this.bufferSize, 1, 1);
        (this.processor as ScriptProcessorNode).onaudioprocess = this.handleAudioProcess.bind(this);
      }

      this.audioInput.connect(this.processor);
      this.processor.connect(this.audioContext.destination); // Connect to destination to keep it alive

      this.ws = new WebSocket(this.websocketUrl);
      this.ws.binaryType = 'arraybuffer';

      this.ws.onopen = () => console.log('WebSocket connected.');
      this.ws.onclose = () => console.log('WebSocket disconnected.');
      this.ws.onerror = (error) => console.error('WebSocket error:', error);

      if (this.processor instanceof AudioWorkletNode) {
        this.processor.port.onmessage = (event) => {
          if (this.ws?.readyState === WebSocket.OPEN) {
            this.ws.send(event.data); // Send raw PCM float32 data
          }
        };
      }

    } catch (error) {
      console.error('Error starting audio stream:', error);
      this.stop();
      throw error;
    }
  }

  private handleAudioProcess(event: AudioProcessingEvent): void {
    if (this.ws?.readyState === WebSocket.OPEN) {
      const inputBuffer = event.inputBuffer.getChannelData(0);
      this.ws.send(inputBuffer.buffer); // Send raw PCM float32 data
    }
  }

  public stop(): void {
    if (this.processor) {
      this.processor.disconnect();
      this.processor = null;
    }
    if (this.audioInput) {
      this.audioInput.disconnect();
      this.audioInput = null;
    }
    if (this.mediaStream) {
      this.mediaStream.getTracks().forEach(track => track.stop());
      this.mediaStream = null;
    }
    if (this.audioContext) {
      this.audioContext.close();
      this.audioContext = null;
    }
    if (this.ws) {
      this.ws.close();
      this.ws = null;
    }
    console.log('AudioStreamer stopped.');
  }
}

// client/public/audio-processor.js (AudioWorklet module)
// This file needs to be served by your web server
class AudioProcessor extends AudioWorkletProcessor {
  process(inputs, outputs, parameters) {
    const input = inputs[0];
    if (input.length > 0) {
      const channelData = input[0]; // Get the first channel (mono)
      this.port.postMessage(channelData.buffer, [channelData.buffer]); // Transfer ownership
    }
    return true; // Keep the processor alive
  }
}
registerProcessor('audio-processor', AudioProcessor);

Rust Backend: WebSockets, Audio Processing, and Inference

The Rust backend will use tokio for asynchronous operations, warp or axum for the WebSocket server, and hound for WAV encoding (if needed for debugging/storage). For Whisper, we'll use whisper-rs (Rust bindings for whisper.cpp) and silero-vad for VAD.

Dependencies

# Cargo.toml
[dependencies]
tokio = { version = "1", features = ["full"] }
warp = "0.3" # Or axum = { version = "0.6", features = ["ws"] }
futures-util = "0.3"
bytes = "1"
log = "0.4"
env_logger = "0.10"
# For whisper.cpp bindings
whisper-rs = { version = "0.1.1", features = ["full"] } # Ensure whisper.cpp is built with GGML_CUDA=1 if using GPU
# For VAD
silero-vad = "0.1.0" # Or a custom VAD implementation
# For audio processing
symphonia = { version = "0.5", features = ["all"] } # For potential audio decoding/resampling if not 16kHz PCM
# For serialization/deserialization
serde = { version = "1", features = ["derive"] }
serde_json = "1"

Backend Structure

The core logic involves:

  1. WebSocket Server: Accepts client connections and handles incoming audio frames.
  2. Audio Buffer Management: Accumulates incoming audio frames into a larger buffer suitable for VAD and Whisper.
  3. VAD: Detects speech segments to trigger Whisper inference.
  4. Whisper Inference: Transcribes detected speech.
  5. Result Streaming: Sends partial and final transcription results back to the client.
// src/main.rs
use tokio::sync::mpsc;
use tokio_stream::wrappers::ReceiverStream;
use futures_util::{StreamExt, SinkExt};
use warp::ws::{Message, WebSocket};
use warp::Filter;
use std::sync::{Arc, Mutex};
use std::collections::VecDeque;
use std::time::{Instant, Duration};

// --- Whisper and VAD related imports ---
use whisper_rs::{FullParams, SamplingStrategy, WhisperContext, WhisperContextParameters};
use silero_vad::{Vad, VadNode};

const SAMPLE_RATE: u32 = 16000;
const WHISPER_MODEL_PATH: &str = "path/to/ggml-medium.en.bin"; // Path to your quantized Whisper model
const VAD_MODEL_PATH: &str = "path/to/silero_vad.onnx"; // Path to your Silero VAD ONNX model

// Audio buffer for a single client
struct ClientAudioBuffer {
    buffer: VecDeque<f32>,
    last_audio_activity: Instant,
    is_speaking: bool,
}

impl ClientAudioBuffer {
    fn new() -> Self {
        ClientAudioBuffer {
            buffer: VecDeque::new(),
            last_audio_activity: Instant::now(),
            is_speaking: false,
        }
    }

    fn push_audio(&mut self, audio_data: &[f32]) {
        self.buffer.extend(audio_data);
        self.last_audio_activity = Instant::now();
    }

    fn get_audio_chunk(&mut self, duration_ms: u64) -> Option<Vec<f32>> {
        let samples_needed = (SAMPLE_RATE as f32 * duration_ms as f32 / 1000.0) as usize;
        if self.buffer.len() >= samples_needed {
            let chunk: Vec<f32> = self.buffer.drain(0..samples_needed).collect();
            Some(chunk)
        } else {
            None
        }
    }

    fn clear_buffer(&mut self) {
        self.buffer.clear();
    }
}

#[tokio::main]
async fn main() {
    env_logger::init();

    // Load Whisper model once
    let whisper_context = Arc::new(
        WhisperContext::new_with_params(
            WHISPER_MODEL_PATH,
            WhisperContextParameters::default(),
        )
        .expect("Failed to load Whisper model"),
    );
    log::info!("Whisper model loaded: {}", WHISPER_MODEL_PATH);

    // Load VAD model once
    let vad_model = Arc::new(
        Vad::builder()
            .set_model_path(VAD_MODEL_PATH)
            .set_sample_rate(SAMPLE_RATE)
            .build()
            .expect("Failed to load Silero VAD model"),
    );
    log::info!("Silero VAD model loaded: {}", VAD_MODEL_PATH);


    let whisper_context_filter = warp::any().map(move || Arc::clone(&whisper_context));
    let vad_model_filter = warp::any().map(move || Arc::clone(&vad_model));

    let ws_route = warp::path("ws")
        .and(warp::ws())
        .and(whisper_context_filter)
        .and(vad_model_filter)
        .map(|ws: warp::ws::Ws, whisper_ctx: Arc<WhisperContext>, vad_model: Arc<Vad>| {
            ws.on_upgrade(move |websocket| handle_websocket(websocket, whisper_ctx, vad_model))
        });

    let routes = ws_route.with(warp::log("websocket_server"));

    log::info!("Server started on 127.0.0.1:8080");
    warp::serve(routes).run(([127, 0, 0, 1], 8080)).await;
}

async fn handle_websocket(
    websocket: WebSocket,
    whisper_ctx: Arc<WhisperContext>,
    vad_model: Arc<Vad>,
) {
    let (mut client_ws_tx, mut client_ws_rx) = websocket.split();

    let (audio_tx, audio_rx) = mpsc::channel::<Vec<f32>>(100); // Channel for raw audio chunks
    let (transcription_tx, mut transcription_rx) = mpsc::channel::<String>(10); // Channel for transcription results

    // Spawn a task to send transcription results back to the client
    tokio::spawn(async move {
        while let Some(transcription) = transcription_rx.recv().await {
            if let Err(e) = client_ws_tx.send(Message::text(transcription)).await {
                log::error!("Failed to send transcription to client: {}", e);
                break;
            }
        }
        log::info!("Transcription sender task terminated.");
    });

    // Spawn a task to process audio and run VAD/Whisper
    let whisper_ctx_clone = Arc::clone(&whisper_ctx);
    let vad_model_clone = Arc::clone(&vad_model);
    tokio::spawn(async move {
        let mut client_audio_buffer = ClientAudioBuffer::new();
        let mut vad_node = VadNode::new(vad_model_clone, SAMPLE_RATE, 512); // VAD processing frame size

        let mut current_speech_buffer: Vec<f32> = Vec::new();
        let mut last_vad_activity = Instant::now();
        let vad_timeout = Duration::from_secs(2); // How long to wait after speech ends before transcribing

        let mut whisper_session = whisper_ctx_clone.create_state().expect("Failed to create Whisper state");

        while let Some(audio_chunk) = audio_rx.recv().await {
            client_audio_buffer.push_audio(&audio_chunk);

            // Process audio in smaller VAD-friendly chunks
            while let Some(vad_chunk) = client_audio_buffer.get_audio_chunk(30) { // 30ms chunks for VAD
                let speech_prob = vad_node.process(&vad_chunk).expect("VAD processing failed");

                if speech_prob > 0.5 { // Threshold for speech detection
                    current_speech_buffer.extend(vad_chunk);
                    last_vad_activity = Instant::now();
                    client_audio_buffer.is_speaking = true;
                } else {
                    // If not speaking, but we were recently, check for timeout
                    if client_audio_buffer.is_speaking && last_vad_activity.elapsed() > vad_timeout {
                        // Speech has ended, transcribe the accumulated buffer
                        if !current_speech_buffer.is_empty() {
                            log::info!("Speech ended, transcribing {} samples.", current_speech_buffer.len());
                            let transcription = run_whisper_inference(
                                &whisper_ctx_clone,
                                &mut whisper_session,
                                &current_speech_buffer,
                            );
                            if let Err(e) = transcription_tx.send(transcription).await {
                                log::error!("Failed to send final transcription: {}", e);
                            }
                            current_speech_buffer.clear();
                        }
                        client_audio_buffer.is_speaking = false;
                    } else if client_audio_buffer.is_speaking {
                        // Still within timeout, keep accumulating non-speech for context
                        current_speech_buffer.extend(vad_chunk);
                    }
                }
            }

            // Periodically transcribe partial results if speaking
            if client_audio_buffer.is_speaking && current_speech_buffer.len() > (SAMPLE_RATE as usize * 1) { // Transcribe every 1 second of speech
                let partial_transcription = run_whisper_inference(
                    &whisper_ctx_clone,
                    &mut whisper_session,
                    &current_speech_buffer,
                );
                if let Err(e) = transcription_tx.send(format!("[partial] {}", partial_transcription)).await {
                    log::error!("Failed to send partial transcription: {}", e);
                }
            }
        }
        // Handle any remaining speech in buffer when audio stream ends
        if !current_speech_buffer.is_empty() {
            log::info!("Stream ended, transcribing remaining {} samples.", current_speech_buffer.len());
            let transcription = run_whisper_inference(
                &whisper_ctx_clone,
                &mut whisper_session,
                &current_speech_buffer,
            );
            if let Err(e) = transcription_tx.send(transcription).await {
                log::error!("Failed to send final transcription on stream end: {}", e);
            }
        }
        log::info!("Audio processor task terminated.");
    });

    // Receive audio data from client
    while let Some(result) = client_ws_rx.next().await {
        match result {
            Ok(msg) => {
                if msg.is_binary() {
                    let audio_bytes = msg.as_bytes();
                    // Assuming client sends f32 raw PCM
                    let audio_data: Vec<f32> = audio_bytes
                        .chunks_exact(4)
                        .map(|chunk| f32::from_le_bytes(chunk.try_into().unwrap()))
                        .collect();

                    if let Err(e) = audio_tx.send(audio_data).await {
                        log::error!("Failed to send audio chunk to processor: {}", e);
                        break;
                    }
                } else if msg.is_text() {
                    log::debug!("Received text message from client: {}", msg.to_str().unwrap_or_default());
                    // Handle control messages if any
                }
            }
            Err(e) => {
                log::error!("WebSocket receive error: {}", e);
                break;
            }
        }
    }
    log::info!("Client WebSocket disconnected.");
}

fn run_whisper_inference(
    ctx: &WhisperContext,
    state: &mut whisper_rs::WhisperState,
    audio_data: &[f32],
) -> String {
    let mut params = FullParams::new(SamplingStrategy::Greedy { best_of: 1 });
    params.set_print_progress(false);
    params.set_print_special(false);
    params.set_print_realtime(false);
    params.set_print_timestamps(false);
    params.set_language(Some("en"));
    params.set_n_threads(4); // Adjust based on CPU cores

    // Run the inference
    state.full(params, audio_data).expect("Failed to run Whisper inference");

    // Iterate over the segments and collect the text
    let mut result = String::new();
    let num_segments = state.full_n_segments().expect("Failed to get number of segments");
    for i in 0..num_segments {
        let text = state.full_get_segment_text(i).expect("Failed to get segment text");
        result.push_str(&text);
    }
    result.trim().to_string()
}

Performance Considerations and Tradeoffs

Feature/MetricWebRTC Data ChannelWebSocket (Raw PCM)whisper.cpp (Quantized)Silero VAD
LatencyLow (P2P)Low (Client-Server)Medium (depends on model size/CPU/GPU)Very Low
ThroughputHighHighN/A (inference time)N/A (inference time)
ComplexityHigh (ICE, SDP, NAT)Medium (Server-Client)Medium (C bindings, model management)Low (ONNX runtime)
Audio QualityNegotiated (Opus)Raw (configurable)Input: 16kHz PCMInput: 16kHz PCM
CPU UsageLow (client-side encoding)Low (client-side encoding)High (CPU/GPU inference)Low
Network OverheadHigher (protocol)Lower (raw data)N/AN/A
ScalabilityHarder (P2P)Easier (load balancing)Vertical scaling (GPU)Horizontal scaling (more instances)
DeploymentComplexSimplerRequires model filesRequires ONNX model

Latency Breakdown:

  • Client Audio Capture & Encoding: ~10-30ms (browser buffer, AudioContext processing).
  • Network Latency (Client to Server): ~10-100ms (depends on network conditions).
  • Server Audio Buffering: ~30-100ms (to accumulate enough audio for VAD/Whisper).
  • VAD Processing: ~5-10ms per chunk.
  • Whisper Inference: ~50-500ms (depends on model size, hardware, audio chunk length). For sub-200ms, this is the critical path. Using ggml-tiny.en or ggml-base.en on a capable CPU or GPU is essential.
  • Network Latency (Server to Client): ~10-100ms.
  • Client Rendering: ~10ms.

Achieving sub-200ms end-to-end latency requires aggressive chunking, fast VAD, and highly optimized Whisper inference (e.g., ggml-tiny.en on a powerful CPU or GPU).

Production Gotchas & Troubleshooting

  1. Whisper Model Loading Failures:
    • Symptom: Failed to load Whisper model or whisper_init_from_file: failed to open
    • Cause: Incorrect path to ggml-*.bin model file, file permissions, or corrupted model.
    • Fix: Double-check WHISPER_MODEL_PATH. Ensure the Rust process has read access. Download the model again if corrupted. Ensure whisper-rs is built against a compatible whisper.cpp version.
  2. VAD Model Loading Failures:
    • Symptom: Failed to load Silero VAD model
    • Cause: Incorrect path to silero_vad.onnx, missing ONNX runtime dependencies, or incompatible ONNX model version.
    • Fix: Verify VAD_MODEL_PATH. Ensure onnxruntime is correctly installed and linked if silero-vad relies on it (it typically bundles a minimal version).
  3. Audio Resampling/Format Mismatch:
    • Symptom: Garbled transcription, no transcription, or whisper_full: invalid audio length errors.
    • Cause: Client sending audio at a different sample rate or format (e.g., 48kHz, int16) than expected (16kHz, f32).
    • Fix: Ensure client AudioContext is configured for 16kHz. Verify f32 conversion on the server. If client sends int16, convert to f32 on the server.
      // Example: converting i16 to f32
      fn convert_i16_to_f32(audio_data_i16: &[i16]) -> Vec<f32> {
          audio_data_i16.iter().map(|&s| s as f32 / 32768.0).collect()
      }
      
  4. High Latency / Slow Transcription:
    • Symptom: Transcription appears significantly delayed.
    • Cause: Large audio chunks sent to Whisper, slow CPU/GPU, large Whisper model (ggml-large), insufficient threads for Whisper, or VAD not triggering inference quickly enough.
    • Fix:
      • Use smaller Whisper models (ggml-tiny.en, ggml-base.en).
      • Increase params.set_n_threads() for Whisper (up to physical core count).
      • Optimize VAD parameters (vad_timeout, get_audio_chunk size) to trigger inference faster.
      • Ensure whisper.cpp is compiled with GPU support (e.g., GGML_CUDA=1) if a GPU is available and whisper-rs is configured to use it.
  5. WebSocket Disconnections:
    • Symptom: Client or server logs show frequent WebSocket close events.
    • Cause: Network instability, server overload, unhandled errors in WebSocket handler, or client-side AudioContext being garbage collected if not connected to destination.
    • Fix: Implement client-side reconnection logic. Ensure server-side error handling is robust. Keep AudioContext connected to destination or a GainNode connected to destination to prevent it from being culled. Implement WebSocket heartbeats (ping/pong) to detect dead connections.
  6. Memory Leaks:
    • Symptom: Server memory usage steadily increases over time.
    • Cause: Audio buffers not being cleared, Whisper WhisperState not being properly managed, or VecDeque growing indefinitely.
    • Fix: Ensure ClientAudioBuffer.clear_buffer() is called when appropriate. whisper-rs WhisperState should be reused per session. Monitor VecDeque size and implement limits if necessary.

Frequently Asked Questions

  1. Why use WebSockets instead of WebRTC for audio streaming? While WebRTC offers peer-to-peer capabilities and built-in media handling, for a server-centric transcription service, WebSockets often simplify the architecture. WebRTC's signaling, ICE negotiation, and NAT traversal add significant complexity. WebSockets provide a full-duplex, low-latency channel suitable for sending raw audio buffers directly to a central server for processing, which is easier to scale horizontally.

  2. How do I achieve sub-200ms end-to-end latency? This is challenging. Key strategies include:

    • Client-side: Capture and send small audio chunks (e.g., 30-50ms) immediately.
    • Network: Minimize network latency (e.g., deploy server geographically close to users).
    • Server-side:
      • Use a highly optimized, quantized Whisper model (e.g., ggml-tiny.en or ggml-base.en).
      • Leverage GPU acceleration for Whisper inference if available.
      • Employ aggressive VAD to only transcribe speech segments, minimizing Whisper calls.
      • Process audio in small, fixed-size chunks for VAD and then accumulate for Whisper.
      • Utilize Rust's performance and whisper.cpp's C bindings for minimal overhead.
      • Stream partial results as soon as they are available.
  3. Can I use a larger Whisper model for better accuracy? Yes, but at the cost of increased latency. Larger models like ggml-medium.en or ggml-large.en provide higher accuracy but require significantly more computational resources and time for inference. For real-time applications with strict latency requirements, a smaller model is often a necessary compromise. Consider using a smaller model for real-time partial results and a larger model for final, post-processed transcription if accuracy is paramount.

  4. How does VAD improve the transcription pipeline? Voice Activity Detection (VAD) is crucial for real-time performance and resource efficiency. It identifies segments of speech within the audio stream, allowing the Whisper model to only process relevant audio. This reduces the number of Whisper inference calls, saves CPU/GPU cycles, and prevents transcribing silence or background noise, leading to faster and cleaner results. It also helps in segmenting continuous speech into meaningful utterances.

  5. What if my client audio is not 16kHz PCM? The Whisper model expects 16kHz mono PCM f32 audio. If your client sends audio in a different format (e.g., 48kHz, stereo, int16), you must resample and convert it on either the client or server side. Performing this on the client (using AudioContext's sampleRate property) offloads work from the server. If done on the server, libraries like symphonia or rubato can handle resampling and format conversion efficiently.

Share this article:

Stay Updated

Get the latest posts delivered straight to your inbox.

Free Developer Utilities

Free In-Browser Developer Tools

Clean AI CLI logs, build cron expressions, decode JWTs, and calculate chmod permissions offline.

Explore Tools
Advertisement