PipelineTask: force constructor keyword arguments

This commit is contained in:
Aleix Conchillo Flaqué
2025-02-25 17:56:01 -08:00
parent 66564392a6
commit 6722aae598
95 changed files with 140 additions and 98 deletions

View File

@@ -42,6 +42,9 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0
### Changed ### Changed
- ⚠️ `PipelineTask` now requires keyword arguments (except for the first one for
the pipeline).
- The base `TTSService` class now strips leading newlines before sending text - The base `TTSService` class now strips leading newlines before sending text
to the TTS provider. This change is to solve issues where some TTS providers, to the TTS provider. This change is to solve issues where some TTS providers,
like Azure, would not output text due to newlines. like Azure, would not output text due to newlines.

View File

@@ -17,7 +17,7 @@ from runner import configure
from pipecat.frames.frames import AudioRawFrame, EndFrame, OutputAudioRawFrame, TTSSpeakFrame from pipecat.frames.frames import AudioRawFrame, EndFrame, OutputAudioRawFrame, TTSSpeakFrame
from pipecat.pipeline.pipeline import Pipeline from pipecat.pipeline.pipeline import Pipeline
from pipecat.pipeline.runner import PipelineRunner from pipecat.pipeline.runner import PipelineRunner
from pipecat.pipeline.task import PipelineParams, PipelineTask from pipecat.pipeline.task import PipelineTask
from pipecat.services.cartesia import CartesiaTTSService from pipecat.services.cartesia import CartesiaTTSService
from pipecat.transports.services.daily import DailyParams, DailyTransport from pipecat.transports.services.daily import DailyParams, DailyTransport

View File

@@ -119,7 +119,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -124,7 +124,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@audiobuffer.event_handler("on_audio_data") @audiobuffer.event_handler("on_audio_data")
async def on_audio_data(buffer, audio, sample_rate, num_channels): async def on_audio_data(buffer, audio, sample_rate, num_channels):

View File

@@ -70,7 +70,7 @@ async def main(room_url: str, token: str):
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -62,7 +62,7 @@ async def main(room_url: str, token: str):
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -44,7 +44,8 @@ async def main():
runner = PipelineRunner() runner = PipelineRunner()
task = PipelineTask( task = PipelineTask(
Pipeline([imagegen, transport.output()]), PipelineParams(enable_metrics=True) Pipeline([imagegen, transport.output()]),
params=PipelineParams(enable_metrics=True),
) )
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")

View File

@@ -105,7 +105,10 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams(enable_metrics=True, enable_usage_metrics=True), params=PipelineParams(
enable_metrics=True,
enable_usage_metrics=True,
),
) )
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")

View File

@@ -127,7 +127,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -76,7 +76,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -74,7 +74,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -79,7 +79,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -103,7 +103,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -81,7 +81,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -74,7 +74,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -74,7 +74,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -75,7 +75,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -77,7 +77,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -83,7 +83,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -81,7 +81,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -81,7 +81,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -75,7 +75,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -80,7 +80,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -71,7 +71,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -88,7 +88,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -81,7 +81,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -79,7 +79,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -80,7 +80,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -76,7 +76,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -74,7 +74,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -74,7 +74,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -251,7 +251,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -74,7 +74,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -78,7 +78,11 @@ async def main():
runner = PipelineRunner() runner = PipelineRunner()
task = PipelineTask( task = PipelineTask(
pipeline, PipelineParams(audio_in_sample_rate=24000, audio_out_sample_rate=24000) pipeline,
params=PipelineParams(
audio_in_sample_rate=24000,
audio_out_sample_rate=24000,
),
) )
await runner.run(task) await runner.run(task)

View File

@@ -82,7 +82,11 @@ async def main():
pipeline = Pipeline([daily_transport.input(), MirrorProcessor(), tk_transport.output()]) pipeline = Pipeline([daily_transport.input(), MirrorProcessor(), tk_transport.output()])
task = PipelineTask( task = PipelineTask(
pipeline, PipelineParams(audio_in_sample_rate=24000, audio_out_sample_rate=24000) pipeline,
params=PipelineParams(
audio_in_sample_rate=24000,
audio_out_sample_rate=24000,
),
) )
async def run_tk(): async def run_tk():

View File

@@ -76,7 +76,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -112,7 +112,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -99,7 +99,13 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True, enable_metrics=True)) task = PipelineTask(
pipeline,
params=PipelineParams(
allow_interruptions=True,
enable_metrics=True,
),
)
@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):

View File

@@ -153,7 +153,13 @@ If you need to use a tool, simply use the tool. Do not tell the user the tool yo
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True, enable_metrics=True)) task = PipelineTask(
pipeline,
params=PipelineParams(
allow_interruptions=True,
enable_metrics=True,
),
)
@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):

View File

@@ -152,7 +152,7 @@ indicate you should use the get_image tool are:
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -116,7 +116,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -113,7 +113,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -117,7 +117,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -116,7 +116,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -116,7 +116,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -123,7 +123,7 @@ Start by asking me for my location. Then, use 'get_weather_current' to give me a
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -123,7 +123,7 @@ Start by asking me for my location. Then, use 'get_weather_current' to give me a
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -117,7 +117,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -83,7 +83,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -133,7 +133,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -126,7 +126,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -85,7 +85,13 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True, enable_metrics=True)) task = PipelineTask(
pipeline,
params=PipelineParams(
allow_interruptions=True,
enable_metrics=True,
),
)
# When a participant joins, start transcription for that participant so the # When a participant joins, start transcription for that participant so the
# bot can "hear" and respond to them. # bot can "hear" and respond to them.

View File

@@ -108,7 +108,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
report_only_initial_ttfb=True, report_only_initial_ttfb=True,

View File

