from __future__ import annotations import asyncio from app.core.logger import get_logger from app.services.cleanup import CleanupService logger = get_logger(__name__) class CleanupWorker: def __init__(self, service: CleanupService, interval_seconds: int) -> None: self.service = service self.interval_seconds = interval_seconds self._task: asyncio.Task[None] | None = None self._stop = asyncio.Event() async def start(self) -> None: self._stop.clear() self._task = asyncio.create_task(self._run(), name="media-cleanup-worker") async def stop(self) -> None: self._stop.set() if self._task: self._task.cancel() try: await self._task except asyncio.CancelledError: logger.debug("cleanup worker task cancelled") self._task = None async def _run(self) -> None: while not self._stop.is_set(): try: removed = await self.service.cleanup_expired() if removed: logger.info("cleanup worker removed workspaces", extra={"count": removed}) except asyncio.CancelledError: raise except Exception: logger.exception("cleanup worker iteration failed") try: await asyncio.wait_for(self._stop.wait(), timeout=self.interval_seconds) except asyncio.TimeoutError: continue