update send_metrics() to support changes introduced in #474

This commit is contained in:
joachimchauvet
2024-09-25 21:38:55 +03:00
parent ec5998bc36
commit 447baad5c3

View File

@@ -16,6 +16,12 @@ from pipecat.frames.frames import (
StartFrame, StartFrame,
TransportMessageFrame, TransportMessageFrame,
) )
from pipecat.metrics.metrics import (
LLMUsageMetricsData,
ProcessingMetricsData,
TTFBMetricsData,
TTSUsageMetricsData,
)
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor from pipecat.processors.frame_processor import FrameDirection, FrameProcessor
from pipecat.transports.base_input import BaseInputTransport from pipecat.transports.base_input import BaseInputTransport
from pipecat.transports.base_output import BaseOutputTransport from pipecat.transports.base_output import BaseOutputTransport
@@ -257,9 +263,9 @@ class LiveKitTransportClient:
async def _async_on_connected(self): async def _async_on_connected(self):
await self._callbacks.on_connected() await self._callbacks.on_connected()
async def _async_on_disconnected(self): async def _async_on_disconnected(self, reason=None):
self._connected = False self._connected = False
logger.info(f"Disconnected from {self._room_name}") logger.info(f"Disconnected from {self._room_name}. Reason: {reason}")
await self._callbacks.on_disconnected() await self._callbacks.on_disconnected()
async def _process_audio_stream(self, audio_stream: rtc.AudioStream, participant_id: str): async def _process_audio_stream(self, audio_stream: rtc.AudioStream, participant_id: str):
@@ -413,14 +419,23 @@ class LiveKitOutputTransport(BaseOutputTransport):
async def send_metrics(self, frame: MetricsFrame): async def send_metrics(self, frame: MetricsFrame):
metrics = {} metrics = {}
if frame.ttfb: for d in frame.data:
metrics["ttfb"] = frame.ttfb if isinstance(d, TTFBMetricsData):
if frame.processing: if "ttfb" not in metrics:
metrics["processing"] = frame.processing metrics["ttfb"] = []
if hasattr(frame, "tokens"): metrics["ttfb"].append(d.model_dump())
metrics["tokens"] = frame.tokens elif isinstance(d, ProcessingMetricsData):
if hasattr(frame, "characters"): if "processing" not in metrics:
metrics["characters"] = frame.characters metrics["processing"] = []
metrics["processing"].append(d.model_dump())
elif isinstance(d, LLMUsageMetricsData):
if "tokens" not in metrics:
metrics["tokens"] = []
metrics["tokens"].append(d.value.model_dump(exclude_none=True))
elif isinstance(d, TTSUsageMetricsData):
if "characters" not in metrics:
metrics["characters"] = []
metrics["characters"].append(d.model_dump())
message = LiveKitTransportMessageFrame( message = LiveKitTransportMessageFrame(
message={"type": "pipecat-metrics", "metrics": metrics} message={"type": "pipecat-metrics", "metrics": metrics}