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

> Learn how the Hugging Face Speech-to-Speech pipeline manages concurrent WebSocket connections with a fixed pool of PipelineUnits for thread-safe, real-time audio processing.

- Repository: [Hugging Face/speech-to-speech](https://github.com/huggingface/speech-to-speech)
- Tags: internals
- Published: 2026-07-30

---

**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`](https://github.com/huggingface/speech-to-speech/blob/main/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`](https://github.com/huggingface/speech-to-speech/blob/main/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).

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

```

2. **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`](https://github.com/huggingface/speech-to-speech/blob/main/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`](https://github.com/huggingface/speech-to-speech/blob/main/scripts/listen_and_play_realtime.py) demonstrates creating a pool of four units:

```python
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:

```python
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:

```python
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`](https://github.com/huggingface/speech-to-speech/blob/main/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`](https://github.com/huggingface/speech-to-speech/blob/main/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.