# How Pipeline Pooling Handles Multiple Concurrent WebSocket Connections in Hugging Face Speech-to-Speech

> Discover how pipeline pooling manages concurrent WebSocket connections in Hugging Face Speech-to-Speech. Learn about isolated PipelineUnit instances and strict concurrency control.

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

---

**The Speech-to-Speech service maintains a fixed-size pool of isolated `PipelineUnit` instances, atomically claiming idle units for new WebSocket connections and rejecting requests with error code 1008 when the pool is exhausted, ensuring strict concurrency limits without shared state.**

The huggingface/speech-to-speech repository implements a robust concurrency model for its OpenAI Realtime API-compatible server. Understanding how pipeline pooling handles multiple concurrent WebSocket connections is essential for deploying production workloads that require predictable resource usage and strict session isolation.

## Architecture of the Pipeline Pool

When the server initializes, `RealtimeServer.run` constructs a pool of independent `PipelineUnit` objects based on the `--num_pipelines` command-line flag (defaulting to 1). This pool is injected into the FastAPI application via `create_app`, making it available to the WebSocket router for the lifetime of the server.

```python

# src/speech_to_speech/api/openai_realtime/server.py

app = create_app(pool=self.pool, stop_event=self.stop_event)

```

Each `PipelineUnit` defined in [`pipeline_unit.py`](https://github.com/huggingface/speech-to-speech/blob/main/pipeline_unit.py) encapsulates its own **input queues**, **output queues**, **control events**, and a dedicated `RealtimeService` instance. This design ensures that no mutable state is shared between concurrent sessions, eliminating race conditions at the architecture level.

## The Claim-and-Release Cycle

The WebSocket endpoint `/v1/realtime` orchestrates a strict claim-and-release lifecycle to manage access to the finite pipeline resources.

### Claiming a Pipeline Unit

When a client connects, the router invokes `_claim_unit` in [`websocket_router.py`](https://github.com/huggingface/speech-to-speech/blob/main/websocket_router.py). This function performs a synchronous scan of the pool (without `await` points) to find the first unit whose `session` attribute is `None`, then atomically reserves it by attaching a new `SessionState` object.

```python

# src/speech_to_speech/api/openai_realtime/websocket_router.py

def _claim_unit(transport: SessionTransport | None) -> PipelineUnit | None:
    for unit in pool:
        if unit.session is None:
            unit.session = SessionState(transport=transport)
            return unit
    return None

```

Because the loop contains no suspension points, Python’s single-threaded async model guarantees that only one coroutine can successfully claim a given unit, preventing double-booking under high concurrency.

### Rejecting Excess Connections

If `_claim_unit` returns `None`, indicating all units are occupied, the router immediately terminates the WebSocket handshake with a **session-limit-reached** error. It sends a JSON error payload and closes the connection with code `1008`, ensuring that the number of simultaneous active sessions never exceeds the configured pool size.

```python

# src/speech_to_speech/api/openai_realtime/websocket_router.py

if unit is None:
    logger.warning(f"Rejected connection: all {len(pool)} pipeline slots in use")
    await send_ws_event(ws, build_error_event(...))
    await ws.close(code=1008, reason="All session slots are in use")
    return

```

### Releasing and Recycling Units

Upon client disconnect, the router’s `finally` block triggers `_release_session`. This initiates a graceful teardown sequence:

1. **`_clean_unit`** clears the unit’s internal queues to drop any pending audio.
2. A **`SESSION_END`** control message is enqueued to signal the handler chain.
3. **`_release_unit_after_drain`** waits asynchronously for the sentinel to propagate through the pipeline before resetting `unit.session` to `None`, making the slot available for the next connection.

```python

# src/speech_to_speech/api/openai_realtime/websocket_router.py

finally:
    _release_session(unit, session_id)

```

## Concurrency Guarantees and Isolation

The pooling mechanism provides several hard guarantees for concurrent WebSocket handling.

### Session Isolation

Each `PipelineUnit` maintains its own `should_listen` and `response_playing` events, along with distinct `input_queue` and `output_queue` instances. Because the `RealtimeService` runs within this isolated container, memory corruption or state leakage between concurrent sessions is architecturally impossible.

### Atomic Claiming Without Locks

The claim loop’s absence of `await` statements acts as a **critical section**. By not yielding control during the scan-and-reserve operation, the system avoids complex locking primitives while still guaranteeing that two simultaneous connections cannot claim the same pipeline index.

### Back-Pressure and Flow Control

The `should_listen` event inside each unit controls whether the audio input pipeline continues processing. When the system is generating a response, this event is cleared, creating natural back-pressure that prevents the client from flooding a busy unit with new audio until the current response completes.

### Quarantine Handling for Stuck Pipelines

If a `SESSION_END` message fails to drain within **`SESSION_END_QUARANTINE_TIMEOUT_S`** (180 seconds), the unit is marked as **stuck**. The `_release_unit_after_drain` task detects this timeout, unregisters the session, and logs an error, keeping the unit occupied until the handler chain finally clears. This prevents partial session data from leaking into a new client connection.

```python

# src/speech_to_speech/api/openai_realtime/websocket_router.py

if session.quarantined_at is None and elapsed >= SESSION_END_QUARANTINE_TIMEOUT_S:
    session.quarantined_at = time.monotonic()
    _safe_unregister(unit, session_id)
    logger.error(f"Pipeline {unit.index} stuck for session {session_id}")

```

## Monitoring Pool Health via the Status Endpoint

The `/v1/pool` HTTP endpoint exposes real-time visibility into the pool’s state. It returns the total size, number of slots in use, and a per-unit breakdown indicating whether each unit is `idle`, `active`, `draining`, or `stuck`, along with timing metadata for debugging.

```python

# src/speech_to_speech/api/openai_realtime/websocket_router.py

async def pool_endpoint() -> dict[str, Any]:
    now = time.monotonic()
    def _state(u: PipelineUnit) -> dict[str, Any]:
        s = u.session
        if s is None:
            return {"index": u.index, "state": "idle", "session_id": None}
        if s.released_at is None:
            return {"index": u.index, "state": "active", "session_id": s.session_id}
        return {"index": u.index, "state": "draining", "session_id": s.session_id}

```

## Practical Implementation Examples

### Launching a Server with a Pool of Three Pipelines

```python
import threading
from speech_to_speech.api.openai_realtime.server import RealtimeServer
from speech_to_speech.api.openai_realtime.pipeline_unit import PipelineUnit
from speech_to_speech.pipeline.cancel_scope import CancelScope
from threading import Event

stop_event = Event()
pool = [
    PipelineUnit(
        index=i,
        service=...,          # RealtimeService instance

        cancel_scope=CancelScope(),
        should_listen=Event(),
        response_playing=Event(),
        input_queue=asyncio.Queue(),
        output_queue=asyncio.Queue(),
        text_output_queue=asyncio.Queue(),
        text_prompt_queue=asyncio.Queue(),
        handlers=...
    )
    for i in range(3)
]

server = RealtimeServer(stop_event=stop_event, pool=pool, host="0.0.0.0", port=8765)
thread = threading.Thread(target=server.run, daemon=True)
thread.start()

```

### Inspecting Pool Status via HTTP

```python
import requests

resp = requests.get("http://localhost:8765/v1/pool")
print(resp.json())

# Output:

# {

#   "size": 3,

#   "in_use": 1,

#   "units": [

#     {"index": 0, "state": "idle", "session_id": null},

#     {"index": 1, "state": "active", "session_id": "abc123"},

#     {"index": 2, "state": "idle", "session_id": null}

#   ]

# }

```

### Handling Rejection Client-Side

```python
import websockets
import json

async def connect():
    try:
        ws = await websockets.connect("ws://localhost:8765/v1/realtime")
    except websockets.exceptions.InvalidStatusCode as e:
        error = json.loads(e.body.decode())
        print("Connection rejected:", error["error"]["message"])
        # Output: All 3 session slots are in use. Disconnect an existing client first.

```

## Summary

- The **pipeline pool** pre-allocates a fixed number of `PipelineUnit` instances at server startup, configured via `--num_pipelines`.
- **Atomic claiming** in `_claim_unit` ensures that only one WebSocket connection can reserve a specific unit, with excess connections receiving a **1008 error code** and immediate closure.
- **Complete isolation** between units (separate queues, events, and services) guarantees that concurrent sessions cannot interfere with each other.
- **Graceful release** involves enqueuing a `SESSION_END` sentinel and draining the pipeline before resetting the unit to `idle`.
- A **180-second quarantine timeout** detects stuck pipelines, logging errors and preventing the unit from serving new clients until fully drained.
- The **`/v1/pool` endpoint** provides observability into slot utilization and unit health for operational monitoring.

## Frequently Asked Questions

### What happens when all pipeline slots are in use?

When the pool is exhausted, the `_claim_unit` function returns `None`, triggering the router to send a JSON error event and close the WebSocket with code `1008`. The client receives a message indicating that all session slots are in use and must disconnect an existing client or wait for a slot to become available.

### How does the system prevent race conditions when claiming units?

The claim loop in `_claim_unit` contains no `await` expressions, executing atomically within Python’s async event loop. This means the scan for an idle unit and the reservation of that unit happen in a single, uninterruptible block of code, eliminating the need for explicit locks while preventing double-booking.

### Can a single pipeline unit handle multiple WebSocket connections simultaneously?

No. Each `PipelineUnit` can be bound to exactly one `SessionState` object at a time. The `session` attribute acts as an occupancy flag; it is set during the claim phase and only reset to `None` after the release sequence completes, ensuring strict one-to-one mapping between units and active connections.

### How long does the system wait before marking a pipeline as stuck?

If a `SESSION_END` control message does not propagate through the handler chain within **180 seconds** (`SESSION_END_QUARANTINE_TIMEOUT_S`), the unit is flagged as stuck. The system logs an error, unregisters the session, and keeps the unit out of the available pool until the pipeline finally drains, preventing cross-session contamination.