PipelineUnit Architecture for Realtime Mode in Hugging Face Speech-to-Speech

The PipelineUnit architecture creates isolated, self-contained processing units with private queues, events, and handler chains, allowing the RealtimeServer to assign dedicated pipelines to each WebSocket client for concurrent, state-separated speech-to-speech conversations.

The huggingface/speech-to-speech library implements a modular PipelineUnit architecture for realtime mode to support OpenAI-compatible realtime API connections. Each PipelineUnit functions as an independent speech-to-speech engine, complete with dedicated queues, threading events, and a chain of handlers spanning VAD, STT, LLM, and TTS stages. This design ensures complete isolation between concurrent client sessions while maintaining low-latency audio streaming through thread-safe queue-based communication.

Core Components of the PipelineUnit Container

According to the source code in src/speech_to_speech/api/openai_realtime/pipeline_unit.py, the PipelineUnit class is a Pydantic BaseModel that encapsulates all resources required for a single realtime conversation session.

Isolated State Management

Each unit maintains strict separation through its own:

  • Input and output queues – input_queue (audio from client), output_queue (synthesized audio to client), text_output_queue (transcription events), and text_prompt_queue (text for LLM processing)
  • Threading events – should_listen controls VAD activation, while response_playing tracks active audio playback
  • Cancellation scope – A CancelScope instance shared across all handlers enables graceful interruption of the entire pipeline chain

The session field (typed as Optional[SessionState]) indicates availability; when None, the unit is free for new WebSocket connections.

RealtimeService Integration

Every PipelineUnit hosts a RealtimeService instance that translates OpenAI-compatible realtime protocol events into internal pipeline messages. This service manages the SpeculativeTurnTracker for handling overlapping user turns and maintains the conversation state according to the OpenAI realtime specification.

Building a Realtime PipelineUnit

The factory function _build_realtime_pipeline_unit in src/speech_to_speech/s2s_pipeline.py constructs fully configured units with isolated resources.

Deep Copy Isolation

The builder creates deep copies of all argument objects to prevent state leakage between pooled units. This ensures that mutations to shared configuration objects (like cancel_scope) remain local to the specific unit.

Queue and Event Initialization

from threading import Event
from queue import Queue
from speech_to_speech.pipeline.cancel_scope import CancelScope

# Per-unit isolation primitives

should_listen = Event()
response_playing = Event()
cancel_scope = CancelScope()

# Communication queues

recv_audio_chunks_queue = Queue()
send_audio_chunks_queue = Queue()
spoken_prompt_queue = Queue()
stt_output_queue = Queue()
text_prompt_queue = Queue()
lm_response_queue = Queue()
lm_processed_queue = Queue()
text_output_queue = Queue()

Service and Handler Chain Construction

The function instantiates RealtimeService with the text_prompt_queue and should_listen event, then calls _build_pipeline_handlers to create the processing chain. The resulting handlers list contains configured instances of VAD, STT, transcription notifier, LLM, output processor, and TTS components.

The Handler Chain Execution Flow

The realtime pipeline processes audio through six distinct stages, each running in its own thread and communicating via the unit's queues:

  1. VADHandler – Detects voice activity in recv_audio_chunks_queue and pushes audio frames to spoken_prompt_queue
  2. STT Handler – Consumes audio from spoken_prompt_queue and outputs text fragments to stt_output_queue (supports WhisperSTTHandler, LightningWhisperSTTHandler, and others)
  3. TranscriptionNotifier – Adds timestamps and metadata, forwarding to text_prompt_queue
  4. LLM Handler – Generates assistant replies using ResponsesApiModelHandler, ChatCompletionsApiModelHandler, or LanguageModelHandler, writing to lm_response_queue
  5. LMOutputProcessor – Converts LLM messages into TTS-ready items and emits text-output events
  6. TTS Handler – Synthesizes audio (via ChatTTSHandler, KokoroTTSHandler, Qwen3TTSHandler, etc.) and pushes PCM chunks to send_audio_chunks_queue

Each handler receives the unit's cancel_scope and relevant events, enabling coordinated interruption across the entire chain.

RealtimeServer and Connection Management

The RealtimeServer class in src/speech_to_speech/api/openai_realtime/server.py manages a pool of PipelineUnit instances and coordinates their assignment to WebSocket clients.

Unit Pool Claiming

