Merge pull request #662 from pipecat-ai/aleix/daily-transport-async-functions

transports(daily): make functions async
This commit is contained in:
Aleix Conchillo Flaqué
2024-10-25 16:14:06 -07:00
committed by GitHub
53 changed files with 112 additions and 93 deletions

View File

@@ -18,6 +18,11 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Changed ### Changed
- The following `DailyTransport` functions are now `async` which means they need
to be awaited: `start_dialout`, `stop_dialout`, `start_recording`,
`stop_recording`, `capture_participant_transcription` and
`capture_participant_video`.
- Changed default output sample rate to 24000. This changes all TTS service to - Changed default output sample rate to 24000. This changes all TTS service to
output to 24000 and also the default output transport sample rate. This output to 24000 and also the default output transport sample rate. This
improves audio quality at the cost of some extra bandwidth. improves audio quality at the cost of some extra bandwidth.

View File

@@ -124,7 +124,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])
@transport.event_handler("on_participant_left") @transport.event_handler("on_participant_left")

View File

@@ -123,7 +123,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])
@transport.event_handler("on_participant_left") @transport.event_handler("on_participant_left")

View File

@@ -75,7 +75,7 @@ async def main(room_url: str, token: str):
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])
@transport.event_handler("on_participant_left") @transport.event_handler("on_participant_left")

View File

@@ -81,7 +81,7 @@ async def main(room_url: str, token: str, callId: str, callDomain: str):
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])
@transport.event_handler("on_participant_left") @transport.event_handler("on_participant_left")

View File

@@ -84,7 +84,7 @@ async def main(room_url: str, token: str, callId: str, sipUri: str):
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])
@transport.event_handler("on_participant_left") @transport.event_handler("on_participant_left")

View File

@@ -110,7 +110,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -127,7 +127,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
participant_name = participant.get("info", {}).get("userName", "") participant_name = participant.get("info", {}).get("userName", "")
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
await task.queue_frames([TextFrame(f"Hi there {participant_name}!")]) await task.queue_frames([TextFrame(f"Hi there {participant_name}!")])
runner = PipelineRunner() runner = PipelineRunner()

View File

@@ -89,7 +89,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -87,7 +87,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -82,7 +82,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -109,7 +109,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
lc.set_participant_id(participant["id"]) lc.set_participant_id(participant["id"])
# Kick off the conversation. # Kick off the conversation.
# the `LLMMessagesFrame` will be picked up by the LangchainProcessor using # the `LLMMessagesFrame` will be picked up by the LangchainProcessor using

View File

@@ -85,7 +85,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -88,7 +88,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -89,7 +89,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -74,7 +74,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -86,7 +86,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -81,7 +81,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -83,7 +83,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -77,7 +77,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -96,7 +96,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -84,7 +84,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -82,7 +82,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -83,7 +83,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -74,7 +74,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_video(participant["id"]) await transport.capture_participant_video(participant["id"])
pipeline = Pipeline([transport.input(), MirrorProcessor(), transport.output()]) pipeline = Pipeline([transport.input(), MirrorProcessor(), transport.output()])

View File

@@ -81,7 +81,7 @@ async def main():
@daily_transport.event_handler("on_first_participant_joined") @daily_transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_video(participant["id"]) await transport.capture_participant_video(participant["id"])
pipeline = Pipeline([daily_transport.input(), MirrorProcessor(), tk_transport.output()]) pipeline = Pipeline([daily_transport.input(), MirrorProcessor(), tk_transport.output()])

View File

@@ -82,7 +82,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
await tts.say("Hi! If you want to talk to me, just say 'Hey Robot'.") await tts.say("Hi! If you want to talk to me, just say 'Hey Robot'.")
runner = PipelineRunner() runner = PipelineRunner()

View File

@@ -134,7 +134,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
await tts.say("Hi, I'm listening!") await tts.say("Hi, I'm listening!")
await transport.send_audio(sounds["ding1.wav"]) await transport.send_audio(sounds["ding1.wav"])

View File

