Handling Concurrent WebSocket Connections in Realtime Mode with PipelineUnit Pools

The huggingface/speech-to-speech library manages concurrent WebSocket connections by maintaining a fixed-size pool of isolated PipelineUnit workers, where each unit owns dedicated message queues and a cancellation scope to ensure complete session isolation.

The realtime API server in the huggingface/speech-to-speech repository processes multiple simultaneous audio streams through OpenAI-compatible WebSocket endpoints at /v1/realtime. To safely handle concurrent WebSocket connections in realtime mode, the architecture pre-allocates a pool of independent PipelineUnit objects that isolate each client's session state, processing queues, and lifecycle from all other active connections.

The PipelineUnit Pool Architecture

At server startup, the system initializes a fixed-size pool of PipelineUnit objects defined in src/speech_to_speech/api/openai_realtime/pipeline_unit.py. According to the source code, each unit is a self-contained worker that maintains its own set of queues for audio input, text prompts, audio output, and text-event output, plus a cancellation scope and a mutable SessionState object holding the active WebSocket connection.

This design ensures that every WebSocket session operates within its own isolated execution context. The pool is typically constructed in the server entry point as demonstrated in scripts/listen_and_play_realtime.py:

from speech_to_speech.api.openai_realtime.pipeline_unit import build_realtime_pipeline_unit

# Create a pool of 4 pipeline units

pipeline_pool = [
    build_realtime_pipeline_unit(index=i) for i in range(4)
]

Claiming Units for Incoming Connections

When a client connects to the /v1/realtime endpoint, the router executes the _claim_unit function in src/speech_to_speech/api/openai_realtime/websocket_router.py. This function iterates over the pool atomically to reserve the first available unit:

for unit in pool:
    if unit.session is None:
        unit.session = SessionState(websocket=ws)
        return unit

If every unit in the pool already hosts an active session (lines 299-308), the connection is immediately rejected with a session_limit_reached error event. This backpressure mechanism ensures that the server only accepts connections it can service with dedicated resources.

Isolating Sessions with Dedicated Send Loops

Each PipelineUnit runs an independent background send loop (_send_loop_for) that starts when the FastAPI application lifespan begins (lines 61-64 of websocket_router.py). This loop continuously pulls items from that specific unit's output queues and forwards them exclusively to the WebSocket stored in that unit's SessionState (lines 48-100).

Because every unit owns its own queues, race conditions are architecturally impossible: one client never receives data that was queued by a different client's handlers. The send loop only processes messages from its assigned unit, creating a hard boundary between concurrent sessions.

Processing Realtime Events Through the Pipeline

Inside the realtime_endpoint function (lines 27-71), the server receives JSON events from the WebSocket and dispatches them to a RealtimeService instance. This service routes the data through the STT, LLM, and TTS handlers attached to the same PipelineUnit claimed for that connection. As these handlers process audio and text, they push results onto the unit’s specific queues:

  • unit.input_queue for incoming audio and text prompts
  • unit.output_queue for generated audio and text responses

The unit-bound routing ensures that processing logic for one connection cannot interfere with another, even during high-load scenarios.

Graceful Teardown and Unit Recycling

When a client disconnects or an error aborts the connection, the architecture executes a three-phase cleanup sequence defined in websocket_router.py:

  1. Immediate cleanup: _clean_unit cancels any in-flight work and flushes all queues (lines 86-94).
  2. End-of-session signal: A special SESSION_END control message is enqueued (unit.input_queue.put(SESSION_END)) to allow the signal to propagate through the entire handler chain.
  3. Drain and release: _release_unit_after_drain (lines 44-53) blocks as a background task until the send loop observes the SESSION_END sentinel. Only then is the unit's session attribute cleared, making it eligible for the next _claim_unit call.

This draining process guarantees that no residual messages from the previous session are transmitted to the next client who claims the unit.

Preventing Cross-Session Data Leakage with Cancellation Scopes

The cancellation scope inside each PipelineUnit (see pipeline_unit.py) provides an additional safety layer. Functions like _generation_is_discardable and _should_discard_audio inspect generation identifiers against the current cancel-scope state (lines 5-20 of websocket_router.py). If a generation was started by a previous session but completes after the unit has been reclaimed, the cancellation scope flags it as stale, ensuring that audio from a previous client cannot be emitted to a new connection.

Configuring the Pipeline Pool Size

The pool size determines the maximum number of concurrent WebSocket connections the server can handle simultaneously. Define the pool when creating the FastAPI application in src/speech_to_speech/api/openai_realtime/server.py:

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

# Configure pool for 2 concurrent sessions

pipeline_pool = [
    build_realtime_pipeline_unit(index=i) for i in range(2)
]

app = create_app(pool=pipeline_pool, stop_event=Event())

Clients can then connect using standard WebSocket libraries:

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

        }))
        async for msg in ws:
            print(json.loads(msg))

asyncio.run(realtime_demo())

Summary

  • The architecture uses a pre-allocated pool of PipelineUnit workers to manage concurrent realtime sessions.
  • Each WebSocket connection is atomically bound to a dedicated unit via the _claim_unit function, with explicit rejection when the pool is exhausted.
  • Isolated queues and send loops per unit prevent data leakage between clients.
  • Graceful teardown uses a SESSION_END sentinel and drain logic to ensure units are only recycled after all session data has cleared.
  • Cancellation scopes prevent stale generations from one session affecting the next.

Frequently Asked Questions

What happens when all PipelineUnits are in use?

When the pool is exhausted, the _claim_unit function in websocket_router.py fails to find a unit with session set to None. The server responds with a session_limit_reached error event (lines 300-307), rejecting the new WebSocket connection until an existing unit is released.

How does the system prevent audio from one client leaking to another?

Each PipelineUnit maintains its own set of queues and runs a dedicated _send_loop_for task that only reads from that unit's queues. Because handlers push output exclusively to their bound unit's queues, and the send loop only forwards to that unit's specific WebSocket, cross-session data leakage is prevented by architectural isolation.

What is the purpose of the SESSION_END sentinel?

The SESSION_END message acts as a tombstone signal that travels through the handler chain when a client disconnects. The _release_unit_after_drain function waits until this sentinel reaches the send loop, ensuring all queued messages are flushed before the unit's session is cleared and the unit returned to the available pool.

How many concurrent connections can the server handle?

The limit equals the size of the PipelineUnit pool configured at startup. For example, a pool initialized with range(4) supports exactly four simultaneous realtime sessions; additional connections receive an immediate error until an existing session terminates and its unit is fully drained and released.

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 →