# Handling Concurrent WebSocket Connections in Realtime Mode with PipelineUnit Pools

> Discover how huggingface speech-to-speech handles concurrent WebSocket connections using isolated PipelineUnit pools for efficient realtime processing and full session isolation.

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

---

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

```python
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`](https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/api/openai_realtime/websocket_router.py). This function iterates over the pool atomically to reserve the first available unit:

```python
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`](https://github.com/huggingface/speech-to-speech/blob/main/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`](https://github.com/huggingface/speech-to-speech/blob/main/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`](https://github.com/huggingface/speech-to-speech/blob/main/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`](https://github.com/huggingface/speech-to-speech/blob/main/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`](https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/api/openai_realtime/server.py):

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

# 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:

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

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