Using the OpenAI Realtime WebSocket Protocol with Custom Clients
The huggingface/speech-to-speech library implements the full OpenAI Realtime API WebSocket protocol, allowing custom clients to connect using standard JSON events while the server handles STT, LLM, and TTS processing through isolated pipeline units.
The speech-to-speech repository provides a complete, self-hosted alternative to OpenAI's Realtime API. By implementing the same WebSocket event schema, it enables developers to build custom clients in any language while leveraging local or remote speech-to-text, language model, and text-to-speech backends. This guide explains how the protocol works, how to connect a custom client, and how the server translates OpenAI-format events into internal pipeline operations.
How the RealtimeProtocol Architecture Works
Every client connection receives a dedicated pipeline unit that maintains complete session isolation. The architecture centers on three core components that handle protocol translation, event routing, and bidirectional streaming.
The PipelineUnit: Per-Session State Container
When a client connects to /v1/realtime, the server claims a PipelineUnit from a shared pool. As defined in [pipeline_unit.py](https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/api/openai_realtime/pipeline_unit.py), this object holds:
- Input and output queues (
input_queue,output_queue,text_output_queue,text_prompt_queue) for audio and text flow - A RealtimeService instance that parses and validates OpenAI protocol events
- A SessionState tracking transport, session ID, and drain signaling
- A cancel scope for clean interruption handling
This design ensures that state, queues, and handlers never leak between sessions.
RealtimeService: Protocol Translation Engine
The [service.py](https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/api/openai_realtime/service.py) file contains RealtimeService, which serves as the single source of truth for OpenAI Realtime protocol handling. Its responsibilities include:
- Event parsing via
parse_client_event()— maps raw JSON to typed Pydantic models - State management through
ConnState— tracks conversation history, audio buffer state, and pending responses - Handler dispatch — routes validated events to
AudioHandler,ConversationHandler,ResponseHandler, orSessionHandler
The _EVENT_TYPE_TO_MODEL dictionary in service.py defines the complete event type mapping. Extending the protocol requires only adding new entries here and implementing corresponding handler methods.
Transport Abstractions: WebSocket and WebRTC
The [transports.py](https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/api/openai_realtime/transports.py) file provides concrete implementations for both transports:
- WebSocketTransport — wraps FastAPI's WebSocket with
send_events(),send_audio_chunk(), anddiscard_pending_audio() - WebRTCSession — manages the SDP exchange and maps media track packets to the same queue interface
Both transports implement identical interfaces, so the pipeline logic remains transport-agnostic.
WebSocket Data Flow: Event-by-Event Breakdown
Understanding the message sequence helps custom clients implement reliable integrations. Here's how data flows through the system:
1. Connection Establishment
In [websocket_router.py](https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/api/openai_realtime/websocket_router.py), the route handler:
@router.websocket("/v1/realtime")
async def realtime_websocket(websocket: WebSocket):
await websocket.accept()
transport = WebSocketTransport(websocket, logger)
unit = await _claim_unit(pipeline_units, stop_event)
session_id = unit.service.register(transport)
# ... starts send loop and enters receive loop
The server immediately emits session.created via service.build_session_created().
2. Client Event Ingestion
Client events arrive as JSON messages. The route's _dispatch_client_event function:
- Calls
unit.service.parse_client_event(raw_message)for validation - Dispatches to handler methods like
handle_audio_append(),handle_conversation_item_create(),handle_response_create() - Pushes data onto appropriate unit queues
3. Pipeline Processing
Background workers consume from queues and produce pipeline events defined in [events.py](https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/pipeline/events.py):
| Pipeline Event | Source | Translation |
|---|---|---|
AssistantTextEvent |
LLM output | response.text.delta |
AudioOutput |
TTS synthesis | response.audio.delta |
SpeechStartedEvent |
VAD trigger | input_audio_buffer.speech_started |
PIPELINE_END |
Handler chain | response.done |
4. The Send Loop: Ordered Delivery Guaranteed
The _send_loop_for function in websocket_router.py continuously drains queues:
async def _send_loop_for(unit: PipelineUnit, session: SessionState, ...):
while not stop_event.is_set():
# Text takes precedence to preserve ordering
if not unit.text_output_queue.empty():
events = await unit.text_output_queue.get()
await transport.send_events(events)
# Audio batches for efficiency
elif not unit.output_queue.empty():
chunk = await unit.output_queue.get()
# Accumulate up to MAX_AUDIO_BATCH_BYTES
await transport.send_audio_chunk(batch)
The loop respects SESSION_END propagation from [control.py](https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/pipeline/control.py), ensuring clean shutdown.
Minimal Python WebSocket Client Example
This runnable example demonstrates the complete protocol flow using the standard websockets library. Any WebSocket-capable language can follow the same JSON schema.
import asyncio
import json
import websockets
API_URL = "ws://localhost:8000/v1/realtime"
def make_event(event_type: str, **kwargs) -> str:
"""Build OpenAI Realtime protocol events."""
return json.dumps({"type": event_type, **kwargs})
async def realtime_client():
async with websockets.connect(API_URL) as ws:
# 1. Receive mandatory session.created
greeting = json.loads(await ws.recv())
print("Server:", greeting)
# 2. Send audio (16-bit little-endian PCM, 16kHz)
# Encode bytes as hex string for JSON transport
silence = b"\x00\x00" * 800 # 0.05s of silence
await ws.send(make_event(
"input_audio_buffer.append",
audio=silence.hex(),
))
# 3. Commit buffer to trigger STT
await ws.send(make_event("input_audio_buffer.commit"))
# 4. Add user message to conversation
await ws.send(make_event(
"conversation.item.create",
item={
"type": "message",
"role": "user",
"content": [
{"type": "input_text", "text": "Hello, how are you?"}
],
},
))
# 5. Request model response with audio and text
await ws.send(make_event(
"response.create",
response={"modalities": ["text", "audio"]},
))
# 6. Stream server events until completion
while True:
raw = await ws.recv()
event = json.loads(raw)
print("←", event["type"])
if event["type"] == "response.done":
break
elif event["type"] == "response.audio.delta":
# Decode base64 audio data and play
audio_bytes = event.get("delta", "")
# ... playback implementation
asyncio.run(realtime_client())
Client Event Effects in the Server
| Step | Client Event | Server Handler Action |
|---|---|---|
| 1 | (connection) | RealtimeService.register() → session.created emitted |
| 2 | input_audio_buffer.append |
AudioHandler.handle_audio_append() resamples and enqueues to unit.input_queue |
| 3 | input_audio_buffer.commit |
Signals VAD/STT pipeline start |
| 4 | conversation.item.create |
ConversationHandler updates per-connection chat history |
| 5 | response.create |
ResponseHandler triggers LLM generation |
| 6 | (receive loop) | Events stream from send loop via transport.send_events() / send_audio_chunk() |
Launching the Server for Client Testing
The repository provides a factory function for server creation. Here's a minimal setup:
from speech_to_speech.api.openai_realtime.websocket_router import create_app
from speech_to_speech.s2s_pipeline import build_pipeline_units
import uvicorn
import threading
# Configure pipeline units with STT, LLM, TTS handlers
# See scripts/listen_and_play_realtime.py for full setup
pipeline_units = build_pipeline_units(
stt="whisper",
llm="huggingface/meta-llama/Meta-Llama-3.1-8B-Instruct",
tts="parler-tts",
num_units=4, # supports 4 concurrent sessions
)
stop_event = threading.Event()
app = create_app(pipeline_units, stop_event)
if __name__ == "__main__":
uvicorn.run(app, host="0.0.0.0", port=8000)
The build_pipeline_units function (from s2s_pipeline.py) constructs handler chains that the RealtimeService orchestrates.
WebRTC Alternative: Same Protocol, Different Transport
For lower-latency applications, the server also exposes the Realtime protocol over WebRTC. The [webrtc_session.py](https://github.com/huggingface/speech-to-speech/blob/main/src/speech_to_speech/api/openai_realtime/webrtc_session.py) file handles:
- SDP offer/answer exchange via the
/v1/realtime/callsendpoint - Audio track callbacks that map incoming Opus packets to
unit.input_queue - Data channel "oai-events" carrying the identical JSON event schema
Custom WebRTC clients must negotiate the peer connection, then send OpenAI-format events through the data channel rather than WebSocket messages. The pipeline processing remains identical.
Reference Implementation: listen_and_play_realtime.py
For a complete working example, study [scripts/listen_and_play_realtime.py](https://github.com/huggingface/speech-to-speech/blob/main/scripts/listen_and_play_realtime.py). This script:
- Launches the server in a background thread
- Opens a local microphone stream
- Connects to
/v1/realtimevia WebSocket - Streams live audio and plays back synthesized responses
It demonstrates proper handling of:
- Audio format conversion (microphone → 16kHz PCM → hex encoding)
- Event sequencing with user interruptions
- Graceful shutdown with queue draining
Summary
-
Session isolation — Each client receives a dedicated
PipelineUnitwith independent queues, state, and cancel scopes, ensuring no cross-session leakage. -
Protocol fidelity —
RealtimeServiceimplements complete OpenAI Realtime event parsing and generation; custom clients speak standard JSON without server modifications. -
Transport flexibility — Identical pipeline logic works over both WebSocket and WebRTC via the
Transportabstraction intransports.py. -
Ordered delivery — The
_send_loop_forfunction guarantees text events precede associated audio and respectsSESSION_ENDfor clean resource release. -
Extensible architecture — Adding new event types requires only updating
_EVENT_TYPE_TO_MODELand implementing a handler method inservice.py.
Frequently Asked Questions
What audio format does the server expect from custom clients?
The speech-to-speech server expects 16-bit little-endian PCM at 16 kHz sample rate. In the WebSocket transport, encode audio bytes as hexadecimal strings for JSON transport: audio_chunk.hex(). The AudioHandler.handle_audio_append() method in handlers/__init__.py automatically resamples and normalizes incoming audio before enqueueing to unit.input_queue.
How do I handle user interruptions in a custom client?
Send input_audio_buffer.clear to discard pending audio, followed by a new response.create event. The server detects this sequence via discard_pending_audio() in the transport layer, which triggers the cancel scope in PipelineUnit to halt in-progress generation. The send loop then flushes stale output before processing the new request.
Can I use the Realtime protocol without the full pipeline?
No — the WebSocket endpoint requires a complete PipelineUnit with STT, LLM, and TTS handlers configured. However, you can substitute components: for example, use a remote Whisper API for STT by implementing a custom handler class. The protocol layer in RealtimeService remains unchanged regardless of backend implementation.
What's the difference between WebSocket and WebRTC performance?
WebRTC offers lower transport latency through UDP-based media channels and Opus compression, making it preferable for production voice applications. WebSocket provides simpler implementation and debugging — the JSON events are human-readable, and standard tools like websocat work for testing. Both transports use identical event schemas and pipeline processing; only the underlying byte transport differs.
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 →