roadrecon / backend /app /api /fleet.py
Deeraj
Feat: add fleet mode and inference telemetry worker
f17ffc5
Raw
History Blame Contribute Delete
4.91 kB
"""Fleet Mode telemetry ingestion.
Mounted phones run inference locally and send only compact hazard telemetry plus an optional
privacy-processed thumbnail. This endpoint intentionally does no model inference and no DBSCAN in
the request path.
"""
from __future__ import annotations
import base64
import io
import uuid
from datetime import datetime, timedelta, timezone
from fastapi import APIRouter, Depends, HTTPException
from PIL import Image, UnidentifiedImageError
from sqlmodel import Session
from app.config import settings
from app.db.repository import Repository
from app.db.session import get_session
from app.models.schemas import (
FleetTelemetryIn,
FleetTelemetryResponse,
GeoSource,
HazardType,
Severity,
Source,
)
from app.models.tables import HazardReport
router = APIRouter(prefix="/fleet", tags=["fleet"])
def _utc(dt: datetime) -> datetime:
if dt.tzinfo is None:
return dt.replace(tzinfo=timezone.utc)
return dt.astimezone(timezone.utc)
def _severity_from_area(area_frac: float | None) -> Severity:
if area_frac is None:
return Severity.low
if area_frac >= settings.SEV_HIGH:
return Severity.high
if area_frac >= settings.SEV_MED:
return Severity.medium
return Severity.low
def _validate_timestamp(ts: datetime) -> datetime:
observed_at = _utc(ts)
now = datetime.now(timezone.utc)
if observed_at < now - timedelta(hours=settings.FLEET_MAX_EVENT_AGE_HOURS):
raise HTTPException(status_code=422, detail="Telemetry timestamp is too old.")
if observed_at > now + timedelta(minutes=settings.FLEET_MAX_CLOCK_SKEW_MINUTES):
raise HTTPException(status_code=422, detail="Telemetry timestamp is in the future.")
return observed_at
def _store_thumbnail(thumbnail_b64: str | None) -> str | None:
if not thumbnail_b64:
return None
raw_b64 = thumbnail_b64.split(",", 1)[1] if "," in thumbnail_b64 else thumbnail_b64
if len(raw_b64) > settings.FLEET_MAX_THUMBNAIL_BYTES * 2:
raise HTTPException(status_code=413, detail="Thumbnail payload is too large.")
try:
data = base64.b64decode(raw_b64, validate=True)
except ValueError as exc:
raise HTTPException(status_code=422, detail="Thumbnail is not valid Base64.") from exc
if len(data) > settings.FLEET_MAX_THUMBNAIL_BYTES:
raise HTTPException(status_code=413, detail="Thumbnail payload is too large.")
try:
img = Image.open(io.BytesIO(data)).convert("RGB")
except (UnidentifiedImageError, OSError) as exc:
raise HTTPException(status_code=422, detail="Thumbnail is not a valid image.") from exc
img.thumbnail((320, 320))
out = io.BytesIO()
img.save(out, format="JPEG", quality=62, optimize=True)
fleet_dir = settings.MEDIA_DIR / "fleet"
fleet_dir.mkdir(parents=True, exist_ok=True)
filename = f"{uuid.uuid4().hex}.jpg"
path = fleet_dir / filename
path.write_bytes(out.getvalue())
return f"media/fleet/{filename}"
@router.post("/telemetry", response_model=FleetTelemetryResponse)
async def ingest_telemetry(
payload: FleetTelemetryIn,
session: Session = Depends(get_session),
) -> FleetTelemetryResponse:
if payload.hazard_type == HazardType.other:
raise HTTPException(status_code=422, detail="Fleet telemetry requires a concrete hazard type.")
if payload.confidence < settings.FLEET_MIN_CONFIDENCE:
raise HTTPException(status_code=422, detail="Detection confidence is below the fleet threshold.")
if payload.bbox is not None and len(payload.bbox) != 4:
raise HTTPException(status_code=422, detail="bbox must contain [x1, y1, x2, y2].")
observed_at = _validate_timestamp(payload.timestamp)
repo = Repository(session)
if payload.client_event_id:
duplicate = repo.get_report_by_client_event(payload.client_event_id)
if duplicate is not None:
return FleetTelemetryResponse(
status="accepted",
report_id=duplicate.id,
queued_for_cluster=False,
)
thumbnail_path = _store_thumbnail(payload.thumbnail_b64)
report = HazardReport(
observed_at=observed_at,
source=Source.fleet_stream,
hazard_type=payload.hazard_type,
severity=_severity_from_area(payload.bbox_area_frac),
confidence=round(payload.confidence, 4),
lat=payload.lat,
lon=payload.lon,
geo_source=GeoSource.device,
media_path=thumbnail_path or "",
thumbnail_path=thumbnail_path,
bbox=payload.bbox or [],
bbox_area_frac=payload.bbox_area_frac,
device_speed=payload.speed,
client_event_id=payload.client_event_id,
auto_validated=True,
)
report = repo.add_report(report)
return FleetTelemetryResponse(status="accepted", report_id=report.id, queued_for_cluster=True)