Fix observer cleanup ordering to stop proxy tasks before closing resources
During pipeline shutdown, proxy tasks must be cancelled before observer resources are cleaned up. Previously, stop() was called inside _cancel_tasks() and start() was called in _start_tasks(), which could lead to proxy tasks still consuming frames after observer resources were closed. Now the lifecycle is explicit in _handle_start_frame: start() after all observers are loaded, and stop() before cleanup() on shutdown. Also fixes misleading variable name in TaskObserver.cleanup() where iterating self._proxies yields observer keys, not Proxy values. Fixes #4195
This commit is contained in:
@@ -613,9 +613,6 @@ class PipelineTask(BasePipelineTask):
|
|||||||
self._process_push_task = self._task_manager.create_task(
|
self._process_push_task = self._task_manager.create_task(
|
||||||
self._process_push_queue(), f"{self}::_process_push_queue"
|
self._process_push_queue(), f"{self}::_process_push_queue"
|
||||||
)
|
)
|
||||||
|
|
||||||
await self._observer.start()
|
|
||||||
|
|
||||||
return self._process_push_task
|
return self._process_push_task
|
||||||
|
|
||||||
def _maybe_start_heartbeat_tasks(self):
|
def _maybe_start_heartbeat_tasks(self):
|
||||||
@@ -637,8 +634,6 @@ class PipelineTask(BasePipelineTask):
|
|||||||
|
|
||||||
async def _cancel_tasks(self):
|
async def _cancel_tasks(self):
|
||||||
"""Cancel all running pipeline tasks."""
|
"""Cancel all running pipeline tasks."""
|
||||||
await self._observer.stop()
|
|
||||||
|
|
||||||
if self._process_push_task:
|
if self._process_push_task:
|
||||||
await self._task_manager.cancel_task(self._process_push_task)
|
await self._task_manager.cancel_task(self._process_push_task)
|
||||||
self._process_push_task = None
|
self._process_push_task = None
|
||||||
@@ -720,12 +715,6 @@ class PipelineTask(BasePipelineTask):
|
|||||||
|
|
||||||
async def _setup(self, params: PipelineTaskParams):
|
async def _setup(self, params: PipelineTaskParams):
|
||||||
"""Set up the pipeline task and all processors."""
|
"""Set up the pipeline task and all processors."""
|
||||||
# Do any additional pipeline task setup externally.
|
|
||||||
await self._load_setup_files()
|
|
||||||
|
|
||||||
# Load additional observers.
|
|
||||||
await self._load_observer_files()
|
|
||||||
|
|
||||||
mgr_params = TaskManagerParams(loop=params.loop)
|
mgr_params = TaskManagerParams(loop=params.loop)
|
||||||
self._task_manager.setup(mgr_params)
|
self._task_manager.setup(mgr_params)
|
||||||
|
|
||||||
@@ -736,14 +725,21 @@ class PipelineTask(BasePipelineTask):
|
|||||||
)
|
)
|
||||||
await self._pipeline.setup(setup)
|
await self._pipeline.setup(setup)
|
||||||
|
|
||||||
|
# Do any additional pipeline task setup externally.
|
||||||
|
await self._load_setup_files()
|
||||||
|
await self._load_observer_files()
|
||||||
|
|
||||||
|
# Start task observer.
|
||||||
|
await self._observer.start()
|
||||||
|
|
||||||
async def _cleanup(self, cleanup_pipeline: bool):
|
async def _cleanup(self, cleanup_pipeline: bool):
|
||||||
"""Clean up the pipeline task and processors."""
|
"""Clean up the pipeline task and processors."""
|
||||||
# Cleanup base object.
|
# Cleanup base object.
|
||||||
await self.cleanup()
|
await self.cleanup()
|
||||||
|
|
||||||
# Cleanup observers.
|
# Cleanup observers.
|
||||||
if self._observer:
|
await self._observer.stop()
|
||||||
await self._observer.cleanup()
|
await self._observer.cleanup()
|
||||||
|
|
||||||
# End conversation tracing if it's active - this will also close any active turn span
|
# End conversation tracing if it's active - this will also close any active turn span
|
||||||
if self._enable_tracing and hasattr(self, "_turn_trace_observer"):
|
if self._enable_tracing and hasattr(self, "_turn_trace_observer"):
|
||||||
|
|||||||
@@ -131,8 +131,8 @@ class TaskObserver(BaseObserver):
|
|||||||
if not self._proxies:
|
if not self._proxies:
|
||||||
return
|
return
|
||||||
|
|
||||||
for proxy in self._proxies:
|
for observer in self._proxies:
|
||||||
await proxy.cleanup()
|
await observer.cleanup()
|
||||||
|
|
||||||
async def on_pipeline_started(self):
|
async def on_pipeline_started(self):
|
||||||
"""Forward pipeline started signal to all managed observers."""
|
"""Forward pipeline started signal to all managed observers."""
|
||||||
|
|||||||
Reference in New Issue
Block a user