77 lines
2.9 KiB
Python
77 lines
2.9 KiB
Python
import asyncio
|
|
from uuid import UUID
|
|
|
|
from .models import CheckResult, Monitor, MonitorCreate, MonitorStatus, MonitorUpdate, now_utc
|
|
|
|
|
|
class NotFoundError(Exception):
|
|
pass
|
|
|
|
|
|
class StaleCheckError(Exception):
|
|
pass
|
|
|
|
|
|
class MonitorStore:
|
|
def __init__(self, default_timeout: float, max_timeout: float) -> None:
|
|
self._items: dict[UUID, Monitor] = {}
|
|
self._lock = asyncio.Lock()
|
|
self._default_timeout = default_timeout
|
|
self._max_timeout = max_timeout
|
|
|
|
async def create(self, data: MonitorCreate) -> Monitor:
|
|
timeout = data.timeout_seconds or self._default_timeout
|
|
self._validate_timeout(timeout)
|
|
monitor = Monitor(name=data.name, url=data.url, timeout_seconds=timeout)
|
|
async with self._lock:
|
|
self._items[monitor.id] = monitor
|
|
return monitor.model_copy(deep=True)
|
|
|
|
async def list(self) -> list[Monitor]:
|
|
async with self._lock:
|
|
values = sorted(self._items.values(), key=lambda item: item.created_at)
|
|
return [item.model_copy(deep=True) for item in values]
|
|
|
|
async def get(self, monitor_id: UUID) -> Monitor:
|
|
async with self._lock:
|
|
item = self._items.get(monitor_id)
|
|
if item is None:
|
|
raise NotFoundError
|
|
return item.model_copy(deep=True)
|
|
|
|
async def update(self, monitor_id: UUID, data: MonitorUpdate) -> Monitor:
|
|
async with self._lock:
|
|
current = self._items.get(monitor_id)
|
|
if current is None:
|
|
raise NotFoundError
|
|
changes = data.model_dump(exclude_unset=True)
|
|
if "timeout_seconds" in changes:
|
|
self._validate_timeout(changes["timeout_seconds"])
|
|
material = "url" in changes or "timeout_seconds" in changes
|
|
changes["updated_at"] = now_utc()
|
|
if material:
|
|
changes["status"] = MonitorStatus()
|
|
updated = current.model_copy(update=changes, deep=True)
|
|
self._items[monitor_id] = updated
|
|
return updated.model_copy(deep=True)
|
|
|
|
async def delete(self, monitor_id: UUID) -> None:
|
|
async with self._lock:
|
|
if self._items.pop(monitor_id, None) is None:
|
|
raise NotFoundError
|
|
|
|
async def apply_check(self, monitor_id: UUID, expected_url: str, result: CheckResult) -> Monitor:
|
|
async with self._lock:
|
|
current = self._items.get(monitor_id)
|
|
if current is None:
|
|
raise NotFoundError
|
|
if str(current.url) != expected_url:
|
|
raise StaleCheckError
|
|
updated = current.model_copy(update={"status": result, "updated_at": now_utc()}, deep=True)
|
|
self._items[monitor_id] = updated
|
|
return updated.model_copy(deep=True)
|
|
|
|
def _validate_timeout(self, timeout: float) -> None:
|
|
if timeout > self._max_timeout:
|
|
raise ValueError(f"timeout_seconds must not exceed {self._max_timeout}")
|