Spaces:
Running
Running
| 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 | |