Add GET /api/map/sentinel1?bbox=... which queries Planetary Computer STAC (sentinel-1-grd, last 7d, most recent), signs the vv/hh COG with a SAS token, and returns a same-origin /titiler/... XYZ tile template. Cache keyed on quantized bbox + UTC day (20 min TTL). 429 -> 429 (Retry-After), no imagery -> 404, other upstream errors -> 502. Self-host TiTiler on the Pi (ghcr.io/developmentseed/titiler, arm64, 1G cap, host loopback 8001 -> container 8000) instead of titiler.xyz. nginx /titiler/ proxy snippet routes browser tiles to it; TITILER_PUBLIC_BASE and TITILER_INTERNAL_URL are env-driven (no hardcoded hostnames).
1269 lines
45 KiB
Python
1269 lines
45 KiB
Python
"""Live map overlays: parsers, TTL cache, and upstream fetchers.
|
||
|
||
Moving objects (aircraft, vessels, trains) and alerts are vectors. Radar /
|
||
GIBS fire tiles are rasters served directly to the browser — this module
|
||
only returns tile *templates* and GeoJSON/JSON for vectors.
|
||
|
||
Third-party APIs that leak keys, lack CORS, or rate-limit by IP are fetched
|
||
here (FastAPI), never from Leaflet. Viewport bbox only; never a global ADS-B
|
||
or AIS poll.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import logging
|
||
import math
|
||
import time
|
||
from datetime import datetime, timedelta, timezone
|
||
from typing import Any, Awaitable, Callable
|
||
|
||
import httpx
|
||
from urllib.parse import quote, urlencode
|
||
|
||
from config import OSINT_USER_AGENT, TITILER_PUBLIC_BASE
|
||
|
||
logger = logging.getLogger("osint.live_layers")
|
||
|
||
MARKER_FIELDS = ("id", "lat", "lon", "heading", "speed", "label", "extra")
|
||
|
||
ADSB_LOL_BASE = "https://api.adsb.lol"
|
||
PLANESPOTTERS_PHOTO = "https://api.planespotters.net/pub/photos"
|
||
# Planespotters ToS cap server-side JSON caching at 24 hours.
|
||
PLANESPOTTERS_CACHE_TTL = 24 * 3600
|
||
AMTRAKER_TRAINS = "https://api.amtraker.com/v3/trains"
|
||
RAINVIEWER_MAPS = "https://api.rainviewer.com/public/weather-maps.json"
|
||
NWS_ALERTS = "https://api.weather.gov/alerts/active"
|
||
IEM_SBW = "https://mesonet.agron.iastate.edu/geojson/sbw.geojson"
|
||
WFIGS_INCIDENTS = (
|
||
"https://services3.arcgis.com/T4QMspbfLg3qTGWY/arcgis/rest/services/"
|
||
"WFIGS_Incident_Locations_Current/FeatureServer/0/query"
|
||
)
|
||
WFIGS_PERIMETERS = (
|
||
"https://services3.arcgis.com/T4QMspbfLg3qTGWY/arcgis/rest/services/"
|
||
"WFIGS_Interagency_Perimeters_Current/FeatureServer/0/query"
|
||
)
|
||
NHC_STORMS = "https://www.nhc.noaa.gov/CurrentStorms.json"
|
||
|
||
PC_STAC_SEARCH = "https://planetarycomputer.microsoft.com/api/stac/v1/search"
|
||
PC_SAS_TOKEN = "https://planetarycomputer.microsoft.com/api/sas/v1/token/sentinel-1-grd"
|
||
# Self-hosted TiTiler on the Pi, exposed same-origin through the osint.rpi.local
|
||
# nginx vhost. Relative template — Leaflet resolves it against the page origin,
|
||
# so the browser never touches a raw loopback port or titiler.xyz.
|
||
TITILER_COG_TILES = f"{TITILER_PUBLIC_BASE}/cog/tiles/WebMercatorQuad/{{z}}/{{x}}/{{y}}@1x"
|
||
SENTINEL1_TTL = 20 * 60 # 15–30 min quota-friendly window
|
||
SENTINEL1_ATTRIBUTION = "Copernicus Sentinel-1 / Microsoft Planetary Computer"
|
||
|
||
IEM_NEXRAD = "https://mesonet.agron.iastate.edu/cache/tile.py/1.0.0/nexrad-n0q/{z}/{x}/{y}.png"
|
||
GIBS_THERMAL = (
|
||
"https://gibs.earthdata.nasa.gov/wmts/epsg3857/best/"
|
||
"VIIRS_SNPP_Thermal_Anomalies_375m_All/default/{time}/"
|
||
"GoogleMapsCompatible_Level9/{z}/{y}/{x}.png"
|
||
)
|
||
|
||
CONUS = (-125.0, 24.0, -66.0, 50.0)
|
||
MAX_RADIUS_NM = 150
|
||
DEFAULT_LIMIT = 2000
|
||
|
||
_HTTP_TIMEOUT = httpx.Timeout(8.0, connect=3.0)
|
||
_http: httpx.AsyncClient | None = None
|
||
_COMPASS = {
|
||
"N": 0, "NE": 45, "E": 90, "SE": 135,
|
||
"S": 180, "SW": 225, "W": 270, "NW": 315,
|
||
"NNE": 22, "ENE": 67, "ESE": 112, "SSE": 157,
|
||
"SSW": 202, "WSW": 247, "WNW": 292, "NNW": 337,
|
||
}
|
||
|
||
_cache: dict[str, tuple[float, Any]] = {}
|
||
_key_locks: dict[str, asyncio.Lock] = {}
|
||
_key_locks_guard = asyncio.Lock()
|
||
_QUANT = 0.25 # degrees — pan jitter inside a cell reuses the TTL entry
|
||
|
||
# Last-known AIS positions (MMSI -> marker). Filled by ais_stream worker.
|
||
vessel_last_known: dict[str, dict] = {}
|
||
vessel_lock = asyncio.Lock()
|
||
# Viewport-following accumulates vessels across every region visited in a
|
||
# session — keep the in-memory store bounded (oldest entries evicted).
|
||
_MAX_VESSELS = 6000
|
||
# Last ADS-B snapshot + WFIGS points for fire↔tanker correlation.
|
||
aircraft_last_known: dict[str, dict] = {}
|
||
fire_last_known: list[dict] = []
|
||
|
||
|
||
def overlay_catalog() -> dict:
|
||
"""Tile templates + vector endpoint map for the layer panel. No secrets."""
|
||
today = datetime.now(timezone.utc).date().isoformat()
|
||
return {
|
||
"radar_iem": {
|
||
"id": "radar_iem",
|
||
"title": "IEM NEXRAD (CONUS)",
|
||
"kind": "raster",
|
||
"tileUrl": IEM_NEXRAD,
|
||
"maxZoom": 18,
|
||
"attribution": "Iowa Environmental Mesonet",
|
||
},
|
||
"radar_rainviewer": {
|
||
"id": "radar_rainviewer",
|
||
"title": "RainViewer (global)",
|
||
"kind": "raster",
|
||
"tileUrl": None, # filled from /api/map/radar frames
|
||
"maxZoom": 7,
|
||
"attribution": '<a href="https://www.rainviewer.com/">Weather data by RainViewer</a>',
|
||
},
|
||
"gibs_thermal": {
|
||
"id": "gibs_thermal",
|
||
"title": "GIBS VIIRS thermal anomalies",
|
||
"kind": "raster",
|
||
"tileUrl": GIBS_THERMAL.replace("{time}", today),
|
||
"timeTemplate": GIBS_THERMAL,
|
||
"maxZoom": 9,
|
||
"attribution": "NASA GIBS / EOSDIS",
|
||
},
|
||
"sentinel1_sar": {
|
||
"id": "sentinel1_sar",
|
||
"title": "Sentinel-1 SAR (Cloud-Penetrating)",
|
||
"kind": "raster",
|
||
"tileUrl": None, # filled from /api/map/sentinel1 (signed COG)
|
||
"endpoint": "/api/map/sentinel1",
|
||
"maxZoom": 14,
|
||
"attribution": SENTINEL1_ATTRIBUTION,
|
||
},
|
||
"nws_alerts": {"id": "nws_alerts", "kind": "geojson", "endpoint": "/api/weather-alerts"},
|
||
"wfigs_incidents": {"id": "wfigs_incidents", "kind": "points", "endpoint": "/api/fire-incidents"},
|
||
"wfigs_perimeters": {"id": "wfigs_perimeters", "kind": "geojson", "endpoint": "/api/fire-perimeters"},
|
||
"aircraft": {"id": "aircraft", "kind": "points", "endpoint": "/api/aircraft"},
|
||
"vessels": {"id": "vessels", "kind": "points", "endpoint": "/api/vessels"},
|
||
"trains": {"id": "trains", "kind": "points", "endpoint": "/api/trains"},
|
||
"storms": {"id": "storms", "kind": "points", "endpoint": "/api/storms"},
|
||
}
|
||
|
||
|
||
def parse_bbox(bbox: str) -> tuple[float, float, float, float]:
|
||
"""Parse 'minlon,minlat,maxlon,maxlat' into four floats."""
|
||
parts = [p.strip() for p in (bbox or "").split(",")]
|
||
if len(parts) != 4:
|
||
raise ValueError("bbox must be 'minlon,minlat,maxlon,maxlat'")
|
||
try:
|
||
minlon, minlat, maxlon, maxlat = (float(p) for p in parts)
|
||
except ValueError as exc:
|
||
raise ValueError("bbox values must be floats") from exc
|
||
return minlon, minlat, maxlon, maxlat
|
||
|
||
|
||
def quantize_bbox(
|
||
minlon: float, minlat: float, maxlon: float, maxlat: float,
|
||
step: float = _QUANT,
|
||
) -> tuple[float, float, float, float]:
|
||
"""Snap a viewport to a coarse cell so nearby pans share a cache key.
|
||
|
||
The returned envelope is expanded to cover the original box.
|
||
"""
|
||
def q_down(v: float, lo: float, hi: float) -> float:
|
||
v = max(lo, min(hi, v))
|
||
return math.floor(v / step) * step
|
||
|
||
return (
|
||
round(q_down(minlon, -180.0, 180.0), 4),
|
||
round(q_down(minlat, -90.0, 90.0), 4),
|
||
round(q_down(maxlon, -180.0, 180.0) + step, 4),
|
||
round(q_down(maxlat, -90.0, 90.0) + step, 4),
|
||
)
|
||
|
||
|
||
def bbox_cell_key(bbox: str | None) -> str:
|
||
"""Stable cache-key fragment for a viewport (or 'all')."""
|
||
if not bbox:
|
||
return "all"
|
||
return ",".join(f"{v:.4f}" for v in quantize_bbox(*parse_bbox(bbox)))
|
||
|
||
|
||
def bbox_center_radius_nm(
|
||
minlon: float, minlat: float, maxlon: float, maxlat: float,
|
||
) -> tuple[float, float, int]:
|
||
"""Viewport center + half-diagonal radius in nautical miles, clamped ≤ 150."""
|
||
lat = (minlat + maxlat) / 2.0
|
||
lon = (minlon + maxlon) / 2.0
|
||
# Half the diagonal of the box, in nm (1 deg lat ≈ 60 nm).
|
||
dlat = abs(maxlat - minlat) / 2.0
|
||
dlon = abs(maxlon - minlon) / 2.0
|
||
km = _haversine_km(lat, lon, lat + dlat, lon + dlon)
|
||
nm = km / 1.852
|
||
radius = max(1, min(MAX_RADIUS_NM, int(math.ceil(nm))))
|
||
return lat, lon, radius
|
||
|
||
|
||
def _haversine_km(lat1: float, lon1: float, lat2: float, lon2: float) -> float:
|
||
r = 6371.0
|
||
p1, p2 = math.radians(lat1), math.radians(lat2)
|
||
dphi = math.radians(lat2 - lat1)
|
||
dlmb = math.radians(lon2 - lon1)
|
||
a = math.sin(dphi / 2) ** 2 + math.cos(p1) * math.cos(p2) * math.sin(dlmb / 2) ** 2
|
||
return 2 * r * math.asin(min(1.0, math.sqrt(a)))
|
||
|
||
|
||
def to_marker(
|
||
id_: str,
|
||
lat: float | None,
|
||
lon: float | None,
|
||
heading: float | None = None,
|
||
speed: float | None = None,
|
||
label: str | None = None,
|
||
extra: dict | None = None,
|
||
) -> dict:
|
||
return {
|
||
"id": str(id_),
|
||
"lat": lat,
|
||
"lon": lon,
|
||
"heading": heading,
|
||
"speed": speed,
|
||
"label": label or str(id_),
|
||
"extra": extra or {},
|
||
}
|
||
|
||
|
||
def filter_points_bbox(
|
||
points: list[dict],
|
||
minlon: float, minlat: float, maxlon: float, maxlat: float,
|
||
limit: int = DEFAULT_LIMIT,
|
||
) -> list[dict]:
|
||
out = []
|
||
for p in points:
|
||
lat, lon = p.get("lat"), p.get("lon")
|
||
if lat is None or lon is None:
|
||
continue
|
||
if minlat <= lat <= maxlat and minlon <= lon <= maxlon:
|
||
out.append(p)
|
||
if len(out) >= limit:
|
||
break
|
||
return out
|
||
|
||
|
||
_ALERT_KEEP = ("event", "severity", "headline", "areaDesc", "wfo", "source")
|
||
|
||
|
||
def slim_alert_properties(props: dict | None) -> dict:
|
||
"""Keep only the fields the map popup reads."""
|
||
src = props or {}
|
||
return {k: src.get(k) for k in _ALERT_KEEP}
|
||
|
||
|
||
def _walk_coords(coords: Any, acc: list[float]) -> None:
|
||
if not coords:
|
||
return
|
||
first = coords[0]
|
||
if isinstance(first, (int, float)):
|
||
lon, lat = float(coords[0]), float(coords[1])
|
||
acc[0] = min(acc[0], lon)
|
||
acc[1] = min(acc[1], lat)
|
||
acc[2] = max(acc[2], lon)
|
||
acc[3] = max(acc[3], lat)
|
||
return
|
||
for child in coords:
|
||
_walk_coords(child, acc)
|
||
|
||
|
||
def _geom_envelope(geom: dict | None) -> tuple[float, float, float, float] | None:
|
||
if not geom or not isinstance(geom, dict):
|
||
return None
|
||
if geom.get("type") == "GeometryCollection":
|
||
env: list[float] | None = None
|
||
for g in geom.get("geometries") or []:
|
||
e = _geom_envelope(g)
|
||
if e is None:
|
||
continue
|
||
if env is None:
|
||
env = list(e)
|
||
else:
|
||
env[0] = min(env[0], e[0])
|
||
env[1] = min(env[1], e[1])
|
||
env[2] = max(env[2], e[2])
|
||
env[3] = max(env[3], e[3])
|
||
return tuple(env) if env else None # type: ignore[return-value]
|
||
coords = geom.get("coordinates")
|
||
if coords is None:
|
||
return None
|
||
acc = [180.0, 90.0, -180.0, -90.0]
|
||
try:
|
||
_walk_coords(coords, acc)
|
||
except (TypeError, ValueError, IndexError):
|
||
return None
|
||
if acc[0] > acc[2]:
|
||
return None
|
||
return acc[0], acc[1], acc[2], acc[3]
|
||
|
||
|
||
def clip_fc_to_bbox(
|
||
fc: dict | None,
|
||
minlon: float, minlat: float, maxlon: float, maxlat: float,
|
||
) -> dict:
|
||
"""Drop features whose geometry envelope misses the viewport. No shapely."""
|
||
box = (minlon, minlat, maxlon, maxlat)
|
||
out = []
|
||
for feat in (fc or {}).get("features") or []:
|
||
if not isinstance(feat, dict):
|
||
continue
|
||
env = _geom_envelope(feat.get("geometry"))
|
||
if env is None:
|
||
continue
|
||
if env[0] <= box[2] and env[2] >= box[0] and env[1] <= box[3] and env[3] >= box[1]:
|
||
out.append(feat)
|
||
return {"type": "FeatureCollection", "features": out}
|
||
|
||
|
||
def _f(value: object) -> float | None:
|
||
if value is None or value == "":
|
||
return None
|
||
try:
|
||
return float(value)
|
||
except (TypeError, ValueError):
|
||
return None
|
||
|
||
|
||
def _heading(value: object) -> float | None:
|
||
if value is None or value == "":
|
||
return None
|
||
if isinstance(value, str):
|
||
key = value.strip().upper()
|
||
if key in _COMPASS:
|
||
return float(_COMPASS[key])
|
||
num = _f(value)
|
||
if num is None:
|
||
return None
|
||
if num < 0 or num > 360:
|
||
return None
|
||
return num
|
||
|
||
|
||
def _s(value: object) -> str | None:
|
||
if value is None:
|
||
return None
|
||
text = str(value).strip()
|
||
return text or None
|
||
|
||
|
||
# ADS-B emitter category (DO-260B). A3 airliners, A5 heavies, A7 helicopters.
|
||
_EMITTER = {
|
||
"A0": "unknown", "A1": "light", "A2": "small", "A3": "large",
|
||
"A4": "high vortex", "A5": "heavy", "A6": "high performance", "A7": "rotorcraft",
|
||
"B0": "unknown", "B1": "glider", "B2": "airship", "B3": "parachute",
|
||
"B4": "ultralight", "B6": "UAV", "B7": "space",
|
||
"C0": "ground unknown", "C1": "emergency vehicle", "C2": "service vehicle",
|
||
"D0": "unknown", "D1": "emergency",
|
||
}
|
||
|
||
# Combat / dedicated-military ICAO types. C-130/C-17 omitted — those also fly
|
||
# as fire tankers and civil contractors; dbFlags/hex catch the real mil ones.
|
||
_MIL_ICAO = frozenset({
|
||
"F15", "F16", "F18", "FA18", "F22", "F35", "F117", "A10", "A10A",
|
||
"B1", "B1B", "B2", "B52", "AV8B", "F4", "F5", "F14",
|
||
"SU27", "SU30", "SU34", "SU35", "SU57",
|
||
"MG29", "MIG29", "MG31", "MIG31", "J10", "J11", "J15", "J16", "J20",
|
||
"EUFI", "RFAL", "TOR", "E3TF", "E3CF", "E6", "E8", "P8",
|
||
"MQ9", "MQ1", "RQ4", "V22", "AH64", "H64",
|
||
})
|
||
|
||
# US DoD Mode-S block AE0000–AEFFFF.
|
||
_US_DOD_HEX_LO, _US_DOD_HEX_HI = 0xAE0000, 0xAEFFFF
|
||
|
||
_MIL_CS_PREFIX = ("RCH", "NAVY", "ARMY", "MARINE", "GOTOF", "REACH")
|
||
|
||
|
||
def classify_adsb(ac: dict) -> tuple[str, str]:
|
||
"""Return (role, role_src). Prefer readsb dbFlags bit0, then type/hex/cs."""
|
||
flags = ac.get("dbFlags")
|
||
try:
|
||
flags_i = int(flags) if flags is not None else 0
|
||
except (TypeError, ValueError):
|
||
flags_i = 0
|
||
if flags_i & 1:
|
||
return "military", "dbFlags"
|
||
icao = str(ac.get("t") or "").strip().upper()
|
||
if icao in _MIL_ICAO:
|
||
return "military", "type"
|
||
try:
|
||
hx = int(str(ac.get("hex") or "").strip(), 16)
|
||
except ValueError:
|
||
hx = -1
|
||
if _US_DOD_HEX_LO <= hx <= _US_DOD_HEX_HI:
|
||
return "military", "hex"
|
||
cs = str(ac.get("flight") or "").strip().upper()
|
||
if cs.startswith(_MIL_CS_PREFIX):
|
||
return "military", "callsign"
|
||
return "civilian", "default"
|
||
|
||
|
||
def _adsb_extra(ac: dict, hex_id: str) -> dict[str, Any]:
|
||
role, src = classify_adsb(ac)
|
||
cat = str(ac.get("category") or "").strip().upper()
|
||
extra: dict[str, Any] = {
|
||
"hex": hex_id,
|
||
"reg": _s(ac.get("r")),
|
||
"type": _s(ac.get("t")),
|
||
"alt_baro": ac.get("alt_baro"),
|
||
"squawk": _s(ac.get("squawk")),
|
||
"emergency": _s(ac.get("emergency")),
|
||
"category": cat or None,
|
||
"emitter": _EMITTER.get(cat),
|
||
"seen_pos": ac.get("seen_pos"),
|
||
"role": role,
|
||
"role_src": src,
|
||
"src": "adsb.lol",
|
||
}
|
||
desc = _s(ac.get("desc"))
|
||
if desc:
|
||
extra["desc"] = desc
|
||
own = _s(ac.get("ownOp") or ac.get("ownOpName") or ac.get("ownop"))
|
||
if own:
|
||
extra["ownOp"] = own
|
||
if ac.get("alt_geom") is not None:
|
||
extra["alt_geom"] = ac.get("alt_geom")
|
||
vs = ac.get("baro_rate")
|
||
if vs is None:
|
||
vs = ac.get("geom_rate")
|
||
if vs is not None:
|
||
extra["vs"] = vs
|
||
try:
|
||
raw_flags = ac.get("dbFlags")
|
||
flags_i = int(raw_flags) if raw_flags is not None else 0
|
||
except (TypeError, ValueError):
|
||
flags_i = 0
|
||
if flags_i:
|
||
extra["dbFlags"] = flags_i
|
||
extra["interesting"] = bool(flags_i & 2)
|
||
return extra
|
||
|
||
|
||
_NAVSTAT = {
|
||
0: "underway",
|
||
1: "at anchor",
|
||
2: "not under command",
|
||
3: "restricted manoeuvre",
|
||
4: "constrained by draught",
|
||
5: "moored",
|
||
6: "aground",
|
||
7: "fishing",
|
||
8: "sailing",
|
||
14: "AIS-SART",
|
||
15: "undefined",
|
||
}
|
||
|
||
# Compact MID → country for the flags that actually show up on AIS.
|
||
_MID_COUNTRY = {
|
||
211: "Germany", 218: "Germany",
|
||
219: "Denmark", 220: "Denmark",
|
||
224: "Spain", 225: "Spain",
|
||
226: "France", 227: "France", 228: "France",
|
||
232: "United Kingdom", 233: "United Kingdom", 234: "United Kingdom", 235: "United Kingdom",
|
||
236: "Gibraltar", 237: "Greece", 239: "Greece", 240: "Greece", 241: "Greece",
|
||
244: "Netherlands", 245: "Netherlands", 246: "Netherlands",
|
||
247: "Italy", 249: "Malta", 250: "Ireland", 251: "Iceland",
|
||
255: "Portugal", 256: "Malta",
|
||
257: "Norway", 258: "Norway", 259: "Norway",
|
||
261: "Poland", 263: "Portugal", 265: "Sweden", 266: "Sweden",
|
||
271: "Turkey", 273: "Russia", 276: "Estonia", 277: "Lithuania",
|
||
301: "Anguilla", 303: "United States", 310: "Bermuda", 316: "Canada",
|
||
319: "Cayman Islands", 338: "United States", 339: "Jamaica",
|
||
345: "Mexico", 352: "Panama", 353: "Panama", 354: "Panama",
|
||
355: "Panama", 356: "Panama", 357: "Panama",
|
||
366: "United States", 367: "United States", 368: "United States", 369: "United States",
|
||
370: "Panama", 371: "Panama", 372: "Panama", 373: "Panama", 374: "Panama",
|
||
375: "St Vincent", 376: "St Vincent", 377: "St Vincent",
|
||
412: "China", 413: "China", 414: "China", 416: "Taiwan",
|
||
419: "India", 431: "Japan", 432: "Japan", 440: "South Korea", 441: "South Korea",
|
||
477: "Hong Kong", 503: "Australia", 525: "Indonesia", 533: "Malaysia",
|
||
538: "Marshall Islands", 548: "Philippines", 563: "Singapore",
|
||
564: "Singapore", 565: "Singapore", 566: "Singapore", 567: "Thailand",
|
||
574: "Vietnam", 636: "Liberia", 637: "Liberia",
|
||
710: "Brazil", 725: "Chile", 730: "Colombia", 760: "Peru",
|
||
}
|
||
|
||
|
||
def _mmsi_country(mmsi: object) -> str | None:
|
||
digits = "".join(ch for ch in str(mmsi or "") if ch.isdigit())
|
||
if len(digits) < 3:
|
||
return None
|
||
try:
|
||
mid = int(digits[:3])
|
||
except ValueError:
|
||
return None
|
||
return _MID_COUNTRY.get(mid)
|
||
|
||
|
||
def classify_ais_type(type_code: int | None) -> tuple[str, str]:
|
||
"""Return (role, kind) from ITU-R M.1371 ship-and-cargo type."""
|
||
if type_code is None:
|
||
return "civilian", "unknown"
|
||
t = int(type_code)
|
||
tens = t // 10
|
||
if t == 35:
|
||
return "military", "military"
|
||
if t == 30:
|
||
return "civilian", "fishing"
|
||
if t in (31, 32, 52):
|
||
return "civilian", "tug"
|
||
if t == 33:
|
||
return "civilian", "dredger"
|
||
if t == 34:
|
||
return "civilian", "diving"
|
||
if t == 36:
|
||
return "civilian", "sailing"
|
||
if t == 37:
|
||
return "civilian", "pleasure"
|
||
if t == 50:
|
||
return "government", "pilot"
|
||
if t == 51:
|
||
return "government", "SAR"
|
||
if t == 55:
|
||
return "government", "law"
|
||
if t == 54:
|
||
return "government", "anti-pollution"
|
||
if t == 58:
|
||
return "government", "medical"
|
||
if tens == 4:
|
||
return "civilian", "HSC"
|
||
if tens == 6:
|
||
return "civilian", "passenger"
|
||
if tens == 7:
|
||
return "civilian", "cargo"
|
||
if tens == 8:
|
||
return "civilian", "tanker"
|
||
if tens in (5, 9) or t in (53, 56, 57, 59):
|
||
return "civilian", "special"
|
||
return "civilian", "other"
|
||
|
||
|
||
def transform_adsb_lol(payload: dict | list | None) -> list[dict]:
|
||
"""Map ADSB.lol v2 aircraft list to shared markers. Dedup on hex."""
|
||
if payload is None:
|
||
return []
|
||
if isinstance(payload, list):
|
||
aircraft = payload
|
||
else:
|
||
aircraft = payload.get("ac") or payload.get("aircraft") or []
|
||
seen: set[str] = set()
|
||
out: list[dict] = []
|
||
for ac in aircraft:
|
||
hex_id = str(ac.get("hex") or "").strip().lower()
|
||
lat, lon = _f(ac.get("lat")), _f(ac.get("lon"))
|
||
if not hex_id or lat is None or lon is None:
|
||
continue
|
||
if hex_id in seen:
|
||
continue
|
||
seen.add(hex_id)
|
||
flight = str(ac.get("flight") or "").strip() or hex_id
|
||
out.append(to_marker(
|
||
hex_id, lat, lon,
|
||
heading=_heading(ac.get("track")),
|
||
speed=_f(ac.get("gs")),
|
||
label=flight,
|
||
extra=_adsb_extra(ac, hex_id),
|
||
))
|
||
return out
|
||
|
||
|
||
def transform_amtraker(payload: dict | None) -> list[dict]:
|
||
"""Flatten Amtraker `{trainNum: [Train, ...]}` to one marker per trainID."""
|
||
if not payload or not isinstance(payload, dict):
|
||
return []
|
||
out: list[dict] = []
|
||
for _num, trains in payload.items():
|
||
if not isinstance(trains, list):
|
||
continue
|
||
for tr in trains:
|
||
if not isinstance(tr, dict):
|
||
continue
|
||
lat, lon = _f(tr.get("lat")), _f(tr.get("lon"))
|
||
if lat is None or lon is None:
|
||
continue
|
||
tid = str(tr.get("trainID") or tr.get("trainId") or "").strip()
|
||
tnum = str(tr.get("trainNum") or _num)
|
||
route = str(tr.get("routeName") or "").strip()
|
||
label = f"{route} #{tnum}".strip() if route else f"Train {tnum}"
|
||
late = tr.get("late")
|
||
if late is None:
|
||
late = tr.get("lateMin")
|
||
out.append(to_marker(
|
||
tid or tnum, lat, lon,
|
||
heading=_heading(tr.get("heading")),
|
||
speed=_f(tr.get("velocity") or tr.get("speed")),
|
||
label=label,
|
||
extra={
|
||
"route": route,
|
||
"trainNum": tnum,
|
||
"late_min": late,
|
||
"iconColor": tr.get("iconColor"),
|
||
"stations": tr.get("stations") or [],
|
||
"src": "amtraker",
|
||
},
|
||
))
|
||
return out
|
||
|
||
|
||
def transform_ais_frame(frame: dict | None) -> dict | None:
|
||
"""Decode one AISStream JSON envelope to a marker (or None if unusable)."""
|
||
if not frame or not isinstance(frame, dict):
|
||
return None
|
||
meta = frame.get("MetaData") or {}
|
||
mmsi = meta.get("MMSI") or meta.get("mmsi")
|
||
if mmsi is None:
|
||
return None
|
||
name = str(meta.get("ShipName") or meta.get("shipName") or "").strip()
|
||
lat = _f(meta.get("Latitude") if "Latitude" in meta else meta.get("latitude"))
|
||
lon = _f(meta.get("Longitude") if "Longitude" in meta else meta.get("longitude"))
|
||
msg = frame.get("Message") or {}
|
||
pos = (
|
||
msg.get("PositionReport")
|
||
or msg.get("StandardClassBPositionReport")
|
||
or msg.get("ExtendedClassBPositionReport")
|
||
or {}
|
||
)
|
||
extra: dict[str, Any] = {"src": "aisstream", "mmsi": mmsi}
|
||
country = _mmsi_country(mmsi)
|
||
if country:
|
||
extra["country"] = country
|
||
if frame.get("MessageType") == "ShipStaticData":
|
||
static = msg.get("ShipStaticData") or {}
|
||
dest = str(static.get("Destination") or static.get("destination") or "").strip()
|
||
extra["dest"] = dest
|
||
extra["static"] = True
|
||
if not name:
|
||
name = str(static.get("Name") or static.get("name") or "").strip()
|
||
cs = _s(static.get("CallSign") or static.get("callSign"))
|
||
if cs:
|
||
extra["callsign"] = cs
|
||
try:
|
||
imo = int(static.get("ImoNumber") or static.get("imoNumber") or 0)
|
||
except (TypeError, ValueError):
|
||
imo = 0
|
||
if imo:
|
||
extra["imo"] = imo
|
||
type_code = static.get("Type") if "Type" in static else static.get("type")
|
||
try:
|
||
type_i = int(type_code) if type_code is not None else None
|
||
except (TypeError, ValueError):
|
||
type_i = None
|
||
if type_i is not None:
|
||
extra["type_code"] = type_i
|
||
role, kind = classify_ais_type(type_i)
|
||
extra["role"] = role
|
||
extra["kind"] = kind
|
||
dim = static.get("Dimension") or static.get("dimension") or {}
|
||
if isinstance(dim, dict):
|
||
a, b = _f(dim.get("A")), _f(dim.get("B"))
|
||
c, d = _f(dim.get("C")), _f(dim.get("D"))
|
||
if a is not None and b is not None:
|
||
extra["length"] = int(round(a + b))
|
||
if c is not None and d is not None:
|
||
extra["beam"] = int(round(c + d))
|
||
draught = _f(static.get("MaximumStaticDraught") or static.get("maximumStaticDraught"))
|
||
if draught is not None:
|
||
extra["draught"] = draught
|
||
eta = static.get("Eta") or static.get("eta") or {}
|
||
if isinstance(eta, dict) and eta.get("Month"):
|
||
extra["eta"] = (
|
||
f"{int(eta.get('Month') or 0):02d}-{int(eta.get('Day') or 0):02d} "
|
||
f"{int(eta.get('Hour') or 0):02d}:{int(eta.get('Minute') or 0):02d}"
|
||
)
|
||
if lat is None or lon is None:
|
||
# Static-only update: caller merges onto last-known by MMSI.
|
||
return to_marker(str(mmsi), None, None, label=name or str(mmsi), extra=extra)
|
||
heading = pos.get("TrueHeading")
|
||
if heading in (511, 511.0, None):
|
||
heading = pos.get("Cog")
|
||
sog = pos.get("Sog")
|
||
navstat = pos.get("NavigationalStatus")
|
||
extra["navstat"] = navstat
|
||
try:
|
||
extra["nav"] = _NAVSTAT.get(int(navstat)) if navstat is not None else None
|
||
except (TypeError, ValueError):
|
||
extra["nav"] = None
|
||
extra["cog"] = pos.get("Cog")
|
||
if not extra.get("dest"):
|
||
extra.pop("dest", None)
|
||
if lat is None or lon is None:
|
||
return None
|
||
return to_marker(
|
||
str(mmsi), lat, lon,
|
||
heading=_heading(heading),
|
||
speed=_f(sog),
|
||
label=name or str(mmsi),
|
||
extra=extra,
|
||
)
|
||
|
||
|
||
def transform_wfigs_incidents(fc: dict | None) -> list[dict]:
|
||
if not fc:
|
||
return []
|
||
out: list[dict] = []
|
||
for feat in fc.get("features") or []:
|
||
geom = feat.get("geometry") or {}
|
||
coords = geom.get("coordinates") or []
|
||
if len(coords) < 2:
|
||
continue
|
||
lon, lat = _f(coords[0]), _f(coords[1])
|
||
if lat is None or lon is None:
|
||
continue
|
||
props = feat.get("properties") or {}
|
||
name = str(props.get("IncidentName") or "Incident").strip()
|
||
out.append(to_marker(
|
||
name, lat, lon,
|
||
label=name,
|
||
extra={
|
||
"acres": props.get("IncidentSize"),
|
||
"contained": props.get("PercentContained"),
|
||
"state": props.get("POOState"),
|
||
"category": props.get("IncidentTypeCategory"),
|
||
"cause": props.get("FireCause"),
|
||
"discovered": props.get("FireDiscoveryDateTime"),
|
||
"src": "wfigs",
|
||
},
|
||
))
|
||
return out
|
||
|
||
|
||
def transform_nhc_storms(payload: dict | None) -> list[dict]:
|
||
if not payload:
|
||
return []
|
||
storms = payload.get("activeStorms") or []
|
||
class_label = {
|
||
"TD": "Tropical Depression",
|
||
"TS": "Tropical Storm",
|
||
"HU": "Hurricane",
|
||
"STD": "Subtropical Depression",
|
||
"STS": "Subtropical Storm",
|
||
"PTC": "Potential Tropical Cyclone",
|
||
"PC": "Post-Tropical Cyclone",
|
||
}
|
||
out: list[dict] = []
|
||
for s in storms:
|
||
lat = _f(s.get("latitudeNumeric") or s.get("lat"))
|
||
lon = _f(s.get("longitudeNumeric") or s.get("lon"))
|
||
if lat is None or lon is None:
|
||
continue
|
||
sid = str(s.get("id") or s.get("binNumber") or "").strip()
|
||
name = str(s.get("name") or "Storm").strip()
|
||
klass = str(s.get("classification") or "").upper()
|
||
kind = class_label.get(klass, klass or "Storm")
|
||
out.append(to_marker(
|
||
sid or name, lat, lon,
|
||
heading=_heading(s.get("movementDir")),
|
||
speed=_f(s.get("movementSpeed")),
|
||
label=f"{kind} {name}".strip(),
|
||
extra={
|
||
"classification": klass,
|
||
"intensity_kt": s.get("intensity"),
|
||
"src": "nhc",
|
||
},
|
||
))
|
||
return out
|
||
|
||
|
||
def rainviewer_tile_url(host: str, path: str, size: int = 256, color: int = 2,
|
||
smooth: int = 1, snow: int = 1) -> str:
|
||
host = host.rstrip("/")
|
||
path = path if path.startswith("/") else f"/{path}"
|
||
return f"{host}{path}/{size}/{{z}}/{{x}}/{{y}}/{color}/{smooth}_{snow}.png"
|
||
|
||
|
||
def _headers() -> dict[str, str]:
|
||
return {"User-Agent": OSINT_USER_AGENT, "Accept": "application/json"}
|
||
|
||
|
||
def _planespotters_headers() -> dict[str, str]:
|
||
"""Planespotters 403s server clients whose UA has no contact URL/email."""
|
||
ua = (OSINT_USER_AGENT or "").strip()
|
||
if ua and ("@" in ua or "http://" in ua.lower() or "https://" in ua.lower()):
|
||
return {"User-Agent": ua, "Accept": "application/json"}
|
||
base = ua or "osint-dashboard/1.0"
|
||
return {
|
||
"User-Agent": f"{base} (lancewalters94@gmail.com)",
|
||
"Accept": "application/json",
|
||
}
|
||
|
||
|
||
async def _lock_for(key: str) -> asyncio.Lock:
|
||
async with _key_locks_guard:
|
||
lock = _key_locks.get(key)
|
||
if lock is None:
|
||
lock = asyncio.Lock()
|
||
_key_locks[key] = lock
|
||
return lock
|
||
|
||
|
||
async def _ttl_get(key: str, ttl: float, factory: Callable[[], Awaitable[Any]]) -> Any:
|
||
now = time.monotonic()
|
||
hit = _cache.get(key)
|
||
if hit and now - hit[0] < ttl:
|
||
return hit[1]
|
||
lock = await _lock_for(key)
|
||
async with lock:
|
||
hit = _cache.get(key)
|
||
if hit and time.monotonic() - hit[0] < ttl:
|
||
return hit[1]
|
||
value = await factory()
|
||
_cache[key] = (time.monotonic(), value)
|
||
return value
|
||
|
||
|
||
async def init_http() -> None:
|
||
"""Shared outbound client — one TLS pool for all overlay upstreams."""
|
||
global _http
|
||
if _http is None:
|
||
_http = httpx.AsyncClient(
|
||
timeout=_HTTP_TIMEOUT,
|
||
follow_redirects=True,
|
||
headers=_headers(),
|
||
limits=httpx.Limits(max_connections=20, max_keepalive_connections=10),
|
||
)
|
||
|
||
|
||
async def close_http() -> None:
|
||
global _http
|
||
if _http is not None:
|
||
await _http.aclose()
|
||
_http = None
|
||
|
||
|
||
async def _get_json(
|
||
url: str, params: dict | None = None, headers: dict | None = None,
|
||
) -> Any:
|
||
if _http is None:
|
||
async with httpx.AsyncClient(
|
||
timeout=_HTTP_TIMEOUT, follow_redirects=True, headers=_headers(),
|
||
) as client:
|
||
resp = await client.get(url, params=params, headers=headers)
|
||
resp.raise_for_status()
|
||
return resp.json()
|
||
resp = await _http.get(url, params=params, headers=headers)
|
||
resp.raise_for_status()
|
||
return resp.json()
|
||
|
||
|
||
def _normalize_planespotter_photo(photo: dict) -> dict | None:
|
||
"""Slim a planespotters.net photo object for the aircraft popup.
|
||
|
||
Image binaries are never proxied/re-hosted: we return the CDN ``src``
|
||
and the browser loads it directly (per ToS).
|
||
"""
|
||
if not isinstance(photo, dict):
|
||
return None
|
||
large = photo.get("thumbnail_large") or {}
|
||
small = photo.get("thumbnail") or {}
|
||
src = large.get("src") or small.get("src")
|
||
if not src:
|
||
return None
|
||
size = large.get("size") or small.get("size") or {}
|
||
return {
|
||
"id": str(photo.get("id") or ""),
|
||
"src": src,
|
||
"width": size.get("width"),
|
||
"height": size.get("height"),
|
||
"link": photo.get("link") or "",
|
||
"photographer": photo.get("photographer") or "",
|
||
}
|
||
|
||
|
||
async def fetch_planespotters_photo(
|
||
hex_code: str | None = None, reg: str | None = None,
|
||
) -> dict | None:
|
||
"""Latest photo for an aircraft from planespotters.net (hex preferred).
|
||
|
||
Server-side only: planespotters 403s any request carrying an ``Origin``
|
||
header (browser fetch() always sends one), so the browser can never reach
|
||
it directly. We fetch here with the identifying UA, cache the JSON ≤24h,
|
||
and hand the CDN image URL back for the browser to load.
|
||
"""
|
||
if hex_code:
|
||
key = f"psp-hex:{hex_code.strip().lower()}"
|
||
url = f"{PLANESPOTTERS_PHOTO}/hex/{quote(hex_code.strip())}"
|
||
elif reg:
|
||
key = f"psp-reg:{reg.strip().lower()}"
|
||
url = f"{PLANESPOTTERS_PHOTO}/reg/{quote(reg.strip())}"
|
||
else:
|
||
return None
|
||
|
||
async def _load() -> dict | None:
|
||
data = await _get_json(url, headers=_planespotters_headers())
|
||
photos = data.get("photos") or []
|
||
return _normalize_planespotter_photo(photos[0]) if photos else None
|
||
|
||
return await _ttl_get(key, PLANESPOTTERS_CACHE_TTL, _load)
|
||
|
||
|
||
async def fetch_aircraft(
|
||
bbox: str, limit: int = DEFAULT_LIMIT, persist: bool = False,
|
||
) -> list[dict]:
|
||
"""Viewport ADS-B.
|
||
|
||
GET path (persist=False) serves in-memory last-known and never writes
|
||
tracks/geofences. Background refresh (persist=True) hits ADSB.lol and
|
||
then persist_aircraft_snapshot.
|
||
"""
|
||
minlon, minlat, maxlon, maxlat = parse_bbox(bbox)
|
||
if not persist:
|
||
cached = [
|
||
dict(v) for v in aircraft_last_known.values()
|
||
if v.get("lat") is not None and v.get("lon") is not None
|
||
]
|
||
filtered = filter_points_bbox(cached, minlon, minlat, maxlon, maxlat, limit)
|
||
if filtered:
|
||
return filtered
|
||
qminlon, qminlat, qmaxlon, qmaxlat = quantize_bbox(minlon, minlat, maxlon, maxlat)
|
||
lat, lon, radius = bbox_center_radius_nm(qminlon, qminlat, qmaxlon, qmaxlat)
|
||
cache_key = f"adsb:{lat:.2f}:{lon:.2f}:{radius}"
|
||
|
||
async def _load():
|
||
url = f"{ADSB_LOL_BASE}/v2/lat/{lat:.4f}/lon/{lon:.4f}/dist/{radius}"
|
||
return transform_adsb_lol(await _get_json(url))
|
||
|
||
rows = await _ttl_get(cache_key, 8.0, _load)
|
||
for m in rows:
|
||
aircraft_last_known[str(m.get("id"))] = m
|
||
if persist:
|
||
await persist_aircraft_snapshot(rows)
|
||
return filter_points_bbox(rows, minlon, minlat, maxlon, maxlat, limit)
|
||
|
||
|
||
async def persist_aircraft_snapshot(rows: list[dict]) -> None:
|
||
"""Track / geofence / WS / fire correlation — never on the GET path."""
|
||
from ws_manager import manager
|
||
from tracks import record_position
|
||
from geofence import record_and_notify
|
||
|
||
for m in rows:
|
||
mlat, mlon = m.get("lat"), m.get("lon")
|
||
if mlat is None or mlon is None:
|
||
continue
|
||
await record_position("aircraft", m)
|
||
if manager.has_clients():
|
||
await manager.publish_point("adsb", m, lat=mlat, lon=mlon)
|
||
await record_and_notify(
|
||
source_kind="adsb", entity_id=str(m.get("id")),
|
||
lat=mlat, lon=mlon, payload=m,
|
||
)
|
||
if fire_last_known:
|
||
from fire_aircraft import correlate_and_notify
|
||
await correlate_and_notify(fire_last_known, rows)
|
||
|
||
|
||
async def fetch_trains(bbox: str | None, limit: int = DEFAULT_LIMIT) -> list[dict]:
|
||
async def _load():
|
||
return transform_amtraker(await _get_json(AMTRAKER_TRAINS))
|
||
|
||
rows = await _ttl_get("amtraker:trains", 20.0, _load)
|
||
if bbox:
|
||
minlon, minlat, maxlon, maxlat = parse_bbox(bbox)
|
||
return filter_points_bbox(rows, minlon, minlat, maxlon, maxlat, limit)
|
||
return rows[:limit]
|
||
|
||
|
||
async def fetch_vessels(bbox: str | None, limit: int = DEFAULT_LIMIT) -> list[dict]:
|
||
async with vessel_lock:
|
||
rows = [dict(v) for v in vessel_last_known.values()
|
||
if v.get("lat") is not None and v.get("lon") is not None]
|
||
if bbox:
|
||
minlon, minlat, maxlon, maxlat = parse_bbox(bbox)
|
||
return filter_points_bbox(rows, minlon, minlat, maxlon, maxlat, limit)
|
||
return rows[:limit]
|
||
|
||
|
||
async def upsert_vessel(marker: dict) -> None:
|
||
"""Merge an AIS marker into last-known by MMSI. Static-only updates names."""
|
||
vid = str(marker.get("id") or "")
|
||
if not vid:
|
||
return
|
||
async with vessel_lock:
|
||
prev = vessel_last_known.get(vid, {})
|
||
extra = {**(prev.get("extra") or {}), **(marker.get("extra") or {})}
|
||
lat = marker.get("lat") if marker.get("lat") is not None else prev.get("lat")
|
||
lon = marker.get("lon") if marker.get("lon") is not None else prev.get("lon")
|
||
label = marker.get("label")
|
||
if not label or label == vid:
|
||
label = prev.get("label") or vid
|
||
stored = {
|
||
**to_marker(
|
||
vid, lat, lon,
|
||
heading=marker.get("heading") if marker.get("heading") is not None else prev.get("heading"),
|
||
speed=marker.get("speed") if marker.get("speed") is not None else prev.get("speed"),
|
||
label=label,
|
||
extra=extra,
|
||
),
|
||
"seen_at": datetime.now(timezone.utc).isoformat(),
|
||
}
|
||
vessel_last_known[vid] = stored
|
||
if len(vessel_last_known) > _MAX_VESSELS:
|
||
excess = len(vessel_last_known) - int(_MAX_VESSELS * 0.9)
|
||
oldest = sorted(
|
||
vessel_last_known,
|
||
key=lambda k: vessel_last_known[k].get("seen_at", ""),
|
||
)[:excess]
|
||
for k in oldest:
|
||
vessel_last_known.pop(k, None)
|
||
if lat is not None and lon is not None:
|
||
from ws_manager import manager
|
||
from tracks import record_position
|
||
from geofence import record_and_notify
|
||
await manager.publish_point("ais", stored, lat=lat, lon=lon)
|
||
await record_position("vessel", stored)
|
||
await record_and_notify(
|
||
source_kind="ais", entity_id=vid, lat=lat, lon=lon, payload=stored,
|
||
)
|
||
|
||
|
||
def _wfigs_params(bbox: str | None, *, offset_m: float = 250.0) -> dict:
|
||
params = {
|
||
"where": "1=1",
|
||
"outSR": "4326",
|
||
"f": "geojson",
|
||
"resultRecordCount": 500,
|
||
"maxAllowableOffset": offset_m / 111_320.0, # metres → degrees
|
||
"geometryPrecision": 5,
|
||
}
|
||
if bbox:
|
||
minlon, minlat, maxlon, maxlat = quantize_bbox(*parse_bbox(bbox))
|
||
params["geometry"] = f"{minlon},{minlat},{maxlon},{maxlat}"
|
||
params["geometryType"] = "esriGeometryEnvelope"
|
||
params["inSR"] = "4326"
|
||
params["spatialRel"] = "esriSpatialRelIntersects"
|
||
return params
|
||
|
||
|
||
async def fetch_fire_incidents(bbox: str | None, limit: int = DEFAULT_LIMIT) -> list[dict]:
|
||
params = _wfigs_params(bbox)
|
||
params["outFields"] = (
|
||
"IncidentName,IncidentSize,FireDiscoveryDateTime,POOState,"
|
||
"PercentContained,IncidentTypeCategory,FireCause"
|
||
)
|
||
|
||
async def _load():
|
||
return transform_wfigs_incidents(await _get_json(WFIGS_INCIDENTS, params))
|
||
|
||
rows = await _ttl_get(f"wfigs:inc:{bbox_cell_key(bbox)}", 600.0, _load)
|
||
fire_last_known[:] = list(rows)
|
||
if rows and aircraft_last_known:
|
||
from fire_aircraft import correlate_and_notify
|
||
await correlate_and_notify(rows, list(aircraft_last_known.values()))
|
||
return rows[:limit]
|
||
|
||
|
||
async def fetch_fire_perimeters(bbox: str | None) -> dict:
|
||
params = _wfigs_params(bbox)
|
||
params["outFields"] = (
|
||
"poly_IncidentName,poly_GISAcres,attr_IncidentSize,"
|
||
"attr_PercentContained,attr_FireDiscoveryDateTime"
|
||
)
|
||
|
||
async def _load():
|
||
return await _get_json(WFIGS_PERIMETERS, params)
|
||
|
||
fc = await _ttl_get(f"wfigs:per:{bbox_cell_key(bbox)}", 600.0, _load)
|
||
if not isinstance(fc, dict):
|
||
return {"type": "FeatureCollection", "features": []}
|
||
return fc
|
||
|
||
|
||
async def fetch_weather_alerts(area: str | None, bbox: str | None) -> dict:
|
||
"""Cached NWS active alerts + IEM storm-based warning polygons."""
|
||
|
||
async def _load_iem():
|
||
try:
|
||
data = await _get_json(IEM_SBW)
|
||
except Exception:
|
||
return {"type": "FeatureCollection", "features": []}
|
||
return data if isinstance(data, dict) else {"type": "FeatureCollection", "features": []}
|
||
|
||
async def _load():
|
||
nws_params: dict[str, str] = {"status": "actual"}
|
||
clip_box = None
|
||
if area:
|
||
nws_params["area"] = area.upper()
|
||
if bbox:
|
||
clip_box = quantize_bbox(*parse_bbox(bbox))
|
||
# NWS /alerts/active 400s on bbox= — never send it; clip locally.
|
||
nws_fc: dict = {"features": []}
|
||
nws_ok = True
|
||
try:
|
||
nws_fc = await _get_json(NWS_ALERTS, nws_params)
|
||
except Exception as exc:
|
||
logger.warning("NWS alerts fetch failed: %s", exc)
|
||
nws_ok = False
|
||
nws_fc = {"features": []}
|
||
sbw_fc = await _ttl_get("iem:sbw", 45.0, _load_iem)
|
||
features = []
|
||
for feat in nws_fc.get("features") or []:
|
||
if not isinstance(feat, dict):
|
||
continue
|
||
props = dict(feat.get("properties") or {})
|
||
props["source"] = "nws"
|
||
features.append({**feat, "properties": slim_alert_properties(props)})
|
||
for feat in (sbw_fc or {}).get("features") or []:
|
||
if not isinstance(feat, dict):
|
||
continue
|
||
props = dict(feat.get("properties") or {})
|
||
props["source"] = "iem-sbw"
|
||
if "event" not in props:
|
||
props["event"] = props.get("ps") or "Storm-based warning"
|
||
features.append({**feat, "properties": slim_alert_properties(props)})
|
||
merged = {"type": "FeatureCollection", "features": features, "nws_ok": nws_ok}
|
||
if clip_box:
|
||
clipped = clip_fc_to_bbox(merged, *clip_box)
|
||
clipped["nws_ok"] = nws_ok
|
||
merged = clipped
|
||
elif bbox:
|
||
clipped = clip_fc_to_bbox(merged, *parse_bbox(bbox))
|
||
clipped["nws_ok"] = nws_ok
|
||
merged = clipped
|
||
return merged
|
||
|
||
key = f"alerts:{area or ''}:{bbox_cell_key(bbox) if bbox else ''}"
|
||
return await _ttl_get(key, 30.0, _load)
|
||
|
||
|
||
async def fetch_radar_meta() -> dict:
|
||
async def _load():
|
||
data = await _get_json(RAINVIEWER_MAPS)
|
||
host = data.get("host") or "https://tilecache.rainviewer.com"
|
||
radar = (data.get("radar") or {})
|
||
past = radar.get("past") or []
|
||
nowcast = radar.get("nowcast") or []
|
||
frames = []
|
||
for fr in past + nowcast:
|
||
path = fr.get("path")
|
||
if not path:
|
||
continue
|
||
frames.append({
|
||
"time": fr.get("time"),
|
||
"path": path,
|
||
"tileUrl": rainviewer_tile_url(host, path),
|
||
})
|
||
latest = frames[-1]["tileUrl"] if frames else None
|
||
return {
|
||
"provider": "rainviewer",
|
||
"host": host,
|
||
"tileUrl": latest,
|
||
"frames": frames,
|
||
"attribution": "Weather data by RainViewer",
|
||
"iemTileUrl": IEM_NEXRAD,
|
||
}
|
||
|
||
return await _ttl_get("radar:rv", 300.0, _load)
|
||
|
||
|
||
async def fetch_storms() -> list[dict]:
|
||
async def _load():
|
||
return transform_nhc_storms(await _get_json(NHC_STORMS))
|
||
|
||
return await _ttl_get("nhc:storms", 300.0, _load)
|
||
|
||
|
||
# ── Sentinel-1 SAR (Planetary Computer STAC → signed COG tile template) ────
|
||
|
||
class UpstreamRateLimited(Exception):
|
||
"""Planetary Computer returned 429. Carries Retry-After for the client."""
|
||
|
||
def __init__(self, retry_after: str | None = None):
|
||
self.retry_after = retry_after
|
||
super().__init__("planetary computer rate limited")
|
||
|
||
|
||
async def _post_json(
|
||
url: str, json: dict | None = None, headers: dict | None = None,
|
||
) -> Any:
|
||
if _http is None:
|
||
async with httpx.AsyncClient(
|
||
timeout=_HTTP_TIMEOUT, follow_redirects=True, headers=_headers(),
|
||
) as client:
|
||
resp = await client.post(url, json=json, headers=headers)
|
||
resp.raise_for_status()
|
||
return resp.json()
|
||
resp = await _http.post(url, json=json, headers=headers)
|
||
resp.raise_for_status()
|
||
return resp.json()
|
||
|
||
|
||
async def _pc_call(coro: Awaitable[Any]) -> Any:
|
||
"""Run a Planetary Computer call, mapping 429 → UpstreamRateLimited."""
|
||
try:
|
||
return await coro
|
||
except httpx.HTTPStatusError as exc:
|
||
if exc.response.status_code == 429:
|
||
raise UpstreamRateLimited(
|
||
exc.response.headers.get("Retry-After")
|
||
) from exc
|
||
raise
|
||
|
||
|
||
def sign_cog_url(href: str, token: str) -> str:
|
||
"""Append a SAS token to a PC blob URL (respect existing query string)."""
|
||
sep = "&" if "?" in href else "?"
|
||
return f"{href}{sep}token={quote(token, safe='')}"
|
||
|
||
|
||
def sentinel1_tile_url(signed_cog: str) -> str:
|
||
"""TiTiler XYZ template for a signed COG (Leaflet substitutes {z}/{x}/{y})."""
|
||
params = urlencode({
|
||
"url": signed_cog,
|
||
"rescale": "0,500",
|
||
"colormap_name": "cfastie",
|
||
})
|
||
return f"{TITILER_COG_TILES}?{params}"
|
||
|
||
|
||
async def fetch_sentinel1(bbox: str) -> dict | None:
|
||
"""Most recent Sentinel-1 GRD COG for a viewport, signed and TiTiler-ready.
|
||
|
||
Returns the overlay tile-template dict, or ``None`` when no GRD imagery
|
||
covers the bbox in the last 7 days (caller maps to 404). Queries Planetary
|
||
Computer only when called; cached per quantized bbox + UTC day.
|
||
"""
|
||
minlon, minlat, maxlon, maxlat = parse_bbox(bbox)
|
||
day = datetime.now(timezone.utc).date().isoformat()
|
||
key = f"sentinel1:{day}:{bbox_cell_key(bbox)}"
|
||
|
||
async def _load() -> dict | None:
|
||
now = datetime.now(timezone.utc)
|
||
week_ago = now - timedelta(days=7)
|
||
payload = {
|
||
"collections": ["sentinel-1-grd"],
|
||
"bbox": [minlon, minlat, maxlon, maxlat],
|
||
"datetime": f"{week_ago.isoformat()}/{now.isoformat()}",
|
||
"limit": 1,
|
||
"sortby": [{"field": "datetime", "direction": "desc"}],
|
||
}
|
||
data = await _pc_call(_post_json(PC_STAC_SEARCH, json=payload))
|
||
features = data.get("features") or []
|
||
if not features:
|
||
return None
|
||
|
||
feat = features[0]
|
||
assets = feat.get("assets") or {}
|
||
chosen_href: str | None = None
|
||
polarization: str | None = None
|
||
for pol in ("vv", "hh"):
|
||
href = (assets.get(pol) or {}).get("href")
|
||
if href:
|
||
chosen_href = href
|
||
polarization = pol
|
||
break
|
||
if not chosen_href:
|
||
return None
|
||
|
||
sas = await _pc_call(_get_json(PC_SAS_TOKEN))
|
||
token = (sas or {}).get("token")
|
||
if not token:
|
||
raise RuntimeError("planetarycomputer SAS token missing")
|
||
signed = sign_cog_url(chosen_href, token)
|
||
|
||
props = feat.get("properties") or {}
|
||
return {
|
||
"id": "sentinel-1-sar",
|
||
"kind": "raster",
|
||
"tileUrl": sentinel1_tile_url(signed),
|
||
"opacity": 0.8,
|
||
"itemId": feat.get("id"),
|
||
"datetime": props.get("datetime") or feat.get("datetime"),
|
||
"polarization": polarization,
|
||
"attribution": SENTINEL1_ATTRIBUTION,
|
||
}
|
||
|
||
return await _ttl_get(key, float(SENTINEL1_TTL), _load)
|