@@ -84,8 +84,8 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
await tts.say("Hi there! Feel free to ask me what I see.") await tts.say("Hi there! Feel free to ask me what I see.")
transport.capture_participant_video(participant["id"], framerate=0) await transport.capture_participant_video(participant["id"], framerate=0)
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
image_requester.set_participant_id(participant["id"]) image_requester.set_participant_id(participant["id"])
pipeline = Pipeline( pipeline = Pipeline(

View File

@@ -86,8 +86,8 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
await tts.say("Hi there! Feel free to ask me what I see.") await tts.say("Hi there! Feel free to ask me what I see.")
transport.capture_participant_video(participant["id"], framerate=0) await transport.capture_participant_video(participant["id"], framerate=0)
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
image_requester.set_participant_id(participant["id"]) image_requester.set_participant_id(participant["id"])
pipeline = Pipeline( pipeline = Pipeline(

View File

@@ -83,8 +83,8 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
await tts.say("Hi there! Feel free to ask me what I see.") await tts.say("Hi there! Feel free to ask me what I see.")
transport.capture_participant_video(participant["id"], framerate=0) await transport.capture_participant_video(participant["id"], framerate=0)
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
image_requester.set_participant_id(participant["id"]) image_requester.set_participant_id(participant["id"])
pipeline = Pipeline( pipeline = Pipeline(

View File

@@ -83,8 +83,8 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
await tts.say("Hi there! Feel free to ask me what I see.") await tts.say("Hi there! Feel free to ask me what I see.")
transport.capture_participant_video(participant["id"], framerate=0) await transport.capture_participant_video(participant["id"], framerate=0)
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
image_requester.set_participant_id(participant["id"]) image_requester.set_participant_id(participant["id"])
pipeline = Pipeline( pipeline = Pipeline(

View File

@@ -127,7 +127,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([context_aggregator.user().get_context_frame()]) await task.queue_frames([context_aggregator.user().get_context_frame()])

View File

@@ -105,7 +105,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([context_aggregator.user().get_context_frame()]) await task.queue_frames([context_aggregator.user().get_context_frame()])

View File

@@ -160,8 +160,8 @@ If you need to use a tool, simply use the tool. Do not tell the user the tool yo
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
global video_participant_id global video_participant_id
video_participant_id = participant["id"] video_participant_id = participant["id"]
transport.capture_participant_transcription(video_participant_id) await transport.capture_participant_transcription(video_participant_id)
transport.capture_participant_video(video_participant_id, framerate=0) await transport.capture_participant_video(video_participant_id, framerate=0)
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([context_aggregator.user().get_context_frame()]) await task.queue_frames([context_aggregator.user().get_context_frame()])

View File

@@ -123,7 +123,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
# await tts.say("Hi! Ask me about the weather in San Francisco.") # await tts.say("Hi! Ask me about the weather in San Francisco.")

View File

@@ -153,8 +153,8 @@ indicate you should use the get_image tool are:
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
global video_participant_id global video_participant_id
video_participant_id = participant["id"] video_participant_id = participant["id"]
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
transport.capture_participant_video(video_participant_id, framerate=0) await transport.capture_participant_video(video_participant_id, framerate=0)
# Kick off the conversation. # Kick off the conversation.
await tts.say("Hi! Ask me about the weather in San Francisco.") await tts.say("Hi! Ask me about the weather in San Francisco.")

View File

@@ -159,8 +159,8 @@ indicate you should use the get_image tool are:
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
global video_participant_id global video_participant_id
video_participant_id = participant["id"] video_participant_id = participant["id"]
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
transport.capture_participant_video(video_participant_id, framerate=0) await transport.capture_participant_video(video_participant_id, framerate=0)
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([context_aggregator.user().get_context_frame()]) await task.queue_frames([context_aggregator.user().get_context_frame()])

View File

@@ -141,7 +141,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append( messages.append(
{ {

View File

@@ -128,7 +128,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append( messages.append(
{ {

View File

@@ -92,7 +92,7 @@ async def main():
# bot can "hear" and respond to them. # bot can "hear" and respond to them.
@transport.event_handler("on_participant_joined") @transport.event_handler("on_participant_joined")
async def on_participant_joined(transport, participant): async def on_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# When the first participant joins, the bot should introduce itself. # When the first participant joins, the bot should introduce itself.
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")

View File

@@ -99,7 +99,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
messages.append({"role": "system", "content": "Please introduce yourself to the user."}) messages.append({"role": "system", "content": "Please introduce yourself to the user."})
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -166,7 +166,7 @@ Remember, your responses should be short. Just one or two sentences, usually."""
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([context_aggregator.user().get_context_frame()]) await task.queue_frames([context_aggregator.user().get_context_frame()])

View File

@@ -223,7 +223,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([context_aggregator.user().get_context_frame()]) await task.queue_frames([context_aggregator.user().get_context_frame()])

View File

@@ -249,7 +249,7 @@ Remember, your responses should be short. Just one or two sentences, usually."""
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([context_aggregator.user().get_context_frame()]) await task.queue_frames([context_aggregator.user().get_context_frame()])

View File

@@ -219,7 +219,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([context_aggregator.user().get_context_frame()]) await task.queue_frames([context_aggregator.user().get_context_frame()])

View File

@@ -276,8 +276,8 @@ async def main():
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
global video_participant_id global video_participant_id
video_participant_id = participant["id"] video_participant_id = participant["id"]
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
transport.capture_participant_video(video_participant_id, framerate=0) await transport.capture_participant_video(video_participant_id, framerate=0)
# Kick off the conversation. # Kick off the conversation.
await task.queue_frames([context_aggregator.user().get_context_frame()]) await task.queue_frames([context_aggregator.user().get_context_frame()])

View File

@@ -203,8 +203,8 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
transport.capture_participant_video(participant["id"], framerate=0) await transport.capture_participant_video(participant["id"], framerate=0)
ir.set_participant_id(participant["id"]) ir.set_participant_id(participant["id"])
await task.queue_frames([LLMMessagesFrame(messages)]) await task.queue_frames([LLMMessagesFrame(messages)])

View File

@@ -352,7 +352,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
print(f"Context is: {context}") print(f"Context is: {context}")
await task.queue_frames([OpenAILLMContextFrame(context)]) await task.queue_frames([OpenAILLMContextFrame(context)])

View File

@@ -102,7 +102,7 @@ async def main(room_url, token=None):
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
logger.debug("Participant joined, storytime commence!") logger.debug("Participant joined, storytime commence!")
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
await intro_task.queue_frames( await intro_task.queue_frames(
[ [
images["book1"], images["book1"],

View File

@@ -165,7 +165,7 @@ Your task is to help the user understand and learn from this article in 2 senten
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
messages.append( messages.append(
{ {
"role": "system", "role": "system",

View File

@@ -121,7 +121,7 @@ async def main():
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")
async def on_first_participant_joined(transport, participant): async def on_first_participant_joined(transport, participant):
transport.capture_participant_transcription(participant["id"]) await transport.capture_participant_transcription(participant["id"])
runner = PipelineRunner() runner = PipelineRunner()

View File

@@ -454,27 +454,37 @@ class DailyTransportClient(EventHandler):
def participant_counts(self): def participant_counts(self):
return self._client.participant_counts() return self._client.participant_counts()
def start_dialout(self, settings): async def start_dialout(self, settings):
self._client.start_dialout(settings) future = self._loop.create_future()
self._client.start_dialout(settings, completion=completion_callback(future))
await future
def stop_dialout(self, participant_id): async def stop_dialout(self, participant_id):
self._client.stop_dialout(participant_id) future = self._loop.create_future()
self._client.stop_dialout(participant_id, completion=completion_callback(future))
await future
def start_recording(self, streaming_settings, stream_id, force_new): async def start_recording(self, streaming_settings, stream_id, force_new):
self._client.start_recording(streaming_settings, stream_id, force_new) future = self._loop.create_future()
self._client.start_recording(
streaming_settings, stream_id, force_new, completion=completion_callback(future)
)
await future
def stop_recording(self, stream_id): async def stop_recording(self, stream_id):
self._client.stop_recording(stream_id) future = self._loop.create_future()
self._client.stop_recording(stream_id, completion=completion_callback(future))
await future
def capture_participant_transcription(self, participant_id: str): async def capture_participant_transcription(self, participant_id: str):
if not self._params.transcription_enabled: if not self._params.transcription_enabled:
return return
self._transcription_ids.append(participant_id) self._transcription_ids.append(participant_id)
if self._joined and self._transcription_status: if self._joined and self._transcription_status:
self.update_transcription(self._transcription_ids) await self.update_transcription(self._transcription_ids)
def capture_participant_video( async def capture_participant_video(
self, self,
participant_id: str, participant_id: str,
callback: Callable, callback: Callable,
@@ -483,7 +493,7 @@ class DailyTransportClient(EventHandler):
color_format: str = "RGB", color_format: str = "RGB",
): ):
# Only enable camera subscription on this participant # Only enable camera subscription on this participant
self._client.update_subscriptions( await self.update_subscriptions(
participant_settings={participant_id: {"media": "subscribed"}} participant_settings={participant_id: {"media": "subscribed"}}
) )
@@ -496,8 +506,12 @@ class DailyTransportClient(EventHandler):
color_format=color_format, color_format=color_format,
) )
def update_transcription(self, participants=None, instance_id=None): async def update_transcription(self, participants=None, instance_id=None):
self._client.update_transcription(participants, instance_id) future = self._loop.create_future()
self._client.update_transcription(
participants, instance_id, completion=completion_callback(future)
)
await future
async def update_subscriptions(self, participant_settings=None, profile_settings=None): async def update_subscriptions(self, participant_settings=None, profile_settings=None):
future = self._loop.create_future() future = self._loop.create_future()
@@ -559,7 +573,7 @@ class DailyTransportClient(EventHandler):
def on_transcription_started(self, status): def on_transcription_started(self, status):
logger.debug(f"Transcription started: {status}") logger.debug(f"Transcription started: {status}")
self._transcription_status = status self._transcription_status = status
self.update_transcription(self._transcription_ids) self._call_async_callback(self.update_transcription, self._transcription_ids)
def on_transcription_stopped(self, stopped_by, stopped_by_error): def on_transcription_stopped(self, stopped_by, stopped_by_error):
logger.debug("Transcription stopped") logger.debug("Transcription stopped")
@@ -662,7 +676,7 @@ class DailyInputTransport(BaseInputTransport):
await super().process_frame(frame, direction) await super().process_frame(frame, direction)
if isinstance(frame, UserImageRequestFrame): if isinstance(frame, UserImageRequestFrame):
self.request_participant_image(frame.user_id) await self.request_participant_image(frame.user_id)
# #
# Frames # Frames
@@ -692,7 +706,7 @@ class DailyInputTransport(BaseInputTransport):
# Camera in # Camera in
# #
def capture_participant_video( async def capture_participant_video(
self, self,
participant_id: str, participant_id: str,
framerate: int = 30, framerate: int = 30,
@@ -705,11 +719,11 @@ class DailyInputTransport(BaseInputTransport):
"render_next_frame": False, "render_next_frame": False,
} }
self._client.capture_participant_video( await self._client.capture_participant_video(
participant_id, self._on_participant_video_frame, framerate, video_source, color_format participant_id, self._on_participant_video_frame, framerate, video_source, color_format
) )
def request_participant_image(self, participant_id: str): async def request_participant_image(self, participant_id: str):
if participant_id in self._video_renderers: if participant_id in self._video_renderers:
self._video_renderers[participant_id]["render_next_frame"] = True self._video_renderers[participant_id]["render_next_frame"] = True
@@ -866,22 +880,22 @@ class DailyTransport(BaseTransport):
def participant_counts(self): def participant_counts(self):
return self._client.participant_counts() return self._client.participant_counts()
def start_dialout(self, settings=None): async def start_dialout(self, settings=None):
self._client.start_dialout(settings) await self._client.start_dialout(settings)
def stop_dialout(self, participant_id): async def stop_dialout(self, participant_id):
self._client.stop_dialout(participant_id) await self._client.stop_dialout(participant_id)
def start_recording(self, streaming_settings=None, stream_id=None, force_new=None): async def start_recording(self, streaming_settings=None, stream_id=None, force_new=None):
self._client.start_recording(streaming_settings, stream_id, force_new) await self._client.start_recording(streaming_settings, stream_id, force_new)
def stop_recording(self, stream_id=None): async def stop_recording(self, stream_id=None):
self._client.stop_recording(stream_id) await self._client.stop_recording(stream_id)
def capture_participant_transcription(self, participant_id: str): async def capture_participant_transcription(self, participant_id: str):
self._client.capture_participant_transcription(participant_id) await self._client.capture_participant_transcription(participant_id)
def capture_participant_video( async def capture_participant_video(
self, self,
participant_id: str, participant_id: str,
framerate: int = 30, framerate: int = 30,
@@ -889,7 +903,7 @@ class DailyTransport(BaseTransport):
color_format: str = "RGB", color_format: str = "RGB",
): ):
if self._input: if self._input:
self._input.capture_participant_video( await self._input.capture_participant_video(
participant_id, framerate, video_source, color_format participant_id, framerate, video_source, color_format
) )