59 lines
2.0 KiB
Python
59 lines
2.0 KiB
Python
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
|
|
from app.core.errors import ConflictError
|
|
from app.rotation.repository import RotationRepository
|
|
from app.rotation.service import RotationService
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
class SchedulerLoop:
|
|
def __init__(
|
|
self,
|
|
repository: RotationRepository,
|
|
rotation_service: RotationService,
|
|
) -> None:
|
|
self.repository = repository
|
|
self.rotation_service = rotation_service
|
|
self._task: asyncio.Task[None] | None = None
|
|
self._stopping = asyncio.Event()
|
|
|
|
def start(self, tick_seconds: int) -> None:
|
|
if self._task and not self._task.done():
|
|
return
|
|
self._task = asyncio.create_task(self._run(tick_seconds), name="fleet-schedule-loop")
|
|
|
|
async def stop(self) -> None:
|
|
self._stopping.set()
|
|
if self._task:
|
|
self._task.cancel()
|
|
await asyncio.gather(self._task, return_exceptions=True)
|
|
|
|
async def _run(self, tick_seconds: int) -> None:
|
|
while not self._stopping.is_set():
|
|
for group_id in self.repository.due_group_ids():
|
|
try:
|
|
self.rotation_service.start_group(group_id, trigger="scheduled")
|
|
break
|
|
except ConflictError:
|
|
break
|
|
except Exception:
|
|
logger.exception(
|
|
"Scheduled group rotation dispatch failed",
|
|
extra={"group_id": group_id},
|
|
)
|
|
try:
|
|
self.repository.defer_group(group_id)
|
|
except Exception:
|
|
logger.exception(
|
|
"Failed to defer group schedule",
|
|
extra={"group_id": group_id},
|
|
)
|
|
try:
|
|
await asyncio.wait_for(self._stopping.wait(), timeout=tick_seconds)
|
|
except TimeoutError:
|
|
continue
|