frames: use InputTransportMessageFrame instead of InputTransportMessageUrgentFrame

By default, input frames are already urgent.
This commit is contained in:
Aleix Conchillo Flaqué
2025-09-25 12:07:31 -07:00
parent 029d76033d
commit c7dc2e886f
7 changed files with 74 additions and 22 deletions

View File

@@ -32,6 +32,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
- Fixed an issue where local SmartTurn was not being ran in a separate thread. - Fixed an issue where local SmartTurn was not being ran in a separate thread.
### Deprecated
- `InputTransportMessageUrgentFrame` is deprecated, use
`InputTransportMessageFrame` instead.
## [0.0.86] - 2025-09-24 ## [0.0.86] - 2025-09-24
### Added ### Added

View File

@@ -1106,20 +1106,43 @@ class TransportMessageUrgentFrame(SystemFrame):
@dataclass @dataclass
class InputTransportMessageUrgentFrame(TransportMessageUrgentFrame): class InputTransportMessageFrame(SystemFrame):
"""Frame for transport messages received from external sources. """Frame for transport messages received from external sources.
This frame wraps incoming transport messages to distinguish them from outgoing Parameters:
urgent transport messages (TransportMessageUrgentFrame), preventing infinite message: The urgent transport message payload.
message loops in the transport layer. It inherits the message payload from
TransportMessageFrame while marking the message as having been received
rather than generated locally.
Used by transport implementations to properly handle bidirectional message
flow without creating feedback loops.
""" """
pass message: Any
def __str__(self):
return f"{self.name}(message: {self.message})"
@dataclass
class InputTransportMessageUrgentFrame(InputTransportMessageFrame):
"""Frame for transport messages received from external sources.
.. deprecated:: 0.0.87
This frame is deprecated and will be removed in a future version.
Instead, use `InputTransportMessageFrame`.
Parameters:
message: The urgent transport message payload.
"""
def __post_init__(self):
super().__post_init__()
import warnings
with warnings.catch_warnings():
warnings.simplefilter("always")
warnings.warn(
"InputTransportMessageUrgentFrame is deprecated and will be removed in a future version. "
"Instead, use InputTransportMessageFrame.",
DeprecationWarning,
stacklevel=2,
)
@dataclass @dataclass

View File

