Logging the last 10 frames received in case idle timeout is detected.
This commit is contained in:
@@ -6,7 +6,8 @@
|
|||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import time
|
import time
|
||||||
from typing import Any, AsyncIterable, Dict, Iterable, List, Optional, Sequence, Tuple, Type
|
from collections import deque
|
||||||
|
from typing import Any, AsyncIterable, Deque, Dict, Iterable, List, Optional, Sequence, Tuple, Type
|
||||||
|
|
||||||
from loguru import logger
|
from loguru import logger
|
||||||
from pydantic import BaseModel, ConfigDict, Field
|
from pydantic import BaseModel, ConfigDict, Field
|
||||||
@@ -23,6 +24,7 @@ from pipecat.frames.frames import (
|
|||||||
ErrorFrame,
|
ErrorFrame,
|
||||||
Frame,
|
Frame,
|
||||||
HeartbeatFrame,
|
HeartbeatFrame,
|
||||||
|
InputAudioRawFrame,
|
||||||
LLMFullResponseEndFrame,
|
LLMFullResponseEndFrame,
|
||||||
MetricsFrame,
|
MetricsFrame,
|
||||||
StartFrame,
|
StartFrame,
|
||||||
@@ -646,12 +648,17 @@ class PipelineTask(BaseTask):
|
|||||||
"""
|
"""
|
||||||
running = True
|
running = True
|
||||||
last_frame_time = 0
|
last_frame_time = 0
|
||||||
|
frame_buffer = deque(maxlen=10) # Store last 10 frames
|
||||||
|
|
||||||
while running:
|
while running:
|
||||||
try:
|
try:
|
||||||
frame = await asyncio.wait_for(
|
frame = await asyncio.wait_for(
|
||||||
self._idle_queue.get(), timeout=self._idle_timeout_secs
|
self._idle_queue.get(), timeout=self._idle_timeout_secs
|
||||||
)
|
)
|
||||||
|
|
||||||
|
if not isinstance(frame, InputAudioRawFrame):
|
||||||
|
frame_buffer.append(frame)
|
||||||
|
|
||||||
if isinstance(frame, StartFrame) or isinstance(frame, self._idle_timeout_frames):
|
if isinstance(frame, StartFrame) or isinstance(frame, self._idle_timeout_frames):
|
||||||
# If we find a StartFrame or one of the frames that prevents a
|
# If we find a StartFrame or one of the frames that prevents a
|
||||||
# time out we update the time.
|
# time out we update the time.
|
||||||
@@ -662,7 +669,7 @@ class PipelineTask(BaseTask):
|
|||||||
# valid frames.
|
# valid frames.
|
||||||
diff_time = time.time() - last_frame_time
|
diff_time = time.time() - last_frame_time
|
||||||
if diff_time >= self._idle_timeout_secs:
|
if diff_time >= self._idle_timeout_secs:
|
||||||
running = await self._idle_timeout_detected()
|
running = await self._idle_timeout_detected(frame_buffer)
|
||||||
# Reset `last_frame_time` so we don't trigger another
|
# Reset `last_frame_time` so we don't trigger another
|
||||||
# immediate idle timeout if we are not cancelling. For
|
# immediate idle timeout if we are not cancelling. For
|
||||||
# example, we might want to force the bot to say goodbye
|
# example, we might want to force the bot to say goodbye
|
||||||
@@ -670,15 +677,20 @@ class PipelineTask(BaseTask):
|
|||||||
last_frame_time = time.time()
|
last_frame_time = time.time()
|
||||||
|
|
||||||
self._idle_queue.task_done()
|
self._idle_queue.task_done()
|
||||||
except asyncio.TimeoutError:
|
|
||||||
running = await self._idle_timeout_detected()
|
|
||||||
|
|
||||||
async def _idle_timeout_detected(self) -> bool:
|
except asyncio.TimeoutError:
|
||||||
|
running = await self._idle_timeout_detected(frame_buffer)
|
||||||
|
|
||||||
|
async def _idle_timeout_detected(self, last_frames: Deque[Frame]) -> bool:
|
||||||
"""Logic for when the pipeline is idle.
|
"""Logic for when the pipeline is idle.
|
||||||
|
|
||||||
Returns:
|
Returns:
|
||||||
bool: Whther the pipeline task is being cancelled or not.
|
bool: Whther the pipeline task is being cancelled or not.
|
||||||
"""
|
"""
|
||||||
|
logger.warning("Idle timeout detected. Last 10 frames received:")
|
||||||
|
for i, frame in enumerate(last_frames, 1):
|
||||||
|
logger.warning(f"Frame {i}: {frame}")
|
||||||
|
|
||||||
await self._call_event_handler("on_idle_timeout")
|
await self._call_event_handler("on_idle_timeout")
|
||||||
if self._cancel_on_idle_timeout:
|
if self._cancel_on_idle_timeout:
|
||||||
logger.warning(f"Idle pipeline detected, cancelling pipeline task...")
|
logger.warning(f"Idle pipeline detected, cancelling pipeline task...")
|
||||||
|
|||||||
Reference in New Issue
Block a user