511NY DOT cameras + traffic events map layers (#183)
GET /traffic/511 serves a bbox slice of the statewide 511NY camera and event feeds from an in-memory cache (cameras 1h, events 2m TTL, fetched lazily) -- public data, so no Firestore writes. A failed refresh keeps the last good data and reports the error; a schema change (nothing parses) is an error, not an empty layer. Frontend: opt-in "DOT Cameras" and "Traffic Events" overlays that fetch only while shown, with an on-map notice when the feed is down or stale. NY511_API_KEY is optional in config (the API answers without one today; the terms require a registered key). Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
1a01497f4d
commit
3347ced863
@@ -26,6 +26,7 @@ NODE_OFFLINE_THRESHOLD=90
|
||||
# Google Maps — for geocoding location strings extracted from transcripts
|
||||
# Enable "Geocoding API" in Cloud Console for this key
|
||||
GOOGLE_MAPS_API_KEY=
|
||||
NY511_API_KEY=
|
||||
|
||||
# OpenAI — for transcription (Whisper), intelligence extraction, embeddings, and summaries
|
||||
OPENAI_API_KEY=
|
||||
|
||||
@@ -32,6 +32,9 @@ class Settings(BaseSettings):
|
||||
|
||||
# Google Maps (geocoding)
|
||||
google_maps_api_key: Optional[str] = None
|
||||
# 511NY developer key (server-26#183). Optional: the API answered without one
|
||||
# as of 2026-09-27, but its terms require a registered key.
|
||||
ny511_api_key: Optional[str] = None
|
||||
|
||||
# Gemini (intelligence extraction, embeddings, incident summaries)
|
||||
gemini_api_key: Optional[str] = None
|
||||
|
||||
@@ -0,0 +1,140 @@
|
||||
"""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]
|
||||
@@ -17,7 +17,7 @@ from app.internal.auth import (
|
||||
require_node_service_or_firebase_token,
|
||||
)
|
||||
from app.routers import nodes, systems, calls, upload, tokens, incidents, alerts, admin, trips, places, links, users
|
||||
from app.routers import enrollment, media, org, waitlist, telemetry, replay
|
||||
from app.routers import enrollment, media, org, waitlist, telemetry, replay, traffic
|
||||
from app.internal import dynsec
|
||||
from app.internal import firestore as fstore
|
||||
|
||||
@@ -127,6 +127,7 @@ app.include_router(incidents.router, dependencies=[Depends(require_service_or_fi
|
||||
app.include_router(alerts.router, dependencies=[Depends(require_service_or_firebase_token)])
|
||||
app.include_router(trips.router, dependencies=[Depends(require_service_or_firebase_token)])
|
||||
app.include_router(places.router, dependencies=[Depends(require_service_or_firebase_token)])
|
||||
app.include_router(traffic.router, dependencies=[Depends(require_service_or_firebase_token)])
|
||||
app.include_router(upload.router) # auth is per-node, handled inline
|
||||
app.include_router(admin.router) # auth is per-endpoint (read: firebase, write: admin)
|
||||
app.include_router(replay.router) # auth: admin only (every route spends or reads a replay run)
|
||||
|
||||
@@ -0,0 +1,38 @@
|
||||
import asyncio
|
||||
from typing import Optional
|
||||
|
||||
from fastapi import APIRouter, HTTPException, Query
|
||||
|
||||
from app.internal import ny511
|
||||
|
||||
router = APIRouter(prefix="/traffic", tags=["traffic"])
|
||||
|
||||
# Per-layer cap per response. A statewide view is ~3k cameras; past this the
|
||||
# map is unreadable anyway, and the client is told the list was cut.
|
||||
MAX_ITEMS = 1500
|
||||
|
||||
|
||||
def _feed_status(feed: "ny511._Feed") -> dict:
|
||||
return {"fetched_at": feed.fetched_at, "error": feed.error}
|
||||
|
||||
|
||||
@router.get("/511")
|
||||
async def get_511(
|
||||
south: float = Query(..., ge=-90, le=90),
|
||||
west: float = Query(..., ge=-180, le=180),
|
||||
north: float = Query(..., ge=-90, le=90),
|
||||
east: float = Query(..., ge=-180, le=180),
|
||||
layers: Optional[str] = Query("cameras,events", description="comma list: cameras, events"),
|
||||
):
|
||||
if south > north or west > east:
|
||||
raise HTTPException(400, "bbox must satisfy south<=north and west<=east")
|
||||
wanted = {s.strip() for s in (layers or "").split(",") if s.strip()}
|
||||
feeds = {name: getattr(ny511, name) for name in ("cameras", "events") if name in wanted}
|
||||
await asyncio.gather(*(f.get() for f in feeds.values()))
|
||||
|
||||
out: dict = {}
|
||||
for name, feed in feeds.items():
|
||||
hits = ny511.in_bbox(feed.items, south, west, north, east)
|
||||
out[name] = hits[:MAX_ITEMS]
|
||||
out[f"{name}_status"] = {**_feed_status(feed), "total_in_bbox": len(hits), "truncated": len(hits) > MAX_ITEMS}
|
||||
return out
|
||||
@@ -0,0 +1,98 @@
|
||||
"""
|
||||
server-26#183 — 511NY cameras/events layer.
|
||||
|
||||
The feed is scraped from a third party, so the tests pin the two failure shapes
|
||||
that would otherwise look like "no traffic right now": a refresh error must keep
|
||||
the last good data AND report the error, and a schema change (every record
|
||||
unparseable) must be an error, not an empty list.
|
||||
"""
|
||||
import asyncio
|
||||
from unittest.mock import AsyncMock, patch
|
||||
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from app.main import app
|
||||
from app.internal import ny511
|
||||
from app.internal.auth import require_service_or_firebase_token
|
||||
|
||||
client = TestClient(app)
|
||||
|
||||
CAM = {"Latitude": 41.03, "Longitude": -73.76, "ID": "NYSDOT-1", "Name": "I-287 at Exit 5",
|
||||
"DirectionOfTravel": "Unknown", "RoadwayName": "I-287", "Url": "https://511ny.org/map/Cctv/1",
|
||||
"VideoUrl": None, "Disabled": False, "Blocked": False}
|
||||
EVENT = {"Latitude": 41.019265, "Longitude": -73.797869, "ID": "TRANSCOM-1", "EventType": "roadwork",
|
||||
"EventSubType": "Gas main repairs", "Severity": "Unknown", "RoadwayName": "NY 100",
|
||||
"DirectionOfTravel": "Both directions", "CountyName": "Westchester", "Description": "x" * 900,
|
||||
"StartDate": "28/09/2026 09:00:00", "PlannedEndDate": "", "LastUpdated": "26/09/2026 14:01:12"}
|
||||
|
||||
|
||||
def setup_function():
|
||||
app.dependency_overrides[require_service_or_firebase_token] = lambda: {"admin": True}
|
||||
for feed in (ny511.cameras, ny511.events):
|
||||
feed.items, feed.fetched_at, feed.error, feed._next_attempt = [], None, None, 0.0
|
||||
|
||||
|
||||
def teardown_function():
|
||||
app.dependency_overrides.pop(require_service_or_firebase_token, None)
|
||||
|
||||
|
||||
def test_normalize_camera_skips_disabled_blocked_and_zero_coords():
|
||||
assert ny511.normalize_camera(CAM)["image_url"] == "https://511ny.org/map/Cctv/1"
|
||||
assert ny511.normalize_camera({**CAM, "Disabled": True}) is None
|
||||
assert ny511.normalize_camera({**CAM, "Blocked": True}) is None
|
||||
assert ny511.normalize_camera({**CAM, "Latitude": 0, "Longitude": 0}) is None
|
||||
|
||||
|
||||
def test_normalize_event_parses_day_first_dates_and_truncates_description():
|
||||
e = ny511.normalize_event(EVENT)
|
||||
assert e["start_local"] == "2026-09-28T09:00:00" # DD/MM, not MM/DD
|
||||
assert e["planned_end_local"] is None
|
||||
assert len(e["description"]) == ny511._DESCRIPTION_MAX + 1
|
||||
|
||||
|
||||
def test_failed_refresh_keeps_last_good_data_and_reports_error():
|
||||
feed = ny511._Feed("getcameras", 3600, ny511.normalize_camera)
|
||||
with patch.object(ny511, "_fetch", AsyncMock(return_value=[ny511.normalize_camera(CAM)])):
|
||||
asyncio.run(feed.get())
|
||||
feed._next_attempt = 0.0
|
||||
with patch.object(ny511, "_fetch", AsyncMock(side_effect=RuntimeError("boom"))):
|
||||
asyncio.run(feed.get())
|
||||
assert len(feed.items) == 1 and feed.fetched_at is not None
|
||||
assert "boom" in feed.error
|
||||
assert feed._next_attempt - feed.fetched_at <= ny511.RETRY_AFTER_FAILURE_S + 5 # retries soon, not after the full TTL
|
||||
|
||||
|
||||
def test_schema_change_is_an_error_not_an_empty_layer():
|
||||
class Resp:
|
||||
def raise_for_status(self): pass
|
||||
def json(self): return [{"lat": 1, "lng": 2}] # renamed fields -> nothing parses
|
||||
|
||||
class Client:
|
||||
async def __aenter__(self): return self
|
||||
async def __aexit__(self, *a): pass
|
||||
async def get(self, *a, **k): return Resp()
|
||||
|
||||
with patch.object(ny511.httpx, "AsyncClient", lambda **k: Client()):
|
||||
try:
|
||||
asyncio.run(ny511._fetch("getcameras", ny511.normalize_camera))
|
||||
assert False, "expected a schema-change error"
|
||||
except ValueError as e:
|
||||
assert "schema" in str(e)
|
||||
|
||||
|
||||
def test_endpoint_filters_to_bbox_and_reports_status():
|
||||
far = {**CAM, "ID": "NYSDOT-2", "Latitude": 42.9, "Longitude": -78.8} # Buffalo
|
||||
for feed, rows, norm in ((ny511.cameras, [CAM, far], ny511.normalize_camera), (ny511.events, [EVENT], ny511.normalize_event)):
|
||||
feed.items = [norm(r) for r in rows]
|
||||
feed.fetched_at, feed._next_attempt = 1.0, float("inf")
|
||||
r = client.get("/traffic/511", params={"south": 40.9, "west": -74.0, "north": 41.4, "east": -73.4})
|
||||
assert r.status_code == 200
|
||||
body = r.json()
|
||||
assert [c["id"] for c in body["cameras"]] == ["NYSDOT-1"]
|
||||
assert body["cameras_status"] == {"fetched_at": 1.0, "error": None, "total_in_bbox": 1, "truncated": False}
|
||||
assert len(body["events"]) == 1
|
||||
|
||||
|
||||
def test_endpoint_rejects_inverted_bbox():
|
||||
r = client.get("/traffic/511", params={"south": 41.4, "west": -74.0, "north": 40.9, "east": -73.4})
|
||||
assert r.status_code == 400
|
||||
Reference in New Issue
Block a user