"""511NY (NYSDOT) traffic cameras and events, cached in memory (server-26#183). Public data, identical for every org, so it is NOT written to Firestore: the statewide feeds are ~3k cameras and ~2.3k events, and re-writing them every poll would be millions of writes a day for data nobody needs history of. The cache is filled lazily on request and refreshed per TTL, so an idle deploy makes no 511 calls at all. A failed refresh keeps serving the last good data and reports the error and its age to the caller -- an empty layer must never be the only symptom of a dead feed (the AI-silent-failures lesson). """ import asyncio import time from datetime import datetime from typing import Any, Dict, List, Optional import httpx from app.config import settings from app.internal.logger import logger _BASE = "https://511ny.org/api" CAMERAS_TTL_S = 60 * 60 # camera list is near-static EVENTS_TTL_S = 2 * 60 # accidents/closures change minute to minute RETRY_AFTER_FAILURE_S = 60 _DESCRIPTION_MAX = 500 class _Feed: def __init__(self, path: str, ttl_s: int, normalize): self.path = path self.ttl_s = ttl_s self.normalize = normalize self.items: List[Dict[str, Any]] = [] self.fetched_at: Optional[float] = None # epoch s of last SUCCESSFUL fetch self.error: Optional[str] = None self._next_attempt = 0.0 self._lock = asyncio.Lock() async def get(self) -> "_Feed": if time.time() < self._next_attempt: return self async with self._lock: if time.time() < self._next_attempt: return self # another request refreshed while we waited try: self.items = await _fetch(self.path, self.normalize) self.fetched_at = time.time() self.error = None self._next_attempt = time.time() + self.ttl_s except Exception as e: self._next_attempt = time.time() + min(self.ttl_s, RETRY_AFTER_FAILURE_S) self.error = f"{type(e).__name__}: {e}"[:300] logger.warning(f"511NY {self.path} refresh failed, serving {len(self.items)} cached: {self.error}") return self async def _fetch(path: str, normalize) -> List[Dict[str, Any]]: params = {"format": "json"} if settings.ny511_api_key: params["key"] = settings.ny511_api_key async with httpx.AsyncClient(timeout=20.0) as client: r = await client.get(f"{_BASE}/{path}", params=params) r.raise_for_status() raw = r.json() if not isinstance(raw, list): raise ValueError(f"expected a JSON list, got {type(raw).__name__}") out = [n for n in (normalize(x) for x in raw) if n is not None] if raw and not out: # every record failed to normalize: the schema changed under us raise ValueError(f"0 of {len(raw)} records parsed -- 511NY schema change?") return out def _coords(x: Dict[str, Any]) -> Optional[tuple]: try: lat, lon = float(x["Latitude"]), float(x["Longitude"]) except (KeyError, TypeError, ValueError): return None if lat == 0 and lon == 0: return None return lat, lon def _local_iso(s: Any) -> Optional[str]: """511NY stamps are 'DD/MM/YYYY HH:MM:SS' New York local time. Returned as a naive ISO string (no offset) -- display-only, never compared to UTC.""" if not s: return None try: return datetime.strptime(s, "%d/%m/%Y %H:%M:%S").isoformat() except (TypeError, ValueError): return None def normalize_camera(x: Dict[str, Any]) -> Optional[Dict[str, Any]]: c = _coords(x) if c is None or x.get("Disabled") or x.get("Blocked") or not x.get("ID"): return None return { "id": x["ID"], "lat": c[0], "lon": c[1], "name": x.get("Name") or "", "roadway": x.get("RoadwayName") or "", "direction": x.get("DirectionOfTravel") or "", "image_url": x.get("Url"), # 511NY serves the current still at this URL "video_url": x.get("VideoUrl"), # HLS playlist, when the camera streams } def normalize_event(x: Dict[str, Any]) -> Optional[Dict[str, Any]]: c = _coords(x) if c is None or not x.get("ID"): return None desc = x.get("Description") or "" return { "id": x["ID"], "lat": c[0], "lon": c[1], "type": x.get("EventType") or "", "subtype": x.get("EventSubType") or "", "severity": x.get("Severity") or "", "roadway": x.get("RoadwayName") or "", "direction": x.get("DirectionOfTravel") or "", "county": x.get("CountyName") or "", "description": desc[:_DESCRIPTION_MAX] + ("…" if len(desc) > _DESCRIPTION_MAX else ""), "start_local": _local_iso(x.get("StartDate")), "planned_end_local": _local_iso(x.get("PlannedEndDate")), "updated_local": _local_iso(x.get("LastUpdated")), } cameras = _Feed("getcameras", CAMERAS_TTL_S, normalize_camera) events = _Feed("getevents", EVENTS_TTL_S, normalize_event) def in_bbox(items: List[Dict[str, Any]], south: float, west: float, north: float, east: float) -> List[Dict[str, Any]]: return [i for i in items if south <= i["lat"] <= north and west <= i["lon"] <= east]