Source code for pipecat.turns.user_stop.eager_user_turn_stop_strategy

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

"""User turn stop strategy that answers an eager end of turn speculatively."""

import asyncio

from loguru import logger

from pipecat.frames.frames import (
    EagerEndOfTurnCancelFrame,
    EagerTranscriptionFrame,
    Frame,
)
from pipecat.turns.types import ProcessFrameResult, UserTurnSpeculation
from pipecat.turns.user_stop.eager_match_policy import EagerMatchPolicy, NormalizedMatch
from pipecat.turns.user_stop.external_user_turn_stop_strategy import ExternalUserTurnStopStrategy


[docs] class EagerUserTurnStopStrategy(ExternalUserTurnStopStrategy): """Answers an eager end of turn while the turn is still open. Some STT services predict the end of a turn before committing to it, and withdraw the prediction if the user turns out to be mid-sentence. This strategy starts generating a response on that prediction, so the gap between it and the committed end of turn is spent generating rather than waiting. The prediction can be wrong in three ways, and each withdraws the response: - the user resumes speaking, and the service withdraws the eager end of turn - the committed transcript differs from the eager one, per ``match_policy`` - the turn never commits within ``speculation_timeout`` Every one of those leaves as an :class:`~pipecat.frames.frames.EagerEndOfTurnCancelFrame` from here, so a response is never dropped somewhere the rest of the pipeline cannot see. Nothing the speculation produces reaches the user or the context. The inference runs against a provisional context, and its response is held by the LLM service's :class:`~pipecat.turns.speculation_gate.SpeculationGate` until the turn is confirmed. The turn ends normally: the user message written to the context is always the committed transcript, never the eager one. Install it with :class:`~pipecat.turns.user_turn_strategies.EagerUserTurnStrategies` rather than directly — the service owns turn detection here, so it replaces the detector chain instead of running alongside it. """
[docs] def __init__( self, *, match_policy: EagerMatchPolicy | None = None, speculation_timeout: float = 5.0, **kwargs, ): """Initialize the eager user turn stop strategy. Args: match_policy: Decides whether the committed transcript is close enough to the eager one to keep the speculative response. Defaults to :class:`~pipecat.turns.user_stop.NormalizedMatch`, which ignores the capitalization and punctuation services commonly add when they commit a transcript. Pass :class:`~pipecat.turns.user_stop.ExactMatch` to require the two to be identical. speculation_timeout: Seconds a prediction may go unresolved before it is withdrawn. A service that stops sending turn signals mid-speculation would otherwise leave the response held and the bot silent for the rest of the session. **kwargs: Additional keyword arguments forwarded to the base class. """ super().__init__(**kwargs) self._match_policy = match_policy or NormalizedMatch() self._speculation_timeout = speculation_timeout self._speculation_timeout_task: asyncio.Task | None = None self._speculation: UserTurnSpeculation | None = None
@property def match_policy(self) -> EagerMatchPolicy: """The policy deciding whether a speculative response still applies.""" return self._match_policy
[docs] async def process_frame(self, frame: Frame) -> ProcessFrameResult: """Start a speculation on an eager end of turn, or withdraw one. Args: frame: The frame to be analyzed. Returns: Always CONTINUE, so subsequent stop strategies are evaluated. """ if isinstance(frame, EagerTranscriptionFrame): await self._speculate(frame) elif isinstance(frame, EagerEndOfTurnCancelFrame): # The service withdrew its prediction, and its frame reaches every # consumer on its own. Only our own state is left to clear. await self._forget() return await super().process_frame(frame)
[docs] async def trigger_user_turn_stopped(self, *, enable_user_speaking_frames: bool | None = None): """End the turn, keeping the speculative response only if it still applies. Args: enable_user_speaking_frames: Whether to emit :class:`~pipecat.frames.frames.UserStoppedSpeakingFrame` for this turn. """ speculation = await self._take_speculation() if not speculation: await super().trigger_user_turn_stopped( enable_user_speaking_frames=enable_user_speaking_frames ) return if self._match_policy.matches(speculation.text, self._text): logger.debug(f"{self}: eager end of turn held, keeping the speculative response") # Inference already ran, on the eager transcript. Only finalize: # the UserStoppedSpeakingFrame it emits is what releases the # response, and the flag stops a second inference answering the # same turn. await self.trigger_user_turn_finalized( enable_user_speaking_frames=enable_user_speaking_frames, confirms_speculation=True, ) return logger.debug( f"{self}: eager end of turn missed, discarding the speculative response " f"(eager: [{speculation.text}], committed: [{self._text}])" ) await self.trigger_user_turn_speculation_cancelled() # Inference has to run again, on the committed transcript, so fire both # events rather than just finalizing. await super().trigger_user_turn_stopped( enable_user_speaking_frames=enable_user_speaking_frames )
async def _reset(self): """Clear per-turn state. Runs at both turn boundaries.""" speculation = await self._take_speculation() await super()._reset() if speculation: # The turn ended without resolving the speculation — the stop # watchdog, an interruption, session end. Nothing else will withdraw # it, so the response would be held until the gate times out. logger.debug(f"{self}: turn ended unresolved, discarding the speculative response") await self.trigger_user_turn_speculation_cancelled() async def _speculate(self, frame: EagerTranscriptionFrame): """Answer an eager end of turn, leaving the turn open.""" # Segments committed earlier in this turn are part of what the LLM will # see, so they're part of what the committed transcript is compared to. speculation = UserTurnSpeculation(text=self._text + frame.text) self._speculation = speculation self._speculation_timeout_task = self.task_manager.create_task( self._speculation_timeout_handler(speculation), f"{self}::_speculation_timeout_handler", ) logger.debug(f"{self}: speculating on eager end of turn: [{speculation.text}]") await self.trigger_user_turn_inference_triggered(speculation=speculation) async def _speculation_timeout_handler(self, speculation: UserTurnSpeculation): """Withdraw a prediction the turn never resolved.""" await asyncio.sleep(self._speculation_timeout) # Cleared here rather than by a cancellation, since this timer fired. self._speculation_timeout_task = None if self._speculation is not speculation: return logger.debug( f"{self}: eager end of turn unresolved after {self._speculation_timeout}s, " "discarding the speculative response" ) self._speculation = None await self.trigger_user_turn_speculation_cancelled() async def _take_speculation(self) -> UserTurnSpeculation | None: """Take the speculation in flight, stopping the clock on it. The single exit, so a resolved speculation can never be withdrawn again by a timer still running for it. """ speculation, self._speculation = self._speculation, None if self._speculation_timeout_task: task, self._speculation_timeout_task = self._speculation_timeout_task, None await self.cancel_task(task) return speculation async def _forget(self): """Drop a speculation the service withdrew.""" await self._take_speculation()