Deprecate passing a worker to PipelineRunner.run()

Register the worker with PipelineRunner.add_workers() before calling
run() instead. The worker argument still works but now emits a
DeprecationWarning and will be removed in a future release.

Update the runner docstrings, the run_test() helper, and all examples
(including the asyncio.gather() forms) to use the new pattern.
This commit is contained in:
Aleix Conchillo Flaqué
2026-05-21 22:38:27 -07:00
parent e8ec7c585f
commit afa880f523
327 changed files with 734 additions and 379 deletions

View File

@@ -129,7 +129,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), handoff_after_start())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), handoff_after_start())
self.assertTrue(worker.active)
self.assertEqual(len(handoff_args_received), 1)
@@ -150,7 +151,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), wait_and_end())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), wait_and_end())
self.assertTrue(worker.active)
self.assertTrue(activated.is_set())
@@ -186,7 +188,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), drive())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), drive())
self.assertTrue(observed_while_active["active"])
self.assertEqual(observed_while_active["args"], args)
@@ -291,7 +294,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), wait_and_end())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), wait_and_end())
async def test_on_pipeline_started_event(self):
"""on_pipeline_started fires after pipeline starts."""
@@ -308,7 +312,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), wait_and_end())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), wait_and_end())
self.assertTrue(started.is_set())
@@ -327,7 +332,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), end_pipeline())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), end_pipeline())
self.assertTrue(finished_fired.is_set())
@@ -361,7 +367,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
)
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), send_end_message())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), send_end_message())
self.assertTrue(finished.is_set())
@@ -375,7 +382,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
try:
await asyncio.gather(runner.run(worker), send_cancel_message())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), send_cancel_message())
except asyncio.CancelledError:
pass
@@ -399,7 +407,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), push_frames())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), push_frames())
self.assertEqual(len(received), 1)
self.assertEqual(received[0].text, "injected")
@@ -422,7 +431,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), push_frames())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), push_frames())
self.assertEqual(len(received), 2)
self.assertEqual(received[0].text, "a")
@@ -447,7 +457,8 @@ class TestPipelineTaskLifecycle(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), self_handoff())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), self_handoff())
self.assertTrue(worker.active)
@@ -564,7 +575,8 @@ class TestEdgeToBus(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), push_frames())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), push_frames())
bus_frame_msgs = [m for m in sent if isinstance(m, BusFrameMessage)]
text_msgs = [m for m in bus_frame_msgs if isinstance(m.frame, TextFrame)]
@@ -592,7 +604,8 @@ class TestEdgeToBus(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), inject_frame())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), inject_frame())
# The frame passes through the identity pipeline and reaches
# EdgeSink, which re-broadcasts with source="worker". That's
@@ -623,7 +636,8 @@ class TestEdgeToBus(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), push_frames())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), push_frames())
bus_frame_msgs = [m for m in sent if isinstance(m, BusFrameMessage)]
self.assertEqual(len(bus_frame_msgs), 0)
@@ -652,7 +666,8 @@ class TestEdgeToBus(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), inject_frame())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), inject_frame())
self.assertEqual(len(received), 1)
self.assertEqual(received[0].text, "from_bus")
@@ -670,7 +685,8 @@ class TestEdgeToBus(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), push_frames())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), push_frames())
bus_frame_msgs = [m for m in sent if isinstance(m, BusFrameMessage)]
text_msgs = [m for m in bus_frame_msgs if isinstance(m.frame, TextFrame)]
@@ -703,7 +719,8 @@ class TestEdgeToBus(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), inject_frame())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), inject_frame())
self.assertEqual(len(received), 1)
self.assertEqual(received[0].text, "voice_frame")
@@ -733,7 +750,8 @@ class TestEdgeToBus(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), inject_frame())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), inject_frame())
self.assertEqual(len(received), 0)
@@ -777,7 +795,8 @@ class TestEdgeToBus(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), inject_frames())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), inject_frames())
self.assertEqual(len(received), 3)
@@ -822,7 +841,8 @@ class TestEdgeToBus(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), inject_frames())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), inject_frames())
texts = sorted([r.text for r in received])
self.assertEqual(texts, ["video", "voice"])
@@ -840,7 +860,8 @@ class TestEdgeToBus(unittest.IsolatedAsyncioTestCase):
await worker.queue_frame(EndFrame())
runner = PipelineRunner(bus=self.bus, handle_sigint=False)
await asyncio.gather(runner.run(worker), push_frames())
await runner.add_workers(worker)
await asyncio.gather(runner.run(), push_frames())
bus_frame_msgs = [m for m in sent if isinstance(m, BusFrameMessage)]
self.assertEqual(len(bus_frame_msgs), 0)