@@ -42,6 +42,7 @@ from pipecat.frames.frames import (
Frame, Frame,
FunctionCallResultFrame, FunctionCallResultFrame,
InputAudioRawFrame, InputAudioRawFrame,
InputTransportMessageUrgentFrame,
InterimTranscriptionFrame, InterimTranscriptionFrame,
LLMConfigureOutputFrame, LLMConfigureOutputFrame,
LLMContextFrame, LLMContextFrame,
@@ -1418,7 +1419,7 @@ class RTVIProcessor(FrameProcessor):
elif isinstance(frame, ErrorFrame): elif isinstance(frame, ErrorFrame):
await self._send_error_frame(frame) await self._send_error_frame(frame)
await self.push_frame(frame, direction) await self.push_frame(frame, direction)
elif isinstance(frame, TransportMessageUrgentFrame): elif isinstance(frame, InputTransportMessageUrgentFrame):
await self._handle_transport_message(frame) await self._handle_transport_message(frame)
# All other system frames # All other system frames
elif isinstance(frame, SystemFrame): elif isinstance(frame, SystemFrame):
@@ -1481,7 +1482,7 @@ class RTVIProcessor(FrameProcessor):
await self._handle_message(message) await self._handle_message(message)
self._message_queue.task_done() self._message_queue.task_done()
async def _handle_transport_message(self, frame: TransportMessageUrgentFrame): async def _handle_transport_message(self, frame: InputTransportMessageUrgentFrame):
"""Handle an incoming transport message frame.""" """Handle an incoming transport message frame."""
try: try:
transport_message = frame.message transport_message = frame.message

View File

@@ -15,6 +15,7 @@ import pipecat.frames.protobufs.frames_pb2 as frame_protos
from pipecat.frames.frames import ( from pipecat.frames.frames import (
Frame, Frame,
InputAudioRawFrame, InputAudioRawFrame,
InputTransportMessageFrame,
OutputAudioRawFrame, OutputAudioRawFrame,
TextFrame, TextFrame,
TranscriptionFrame, TranscriptionFrame,
@@ -138,7 +139,7 @@ class ProtobufFrameSerializer(FrameSerializer):
if class_name == MessageFrame: if class_name == MessageFrame:
try: try:
msg = json.loads(args_dict["data"]) msg = json.loads(args_dict["data"])
instance = TransportMessageUrgentFrame(message=msg) instance = InputTransportMessageFrame(message=msg)
logger.debug(f"ProtobufFrameSerializer: Transport message {instance}") logger.debug(f"ProtobufFrameSerializer: Transport message {instance}")
except Exception as e: except Exception as e:
logger.error(f"Error parsing MessageFrame data: {e}") logger.error(f"Error parsing MessageFrame data: {e}")

View File

@@ -29,7 +29,6 @@ from pipecat.frames.frames import (
CancelFrame, CancelFrame,
EndFrame, EndFrame,
Frame, Frame,
InputTransportMessageUrgentFrame,
InterruptionFrame, InterruptionFrame,
MixerControlFrame, MixerControlFrame,
OutputAudioRawFrame, OutputAudioRawFrame,
@@ -307,9 +306,7 @@ class BaseOutputTransport(FrameProcessor):
elif isinstance(frame, InterruptionFrame): elif isinstance(frame, InterruptionFrame):
await self.push_frame(frame, direction) await self.push_frame(frame, direction)
await self._handle_frame(frame) await self._handle_frame(frame)
elif isinstance(frame, TransportMessageUrgentFrame) and not isinstance( elif isinstance(frame, TransportMessageUrgentFrame):
frame, InputTransportMessageUrgentFrame
):
await self.send_message(frame) await self.send_message(frame)
elif isinstance(frame, OutputDTMFUrgentFrame): elif isinstance(frame, OutputDTMFUrgentFrame):
await self.write_dtmf(frame) await self.write_dtmf(frame)

View File

@@ -30,7 +30,7 @@ from pipecat.frames.frames import (
ErrorFrame, ErrorFrame,
Frame, Frame,
InputAudioRawFrame, InputAudioRawFrame,
InputTransportMessageUrgentFrame, InputTransportMessageFrame,
InterimTranscriptionFrame, InterimTranscriptionFrame,
OutputAudioRawFrame, OutputAudioRawFrame,
OutputImageRawFrame, OutputImageRawFrame,
@@ -96,7 +96,7 @@ class DailyTransportMessageUrgentFrame(TransportMessageUrgentFrame):
@dataclass @dataclass
class DailyInputTransportMessageUrgentFrame(InputTransportMessageUrgentFrame): class DailyInputTransportMessageFrame(InputTransportMessageFrame):
"""Frame for input urgent transport messages in Daily calls. """Frame for input urgent transport messages in Daily calls.
Parameters: Parameters:
@@ -106,6 +106,31 @@ class DailyInputTransportMessageUrgentFrame(InputTransportMessageUrgentFrame):
participant_id: Optional[str] = None participant_id: Optional[str] = None
class DailyInputTransportMessageUrgentFrame(DailyInputTransportMessageFrame):
"""Frame for input urgent transport messages in Daily calls.
.. deprecated:: 0.0.87
This frame is deprecated and will be removed in a future version.
Instead, use `DailyInputTransportMessageFrame`.
Parameters:
participant_id: Optional ID of the participant this message is for/from.
"""
def __post_init__(self):
super().__post_init__()
import warnings
with warnings.catch_warnings():
warnings.simplefilter("always")
warnings.warn(
"DailyInputTransportMessageUrgentFrame is deprecated and will be removed in a future version. "
"Instead, use DailyInputTransportMessageFrame.",
DeprecationWarning,
stacklevel=2,
)
@dataclass @dataclass
class DailyUpdateRemoteParticipantsFrame(ControlFrame): class DailyUpdateRemoteParticipantsFrame(ControlFrame):
"""Frame to update remote participants in Daily calls. """Frame to update remote participants in Daily calls.
@@ -1621,7 +1646,7 @@ class DailyInputTransport(BaseInputTransport):
message: The message data to send. message: The message data to send.
sender: ID of the message sender. sender: ID of the message sender.
""" """
frame = DailyInputTransportMessageUrgentFrame(message=message, participant_id=sender) frame = DailyInputTransportMessageFrame(message=message, participant_id=sender)
await self.push_frame(frame) await self.push_frame(frame)
# #

View File

@@ -26,7 +26,7 @@ from pipecat.frames.frames import (
EndFrame, EndFrame,
Frame, Frame,
InputAudioRawFrame, InputAudioRawFrame,
InputTransportMessageUrgentFrame, InputTransportMessageFrame,
OutputAudioRawFrame, OutputAudioRawFrame,
OutputImageRawFrame, OutputImageRawFrame,
SpriteFrame, SpriteFrame,
@@ -683,7 +683,7 @@ class SmallWebRTCInputTransport(BaseInputTransport):
message: The application message to process. message: The application message to process.
""" """
logger.debug(f"Received app message inside SmallWebRTCInputTransport {message}") logger.debug(f"Received app message inside SmallWebRTCInputTransport {message}")
frame = InputTransportMessageUrgentFrame(message=message) frame = InputTransportMessageFrame(message=message)
await self.push_frame(frame) await self.push_frame(frame)
# Add this method similar to DailyInputTransport.request_participant_image # Add this method similar to DailyInputTransport.request_participant_image