When a client connects, the WebSocket router (websocket_router.py) searches the pool for the first unit where unit.session is None. Upon finding a free unit, it stores the WebSocket connection in unit.session and initiates send/receive loops that drive the unit's queues.

Graceful Release and Drain

Upon disconnection, the router calls _release_unit_after_drain, which waits for the output_queue to empty before resetting unit.session to None. This prevents audio truncation while ensuring the unit returns to the available pool.


# Conceptual pool management

server = RealtimeServer(
    stop_event=stop_event,
    pool=[unit_0, unit_1, unit_2],  # Multiple isolated units

    host="0.0.0.0",
    port=8765
)

Implementation Examples

Starting Realtime Mode from Command Line

python -m speech_to_speech.main \
    --mode realtime \
    --num_pipelines 2 \
    --stt whisper \
    --llm_backend chat-completions \
    --tts qwen3

The --mode realtime flag triggers the build_pipeline function to construct a pool of PipelineUnits via _build_realtime_pipeline_unit.

Python API Integration

from speech_to_speech.s2s_pipeline import build_pipeline, ParsedArguments
from speech_to_speech.arguments_classes.module_arguments import ModuleArguments
from threading import Event

args = ParsedArguments(
    module_kwargs=ModuleArguments(
        mode="realtime",
        num_pipelines=1,
        stt="whisper",
        llm_backend="chat-completions",
        tts="qwen3"
    ),
    # ... other handler kwargs

)

queues = {
    "stop_event": Event(),
    "recv_audio_chunks_queue": Queue(),
    "send_audio_chunks_queue": Queue(),
    # ... remaining queues

}

pipeline_manager = build_pipeline(**args.__dict__, queues_and_events=queues)
pipeline_manager.start()  # Launches RealtimeServer and handler threads

Manual PipelineUnit Construction

from speech_to_speech.api.openai_realtime.pipeline_unit import PipelineUnit
from speech_to_speech.api.openai_realtime.service import RealtimeService
from speech_to_speech.pipeline.cancel_scope import CancelScope
from threading import Event
from queue import Queue

unit = PipelineUnit(
    index=0,
    service=RealtimeService(
        text_prompt_queue=Queue(),
        should_listen=Event(),
        chat_size=10,
        speculative_turns=None
    ),
    cancel_scope=CancelScope(),
    should_listen=Event(),
    response_playing=Event(),
    input_queue=Queue(),
    output_queue=Queue(),
    text_output_queue=Queue(),
    text_prompt_queue=Queue(),
    handlers=[]  # Populate via _build_pipeline_handlers

)

Summary

  • PipelineUnit is an isolated container bundling queues, events, RealtimeService, and handlers for a single client session, defined in src/speech_to_speech/api/openai_realtime/pipeline_unit.py.
  • Complete isolation is achieved through deep-copied arguments and private Queue instances, preventing cross-talk between concurrent WebSocket connections.
  • Six-stage handler chain (VAD → STT → TranscriptionNotifier → LLM → LMOutputProcessor → TTS) processes audio through thread-safe queue handoffs.
  • RealtimeServer manages a configurable pool of units, claiming free instances for new connections and releasing them after output queue drainage.
  • OpenAI compatibility is handled by the RealtimeService, which translates protocol events to pipeline messages while managing speculative turn tracking.

Frequently Asked Questions

How does PipelineUnit prevent state leakage between concurrent clients?

Each PipelineUnit receives deep copies of all configuration arguments and instantiates fresh Queue and Event objects during construction. The session field tracks assignment status, ensuring that a unit's internal state (including conversation history and audio buffers) never shares memory with other units in the pool.

What is the purpose of the should_listen event in the PipelineUnit?

The should_listen threading Event controls Voice Activity Detection (VAD) gating. When set, the VADHandler processes incoming audio; when cleared (during TTS playback or system processing), the pipeline ignores audio input to prevent the system from hearing its own output or processing during interruptions.

How many simultaneous connections can RealtimeServer support?

The server supports any number of concurrent sessions up to the --num_pipelines limit specified at startup. Each WebSocket connection requires one dedicated PipelineUnit, so hardware constraints (GPU memory for TTS/LLM, CPU cores) typically bound the practical maximum rather than architectural limits.

What happens to audio in the queue when a client disconnects?

The router invokes _release_unit_after_drain to wait for the output_queue to empty before releasing the unit. This ensures the client receives all synthesized audio fragments before the session terminates, preventing truncation of the final assistant response.

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 →