Merge pull request #4267 from pipecat-ai/ac/fix-observer-cleanup-ordering
Fix observer cleanup ordering to stop proxy tasks before closing resources
This commit is contained in:
1
changelog/4267.fixed.md
Normal file
1
changelog/4267.fixed.md
Normal file
@@ -0,0 +1 @@
|
|||||||
|
- Fixed `ValueError: write to closed file` during pipeline shutdown when observers were active. Observer proxy tasks are now cancelled before observer resources are cleaned up.
|
||||||
1
changelog/4267.removed.md
Normal file
1
changelog/4267.removed.md
Normal file
@@ -0,0 +1 @@
|
|||||||
|
- ⚠️ Removed deprecated `PIPECAT_OBSERVER_FILES` environment variable support. Use `PIPECAT_SETUP_FILES` instead.
|
||||||
@@ -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,20 @@ 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()
|
||||||
|
|
||||||
|
# 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"):
|
||||||
@@ -971,37 +966,6 @@ class PipelineTask(BasePipelineTask):
|
|||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"{self} error running external setup from {f}: {e}")
|
logger.error(f"{self} error running external setup from {f}: {e}")
|
||||||
|
|
||||||
async def _load_observer_files(self):
|
|
||||||
"""Dynamically load observers from files listed in PIPECAT_OBSERVER_FILES."""
|
|
||||||
observer_files = [f for f in os.environ.get("PIPECAT_OBSERVER_FILES", "").split(":") if f]
|
|
||||||
for f in observer_files:
|
|
||||||
import warnings
|
|
||||||
|
|
||||||
with warnings.catch_warnings():
|
|
||||||
warnings.simplefilter("always")
|
|
||||||
warnings.warn(
|
|
||||||
"Observer files (and environment variable `PIPECAT_OBSERVER_FILES`) is deprecated, use setup files instead (and `PIPECAT_SETUP_FILES`) instead.",
|
|
||||||
DeprecationWarning,
|
|
||||||
)
|
|
||||||
|
|
||||||
try:
|
|
||||||
path = Path(f).resolve()
|
|
||||||
module_name = path.stem
|
|
||||||
spec = importlib.util.spec_from_file_location(module_name, str(path))
|
|
||||||
if spec:
|
|
||||||
logger.debug(f"{self} loading observers from {path}")
|
|
||||||
|
|
||||||
# Load module.
|
|
||||||
module = importlib.util.module_from_spec(spec)
|
|
||||||
spec.loader.exec_module(module)
|
|
||||||
|
|
||||||
# Create observers.
|
|
||||||
observers = await module.create_observers(self)
|
|
||||||
for observer in observers:
|
|
||||||
self.add_observer(observer)
|
|
||||||
except Exception as e:
|
|
||||||
logger.error(f"{self} error loading external observers from {f}: {e}")
|
|
||||||
|
|
||||||
def _print_dangling_tasks(self):
|
def _print_dangling_tasks(self):
|
||||||
"""Log any dangling tasks that haven't been properly cleaned up."""
|
"""Log any dangling tasks that haven't been properly cleaned up."""
|
||||||
tasks = [t.get_name() for t in self._task_manager.current_tasks()]
|
tasks = [t.get_name() for t in self._task_manager.current_tasks()]
|
||||||
|
|||||||
@@ -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