diff --git a/examples/foundational/49-thinking.py b/examples/foundational/49a-thinking-anthropic.py similarity index 64% rename from examples/foundational/49-thinking.py rename to examples/foundational/49a-thinking-anthropic.py index 74da17c66..6017f335e 100644 --- a/examples/foundational/49-thinking.py +++ b/examples/foundational/49a-thinking-anthropic.py @@ -4,10 +4,7 @@ # SPDX-License-Identifier: BSD 2-Clause License # -import argparse import os -import random -import sys from dotenv import load_dotenv from loguru import logger @@ -28,18 +25,12 @@ from pipecat.runner.utils import create_transport from pipecat.services.anthropic.llm import AnthropicLLMService from pipecat.services.cartesia.tts import CartesiaTTSService from pipecat.services.deepgram.stt import DeepgramSTTService -from pipecat.services.google.llm import GoogleLLMService from pipecat.transports.base_transport import BaseTransport, TransportParams from pipecat.transports.daily.transport import DailyParams from pipecat.transports.websocket.fastapi import FastAPIWebsocketParams load_dotenv(override=True) -# LLM provider constants -LLM_ANTHROPIC = "anthropic" -LLM_GOOGLE = "google" -LLM_DEFAULT = LLM_GOOGLE - # We store functions so objects (e.g. SileroVADAnalyzer) don't get # instantiated. The function will be called when the desired transport gets # selected. @@ -65,10 +56,8 @@ transport_params = { } -async def run_bot( - transport: BaseTransport, runner_args: RunnerArguments, llm_provider: str = LLM_DEFAULT -): - logger.info(f"Starting bot with {llm_provider.capitalize()} LLM") +async def run_bot(transport: BaseTransport, runner_args: RunnerArguments): + logger.info(f"Starting bot") stt = DeepgramSTTService(api_key=os.getenv("DEEPGRAM_API_KEY")) @@ -77,27 +66,12 @@ async def run_bot( voice_id="71a7ad14-091c-4e8e-a314-022ece01c121", # British Reading Lady ) - if llm_provider == LLM_ANTHROPIC: - llm = AnthropicLLMService( - api_key=os.getenv("ANTHROPIC_API_KEY"), - params=AnthropicLLMService.InputParams( - thinking=AnthropicLLMService.ThinkingConfig(type="enabled", budget_tokens=2048) - ), - ) - elif llm_provider == LLM_GOOGLE: - llm = GoogleLLMService( - api_key=os.getenv("GOOGLE_API_KEY"), - # model="gemini-3-pro-preview", # A more powerful reasoning model, but slower - params=GoogleLLMService.InputParams( - thinking=GoogleLLMService.ThinkingConfig( - # thinking_level="low", # Use this field instead of thinking_budget for Gemini 3 Pro. Defaults to "high". - thinking_budget=-1, # Dynamic thinking - include_thoughts=True, - ) - ), - ) - else: - raise ValueError(f"Unsupported LLM provider: {llm_provider}") + llm = AnthropicLLMService( + api_key=os.getenv("ANTHROPIC_API_KEY"), + params=AnthropicLLMService.InputParams( + thinking=AnthropicLLMService.ThinkingConfig(type="enabled", budget_tokens=2048) + ), + ) transcript = TranscriptProcessor(process_thoughts=True) @@ -137,15 +111,16 @@ async def run_bot( @transport.event_handler("on_client_connected") async def on_client_connected(transport, client): logger.info(f"Client connected") - # Choose a random prompt to demonstrate thinking capabilities. - # These prompts were chosen from Google and Anthropic docs. - thinking_prompt_1 = "Analogize photosynthesis and growing up." - thinking_prompt_2 = "Compare and contrast electric cars and hybrid cars." - thinking_prompt_3 = "Are there an infinite number of prime numbers such that n mod 4 == 3?" - selected_prompt = random.choice([thinking_prompt_1, thinking_prompt_2, thinking_prompt_3]) - - # Kick off the conversation. - messages.append({"role": "user", "content": selected_prompt}) + # Kick off the conversation, using a prompt conducive to demonstrating + # thinking (chosen from Google and Anthropic docs). + messages.append( + { + "role": "user", + "content": "Analogize photosynthesis and growing up.", + # "content": "Compare and contrast electric cars and hybrid cars." + # "content": "Are there an infinite number of prime numbers such that n mod 4 == 3?" + } + ) await task.queue_frames([LLMRunFrame()]) @transport.event_handler("on_client_disconnected") @@ -169,31 +144,11 @@ async def run_bot( async def bot(runner_args: RunnerArguments): """Main bot entry point compatible with Pipecat Cloud.""" - # Get llm_provider from module attribute set in __main__ - llm_provider = getattr(sys.modules[__name__], "llm_provider", LLM_DEFAULT) transport = await create_transport(runner_args, transport_params) - await run_bot(transport, runner_args, llm_provider) + await run_bot(transport, runner_args) if __name__ == "__main__": - # Parse custom arguments before calling runner main() - parser = argparse.ArgumentParser(description="Thinking LLM Bot") - parser.add_argument( - "--llm", - type=str, - choices=[LLM_ANTHROPIC, LLM_GOOGLE], - default=LLM_DEFAULT, - help=f"LLM provider to use (default: {LLM_DEFAULT})", - ) - # Parse only known args to allow runner's main() to handle its own args - args, remaining = parser.parse_known_args() - - # Store the llm_provider in sys.modules for bot() function to access - sys.modules[__name__].llm_provider = args.llm - - # Restore sys.argv with remaining args for runner's main() - sys.argv[1:] = remaining - from pipecat.runner.run import main main() diff --git a/examples/foundational/49b-thinking-google.py b/examples/foundational/49b-thinking-google.py new file mode 100644 index 000000000..85df6da39 --- /dev/null +++ b/examples/foundational/49b-thinking-google.py @@ -0,0 +1,159 @@ +# +# Copyright (c) 2024–2025, Daily +# +# SPDX-License-Identifier: BSD 2-Clause License +# + +import os + +from dotenv import load_dotenv +from loguru import logger + +from pipecat.audio.turn.smart_turn.base_smart_turn import SmartTurnParams +from pipecat.audio.turn.smart_turn.local_smart_turn_v3 import LocalSmartTurnAnalyzerV3 +from pipecat.audio.vad.silero import SileroVADAnalyzer +from pipecat.audio.vad.vad_analyzer import VADParams +from pipecat.frames.frames import LLMRunFrame, ThoughtTranscriptionMessage, TranscriptionMessage +from pipecat.pipeline.pipeline import Pipeline +from pipecat.pipeline.runner import PipelineRunner +from pipecat.pipeline.task import PipelineParams, PipelineTask +from pipecat.processors.aggregators.llm_context import LLMContext +from pipecat.processors.aggregators.llm_response_universal import LLMContextAggregatorPair +from pipecat.processors.transcript_processor import TranscriptProcessor +from pipecat.runner.types import RunnerArguments +from pipecat.runner.utils import create_transport +from pipecat.services.cartesia.tts import CartesiaTTSService +from pipecat.services.deepgram.stt import DeepgramSTTService +from pipecat.services.google.llm import GoogleLLMService +from pipecat.transports.base_transport import BaseTransport, TransportParams +from pipecat.transports.daily.transport import DailyParams +from pipecat.transports.websocket.fastapi import FastAPIWebsocketParams + +load_dotenv(override=True) + +# We store functions so objects (e.g. SileroVADAnalyzer) don't get +# instantiated. The function will be called when the desired transport gets +# selected. +transport_params = { + "daily": lambda: DailyParams( + audio_in_enabled=True, + audio_out_enabled=True, + vad_analyzer=SileroVADAnalyzer(params=VADParams(stop_secs=0.2)), + turn_analyzer=LocalSmartTurnAnalyzerV3(params=SmartTurnParams()), + ), + "twilio": lambda: FastAPIWebsocketParams( + audio_in_enabled=True, + audio_out_enabled=True, + vad_analyzer=SileroVADAnalyzer(params=VADParams(stop_secs=0.2)), + turn_analyzer=LocalSmartTurnAnalyzerV3(params=SmartTurnParams()), + ), + "webrtc": lambda: TransportParams( + audio_in_enabled=True, + audio_out_enabled=True, + vad_analyzer=SileroVADAnalyzer(params=VADParams(stop_secs=0.2)), + turn_analyzer=LocalSmartTurnAnalyzerV3(params=SmartTurnParams()), + ), +} + + +async def run_bot(transport: BaseTransport, runner_args: RunnerArguments): + logger.info(f"Starting bot") + + stt = DeepgramSTTService(api_key=os.getenv("DEEPGRAM_API_KEY")) + + tts = CartesiaTTSService( + api_key=os.getenv("CARTESIA_API_KEY"), + voice_id="71a7ad14-091c-4e8e-a314-022ece01c121", # British Reading Lady + ) + + llm = GoogleLLMService( + api_key=os.getenv("GOOGLE_API_KEY"), + # model="gemini-3-pro-preview", # A more powerful reasoning model, but slower + params=GoogleLLMService.InputParams( + thinking=GoogleLLMService.ThinkingConfig( + # thinking_level="low", # Use this field instead of thinking_budget for Gemini 3 Pro. Defaults to "high". + thinking_budget=-1, # Dynamic thinking + include_thoughts=True, + ) + ), + ) + + transcript = TranscriptProcessor(process_thoughts=True) + + messages = [ + { + "role": "system", + "content": "You are a helpful LLM in a WebRTC call. Your goal is to demonstrate your capabilities in a succinct way. Your output will be spoken aloud, so avoid special characters that can't easily be spoken, such as emojis or bullet points. Respond to what the user said in a creative and helpful way.", + }, + ] + + context = LLMContext(messages) + context_aggregator = LLMContextAggregatorPair(context) + + pipeline = Pipeline( + [ + transport.input(), # Transport user input + stt, + transcript.user(), # User transcripts + context_aggregator.user(), # User responses + llm, # LLM + tts, # TTS + transport.output(), # Transport bot output + transcript.assistant(), # Assistant transcripts (including thoughts) + context_aggregator.assistant(), # Assistant spoken responses + ] + ) + + task = PipelineTask( + pipeline, + params=PipelineParams( + enable_metrics=True, + enable_usage_metrics=True, + ), + idle_timeout_secs=runner_args.pipeline_idle_timeout_secs, + ) + + @transport.event_handler("on_client_connected") + async def on_client_connected(transport, client): + logger.info(f"Client connected") + # Kick off the conversation, using a prompt conducive to demonstrating + # thinking (chosen from Google and Anthropic docs). + messages.append( + { + "role": "user", + "content": "Analogize photosynthesis and growing up.", + # "content": "Compare and contrast electric cars and hybrid cars." + # "content": "Are there an infinite number of prime numbers such that n mod 4 == 3?" + } + ) + await task.queue_frames([LLMRunFrame()]) + + @transport.event_handler("on_client_disconnected") + async def on_client_disconnected(transport, client): + logger.info(f"Client disconnected") + await task.cancel() + + # Register event handler for transcript updates + @transcript.event_handler("on_transcript_update") + async def on_transcript_update(processor, frame): + for msg in frame.messages: + if isinstance(msg, (ThoughtTranscriptionMessage, TranscriptionMessage)): + timestamp = f"[{msg.timestamp}] " if msg.timestamp else "" + role = "THOUGHT" if isinstance(msg, ThoughtTranscriptionMessage) else msg.role + logger.info(f"Transcript: {timestamp}{role}: {msg.content}") + + runner = PipelineRunner(handle_sigint=runner_args.handle_sigint) + + await runner.run(task) + + +async def bot(runner_args: RunnerArguments): + """Main bot entry point compatible with Pipecat Cloud.""" + transport = await create_transport(runner_args, transport_params) + await run_bot(transport, runner_args) + + +if __name__ == "__main__": + from pipecat.runner.run import main + + main() diff --git a/examples/foundational/49-thinking-functions.py b/examples/foundational/49c-thinking-functions-anthropic.py similarity index 74% rename from examples/foundational/49-thinking-functions.py rename to examples/foundational/49c-thinking-functions-anthropic.py index 411316efe..3d71f2c47 100644 --- a/examples/foundational/49-thinking-functions.py +++ b/examples/foundational/49c-thinking-functions-anthropic.py @@ -4,10 +4,7 @@ # SPDX-License-Identifier: BSD 2-Clause License # -import argparse import os -import random -import sys from dotenv import load_dotenv from loguru import logger @@ -29,7 +26,6 @@ from pipecat.runner.utils import create_transport from pipecat.services.anthropic.llm import AnthropicLLMService from pipecat.services.cartesia.tts import CartesiaTTSService from pipecat.services.deepgram.stt import DeepgramSTTService -from pipecat.services.google.llm import GoogleLLMService from pipecat.services.llm_service import FunctionCallParams from pipecat.transports.base_transport import BaseTransport, TransportParams from pipecat.transports.daily.transport import DailyParams @@ -56,11 +52,6 @@ async def book_taxi(params: FunctionCallParams, time: str): await params.result_callback({"status": "done"}) -# LLM provider constants -LLM_ANTHROPIC = "anthropic" -LLM_GOOGLE = "google" -LLM_DEFAULT = LLM_GOOGLE - # We store functions so objects (e.g. SileroVADAnalyzer) don't get # instantiated. The function will be called when the desired transport gets # selected. @@ -86,10 +77,8 @@ transport_params = { } -async def run_bot( - transport: BaseTransport, runner_args: RunnerArguments, llm_provider: str = LLM_DEFAULT -): - logger.info(f"Starting bot with {llm_provider.capitalize()} LLM") +async def run_bot(transport: BaseTransport, runner_args: RunnerArguments): + logger.info(f"Starting bot") stt = DeepgramSTTService(api_key=os.getenv("DEEPGRAM_API_KEY")) @@ -98,27 +87,12 @@ async def run_bot( voice_id="71a7ad14-091c-4e8e-a314-022ece01c121", # British Reading Lady ) - if llm_provider == LLM_ANTHROPIC: - llm = AnthropicLLMService( - api_key=os.getenv("ANTHROPIC_API_KEY"), - params=AnthropicLLMService.InputParams( - thinking=AnthropicLLMService.ThinkingConfig(type="enabled", budget_tokens=2048) - ), - ) - elif llm_provider == LLM_GOOGLE: - llm = GoogleLLMService( - api_key=os.getenv("GOOGLE_API_KEY"), - # model="gemini-3-pro-preview", # A more powerful reasoning model, but slower - params=GoogleLLMService.InputParams( - thinking=GoogleLLMService.ThinkingConfig( - # thinking_level="low", # Use this field instead of thinking_budget for Gemini 3 Pro. Defaults to "high". - thinking_budget=-1, # Dynamic thinking - include_thoughts=True, - ) - ), - ) - else: - raise ValueError(f"Unsupported LLM provider: {llm_provider}") + llm = AnthropicLLMService( + api_key=os.getenv("ANTHROPIC_API_KEY"), + params=AnthropicLLMService.InputParams( + thinking=AnthropicLLMService.ThinkingConfig(type="enabled", budget_tokens=2048) + ), + ) llm.register_direct_function(check_flight_status) llm.register_direct_function(book_taxi) @@ -193,31 +167,11 @@ async def run_bot( async def bot(runner_args: RunnerArguments): """Main bot entry point compatible with Pipecat Cloud.""" - # Get llm_provider from module attribute set in __main__ - llm_provider = getattr(sys.modules[__name__], "llm_provider", LLM_DEFAULT) transport = await create_transport(runner_args, transport_params) - await run_bot(transport, runner_args, llm_provider) + await run_bot(transport, runner_args) if __name__ == "__main__": - # Parse custom arguments before calling runner main() - parser = argparse.ArgumentParser(description="Thinking LLM Bot") - parser.add_argument( - "--llm", - type=str, - choices=[LLM_ANTHROPIC, LLM_GOOGLE], - default=LLM_DEFAULT, - help=f"LLM provider to use (default: {LLM_DEFAULT})", - ) - # Parse only known args to allow runner's main() to handle its own args - args, remaining = parser.parse_known_args() - - # Store the llm_provider in sys.modules for bot() function to access - sys.modules[__name__].llm_provider = args.llm - - # Restore sys.argv with remaining args for runner's main() - sys.argv[1:] = remaining - from pipecat.runner.run import main main() diff --git a/examples/foundational/49d-thinking-functions-google.py b/examples/foundational/49d-thinking-functions-google.py new file mode 100644 index 000000000..3ec2b62d8 --- /dev/null +++ b/examples/foundational/49d-thinking-functions-google.py @@ -0,0 +1,182 @@ +# +# Copyright (c) 2024–2025, Daily +# +# SPDX-License-Identifier: BSD 2-Clause License +# + +import os + +from dotenv import load_dotenv +from loguru import logger + +from pipecat.adapters.schemas.tools_schema import ToolsSchema +from pipecat.audio.turn.smart_turn.base_smart_turn import SmartTurnParams +from pipecat.audio.turn.smart_turn.local_smart_turn_v3 import LocalSmartTurnAnalyzerV3 +from pipecat.audio.vad.silero import SileroVADAnalyzer +from pipecat.audio.vad.vad_analyzer import VADParams +from pipecat.frames.frames import LLMRunFrame, ThoughtTranscriptionMessage, TranscriptionMessage +from pipecat.pipeline.pipeline import Pipeline +from pipecat.pipeline.runner import PipelineRunner +from pipecat.pipeline.task import PipelineParams, PipelineTask +from pipecat.processors.aggregators.llm_context import LLMContext +from pipecat.processors.aggregators.llm_response_universal import LLMContextAggregatorPair +from pipecat.processors.transcript_processor import TranscriptProcessor +from pipecat.runner.types import RunnerArguments +from pipecat.runner.utils import create_transport +from pipecat.services.cartesia.tts import CartesiaTTSService +from pipecat.services.deepgram.stt import DeepgramSTTService +from pipecat.services.google.llm import GoogleLLMService +from pipecat.services.llm_service import FunctionCallParams +from pipecat.transports.base_transport import BaseTransport, TransportParams +from pipecat.transports.daily.transport import DailyParams +from pipecat.transports.websocket.fastapi import FastAPIWebsocketParams + +load_dotenv(override=True) + + +async def check_flight_status(params: FunctionCallParams, flight_number: str): + """Check the status of a flight. Returns status (e.g., "on time", "delayed") and departure time. + + Args: + flight_number (str): The flight number, e.g. "AA100". + """ + await params.result_callback({"status": "delayed", "departure_time": "14:30"}) + + +async def book_taxi(params: FunctionCallParams, time: str): + """Book a taxi for a given time. Returns status (e.g., "done"). + + Args: + time (str): The time to book the taxi for, e.g. "15:00". + """ + await params.result_callback({"status": "done"}) + + +# We store functions so objects (e.g. SileroVADAnalyzer) don't get +# instantiated. The function will be called when the desired transport gets +# selected. +transport_params = { + "daily": lambda: DailyParams( + audio_in_enabled=True, + audio_out_enabled=True, + vad_analyzer=SileroVADAnalyzer(params=VADParams(stop_secs=0.2)), + turn_analyzer=LocalSmartTurnAnalyzerV3(params=SmartTurnParams()), + ), + "twilio": lambda: FastAPIWebsocketParams( + audio_in_enabled=True, + audio_out_enabled=True, + vad_analyzer=SileroVADAnalyzer(params=VADParams(stop_secs=0.2)), + turn_analyzer=LocalSmartTurnAnalyzerV3(params=SmartTurnParams()), + ), + "webrtc": lambda: TransportParams( + audio_in_enabled=True, + audio_out_enabled=True, + vad_analyzer=SileroVADAnalyzer(params=VADParams(stop_secs=0.2)), + turn_analyzer=LocalSmartTurnAnalyzerV3(params=SmartTurnParams()), + ), +} + + +async def run_bot(transport: BaseTransport, runner_args: RunnerArguments): + logger.info(f"Starting bot") + + stt = DeepgramSTTService(api_key=os.getenv("DEEPGRAM_API_KEY")) + + tts = CartesiaTTSService( + api_key=os.getenv("CARTESIA_API_KEY"), + voice_id="71a7ad14-091c-4e8e-a314-022ece01c121", # British Reading Lady + ) + + llm = GoogleLLMService( + api_key=os.getenv("GOOGLE_API_KEY"), + # model="gemini-3-pro-preview", # A more powerful reasoning model, but slower + params=GoogleLLMService.InputParams( + thinking=GoogleLLMService.ThinkingConfig( + # thinking_level="low", # Use this field instead of thinking_budget for Gemini 3 Pro. Defaults to "high". + thinking_budget=-1, # Dynamic thinking + include_thoughts=True, + ) + ), + ) + + llm.register_direct_function(check_flight_status) + llm.register_direct_function(book_taxi) + + tools = ToolsSchema(standard_tools=[check_flight_status, book_taxi]) + + transcript = TranscriptProcessor(process_thoughts=True) + + messages = [ + { + "role": "system", + "content": "You are a helpful LLM in a WebRTC call. Your goal is to demonstrate your capabilities in a succinct way. Your output will be spoken aloud, so avoid special characters that can't easily be spoken, such as emojis or bullet points. Respond to what the user said in a creative and helpful way.", + }, + ] + + context = LLMContext(messages, tools) + context_aggregator = LLMContextAggregatorPair(context) + + pipeline = Pipeline( + [ + transport.input(), # Transport user input + stt, + transcript.user(), # User transcripts + context_aggregator.user(), # User responses + llm, # LLM + tts, # TTS + transport.output(), # Transport bot output + transcript.assistant(), # Assistant transcripts (including thoughts) + context_aggregator.assistant(), # Assistant spoken responses + ] + ) + + task = PipelineTask( + pipeline, + params=PipelineParams( + enable_metrics=True, + enable_usage_metrics=True, + ), + idle_timeout_secs=runner_args.pipeline_idle_timeout_secs, + ) + + @transport.event_handler("on_client_connected") + async def on_client_connected(transport, client): + logger.info(f"Client connected") + # Kick off the conversation. + # This example comes from Gemini docs. + messages.append( + { + "role": "user", + "content": "Check the status of flight AA100 and, if it's delayed, book me a taxi 2 hours before its departure time.", + } + ) + await task.queue_frames([LLMRunFrame()]) + + @transport.event_handler("on_client_disconnected") + async def on_client_disconnected(transport, client): + logger.info(f"Client disconnected") + await task.cancel() + + @transcript.event_handler("on_transcript_update") + async def on_transcript_update(processor, frame): + for msg in frame.messages: + if isinstance(msg, (ThoughtTranscriptionMessage, TranscriptionMessage)): + timestamp = f"[{msg.timestamp}] " if msg.timestamp else "" + role = "THOUGHT" if isinstance(msg, ThoughtTranscriptionMessage) else msg.role + logger.info(f"Transcript: {timestamp}{role}: {msg.content}") + + runner = PipelineRunner(handle_sigint=runner_args.handle_sigint) + + await runner.run(task) + + +async def bot(runner_args: RunnerArguments): + """Main bot entry point compatible with Pipecat Cloud.""" + transport = await create_transport(runner_args, transport_params) + await run_bot(transport, runner_args) + + +if __name__ == "__main__": + from pipecat.runner.run import main + + main()