"""Periodic retention sweep for workflow runs, events, and artifacts. Consumes ``workflows.retention.event_retention_days`` / ``run_retention_days``. Events older than the event window are dropped first; terminal runs older than the run window are deleted with their node-run and artifact rows, and any leftover events for those run ids are cleaned up so the event log cannot outlive the run it describes. Artifact *files* stay on the existing sandbox/artifact lifecycle — this job only removes the workflow catalog rows. """ from __future__ import annotations import asyncio import logging from deerflow.config.workflow_config import WorkflowConfig from deerflow.persistence.workflow_events.base import WorkflowEventStore from deerflow.persistence.workflow_runs.base import WorkflowRunStore logger = logging.getLogger(__name__) _INTERVAL_SECONDS = 60 * 60 _MAX_BATCHES_PER_SWEEP = 8 _SQL_PURGE_BATCH = 500 class WorkflowRetentionCleaner: def __init__( self, runs: WorkflowRunStore, events: WorkflowEventStore, *, config: WorkflowConfig | None = None, ) -> None: self._runs = runs self._events = events self._config = config or WorkflowConfig() self._task: asyncio.Task | None = None self._stopping = False def start(self) -> None: if self._task is not None and not self._task.done(): return self._stopping = False self._task = asyncio.create_task(self._loop(), name="workflow-retention") logger.info( "workflow retention cleaner started (events=%sd runs=%sd)", self._config.retention.event_retention_days, self._config.retention.run_retention_days, ) async def stop(self) -> None: self._stopping = True if self._task is not None: self._task.cancel() try: await self._task except (asyncio.CancelledError, Exception): # noqa: BLE001 pass self._task = None async def sweep(self) -> dict[str, int]: """One pass. Safe to call from tests without starting the loop.""" event_days = max(1, self._config.retention.event_retention_days) run_days = max(1, self._config.retention.run_retention_days) events_deleted = await self._events.purge_before(older_than_days=event_days) runs_deleted = 0 for _ in range(_MAX_BATCHES_PER_SWEEP): ids = await self._runs.purge_terminal_before(older_than_days=run_days) if not ids: break runs_deleted += len(ids) events_deleted += await self._events.purge_for_runs(ids) if len(ids) < _SQL_PURGE_BATCH: break if events_deleted or runs_deleted: logger.info("workflow retention sweep: events=%s runs=%s", events_deleted, runs_deleted) return {"events": events_deleted, "runs": runs_deleted} async def _loop(self) -> None: try: await self.sweep() except asyncio.CancelledError: raise except Exception: # noqa: BLE001 - the loop must survive a bad pass logger.exception("workflow retention sweep failed") while not self._stopping: try: await asyncio.sleep(_INTERVAL_SECONDS) except asyncio.CancelledError: raise if self._stopping: return try: await self.sweep() except asyncio.CancelledError: raise except Exception: # noqa: BLE001 logger.exception("workflow retention sweep failed") __all__ = ["WorkflowRetentionCleaner"]