@@ -154,7 +154,7 @@ Remember, your responses should be short. Just one or two sentences, usually."""
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -212,7 +212,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -237,7 +237,7 @@ Remember, your responses should be short. Just one or two sentences, usually."""
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -209,7 +209,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -263,7 +263,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -87,7 +87,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
# We just use 16000 because that's what Tavus is expecting and # We just use 16000 because that's what Tavus is expecting and
# we avoid resampling. # we avoid resampling.
audio_in_sample_rate=16000, audio_in_sample_rate=16000,

View File

@@ -145,7 +145,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -355,7 +355,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -564,7 +564,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -742,7 +742,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -87,7 +87,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -122,7 +122,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -354,7 +354,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -63,7 +63,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -89,7 +89,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -120,7 +120,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -79,7 +79,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -106,7 +106,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -1,5 +1,5 @@
# #
# Copyright (c) 2024, Daily # Copyright (c) 2024-2025, Daily
# #
# SPDX-License-Identifier: BSD 2-Clause License # SPDX-License-Identifier: BSD 2-Clause License
# #
@@ -93,7 +93,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -83,7 +83,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
), ),

View File

@@ -150,7 +150,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -150,7 +150,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -178,7 +178,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -117,7 +117,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -32,7 +32,7 @@ async def main():
pipeline = Pipeline([NullProcessor()]) pipeline = Pipeline([NullProcessor()])
task = PipelineTask(pipeline, PipelineParams(enable_heartbeats=True)) task = PipelineTask(pipeline, params=PipelineParams(enable_heartbeats=True))
runner = PipelineRunner() runner = PipelineRunner()

View File

@@ -1,5 +1,5 @@
# #
# Copyright (c) 2024, Daily # Copyright (c) 2024-2025, Daily
# #
# SPDX-License-Identifier: BSD 2-Clause License # SPDX-License-Identifier: BSD 2-Clause License
# #
@@ -117,7 +117,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -230,7 +230,7 @@ Your response will be turned into speech so use only simple words and punctuatio
) )
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -140,7 +140,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams(allow_interruptions=True), params=PipelineParams(allow_interruptions=True),
observers=[GoogleRTVIObserver(rtvi)], observers=[GoogleRTVIObserver(rtvi)],
) )

View File

@@ -346,7 +346,7 @@ async def main():
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=False)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=False))
@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):

View File

@@ -150,7 +150,7 @@ async def main(
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
if dialout_number: if dialout_number:
logger.debug("dialout number detected; doing dialout") logger.debug("dialout number detected; doing dialout")

View File

@@ -271,7 +271,7 @@ DO NOT say anything until you've determined if this is a voicemail or human."""
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams(allow_interruptions=True), params=PipelineParams(allow_interruptions=True),
) )
if dialout_number: if dialout_number:

View File

@@ -77,7 +77,7 @@ async def main(room_url: str, token: str, callId: str, sipUri: str):
] ]
) )
task = PipelineTask(pipeline, PipelineParams(allow_interruptions=True)) task = PipelineTask(pipeline, params=PipelineParams(allow_interruptions=True))
@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):

View File

@@ -90,7 +90,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams(allow_interruptions=True, enable_metrics=True), params=PipelineParams(allow_interruptions=True, enable_metrics=True),
) )
@transport.event_handler("on_first_participant_joined") @transport.event_handler("on_first_participant_joined")

View File

@@ -172,7 +172,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -198,7 +198,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -104,7 +104,7 @@ async def main(room_url, token=None):
main_task = PipelineTask( main_task = PipelineTask(
main_pipeline, main_pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=True, allow_interruptions=True,
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -155,8 +155,10 @@ Your task is to help the user understand and learn from this article in 2 senten
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
audio_out_sample_rate=44100, allow_interruptions=True, enable_metrics=True audio_out_sample_rate=44100,
allow_interruptions=True,
enable_metrics=True,
), ),
) )

View File

@@ -183,7 +183,7 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
PipelineParams( params=PipelineParams(
allow_interruptions=False, # We don't want to interrupt the translator bot allow_interruptions=False, # We don't want to interrupt the translator bot
enable_metrics=True, enable_metrics=True,
enable_usage_metrics=True, enable_usage_metrics=True,

View File

@@ -108,7 +108,9 @@ async def run_bot(websocket_client: WebSocket, stream_sid: str, testing: bool):
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
params=PipelineParams( params=PipelineParams(
audio_in_sample_rate=8000, audio_out_sample_rate=8000, allow_interruptions=True audio_in_sample_rate=8000,
audio_out_sample_rate=8000,
allow_interruptions=True,
), ),
) )

View File

@@ -142,7 +142,9 @@ async def run_client(client_name: str, server_url: str, duration_secs: int):
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
params=PipelineParams( params=PipelineParams(
audio_in_sample_rate=8000, audio_out_sample_rate=8000, allow_interruptions=True audio_in_sample_rate=8000,
audio_out_sample_rate=8000,
allow_interruptions=True,
), ),
) )

View File

@@ -125,7 +125,9 @@ async def main():
task = PipelineTask( task = PipelineTask(
pipeline, pipeline,
params=PipelineParams( params=PipelineParams(
audio_in_sample_rate=16000, audio_out_sample_rate=16000, allow_interruptions=True audio_in_sample_rate=16000,
audio_out_sample_rate=16000,
allow_interruptions=True,
), ),
) )

View File

@@ -131,6 +131,7 @@ class PipelineTask(BaseTask):
def __init__( def __init__(
self, self,
pipeline: BasePipeline, pipeline: BasePipeline,
*,
params: PipelineParams = PipelineParams(), params: PipelineParams = PipelineParams(),
observers: List[BaseObserver] = [], observers: List[BaseObserver] = [],
clock: BaseClock = SystemClock(), clock: BaseClock = SystemClock(),