missing no longer necessary to call super().process_frame(frame, direction)
This commit is contained in:
@@ -292,8 +292,6 @@ class TTSService(AIService):
|
|||||||
await self.queue_frame(TTSSpeakFrame(text))
|
await self.queue_frame(TTSSpeakFrame(text))
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
if isinstance(frame, TextFrame):
|
if isinstance(frame, TextFrame):
|
||||||
await self._process_text_frame(frame)
|
await self._process_text_frame(frame)
|
||||||
elif isinstance(frame, StartInterruptionFrame):
|
elif isinstance(frame, StartInterruptionFrame):
|
||||||
@@ -410,8 +408,6 @@ class WordTTSService(TTSService):
|
|||||||
await self._stop_words_task()
|
await self._stop_words_task()
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
if isinstance(frame, (LLMFullResponseEndFrame, EndFrame)):
|
if isinstance(frame, (LLMFullResponseEndFrame, EndFrame)):
|
||||||
await self.flush_audio()
|
await self.flush_audio()
|
||||||
|
|
||||||
@@ -498,8 +494,6 @@ class STTService(AIService):
|
|||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
"""Processes a frame of audio data, either buffering or transcribing it."""
|
"""Processes a frame of audio data, either buffering or transcribing it."""
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
if isinstance(frame, AudioRawFrame):
|
if isinstance(frame, AudioRawFrame):
|
||||||
# In this service we accumulate audio internally and at the end we
|
# In this service we accumulate audio internally and at the end we
|
||||||
# push a TextFrame. We also push audio downstream in case someone
|
# push a TextFrame. We also push audio downstream in case someone
|
||||||
@@ -597,8 +591,6 @@ class ImageGenService(AIService):
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
if isinstance(frame, TextFrame):
|
if isinstance(frame, TextFrame):
|
||||||
await self.push_frame(frame, direction)
|
await self.push_frame(frame, direction)
|
||||||
await self.start_processing_metrics()
|
await self.start_processing_metrics()
|
||||||
@@ -620,8 +612,6 @@ class VisionService(AIService):
|
|||||||
pass
|
pass
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
if isinstance(frame, VisionImageRawFrame):
|
if isinstance(frame, VisionImageRawFrame):
|
||||||
await self.start_processing_metrics()
|
await self.start_processing_metrics()
|
||||||
await self.process_generator(self.run_vision(frame))
|
await self.process_generator(self.run_vision(frame))
|
||||||
|
|||||||
@@ -270,8 +270,6 @@ class AnthropicLLMService(LLMService):
|
|||||||
)
|
)
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
context = None
|
context = None
|
||||||
if isinstance(frame, OpenAILLMContextFrame):
|
if isinstance(frame, OpenAILLMContextFrame):
|
||||||
context: "AnthropicLLMContext" = AnthropicLLMContext.upgrade_to_anthropic(frame.context)
|
context: "AnthropicLLMContext" = AnthropicLLMContext.upgrade_to_anthropic(frame.context)
|
||||||
@@ -611,7 +609,6 @@ class AnthropicUserContextAggregator(LLMUserContextAggregator):
|
|||||||
self._context = AnthropicLLMContext.from_openai_context(context)
|
self._context = AnthropicLLMContext.from_openai_context(context)
|
||||||
|
|
||||||
async def process_frame(self, frame, direction):
|
async def process_frame(self, frame, direction):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
# Our parent method has already called push_frame(). So we can't interrupt the
|
# Our parent method has already called push_frame(). So we can't interrupt the
|
||||||
# flow here and we don't need to call push_frame() ourselves. Possibly something
|
# flow here and we don't need to call push_frame() ourselves. Possibly something
|
||||||
# to talk through (tagging @aleix). At some point we might need to refactor these
|
# to talk through (tagging @aleix). At some point we might need to refactor these
|
||||||
@@ -664,7 +661,6 @@ class AnthropicAssistantContextAggregator(LLMAssistantContextAggregator):
|
|||||||
self._pending_image_frame_message = None
|
self._pending_image_frame_message = None
|
||||||
|
|
||||||
async def process_frame(self, frame, direction):
|
async def process_frame(self, frame, direction):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
# See note above about not calling push_frame() here.
|
# See note above about not calling push_frame() here.
|
||||||
if isinstance(frame, StartInterruptionFrame):
|
if isinstance(frame, StartInterruptionFrame):
|
||||||
self._function_call_in_progress = None
|
self._function_call_in_progress = None
|
||||||
|
|||||||
@@ -90,7 +90,6 @@ class CanonicalMetricsService(AIService):
|
|||||||
await self._process_audio()
|
await self._process_audio()
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
await self.push_frame(frame, direction)
|
await self.push_frame(frame, direction)
|
||||||
|
|
||||||
async def _process_audio(self):
|
async def _process_audio(self):
|
||||||
|
|||||||
@@ -287,8 +287,6 @@ class CartesiaTTSService(WordTTSService):
|
|||||||
await self._connect_websocket()
|
await self._connect_websocket()
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
# If we received a TTSSpeakFrame and the LLM response included text (it
|
# If we received a TTSSpeakFrame and the LLM response included text (it
|
||||||
# might be that it's only a function calling response) we pause
|
# might be that it's only a function calling response) we pause
|
||||||
# processing more frames until we receive a BotStoppedSpeakingFrame.
|
# processing more frames until we receive a BotStoppedSpeakingFrame.
|
||||||
|
|||||||
@@ -272,8 +272,6 @@ class ElevenLabsTTSService(WordTTSService):
|
|||||||
await self.add_word_timestamps([("LLMFullResponseEndFrame", 0), ("Reset", 0)])
|
await self.add_word_timestamps([("LLMFullResponseEndFrame", 0), ("Reset", 0)])
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
# If we received a TTSSpeakFrame and the LLM response included text (it
|
# If we received a TTSSpeakFrame and the LLM response included text (it
|
||||||
# might be that it's only a function calling response) we pause
|
# might be that it's only a function calling response) we pause
|
||||||
# processing more frames until we receive a BotStoppedSpeakingFrame.
|
# processing more frames until we receive a BotStoppedSpeakingFrame.
|
||||||
|
|||||||
@@ -107,7 +107,6 @@ class GeminiMultimodalLiveContext(OpenAILLMContext):
|
|||||||
|
|
||||||
class GeminiMultimodalLiveUserContextAggregator(OpenAIUserContextAggregator):
|
class GeminiMultimodalLiveUserContextAggregator(OpenAIUserContextAggregator):
|
||||||
async def process_frame(self, frame, direction):
|
async def process_frame(self, frame, direction):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
# kind of a hack just to pass the LLMMessagesAppendFrame through, but it's fine for now
|
# kind of a hack just to pass the LLMMessagesAppendFrame through, but it's fine for now
|
||||||
if isinstance(frame, LLMMessagesAppendFrame):
|
if isinstance(frame, LLMMessagesAppendFrame):
|
||||||
await self.push_frame(frame, direction)
|
await self.push_frame(frame, direction)
|
||||||
@@ -305,8 +304,6 @@ class GeminiMultimodalLiveLLMService(LLMService):
|
|||||||
#
|
#
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
# logger.debug(f"Processing frame: {frame}")
|
# logger.debug(f"Processing frame: {frame}")
|
||||||
|
|
||||||
if isinstance(frame, TranscriptionFrame):
|
if isinstance(frame, TranscriptionFrame):
|
||||||
|
|||||||
@@ -652,8 +652,6 @@ class GoogleLLMService(LLMService):
|
|||||||
await self.push_frame(LLMFullResponseEndFrame())
|
await self.push_frame(LLMFullResponseEndFrame())
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
context = None
|
context = None
|
||||||
|
|
||||||
if isinstance(frame, OpenAILLMContextFrame):
|
if isinstance(frame, OpenAILLMContextFrame):
|
||||||
|
|||||||
@@ -286,8 +286,6 @@ class BaseOpenAILLMService(LLMService):
|
|||||||
)
|
)
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
context = None
|
context = None
|
||||||
if isinstance(frame, OpenAILLMContextFrame):
|
if isinstance(frame, OpenAILLMContextFrame):
|
||||||
context: OpenAILLMContext = frame.context
|
context: OpenAILLMContext = frame.context
|
||||||
@@ -475,7 +473,6 @@ class OpenAIUserContextAggregator(LLMUserContextAggregator):
|
|||||||
super().__init__(context=context)
|
super().__init__(context=context)
|
||||||
|
|
||||||
async def process_frame(self, frame, direction):
|
async def process_frame(self, frame, direction):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
# Our parent method has already called push_frame(). So we can't interrupt the
|
# Our parent method has already called push_frame(). So we can't interrupt the
|
||||||
# flow here and we don't need to call push_frame() ourselves.
|
# flow here and we don't need to call push_frame() ourselves.
|
||||||
try:
|
try:
|
||||||
@@ -516,7 +513,6 @@ class OpenAIAssistantContextAggregator(LLMAssistantContextAggregator):
|
|||||||
self._pending_image_frame_message = None
|
self._pending_image_frame_message = None
|
||||||
|
|
||||||
async def process_frame(self, frame, direction):
|
async def process_frame(self, frame, direction):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
# See note above about not calling push_frame() here.
|
# See note above about not calling push_frame() here.
|
||||||
if isinstance(frame, StartInterruptionFrame):
|
if isinstance(frame, StartInterruptionFrame):
|
||||||
self._function_calls_in_progress.clear()
|
self._function_calls_in_progress.clear()
|
||||||
|
|||||||
@@ -148,7 +148,6 @@ class OpenAIRealtimeUserContextAggregator(OpenAIUserContextAggregator):
|
|||||||
async def process_frame(
|
async def process_frame(
|
||||||
self, frame: Frame, direction: FrameDirection = FrameDirection.DOWNSTREAM
|
self, frame: Frame, direction: FrameDirection = FrameDirection.DOWNSTREAM
|
||||||
):
|
):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
# Parent does not push LLMMessagesUpdateFrame. This ensures that in a typical pipeline,
|
# Parent does not push LLMMessagesUpdateFrame. This ensures that in a typical pipeline,
|
||||||
# messages are only processed by the user context aggregator, which is generally what we want. But
|
# messages are only processed by the user context aggregator, which is generally what we want. But
|
||||||
# we also need to send new messages over the websocket, so the openai realtime API has them
|
# we also need to send new messages over the websocket, so the openai realtime API has them
|
||||||
|
|||||||
@@ -170,8 +170,6 @@ class OpenAIRealtimeBetaLLMService(LLMService):
|
|||||||
#
|
#
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
if isinstance(frame, TranscriptionFrame):
|
if isinstance(frame, TranscriptionFrame):
|
||||||
pass
|
pass
|
||||||
elif isinstance(frame, OpenAILLMContextFrame):
|
elif isinstance(frame, OpenAILLMContextFrame):
|
||||||
|
|||||||
@@ -265,8 +265,6 @@ class PlayHTTTSService(TTSService):
|
|||||||
await self._connect_websocket()
|
await self._connect_websocket()
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
# If we received a TTSSpeakFrame and the LLM response included text (it
|
# If we received a TTSSpeakFrame and the LLM response included text (it
|
||||||
# might be that it's only a function calling response) we pause
|
# might be that it's only a function calling response) we pause
|
||||||
# processing more frames until we receive a BotStoppedSpeakingFrame.
|
# processing more frames until we receive a BotStoppedSpeakingFrame.
|
||||||
|
|||||||
@@ -92,7 +92,6 @@ class TavusVideoService(AIService):
|
|||||||
await self._send_audio_message(audio_base64, done=done)
|
await self._send_audio_message(audio_base64, done=done)
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
if isinstance(frame, TTSStartedFrame):
|
if isinstance(frame, TTSStartedFrame):
|
||||||
await self.start_processing_metrics()
|
await self.start_processing_metrics()
|
||||||
await self.start_ttfb_metrics()
|
await self.start_ttfb_metrics()
|
||||||
|
|||||||
@@ -101,8 +101,6 @@ class FastAPIWebsocketOutputTransport(BaseOutputTransport):
|
|||||||
self._next_send_time = 0
|
self._next_send_time = 0
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
if isinstance(frame, StartInterruptionFrame):
|
if isinstance(frame, StartInterruptionFrame):
|
||||||
await self._write_frame(frame)
|
await self._write_frame(frame)
|
||||||
self._next_send_time = 0
|
self._next_send_time = 0
|
||||||
|
|||||||
@@ -139,8 +139,6 @@ class WebsocketServerOutputTransport(BaseOutputTransport):
|
|||||||
self._websocket = websocket
|
self._websocket = websocket
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
if isinstance(frame, StartInterruptionFrame):
|
if isinstance(frame, StartInterruptionFrame):
|
||||||
await self._write_frame(frame)
|
await self._write_frame(frame)
|
||||||
self._next_send_time = 0
|
self._next_send_time = 0
|
||||||
|
|||||||
@@ -727,8 +727,6 @@ class DailyInputTransport(BaseInputTransport):
|
|||||||
#
|
#
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
if isinstance(frame, UserImageRequestFrame):
|
if isinstance(frame, UserImageRequestFrame):
|
||||||
await self.request_participant_image(frame.user_id)
|
await self.request_participant_image(frame.user_id)
|
||||||
|
|
||||||
|
|||||||
@@ -16,7 +16,6 @@ from pipecat.frames.frames import (
|
|||||||
AudioRawFrame,
|
AudioRawFrame,
|
||||||
CancelFrame,
|
CancelFrame,
|
||||||
EndFrame,
|
EndFrame,
|
||||||
Frame,
|
|
||||||
InputAudioRawFrame,
|
InputAudioRawFrame,
|
||||||
OutputAudioRawFrame,
|
OutputAudioRawFrame,
|
||||||
StartFrame,
|
StartFrame,
|
||||||
@@ -334,12 +333,6 @@ class LiveKitInputTransport(BaseInputTransport):
|
|||||||
await self._client.disconnect()
|
await self._client.disconnect()
|
||||||
logger.info("LiveKitInputTransport stopped")
|
logger.info("LiveKitInputTransport stopped")
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
|
||||||
if isinstance(frame, EndFrame):
|
|
||||||
await self.stop(frame)
|
|
||||||
else:
|
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
async def cancel(self, frame: CancelFrame):
|
async def cancel(self, frame: CancelFrame):
|
||||||
await super().cancel(frame)
|
await super().cancel(frame)
|
||||||
await self._client.disconnect()
|
await self._client.disconnect()
|
||||||
@@ -411,12 +404,6 @@ class LiveKitOutputTransport(BaseOutputTransport):
|
|||||||
await self._client.disconnect()
|
await self._client.disconnect()
|
||||||
logger.info("LiveKitOutputTransport stopped")
|
logger.info("LiveKitOutputTransport stopped")
|
||||||
|
|
||||||
async def process_frame(self, frame: Frame, direction: FrameDirection):
|
|
||||||
if isinstance(frame, EndFrame):
|
|
||||||
await self.stop(frame)
|
|
||||||
else:
|
|
||||||
await super().process_frame(frame, direction)
|
|
||||||
|
|
||||||
async def cancel(self, frame: CancelFrame):
|
async def cancel(self, frame: CancelFrame):
|
||||||
await super().cancel(frame)
|
await super().cancel(frame)
|
||||||
await self._client.disconnect()
|
await self._client.disconnect()
|
||||||
|
|||||||
Reference in New Issue
Block a user