How Realtime Turn Context Manages Speech-to-Speech Sessions in Dograh
The realtime turn context in Dograh uses a hybrid approach combining a global TurnContextManager with Pipecat's turn_var context variable to track conversation turns across asynchronous task boundaries, enabling precise logging and analytics for speech-to-speech AI workflows.
Managing conversation turns in realtime speech-to-speech systems presents unique challenges because traditional request-response boundaries disappear when the LLM handles audio generation internally. The Dograh platform (dograh-hq/dograh) solves this by propagating turn state through asyncio's execution model while integrating with Pipecat's event system. This architecture ensures that every component—from audio processors to analytics loggers—knows exactly which conversational turn it is processing.
The Challenge of Tracking Turns in Async Speech-to-Speech Workflows
Speech-to-speech services like OpenAI Realtime and Gemini Live maintain their own conversation state server-side, but the Dograh platform must still track turns for logging, UI feedback, and callback payloads. Standard context variables (turn_var) do not survive asyncio.create_task() boundaries, meaning sub-tasks lose track of which turn they belong to. The realtime turn context bridges this gap by storing turn numbers in a manager keyed by asyncio.Task objects, allowing retrieval even when execution crosses task boundaries.
Architectural Components of the Realtime Turn Context
TurnContextManager
The TurnContextManager class in api/services/pipecat/turn_context.py serves as the central registry for turn state. It maintains three key data structures:
_task_turns: A mapping ofasyncio.Taskobjects to turn numbers_current_turn: A global fallback storing the last known turn- Pipeline task reference: Tracks the main pipeline execution context
The manager exposes set_turn() and get_turn() methods. When get_turn() is called, it attempts retrieval from the context variable first (fast path), falls back to the task mapping (cross-task lookup), and finally returns the global current turn as a last resort.
Context Variables and Async Propagation
Pipecat provides turn_var (a contextvars.ContextVar) imported from pipecat.utils.run_context for low-latency turn lookup within the same task. However, when Pipecat spawns sub-tasks for audio processing or LLM inference, these context values disappear. The TurnContextManager compensates by hooking into task creation events and maintaining the parallel mapping structure.
Pipeline Integration
In api/services/pipecat/pipeline_builder.py (lines 86‑100), the pipeline builder checks the ENABLE_TURN_LOGGING environment variable. When enabled, it extracts the turn_observer from the PipelineTask and registers an on_turn_started handler:
async def _on_turn_started(observer, turn_number: int):
turn_var.set(turn_number) # Pipecat’s fast path
turn_manager = get_turn_context_manager() # Global manager
turn_manager.set_turn(turn_number) # Persist across tasks
turn_observer.add_event_handler("on_turn_started", _on_turn_started)
This dual-write strategy ensures that both the context variable and the global manager stay synchronized whenever a new turn begins.
How Realtime Mode Configures Turn Management
Detecting Realtime Sessions
The run_pipeline function in api/services/pipecat/run_pipeline.py detects realtime mode by checking user_config.is_realtime (lines 75‑85). When true, it instantiates a speech-to-speech LLM service via create_realtime_llm_service and disables context compaction (lines 101‑107) because the external service maintains its own conversation history:
if is_realtime:
context_compaction_enabled = False
# ... service instantiation logic
Provider-Specific Turn Strategies
The _create_realtime_user_turn_config function (lines 89‑124) returns provider-specific UserTurnStrategies:
- Google Gemini: Returns a VAD-based start strategy combined with a timeout stop strategy, while keeping a local Silero VAD for early user-turn detection
- OpenAI Realtime: Returns external strategies that rely entirely on provider-generated speaking-state frames, disabling local VAD processing
This differentiation ensures that turn detection aligns with each provider's native interruption handling.
Transport Parameter Adjustments
Realtime sessions require different Voice Activity Detection (VAD) behavior. The realtime_param_overrides function in api/services/pipecat/transport_params.py (lines 9‑25) injects bot_vad_stop_secs = 0.5 to ensure the platform quickly detects when the LLM stops speaking, preventing long silent gaps in the conversation.
Step-by-Step Turn Tracking Flow
- Configuration Resolution:
run_pipelinesetsis_realtime = Truewhen the user config contains arealtimeblock - Service Instantiation: Creates the speech-to-speech LLM and a side-channel text LLM for out-of-band tasks like voicemail detection
- Strategy Selection:
_create_realtime_user_turn_configselects appropriate turn detection strategies based on the provider - Observer Registration: If
ENABLE_TURN_LOGGINGis enabled,create_pipeline_taskattaches the turn event handler to syncturn_varandTurnContextManager - Cross-Task Propagation: Sub-tasks retrieve the correct turn number via
get_turn_context_manager().get_turn()even whenturn_varis unavailable - Analytics Emission: The Pipecat engine queries the manager to embed turn numbers in realtime feedback payloads sent to the UI
Implementation Example: Capturing Turn Events
The following example demonstrates how to enable turn logging and retrieve the current turn within a realtime pipeline:
import os
from api.services.pipecat.pipeline_builder import create_pipeline_task
from api.services.pipecat.turn_context import get_turn_context_manager
from pipecat.utils.run_context import turn_var
# Enable turn tracking
os.environ["ENABLE_TURN_LOGGING"] = "true"
# Create the pipeline task for a realtime LLM
task = create_pipeline_task(
pipeline,
workflow_run_id=12345,
audio_config=my_audio_config,
)
# The on_turn_started handler is automatically attached.
# Retrieve the current turn anywhere in your async stack:
current_turn = get_turn_context_manager().get_turn()
print(f"Processing realtime turn: {current_turn}")
Summary
- The realtime turn context combines
TurnContextManagerwith Pipecat'sturn_varto solve async task boundary issues - Turn propagation works via a three-tier fallback system: context variable → task mapping → global current turn
- Realtime detection in
run_pipeline.pydisables context compaction and configures provider-specific strategies via_create_realtime_user_turn_config - Transport parameters adjust VAD timeouts to 0.5 seconds for rapid bot-speech detection
- Logging integration occurs through an
on_turn_startedobserver registered inpipeline_builder.pywhenENABLE_TURN_LOGGINGis active
Frequently Asked Questions
How does Dograh handle turn tracking when async tasks are created?
Dograh handles cross-task turn tracking through the TurnContextManager in api/services/pipecat/turn_context.py. While turn_var (a context variable) provides fast lookup within a single task, the manager maintains a parallel mapping of asyncio.Task objects to turn numbers. When a sub-task calls get_turn(), the manager checks the context variable first, then falls back to the task mapping, ensuring accurate turn identification regardless of task boundaries.
Why is context compaction disabled for speech-to-speech sessions?
Context compaction is disabled because speech-to-speech services like OpenAI Realtime and Gemini Live maintain their own conversation history server-side. As implemented in api/services/pipecat/run_pipeline.py (lines 101‑107), forcing context_compaction_enabled = False prevents Dograh from truncating the context window, which would interfere with the external service's internal state management and potentially break the conversation flow.
What is the difference between Gemini and OpenAI Realtime turn strategies?
The _create_realtime_user_turn_config function returns different UserTurnStrategies based on the provider. For Google Gemini, it configures local VAD-based start detection with timeout-based stop detection, keeping a Silero VAD for early user-turn detection. For OpenAI Realtime, it uses external strategies that rely entirely on the provider's speaking-state frames, disabling local VAD because OpenAI handles interruption detection internally.
How can I access the current turn number in my pipeline code?
Access the current turn by importing get_turn_context_manager from api/services/pipecat/turn_context and calling get_turn(). This works anywhere in the async call stack after the pipeline starts and the on_turn_started observer has fired at least once. The method handles the complexity of checking context variables, task mappings, and global fallbacks automatically.
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 →