deerflow-code/offline-backend-20260512/backend/app/gateway/workflow_retention.py
2026-09-07 18:24:55 +08:00

105 lines
3.7 KiB
Python

"""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"]