How the Hugging Face Speech-to-Speech Pipeline Handles Concurrent WebSocket Connections in Realtime Mode

The Speech-to-Speech library manages concurrent WebSocket connections by pre-allocating a fixed-size pool of independent PipelineUnit objects, binding each new client to a dedicated unit with isolated queues and cancellation scopes to guarantee thread-safe, race-condition-free real-time audio processing.

The huggingface/speech-to-speech repository implements a scalable real-time API using an object-pool architecture that prevents resource contention and cross-session data leakage. When operating in realtime mode, the system handles multiple simultaneous WebSocket clients on the /v1/realtime endpoint by assigning each connection to a dedicated PipelineUnit rather than spawning threads per request. This design ensures that audio streams, text prompts, and generation state remain strictly isolated while the system accepts concurrent sessions up to a configurable pool limit.

Architecture of the Pipeline Unit Pool

At startup, the server constructs a list of PipelineUnit instances defined in src/speech_to_speech/api/openai_realtime/pipeline_unit.py. Each unit is a self-contained worker that encapsulates:

  • Four isolated queues: audio input, text prompts, audio output, and text-event output.
  • A cancellation scope to manage in-flight generations.
  • A mutable SessionState object that holds the active WebSocket connection.

Because each unit owns its own queues, no two clients can read from or write to the same memory buffer, eliminating race conditions by design.

Claiming and Binding Units per Connection

When a client opens a WebSocket, the router in src/speech_to_speech/api/openai_realtime/websocket_router.py executes an atomic claim process:

  1. Iterate the pool: The _claim_unit function scans the pool for the first unit whose session attribute is None (lines 299–308).
for unit in pool:
    if unit.session is None:
        unit.session = SessionState(websocket=ws)
        return unit
  1. Reservation or rejection: If a free unit is found, it is immediately bound to the new WebSocket via SessionState. If every unit is occupied, the server rejects the connection with a session_limit_reached error event (lines 300–307).

This greedy allocation strategy ensures that the system never over-commits resources; the hard pool size acts as a backpressure mechanism.

Per-Unit Send Loops and Event Isolation

Communication from the server to the client happens through dedicated background tasks:

  • Send loop startup: When the FastAPI application lifespan begins, the router starts a background task _send_loop_for for every unit in the pool (lines 61–64).
  • Queue consumption: Each loop continuously pulls items from its unit’s output queues and forwards them exclusively to the WebSocket stored in that unit’s SessionState (lines 48–100).

Because send loops are scoped to individual units, one client never receives data queued by a different session, even under high concurrency.

Processing Inbound Events Through the Pipeline

Inside the realtime_endpoint handler (lines 27–71), the server receives JSON events from the client and dispatches them to a RealtimeService instance. This service routes data through the STT, LLM, and TTS handlers attached to the same PipelineUnit.

Handlers push their results onto the unit’s specific queues (unit.input_queue, unit.output_queue, etc.). The architecture guarantees that the STT output from Client A flows only into Client A’s LLM handler because both are pinned to the same unit object.

Graceful Session Teardown and Unit Reclamation

When a client disconnects or an error aborts the connection, the system performs a three-stage cleanup to prevent data leakage:

  1. Immediate cancellation: _clean_unit cancels any in-flight work and flushes all queues (lines 86–94).
  2. Sentinel propagation: The code enqueues a special SESSION_END control message (unit.input_queue.put(SESSION_END)) so the signal travels through the entire handler chain, ensuring all stages observe the termination.
  3. Drain and release: _release_unit_after_drain blocks until the send loop observes the SESSION_END sentinel (lines 44–53). Only then is the unit’s session cleared, making the unit eligible for the next claim.

The cancellation scope inside each PipelineUnit guarantees that stale generations from a previous session cannot be emitted after a new client has claimed the unit. Helper functions _generation_is_discardable and _should_discard_audio inspect generation identifiers against the current cancel-scope state (lines 5–20 of websocket_router.py) to filter out orphaned audio chunks.

Configuring the Pool Size

The pool size is defined statically when the server launches. The example script scripts/listen_and_play_realtime.py demonstrates creating a pool of four units:

from speech_to_speech.api.openai_realtime.server import create_app
from speech_to_speech.api.openai_realtime.pipeline_unit import build_realtime_pipeline_unit
from threading import Event

# Create a pool of 4 pipeline units

pipeline_pool = [
    build_realtime_pipeline_unit(index=i) for i in range(4)
]
app = create_app(pool=pipeline_pool, stop_event=Event())

Attempting to connect a fifth simultaneous client results in an immediate error response, protecting the server from resource exhaustion.

Practical Client-Server Example

To start a server with a 2-unit pool:

from speech_to_speech.api.openai_realtime.server import create_app
from speech_to_speech.api.openai_realtime.pipeline_unit import build_realtime_pipeline_unit
from threading import Event

pool = [build_realtime_pipeline_unit(index=i) for i in range(2)]
app = create_app(pool=pool, stop_event=Event())

Connecting with a Python client:

import websockets
import json
import asyncio

async def realtime_demo():
    async with websockets.connect("ws://localhost:8000/v1/realtime") as ws:
        await ws.send(json.dumps({
            "type": "input_audio_buffer.append",
            "audio": "base64_encoded_pcm_data...",
        }))
        async for msg in ws:
            event = json.loads(msg)
            print(event)

asyncio.run(realtime_demo())

Summary

  • Pre-allocated pool: The system initializes a fixed-size list of PipelineUnit workers at startup, preventing runtime allocation overhead and capping resource usage.
  • Atomic unit claiming: Each WebSocket connection is bound to a dedicated unit via SessionState assignment, rejecting excess connections with a session_limit_reached error.
  • Isolated I/O loops: Independent send loops per unit read only from that unit’s queues, ensuring no cross-contamination of audio or text events between clients.
  • Draining semantics: The SESSION_END sentinel guarantees that units are released only after all pending data has been flushed, preventing truncation or leakage into subsequent sessions.
  • Cancellation safety: Per-unit cancellation scopes and discard logic ensure that stale generations from disconnected clients cannot corrupt active sessions.

Frequently Asked Questions

What happens when the pool is exhausted?

When all PipelineUnit instances in the pool have an active SessionState, the _claim_unit function in websocket_router.py returns no available unit. The router responds to the new WebSocket request with a session_limit_reached error event and closes the connection immediately, enforcing the concurrency limit without degrading existing sessions.

How does the system prevent audio from one client bleeding into another?

Each PipelineUnit maintains its own set of queues for audio input, text prompts, and output events. The send loop for a unit only reads from that unit’s queues and only writes to the WebSocket stored in its SessionState. Because these objects are never shared between units, there is no shared state that could cause cross-talk between concurrent connections.

Can the pool size be changed dynamically?

No. The pool is defined as a Python list passed to create_app() at server startup, as shown in scripts/listen_and_play_realtime.py. To resize the pool, you must restart the server with a new list of PipelineUnit objects. This static allocation simplifying memory management and eliminates complex auto-scaling logic in the realtime path.

Why is the SESSION_END sentinel necessary for unit release?

The SESSION_END sentinel ensures that the send loop has fully processed and transmitted any queued audio chunks before the unit is marked available. Without this blocking drain mechanism implemented in _release_unit_after_drain, a newly connected client could claim a unit while the previous client’s audio was still queued in the output buffer, resulting in truncated playback or data leakage between sessions.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →