Spawn the main task before setup so attach happens uniformly
`PipelineRunner.run(task)` now calls `spawn(task)` first (which runs `task.attach()`) and lets `_setup_session` start every registered entry — main and pre-spawned — through the same path, instead of relying on `spawn`'s post-running fast-path to start the main task after setup. The two-branch wait stays for the `task is None` case but reads the runner_task directly off the freshly-spawned entry.
This commit is contained in:
@@ -183,25 +183,24 @@ class PipelineRunner(BaseObject, BusSubscriber):
|
|||||||
task: The pipeline task to run, or None.
|
task: The pipeline task to run, or None.
|
||||||
"""
|
"""
|
||||||
logger.debug(f"PipelineRunner '{self}': started running {task}")
|
logger.debug(f"PipelineRunner '{self}': started running {task}")
|
||||||
|
|
||||||
await self._setup_session()
|
|
||||||
|
|
||||||
self._shutdown_event.clear()
|
self._shutdown_event.clear()
|
||||||
|
|
||||||
# Register the main task like any spawned task so it shares the
|
# Treat the main task as a spawned task: ``spawn`` attaches it
|
||||||
# same start/cancel path.
|
# to the bus and registry, and ``_setup_session`` then starts
|
||||||
|
# every entry (main and pre-spawned) through the same code path.
|
||||||
if task is not None:
|
if task is not None:
|
||||||
await self.spawn(task)
|
await self.spawn(task)
|
||||||
|
|
||||||
|
await self._setup_session()
|
||||||
await self._call_event_handler("on_ready")
|
await self._call_event_handler("on_ready")
|
||||||
|
|
||||||
# Wait for the main task's background runner task to finish
|
# Wait for the main task's background runner task to finish
|
||||||
# (or for an explicit shutdown when there's no main task).
|
# (or for an explicit shutdown when there's no main task).
|
||||||
try:
|
try:
|
||||||
if task is not None:
|
if task is not None:
|
||||||
main_entry = self._entries.get(task.name)
|
runner_task = self._entries[task.name].runner_task
|
||||||
if main_entry and main_entry.runner_task:
|
if runner_task is not None:
|
||||||
await main_entry.runner_task
|
await runner_task
|
||||||
else:
|
else:
|
||||||
await self._shutdown_event.wait()
|
await self._shutdown_event.wait()
|
||||||
except asyncio.CancelledError:
|
except asyncio.CancelledError:
|
||||||
|
|||||||
Reference in New Issue
Block a user