How the VAD → STT → LLM → TTS Pipeline Uses Typed Queues for Inter-Component Communication
The VAD → STT → LLM → TTS architecture relies on thread-safe typed queues to pass structured Pydantic messages between decoupled handler components, with each stage consuming from a dedicated input queue and producing well-defined output types declared in queue_types.py.
The huggingface/speech-to-speech repository implements a modular speech-to-speech pipeline where Voice Activity Detection (VAD), Speech-to-Text (STT), Large Language Models (LLM), and Text-to-Speech (TTS) processing stages run as independent threads. The architecture achieves loose coupling through a centralized queue system that transports structured data payloads between components according to strict type contracts.
Typed Queue Aliases and Message Contracts
All inter-component communication is governed by strict type definitions that ensure compile-time and runtime validation of queue payloads.
Centralized Type Definitions
The file src/speech_to_speech/pipeline/queue_types.py declares type aliases for every queue in the system using Python's union syntax:
# https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/pipeline/queue_types.py
PipelineInternalItem = PipelineControlMessage | bytes
AudioInItem = VADIn | PipelineControlMessage
VADOutItem = VADOut | PipelineInternalItem
STTOutItem = STTOut | PipelineInternalItem
TextPromptItem = LLMIn | PipelineInternalItem
LMOutItem = LLMOut | PipelineInternalItem
TTSInItem = TTSIn | PipelineInternalItem
AudioOutItem = bytes | np.ndarray | AudioOutput | PipelineControlMessage
TextEventItem = PipelineEvent | PipelineInternalItem
These aliases guarantee that each queue accepts only expected types, preventing type mismatches between pipeline stages according to the source code definitions.
Pydantic Message Models
Concrete message implementations reside in messages.py, where each class inherits from PipelineMessage and uses a tag field for discriminated unions:
# https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/pipeline/messages.py
class VADAudio(PipelineMessage):
tag: Literal["vad_audio"] = "vad_audio"
audio: np.ndarray
mode: Literal["progressive", "final"] | None = None
turn_id: str | None = None
turn_revision: int | None = None
Similar structures define Transcription, LLMResponseChunk, TTSInput, and AudioOutput, creating a single source of truth for cross-component contracts.
Queue Initialization and Lifecycle
The pipeline instantiates all communication channels during startup through a centralized factory function.
Initializing Queue Objects
The initialize_queues_and_events() function in s2s_pipeline.py constructs every queue with explicit generic type parameters:
# https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/s2s_pipeline.py
def initialize_queues_and_events() -> dict[str, Any]:
return {
"stop_event": Event(),
"should_listen": Event(),
"response_playing": Event(),
"cancel_scope": CancelScope(),
"recv_audio_chunks_queue": Queue[AudioInItem](),
"send_audio_chunks_queue": Queue[AudioOutItem](),
"spoken_prompt_queue": Queue[VADOutItem](),
"stt_output_queue": Queue[STTOutItem](),
"text_prompt_queue": Queue[TextPromptItem](),
"lm_response_queue": Queue[LMOutItem](),
"lm_processed_queue": Queue[TTSInItem](),
"text_output_queue": Queue[TextEventItem](),
}
These queues are injected into handlers during the construction phase, ensuring each component receives its designated input and output channels.
Handler Communication Patterns
All processing units inherit from BaseHandler and operate by reading from queue_in and writing to queue_out. The ThreadManager orchestrates these handlers as separate threads.
VAD to STT Communication
The VADHandler consumes raw audio bytes and produces voice activity segments:
- Input:
recv_audio_chunks_queue(AudioInItem) - Output:
spoken_prompt_queue(VADOutItemcontainingVADAudiomessages)
# https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/VAD/vad_handler.py
vad = VADHandler(stop_event, queue_in=recv_audio_chunks_queue,
queue_out=spoken_prompt_queue, ...)
STT handlers (such as Whisper) subscribe to spoken_prompt_queue and emit Transcription objects to stt_output_queue.
STT to LLM Communication
The TranscriptionNotifier acts as a bridge between speech recognition and language processing:
- Input:
stt_output_queue(STTOutItem) - Output:
text_prompt_queue(TextPromptItem)
This intermediate step filters partial transcriptions and forwards only finalized text to the LLM stage.
LLM to TTS Communication
Language model handlers consume text prompts and stream response chunks:
- Input:
text_prompt_queue(TextPromptItemcontainingTranscription) - Output:
lm_response_queue(LMOutItemcontainingLLMResponseChunkorEndOfResponse)
The LMOutputProcessor then transforms these chunks into TTSInput objects, pushing them to lm_processed_queue for the synthesis stage.
TTS to Audio Output
Text-to-speech handlers perform the final conversion:
- Input:
lm_processed_queue(TTSInItem) - Output:
send_audio_chunks_queue(AudioOutItemcontainingAudioOutputor raw bytes)
This final queue carries completed audio back to the network layer or playback system.
Side-Channel Event Communication
Beyond the primary data flow, the pipeline maintains a separate event queue for synchronization signals. The text_output_queue (type Queue[TextEventItem]) transports PipelineEvent objects such as SpeechStartedEvent and SpeechStoppedEvent. Generated by VADHandler and TranscriptionNotifier, these events allow WebSocket clients to track pipeline state without parsing audio data.
Pipeline Construction Example
The _build_pipeline_handlers function wires queues to concrete handler instances:
# https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/s2s_pipeline.py
def _build_pipeline_handlers(...):
vad = VADHandler(..., queue_in=recv_audio_chunks_queue,
queue_out=spoken_prompt_queue, ...)
transcription_notifier = TranscriptionNotifier(...,
queue_in=stt_output_queue,
queue_out=text_prompt_queue, ...)
stt = get_stt_handler(...)
lm = get_llm_handler(...)
lm_processor = LMOutputProcessor(...,
queue_in=lm_response_queue,
queue_out=lm_processed_queue, ...)
tts = get_tts_handler(...)
return [vad, stt, transcription_notifier, lm, lm_processor, tts]
To launch the pipeline:
from speech_to_speech.s2s_pipeline import (
initialize_queues_and_events,
build_pipeline,
)
queues_and_events = initialize_queues_and_events()
pipeline_manager = build_pipeline(
module_kwargs=args.module_kwargs,
queues_and_events=queues_and_events,
# ... additional kwargs per handler
)
pipeline_manager.start() # Starts all handler threads
pipeline_manager.wait() # Blocks until completion
Summary
- Type-safe queues: All queue payloads are defined as union types in
queue_types.py, enabling static analysis and runtime validation. - Pydantic message contracts: Every data packet inherits from
PipelineMessagewith discriminant tags, ensuring version compatibility across the VAD → STT → LLM → TTS chain. - Decoupled handler architecture: Components communicate exclusively through injected
queue_inandqueue_outobjects, with no direct method calls between stages. - Thread safety: Python's standard
Queueclass provides built-in synchronization for concurrent access across the VAD, STT, LLM, and TTS threads. - Intermediate processors: Specialized handlers like
TranscriptionNotifierandLMOutputProcessormanage queue transitions and data transformation without blocking primary compute threads.
Frequently Asked Questions
How does the pipeline ensure type safety between distributed components?
The architecture enforces type safety through generic type aliases declared in queue_types.py. Each queue is instantiated with a specific payload type (e.g., Queue[VADOutItem]), and messages are validated Pydantic models from messages.py. This prevents runtime errors from mismatched data structures between the VAD, STT, LLM, and TTS stages.
What is the purpose of the TranscriptionNotifier in the queue chain?
The TranscriptionNotifier serves as a routing filter between stt_output_queue and text_prompt_queue. It consumes transcription events from the STT handler, distinguishes between partial and final results, and forwards only complete utterances to the LLM handler. This prevents the language model from processing incomplete speech fragments while maintaining asynchronous flow.
How are control messages and shutdown signals propagated across the pipeline?
Control messages travel through the same queue infrastructure as data payloads via the PipelineInternalItem union type, which includes PipelineControlMessage. The stop_event and should_listen synchronization primitives coordinate with these messages to trigger graceful shutdowns or interrupt ongoing synthesis without corrupting the byte streams in the audio queues.
Can different backend implementations be swapped without modifying the queue architecture?
Yes. Because handlers only depend on the abstract queue interface and typed message contracts defined in queue_types.py, you can substitute Whisper with Paraformer for STT, or ChatTTS with Kokoro for TTS, without changing queue definitions. As long as the new handler respects the input/output type contracts (e.g., consuming VADOutItem and emitting STTOutItem), the pipeline requires no structural changes.
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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →