Source code for pipecat.services.deepgram.flux.stt_base

#
# Copyright (c) 2024-2026, Daily
#
# SPDX-License-Identifier: BSD 2-Clause License
#

"""Deepgram Flux STT base class shared across transports (WebSocket, SageMaker, etc.)."""

import asyncio
import time
from abc import abstractmethod
from dataclasses import dataclass, field
from enum import StrEnum
from typing import Any
from urllib.parse import urlencode

from loguru import logger
from typing_extensions import override

from pipecat.frames.frames import (
    CancelFrame,
    EndFrame,
    ProposedUserStartedSpeakingFrame,
    ProposedUserStoppedSpeakingFrame,
    STTMetadataFrame,
    TranscriptionFrame,
)
from pipecat.processors.frame_processor import FrameProcessorSetup
from pipecat.services.settings import STTSettings
from pipecat.services.stt_service import STTService
from pipecat.transcriptions.language import Language, resolve_language
from pipecat.turns.eager_end_of_turn_mixin import EagerEndOfTurnSTTServiceMixin
from pipecat.utils.errors import ErrorCategory
from pipecat.utils.time import time_now_iso8601
from pipecat.utils.tracing.service_decorators import traced_stt
from pipecat.utils.types import NOT_GIVEN, NotGiven, assert_given, is_given


[docs] class FluxConnectionNotConfirmedError(Exception): """Flux accepted the connection but never confirmed it was ready."""
[docs] class FluxFatalError(Exception): """Flux reported a fatal error and terminated the connection. Attributes: code: The error code Flux sent, e.g. ``UNPARSABLE_CLIENT_MESSAGE``. """
[docs] def __init__(self, message: str, code: str): """Initialize the error. Args: message: The formatted error message. code: The error code Flux sent. """ super().__init__(message) self.code = code
[docs] def language_to_deepgram_flux_language(language: Language) -> str: """Convert a Pipecat Language to a Deepgram Flux language code. Only honored by the ``flux-general-multi`` model. Locale variants (e.g. ``Language.EN_GB``) fall back to the base code. """ LANGUAGE_MAP = { Language.DE: "de", Language.EN: "en", Language.ES: "es", Language.FR: "fr", Language.HI: "hi", Language.IT: "it", Language.JA: "ja", Language.NL: "nl", Language.PT: "pt", Language.RU: "ru", } return resolve_language(language, LANGUAGE_MAP, use_base_code=True)
def _prepare_language_hints(hints: list[Language] | None) -> list[str]: """Convert a list of Pipecat Languages to Deepgram Flux codes. Drops entries that can't be mapped and deduplicates while preserving order. """ if not hints: return [] seen: set[str] = set() prepared: list[str] = [] for hint in hints: code = language_to_deepgram_flux_language(hint) if code is None or code in seen: continue seen.add(code) prepared.append(code) return prepared def _code_to_pipecat_language(code: str) -> Language | None: """Convert a Deepgram-returned language code to a Pipecat Language.""" try: return Language(code) except ValueError: logger.debug(f"Unmapped Deepgram Flux detected language code: {code}") return None
[docs] class FluxMessageType(StrEnum): """Deepgram Flux WebSocket message types. These are the top-level message types that can be received from the Deepgram Flux WebSocket connection. """ RECEIVE_CONNECTED = "Connected" RECEIVE_FATAL_ERROR = "Error" TURN_INFO = "TurnInfo" CONFIGURE_SUCCESS = "ConfigureSuccess" CONFIGURE_FAILURE = "ConfigureFailure"
[docs] class FluxEventType(StrEnum): """Deepgram Flux TurnInfo event types. These events are contained within TurnInfo messages and indicate different stages of speech processing and turn detection. """ START_OF_TURN = "StartOfTurn" TURN_RESUMED = "TurnResumed" END_OF_TURN = "EndOfTurn" EAGER_END_OF_TURN = "EagerEndOfTurn" UPDATE = "Update"
[docs] @dataclass class DeepgramFluxSTTSettings(STTSettings): """Settings for DeepgramFluxSTTService. Parameters: eager_eot_threshold: EagerEndOfTurn/TurnResumed threshold. Off by default. Lower values = more aggressive (faster response, more LLM calls). Higher values = more conservative (slower response, fewer LLM calls). eot_threshold: End-of-turn confidence required to finish a turn (default 0.7). eot_timeout_ms: Time in ms after speech to finish a turn regardless of EOT confidence (default 5000). keyterm: Keyterms to boost recognition accuracy for specialized terminology. min_confidence: Minimum confidence required to create a TranscriptionFrame. numerals: Convert spoken numbers to numeral form (e.g. "twenty three" → "23"). Read only from the connection URL, so an update is applied by reconnecting. profanity_filter: Mask recognized profanity in the transcript. Can be updated mid-stream via ``STTUpdateSettingsFrame``. redact: Remove sensitive numbers from the transcript: ``"numbers"`` or ``"aggressive_numbers"``. Read only from the connection URL, so an update is applied by reconnecting. language_hints: Languages to bias transcription toward. Only honored by the ``flux-general-multi`` model. An empty list clears any active hints; ``None``/``NOT_GIVEN`` means no hints (auto-detect). Can be updated mid-stream via ``STTUpdateSettingsFrame``. """ eager_eot_threshold: float | None | NotGiven = field(default_factory=lambda: NOT_GIVEN) eot_threshold: float | None | NotGiven = field(default_factory=lambda: NOT_GIVEN) eot_timeout_ms: int | None | NotGiven = field(default_factory=lambda: NOT_GIVEN) keyterm: list | NotGiven = field(default_factory=lambda: NOT_GIVEN) min_confidence: float | None | NotGiven = field(default_factory=lambda: NOT_GIVEN) numerals: bool | None | NotGiven = field(default_factory=lambda: NOT_GIVEN) profanity_filter: bool | None | NotGiven = field(default_factory=lambda: NOT_GIVEN) redact: str | None | NotGiven = field(default_factory=lambda: NOT_GIVEN) language_hints: list[Language] | None | NotGiven = field(default_factory=lambda: NOT_GIVEN)
[docs] class DeepgramFluxSTTBase(EagerEndOfTurnSTTServiceMixin, STTService): """Base class for Deepgram Flux STT services across transports. Contains all shared Flux protocol logic (message handling, turn detection, metrics, settings). Concrete subclasses implement the transport layer by providing three abstract primitives: ``_transport_send_audio``, ``_transport_send_json``, and ``_transport_is_active``. """ Settings = DeepgramFluxSTTSettings _settings: Settings _CONFIGURE_FIELDS = { "keyterm", "eot_threshold", "eager_eot_threshold", "eot_timeout_ms", "language_hints", "profanity_filter", } # Fields Flux only accepts in the connection URL, so changing them reconnects. _CONNECTION_FIELDS = {"model", "numerals", "redact"} # Fields applied to results as they arrive, so no connection change is needed. _LOCAL_FIELDS = {"min_confidence"} _MULTILINGUAL_MODEL = "flux-general-multi" # How long to wait for Flux to confirm a new connection. An endpoint that # rejects a connection parameter sends neither a Connected message nor an # error, so the wait needs a bound. _CONNECTION_TIMEOUT = 10.0 # Flux error codes whose cause a retry cannot clear. Rejected credentials # don't appear here: those fail the HTTP handshake and are classified from # its status code, never reaching a Flux error message. Anything not listed # falls back to the default classification, leaving recovery to the # service's own reconnect handling. _ERROR_CODE_CATEGORIES = { "UNPARSABLE_CLIENT_MESSAGE": ErrorCategory.INVALID_REQUEST, } # Threshold applied when eager end of turn is enabled without one. Flux # reports a prediction only above this confidence: lower is more eager # (faster responses, more discarded inferences), higher more conservative. _DEFAULT_EAGER_EOT_THRESHOLD = 0.5 # How long an in-flight Configure is trusted before a new update supersedes # it outright. Flux caps the number of un-acked Configure messages, so at # most one is ever in flight; this bounds how long a missing ack can block # later updates from ever being sent. _CONFIGURE_ACK_TIMEOUT = 5.0
[docs] def __init__( self, *, encoding: str = "linear16", mip_opt_out: bool | None = None, tag: list | None = None, should_interrupt: bool = True, watchdog_min_timeout: float = 0.5, enable_eager_end_of_turn: bool = False, settings: Settings, **kwargs, ): """Initialize the Deepgram Flux STT base service. Args: encoding: Audio encoding format. Must be "linear16". mip_opt_out: Opt out of the Deepgram Model Improvement Program. tag: Tags to label requests for identification during usage reporting. should_interrupt: Whether to interrupt the bot when Flux detects that the user is speaking. Passed along to the user turn strategies this service recommends, which own the interruption; a user-supplied ``user_turn_strategies`` overrides the recommendation and this setting with it. watchdog_min_timeout: minimum idle timeout before sending silence to prevent dangling turns. The actual threshold is ``max(chunk_duration * 2, watchdog_min_timeout)``. Defaults to 0.5. enable_eager_end_of_turn: Whether to answer Flux's predicted end of turn ahead of the committed one, so the gap between the two is spent generating a response rather than waiting. Off by default: it spends an inference on every prediction, including the ones Flux withdraws. Turning it on sets ``eager_eot_threshold`` to ``_DEFAULT_EAGER_EOT_THRESHOLD`` when the settings leave it unset, since Flux reports no prediction without it. settings: Fully resolved settings instance (built by concrete subclass). **kwargs: Additional arguments passed to the parent STTService (e.g. ``sample_rate``, ``reconnect_on_error``). """ super().__init__( settings=settings, enable_eager_end_of_turn=enable_eager_end_of_turn, **kwargs ) # Flux reports no prediction unless a threshold asks for one, so an # unconfigured threshold would leave the feature silently inert. if self.eager_end_of_turn_enabled and self._settings.eager_eot_threshold is None: self._settings.eager_eot_threshold = self._DEFAULT_EAGER_EOT_THRESHOLD self._encoding = encoding self._mip_opt_out = mip_opt_out self._tag = tag or [] self._should_interrupt = should_interrupt self._watchdog_min_timeout = watchdog_min_timeout # Connection readiness: Flux sends a "Connected" message when ready self._connection_established_event = asyncio.Event() # Configure serialization: Flux caps the number of un-acked Configure # messages, so we only allow one in flight at a time. A Configure sent # while one is already in flight is coalesced into # `_configure_pending_fields` instead, and sent once the in-flight one # is acked (see `_on_configure_acked`) — only the latest settings # matter, so there's no need to replay every intermediate update. self._configure_in_flight = False self._configure_sent_at: float | None = None self._configure_pending_fields: set[str] | None = None # Watchdog state — see _watchdog_task_handler for details self._last_stt_time: float | None = None self._watchdog_task: asyncio.Task | None = None self._user_is_speaking = False self._last_audio_chunk_duration: float = 0.0 # Flux event handlers self._register_event_handler("on_start_of_turn") self._register_event_handler("on_turn_resumed") self._register_event_handler("on_end_of_turn") self._register_event_handler("on_eager_end_of_turn") self._register_event_handler("on_update")
[docs] def can_generate_metrics(self) -> bool: """Check if this service can generate processing metrics. Returns: True, as Deepgram Flux service supports metrics generation. """ return True
[docs] def service_metadata_frame(self) -> STTMetadataFrame: """Recommend turn strategies that leave turn detection to Flux. Flux emits its own start-of-turn and end-of-turn events (as ``ProposedUserStarted/StoppedSpeakingFrame``), so the user aggregator resolves those rather than running local VAD/smart-turn. With ``enable_eager_end_of_turn``, the recommendation also answers Flux's predicted end of turn. Applied unless the user passed their own ``user_turn_strategies``. """ frame = super().service_metadata_frame() frame.user_turn_strategies = self.recommended_user_turn_strategies( enable_interruptions=self._should_interrupt, ) return frame
# ------------------------------------------------------------------ # Abstract transport interface — implemented by each concrete subclass # ------------------------------------------------------------------ @abstractmethod async def _transport_send_audio(self, audio: bytes): """Send raw audio bytes over the transport.""" pass @abstractmethod async def _transport_send_json(self, message: dict): """Serialize and send a JSON control message over the transport.""" pass @abstractmethod def _transport_is_active(self) -> bool: """Return True if the transport connection is currently active.""" pass @abstractmethod async def _connect(self): """Establish the transport connection.""" pass @abstractmethod async def _disconnect(self): """Tear down the transport connection.""" pass # ------------------------------------------------------------------ # Connection helpers # ------------------------------------------------------------------ @override async def _do_reconnect(self): """Tear down the transport connection and re-establish it. Called by ``STTService._reconnect()`` inside the reconnecting guard. """ await self._disconnect() await self._connect() def _classify_error(self, exception: Exception) -> ErrorCategory | None: """Classify the failures Flux signals in its own protocol. Flux reports these over the connection rather than as an HTTP status, so they carry nothing the default classification can read. Args: exception: The exception to classify. Returns: The category, or None to fall back to the default classification. """ if isinstance(exception, FluxConnectionNotConfirmedError): # Flux stays silent rather than refusing the connection when a # setting is unsupported, so an unconfirmed connection means the # request was rejected, not that the network was slow. return ErrorCategory.INVALID_REQUEST if isinstance(exception, FluxFatalError): return self._ERROR_CODE_CATEGORIES.get(exception.code) return None async def _await_connection_established(self): """Wait for Flux to confirm the connection is ready. Raises: FluxConnectionNotConfirmedError: If no confirmation arrives within ``_CONNECTION_TIMEOUT``. """ try: await asyncio.wait_for( self._connection_established_event.wait(), timeout=self._CONNECTION_TIMEOUT ) except TimeoutError: raise FluxConnectionNotConfirmedError( f"Flux did not confirm the connection within {self._CONNECTION_TIMEOUT}s; " "the endpoint may not accept the current connection settings" ) from None def _build_query_string(self) -> str: """Build query string from current settings and init-only connection config.""" params = [ f"model={self._settings.model}", f"sample_rate={self.sample_rate}", f"encoding={self._encoding}", ] if self._settings.eager_eot_threshold is not None: params.append(f"eager_eot_threshold={self._settings.eager_eot_threshold}") if self._settings.eot_threshold is not None: params.append(f"eot_threshold={self._settings.eot_threshold}") if self._settings.eot_timeout_ms is not None: params.append(f"eot_timeout_ms={self._settings.eot_timeout_ms}") if self._settings.numerals is not None: params.append(f"numerals={str(self._settings.numerals).lower()}") if self._settings.profanity_filter is not None: params.append(f"profanity_filter={str(self._settings.profanity_filter).lower()}") if self._settings.redact is not None: params.append(urlencode({"redact": self._settings.redact})) if self._mip_opt_out is not None: params.append(f"mip_opt_out={str(self._mip_opt_out).lower()}") # Add keyterm parameters (can have multiple) for keyterm in assert_given(self._settings.keyterm): params.append(urlencode({"keyterm": keyterm})) # Add tag parameters (can have multiple) for tag_value in self._tag: params.append(urlencode({"tag": tag_value})) # Add language_hint parameters (only valid on flux-general-multi) hints = self._settings.language_hints if hints and is_given(hints): if self._settings.model == self._MULTILINGUAL_MODEL: for code in _prepare_language_hints(hints): params.append(urlencode({"language_hint": code})) else: logger.warning( f"language_hints only supported on {self._MULTILINGUAL_MODEL}; " f"ignoring hints for model {self._settings.model!r}" ) return "&".join(params) async def _send_silence(self, duration_secs: float = 0.5): """Send a block of silence of the specified duration (default 500 ms).""" sample_width = 2 # bytes per sample for 16-bit PCM num_channels = 1 # mono num_samples = int(self.sample_rate * duration_secs) silence = b"\x00" * (num_samples * sample_width * num_channels) await self._transport_send_audio(silence) # Watchdog silence is real audio submitted to the service, so it # counts toward usage. self._record_stt_audio_usage(silence) async def _watchdog_task_handler(self): """Prevent dangling turns by sending silence when audio stops flowing. If we stop sending audio to Flux after receiving a StartOfTurn, we never receive the UserStoppedSpeaking event unless we resume sending audio. """ while self._transport_is_active(): now = time.monotonic() # Send silence if we go more than 500 ms or twice the chunk size # without sending new audio to Flux. threshold = max(self._last_audio_chunk_duration * 2, self._watchdog_min_timeout) if ( self._user_is_speaking and self._last_stt_time and now - self._last_stt_time > threshold ): logger.warning( f"No audio received for {threshold * 1000:.0f} ms. Sending silence to Flux to prevent a dangling task" ) try: await self._send_silence() except Exception as e: logger.warning(f"Failed to send silence: {e}") self._last_stt_time = time.monotonic() # check every 100ms await asyncio.sleep(0.1) async def _send_close_stream(self) -> None: """Sends a CloseStream control message to Deepgram Flux. This signals to the server that no more audio data will be sent. """ try: if self._transport_is_active(): logger.debug("Sending CloseStream message to Deepgram Flux") await self._transport_send_json({"type": "CloseStream"}) except Exception as e: await self.push_error(error_msg=f"Error sending CloseStream: {e}", exception=e) # ------------------------------------------------------------------ # Lifecycle # ------------------------------------------------------------------
[docs] async def setup(self, setup: FrameProcessorSetup): """Set up the service and connect. Args: setup: Configuration object containing setup parameters. """ await super().setup(setup) await self._connect()
[docs] async def stop(self, frame: EndFrame): """Stop the Deepgram Flux STT service. Args: frame: The end frame. """ await super().stop(frame) await self._disconnect()
[docs] async def cancel(self, frame: CancelFrame): """Cancel the Deepgram Flux STT service. Args: frame: The cancel frame. """ await super().cancel(frame) await self._disconnect()
[docs] async def cleanup(self): """Release Deepgram Flux STT resources at teardown.""" await super().cleanup() await self._disconnect()
@traced_stt async def _handle_transcription( self, transcript: str, is_final: bool, language: Language | None = None ): """Handle a transcription result with tracing.""" pass async def _send_configure(self, fields: set[str]): """Send a Configure control message to update settings mid-stream. Builds a Configure JSON message containing only the fields that changed and sends it over the existing connection. At most one Configure is ever in flight, since Flux caps the number of un-acked Configure messages. If one is already in flight, ``fields`` is merged into the pending set and sent once the in-flight one is acked (see ``_on_configure_acked``) instead of being sent now — the message is always built from the current settings, so only the latest values ever need to go out. An in-flight Configure older than ``_CONFIGURE_ACK_TIMEOUT`` is treated as lost so a missing ack can't permanently block later updates. Args: fields: Set of changed field names to include in the message. """ if self._configure_in_flight: assert self._configure_sent_at is not None if time.monotonic() - self._configure_sent_at < self._CONFIGURE_ACK_TIMEOUT: self._configure_pending_fields = (self._configure_pending_fields or set()) | fields return logger.warning( f"{self}: timed out after {self._CONFIGURE_ACK_TIMEOUT}s waiting for " "Configure ack; sending the next Configure anyway" ) message: dict[str, Any] = {"type": "Configure"} if "keyterm" in fields: message["keyterms"] = self._settings.keyterm if "profanity_filter" in fields: message["profanity_filter"] = self._settings.profanity_filter thresholds: dict[str, Any] = {} if "eot_threshold" in fields: thresholds["eot_threshold"] = self._settings.eot_threshold if "eager_eot_threshold" in fields: thresholds["eager_eot_threshold"] = self._settings.eager_eot_threshold if "eot_timeout_ms" in fields: thresholds["eot_timeout_ms"] = self._settings.eot_timeout_ms if thresholds: message["thresholds"] = thresholds if "language_hints" in fields: if self._settings.model != self._MULTILINGUAL_MODEL: logger.warning( f"language_hints only supported on {self._MULTILINGUAL_MODEL}; " f"skipping Configure update for model {self._settings.model!r}" ) else: hints = self._settings.language_hints # Empty list clears hints; NOT_GIVEN/None also treated as clear # since we only reach this branch when the user set the field. if hints is None or not is_given(hints): message["language_hints"] = [] else: message["language_hints"] = _prepare_language_hints(hints) self._configure_in_flight = True self._configure_sent_at = time.monotonic() logger.debug(f"{self}: sending Configure message: {message}") await self._transport_send_json(message) async def _on_configure_acked(self): """Mark the in-flight Configure as acked and flush any pending update. Called when a ConfigureSuccess/ConfigureFailure arrives. If fields were coalesced into ``_configure_pending_fields`` while this Configure was in flight, immediately sends a follow-up Configure covering all of them — unless the transport has since gone inactive, in which case the pending fields are simply dropped, since a reconnect re-applies current settings via the connection URL anyway. Safe to call with nothing in flight (e.g. a stray/duplicate ack), which is a no-op. """ self._configure_in_flight = False self._configure_sent_at = None if self._configure_pending_fields is not None: fields = self._configure_pending_fields self._configure_pending_fields = None if self._transport_is_active(): await self._send_configure(fields) def _reset_configure_state(self): """Clear Configure-serialization state during teardown. Called when the connection is torn down (including ahead of a reconnect), since any in-flight or pending Configure can no longer be acked or sent on a dead connection. A reconnect re-applies the current settings via the connection URL, so nothing needs to be replayed. """ self._configure_in_flight = False self._configure_sent_at = None self._configure_pending_fields = None async def _update_settings(self, delta: Settings) -> dict[str, Any]: """Apply a settings delta. Configure-able fields (keyterm, eot_threshold, eager_eot_threshold, eot_timeout_ms, language_hints) are sent to Deepgram via a Configure message. Fields Flux only reads from the connection URL trigger a reconnect, which waits until the user stops speaking. """ changed = await super()._update_settings(delta) if not changed: return changed configure_fields = changed.keys() & self._CONFIGURE_FIELDS if configure_fields and self._transport_is_active(): await self._send_configure(configure_fields) if changed.keys() & self._CONNECTION_FIELDS: await self._request_reconnect() self._warn_unhandled_updated_settings( changed.keys() - self._CONFIGURE_FIELDS - self._CONNECTION_FIELDS - self._LOCAL_FIELDS ) return changed # ------------------------------------------------------------------ # Message handling # ------------------------------------------------------------------ def _validate_message(self, data: dict[str, Any]) -> bool: """Validate basic message structure from Deepgram Flux. Ensures the received message has the expected structure before processing. Args: data: The parsed JSON message data to validate. Returns: True if the message structure is valid, False otherwise. """ if not isinstance(data, dict): logger.warning("Message is not a dictionary") return False if "type" not in data: logger.warning("Message missing 'type' field") return False return True async def _handle_message(self, data: dict[str, Any]): """Handle a parsed message from Deepgram Flux. Routes messages to appropriate handlers based on their type. Validates message structure before processing. Args: data: The parsed JSON message data. """ if not self._validate_message(data): return message_type = data.get("type") try: flux_message_type = FluxMessageType(message_type) except ValueError: logger.debug(f"Unhandled message type: {message_type or 'unknown'}") return match flux_message_type: case FluxMessageType.RECEIVE_CONNECTED: await self._handle_connection_established(data) case FluxMessageType.RECEIVE_FATAL_ERROR: await self._handle_fatal_error(data) case FluxMessageType.TURN_INFO: await self._handle_turn_info(data) case FluxMessageType.CONFIGURE_SUCCESS: logger.info(f"{self}: Configure accepted: {data}") await self._on_configure_acked() case FluxMessageType.CONFIGURE_FAILURE: error_code = data.get("error_code", "unknown") description = data.get("description", "no description") error_msg = f"Configure rejected: [{error_code}] {description}" logger.warning(f"{self}: {error_msg}") await self._on_configure_acked() await self.push_error(error_msg=error_msg) async def _handle_connection_established(self, data: dict[str, Any]): """Handle successful connection establishment to Deepgram Flux. This event is fired when the connection to Deepgram Flux is successfully established and ready to receive audio data for transcription processing. """ request_id = data.get("request_id") logger.info(f"{self}: Connected to Flux - ready to stream audio ({request_id=})") # Notify connection is established self._connection_established_event.set() async def _handle_fatal_error(self, data: dict[str, Any]): """Handle fatal error messages from Deepgram Flux. Fatal errors indicate unrecoverable issues with the connection or configuration that require intervention. These errors will cause the connection to be terminated. Args: data: The error message data containing error details. Raises: FluxFatalError: Always raises to trigger error handling in the transport layer. """ error_code = data.get("code", "unknown") description = data.get("description", "no description") deepgram_error = f"{self}: Fatal error [{error_code}] {description}" logger.error(deepgram_error) # Error will be handled by the transport's receive loop error handler raise FluxFatalError(deepgram_error, code=error_code) async def _handle_turn_info(self, data: dict[str, Any]): """Handle TurnInfo events from Deepgram Flux. TurnInfo messages contain various turn-based events that indicate the state of speech processing, including turn boundaries, interim results, and turn finalization events. Args: data: The TurnInfo message data containing event type, transcript and some extra metadata. """ event = data.get("event") transcript = data.get("transcript", "") if not isinstance(event, str): logger.debug(f"Unhandled TurnInfo event (not a string): {event}") return try: flux_event_type = FluxEventType(event) except ValueError: logger.debug(f"Unhandled TurnInfo event: {event}") return match flux_event_type: case FluxEventType.START_OF_TURN: await self._handle_start_of_turn(transcript) case FluxEventType.TURN_RESUMED: await self._handle_turn_resumed(event) case FluxEventType.END_OF_TURN: await self._handle_end_of_turn(transcript, data) case FluxEventType.EAGER_END_OF_TURN: await self._handle_eager_end_of_turn(transcript, data) case FluxEventType.UPDATE: await self._handle_update(transcript) async def _handle_start_of_turn(self, transcript: str): """Handle StartOfTurn events from Deepgram Flux. StartOfTurn events are fired when Deepgram Flux detects the beginning of a new speaking turn. The service will: - Propose a turn start, which the user turn strategies resolve into a UserStartedSpeakingFrame and an interruption Args: transcript: maybe the first few words of the turn. """ logger.debug("User started speaking") self._user_is_speaking = True await self.broadcast_frame(ProposedUserStartedSpeakingFrame) await self._call_event_handler("on_start_of_turn", transcript) if transcript: logger.trace(f"Start of turn transcript: {transcript}") async def _handle_turn_resumed(self, event: str): """Handle TurnResumed events from Deepgram Flux. TurnResumed events indicate that speech has resumed after a brief pause within the same turn, which withdraws the EagerEndOfTurn that preceded it: whatever was generated from that prediction no longer answers the turn the user is still speaking. Args: event: The event type string for logging purposes. """ logger.trace(f"Received event TurnResumed: {event}") await self._cancel_eager_end_of_turn() await self._call_event_handler("on_turn_resumed") def _calculate_average_confidence(self, transcript_data) -> float | None: """Calculate the average confidence from transcript data. Return None if the data is missing or invalid. """ # Example: Assume transcript_data has a list of words with confidence words = transcript_data.get("words") if not words or not isinstance(words, list): return None confidences = [ w.get("confidence") for w in words if isinstance(w.get("confidence"), (float, int)) ] if not confidences: return None return sum(confidences) / len(confidences) def _primary_detected_language(self, data: dict[str, Any]) -> Language | None: """Extract the primary detected language from a TurnInfo payload. On ``flux-general-multi`` the language is read from TurnInfo's ``languages`` field. On ``flux-general-en`` the field is absent, so we fall back to ``Language.EN`` to match the model's fixed language. """ codes = data.get("languages") or [] if codes: return _code_to_pipecat_language(codes[0]) if self._settings.model == "flux-general-en": return Language.EN return None async def _handle_end_of_turn(self, transcript: str, data: dict[str, Any]): """Handle EndOfTurn events from Deepgram Flux. EndOfTurn events are fired when Deepgram Flux determines that a speaking turn has concluded, either due to sufficient silence or end-of-turn confidence thresholds being met. This provides the final transcript for the completed turn. The service will: - Create and send a final TranscriptionFrame with the complete transcript - Trigger transcription handling with tracing for metrics - Propose a turn stop, which the user turn strategies resolve into a UserStoppedSpeakingFrame Args: transcript: The final transcript text for the completed turn. data: The TurnInfo message data containing event type, transcript and some extra metadata. """ logger.debug("User stopped speaking") self._user_is_speaking = False # The turn is committed, so any eager prediction it followed is resolved. self._clear_eager_end_of_turn() # Compute the average confidence average_confidence = self._calculate_average_confidence(data) detected_language = self._primary_detected_language(data) min_confidence = assert_given(self._settings.min_confidence) # No threshold (None or 0.0) → accept. Otherwise require confidence # data and compare; drop if data is missing. if not min_confidence or ( average_confidence is not None and average_confidence > min_confidence ): # Report usage before the transcription frame so tracing can # attach it to the STT span the frame closes. await self.emit_stt_usage_metrics() # EndOfTurn means Flux has determined the turn is complete, # so this TranscriptionFrame is always finalized await self.push_frame( TranscriptionFrame( transcript, self._user_id, time_now_iso8601(), detected_language, result=data, finalized=True, ) ) else: logger.warning( f"Transcription confidence below min_confidence threshold: {average_confidence}" ) await self._handle_transcription(transcript, True, detected_language) await self.broadcast_frame(ProposedUserStoppedSpeakingFrame) await self._call_event_handler("on_end_of_turn", transcript) async def _handle_eager_end_of_turn(self, transcript: str, data: dict[str, Any]): """Handle EagerEndOfTurn events from Deepgram Flux. EagerEndOfTurn events are fired when the end-of-turn confidence reaches the EagerEndOfTurn threshold but hasn't yet reached the full end-of-turn threshold, so a response can be generated during the gap. The prediction may not hold: the user may resume speaking, or the committed transcript may differ from this one. Pair the service with :class:`~pipecat.turns.user_turn_strategies.EagerUserTurnStrategies` to have a response generated here and discarded if either happens. Args: transcript: The predicted transcript for the turn. data: The TurnInfo message data containing event type, transcript and some extra metadata. """ await self._push_eager_end_of_turn( transcript, user_id=self._user_id, language=self._primary_detected_language(data), result=data, ) await self._call_event_handler("on_eager_end_of_turn", transcript) async def _handle_update(self, transcript: str): """Handle Update events from Deepgram Flux. Update events provide incremental transcript updates during an ongoing turn. These events allow for real-time display of transcription progress and can be used to provide visual feedback to users about what's being recognized. Args: transcript: The current partial transcript text for the ongoing turn. """ if transcript: logger.trace(f"Update event: {transcript}") # TTFB (Time To First Byte) metrics are currently disabled for Deepgram Flux. # Ideally, TTFB should measure the time from when a user starts speaking # until we receive the first transcript. However, Deepgram Flux delivers # both the "user started speaking" event and the first transcript simultaneously, # making this timing measurement meaningless in this context. # await self.stop_ttfb_metrics() await self._call_event_handler("on_update", transcript)