| """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) |
|
|