examples(simple-chatbot): use RTVIObserver for server-client messages

This commit is contained in:
Aleix Conchillo Flaqué
2025-01-15 09:50:04 -08:00
parent dd9f9179cc
commit e50c76d075
2 changed files with 4 additions and 54 deletions

View File

@@ -41,17 +41,8 @@ from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineParams, PipelineTask from pipecat.pipeline.task import PipelineParams, PipelineTask
from pipecat.processors.aggregators.openai_llm_context import OpenAILLMContext from pipecat.processors.aggregators.openai_llm_context import OpenAILLMContext
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor from pipecat.processors.frame_processor import FrameDirection, FrameProcessor
from pipecat.processors.frameworks.rtvi import ( from pipecat.processors.frameworks.rtvi import RTVIConfig, RTVIProcessor
RTVIBotTranscriptionProcessor,
RTVIConfig,
RTVIMetricsProcessor,
RTVIProcessor,
RTVISpeakingProcessor,
RTVIUserTranscriptionProcessor,
)
from pipecat.services.elevenlabs import ElevenLabsTTSService
from pipecat.services.gemini_multimodal_live.gemini import GeminiMultimodalLiveLLMService from pipecat.services.gemini_multimodal_live.gemini import GeminiMultimodalLiveLLMService
from pipecat.services.openai import OpenAILLMService
from pipecat.transports.services.daily import DailyParams, DailyTransport from pipecat.transports.services.daily import DailyParams, DailyTransport
load_dotenv(override=True) load_dotenv(override=True)
@@ -168,20 +159,6 @@ async def main():
# #
# RTVI events for Pipecat client UI # RTVI events for Pipecat client UI
# #
# This will send `user-*-speaking` and `bot-*-speaking` messages.
rtvi_speaking = RTVISpeakingProcessor()
# This will emit UserTranscript events.
rtvi_user_transcription = RTVIUserTranscriptionProcessor()
# This will emit BotTranscript events.
rtvi_bot_transcription = RTVIBotTranscriptionProcessor()
# This will send `metrics` messages.
rtvi_metrics = RTVIMetricsProcessor()
# Handles RTVI messages from the client
rtvi = RTVIProcessor(config=RTVIConfig(config=[])) rtvi = RTVIProcessor(config=RTVIConfig(config=[]))
pipeline = Pipeline( pipeline = Pipeline(
@@ -190,11 +167,7 @@ async def main():
rtvi, rtvi,
context_aggregator.user(), context_aggregator.user(),
llm, llm,
rtvi_speaking,
rtvi_user_transcription,
rtvi_bot_transcription,
ta, ta,
rtvi_metrics,
transport.output(), transport.output(),
context_aggregator.assistant(), context_aggregator.assistant(),
] ]
@@ -206,6 +179,7 @@ async def main():
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,
observers=[rtvi.observer()],
), ),
) )
await task.queue_frame(quiet_frame) await task.queue_frame(quiet_frame)

View File

@@ -41,14 +41,7 @@ from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineParams, PipelineTask from pipecat.pipeline.task import PipelineParams, PipelineTask
from pipecat.processors.aggregators.openai_llm_context import OpenAILLMContext from pipecat.processors.aggregators.openai_llm_context import OpenAILLMContext
from pipecat.processors.frame_processor import FrameDirection, FrameProcessor from pipecat.processors.frame_processor import FrameDirection, FrameProcessor
from pipecat.processors.frameworks.rtvi import ( from pipecat.processors.frameworks.rtvi import RTVIConfig, RTVIProcessor
RTVIBotTranscriptionProcessor,
RTVIConfig,
RTVIMetricsProcessor,
RTVIProcessor,
RTVISpeakingProcessor,
RTVIUserTranscriptionProcessor,
)
from pipecat.services.elevenlabs import ElevenLabsTTSService from pipecat.services.elevenlabs import ElevenLabsTTSService
from pipecat.services.openai import OpenAILLMService from pipecat.services.openai import OpenAILLMService
from pipecat.transports.services.daily import DailyParams, DailyTransport from pipecat.transports.services.daily import DailyParams, DailyTransport
@@ -189,34 +182,16 @@ async def main():
# #
# RTVI events for Pipecat client UI # RTVI events for Pipecat client UI
# #
# This will send `user-*-speaking` and `bot-*-speaking` messages.
rtvi_speaking = RTVISpeakingProcessor()
# This will emit UserTranscript events.
rtvi_user_transcription = RTVIUserTranscriptionProcessor()
# This will emit BotTranscript events.
rtvi_bot_transcription = RTVIBotTranscriptionProcessor()
# This will send `metrics` messages.
rtvi_metrics = RTVIMetricsProcessor()
# Handles RTVI messages from the client
rtvi = RTVIProcessor(config=RTVIConfig(config=[])) rtvi = RTVIProcessor(config=RTVIConfig(config=[]))
pipeline = Pipeline( pipeline = Pipeline(
[ [
transport.input(), transport.input(),
rtvi, rtvi,
rtvi_speaking,
rtvi_user_transcription,
context_aggregator.user(), context_aggregator.user(),
llm, llm,
rtvi_bot_transcription,
tts, tts,
ta, ta,
rtvi_metrics,
transport.output(), transport.output(),
context_aggregator.assistant(), context_aggregator.assistant(),
] ]
@@ -228,6 +203,7 @@ async def main():
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,
observers=[rtvi.observer()],
), ),
) )
await task.queue_frame(quiet_frame) await task.queue_frame(quiet_frame)