osint-dashboard/app/geofence.py
Sirius DevOps f6c1cfc454 feat(geofence): watch list, filtered hit log, fence snapshot
Off-viewport clients that send {type:watch_geofences,ids} on /ws/live
still receive geofence_alert; AIS/ADS-B stay viewport-only.

GET /api/geofence-alerts accepts geofence_id/since/until/source_kind.
GET /api/geofences/{id}/at returns CAGG+FIRMS inside the fence at T
(404 if missing, empty lists if DB down). DELETE missing fences 404s.
2026-09-01 01:24:36 -04:00

472 lines
15 KiB
Python

"""Geofences: GeoJSON polygons, ST_Intersects on ingest, WS alerts.
``/api/alerts`` is the dashboard entity/keyword table — geofence hits live
in ``geofence_alerts`` and fan out as WS type ``geofence_alert``.
"""
from __future__ import annotations
import json
from datetime import datetime, timedelta, timezone
from typing import Any
from uuid import uuid4
from sqlalchemy import text
from database import async_session
def _rings_from_geojson(geojson: dict) -> list[list[list[float]]]:
if not isinstance(geojson, dict):
raise ValueError("geojson must be an object")
gj = geojson
if gj.get("type") == "Feature":
gj = gj.get("geometry") or {}
if gj.get("type") == "FeatureCollection":
raise ValueError("FeatureCollection is not a single polygon")
if gj.get("type") != "Polygon":
raise ValueError("geojson must be a Polygon")
coords = gj.get("coordinates")
if not isinstance(coords, list) or not coords:
raise ValueError("polygon has no rings")
rings: list[list[list[float]]] = []
for ring in coords:
if not isinstance(ring, list) or len(ring) < 4:
raise ValueError("polygon ring needs ≥4 positions (closed)")
pts = []
for pt in ring:
if not isinstance(pt, (list, tuple)) or len(pt) < 2:
raise ValueError("position must be [lon, lat]")
pts.append([float(pt[0]), float(pt[1])])
rings.append(pts)
return rings
def validate_polygon_geojson(geojson: dict) -> dict:
"""Return a canonical Polygon GeoJSON or raise ValueError."""
rings = _rings_from_geojson(geojson)
return {"type": "Polygon", "coordinates": rings}
def _ring_contains(lon: float, lat: float, ring: list[list[float]]) -> bool:
"""Ray-cast even-odd rule. Ring is [lon, lat] positions."""
inside = False
n = len(ring)
if n < 4:
return False
j = n - 1
for i in range(n):
xi, yi = ring[i][0], ring[i][1]
xj, yj = ring[j][0], ring[j][1]
intersects = ((yi > lat) != (yj > lat)) and (
lon < (xj - xi) * (lat - yi) / ((yj - yi) or 1e-16) + xi
)
if intersects:
inside = not inside
j = i
return inside
def point_in_geojson(lon: float, lat: float, geojson: dict) -> bool:
"""True if (lon, lat) is inside the outer ring and outside holes."""
try:
rings = _rings_from_geojson(geojson)
except (ValueError, TypeError, KeyError):
return False
if not _ring_contains(lon, lat, rings[0]):
return False
for hole in rings[1:]:
if _ring_contains(lon, lat, hole):
return False
return True
def matching_geofences(lon: float, lat: float, fences: list[dict]) -> list[dict]:
hits = []
for fence in fences:
if not fence.get("active", True):
continue
gj = fence.get("geojson") or {}
if point_in_geojson(lon, lat, gj):
hits.append(fence)
return hits
# In-process copy of active fences so ingest does not round-trip Postgres
# on every AIS frame. CRUD endpoints refresh this list.
_cache: list[dict] = []
_recent_hits: dict[tuple[str, str], datetime] = {}
_HIT_COOLDOWN = timedelta(minutes=5)
async def refresh_cache() -> list[dict]:
global _cache
async with async_session() as session:
rows = (await session.execute(text(
"SELECT id::text, name, geojson, active FROM geofences"
))).mappings().all()
_cache = [
{
"id": r["id"],
"name": r["name"],
"geojson": r["geojson"] if isinstance(r["geojson"], dict)
else json.loads(r["geojson"] or "{}"),
"active": bool(r["active"]),
}
for r in rows
]
return _cache
def cached_fences() -> list[dict]:
return list(_cache)
async def list_geofences() -> list[dict]:
if not _cache:
try:
await refresh_cache()
except Exception:
return []
return cached_fences()
async def create_geofence(name: str, geojson: dict, active: bool = True) -> dict:
polygon = validate_polygon_geojson(geojson)
gid = str(uuid4())
gj = json.dumps(polygon)
async with async_session() as session:
await session.execute(
text(
"""
INSERT INTO geofences (id, name, geojson, geom, active)
VALUES (
:id, :name, CAST(:geojson AS jsonb),
ST_SetSRID(ST_GeomFromGeoJSON(:geojson), 4326),
:active
)
"""
),
{"id": gid, "name": name, "geojson": gj, "active": 1 if active else 0},
)
await session.commit()
row = {"id": gid, "name": name, "geojson": polygon, "active": active}
_cache.append(row)
return row
async def update_geofence(gid: str, *, name: str | None = None,
geojson: dict | None = None,
active: bool | None = None) -> dict | None:
current = next((f for f in _cache if f["id"] == gid), None)
if current is None:
await refresh_cache()
current = next((f for f in _cache if f["id"] == gid), None)
if current is None:
return None
if name is not None:
current["name"] = name
if geojson is not None:
current["geojson"] = validate_polygon_geojson(geojson)
if active is not None:
current["active"] = active
gj = json.dumps(current["geojson"])
async with async_session() as session:
await session.execute(
text(
"""
UPDATE geofences SET
name = :name,
geojson = CAST(:geojson AS jsonb),
geom = ST_SetSRID(ST_GeomFromGeoJSON(:geojson), 4326),
active = :active,
updated_at = now()
WHERE id = CAST(:id AS uuid)
"""
),
{
"id": gid,
"name": current["name"],
"geojson": gj,
"active": 1 if current["active"] else 0,
},
)
await session.commit()
return current
async def delete_geofence(gid: str) -> bool:
async with async_session() as session:
result = await session.execute(
text("DELETE FROM geofences WHERE id = CAST(:id AS uuid)"),
{"id": gid},
)
await session.commit()
_cache[:] = [f for f in _cache if f["id"] != gid]
return bool(result.rowcount)
async def st_intersects(lon: float, lat: float) -> list[dict]:
"""PostGIS ST_Intersects against active geofences.
Falls back to the in-memory GeoJSON test if the DB is unreachable so
ingest never dies because a fence check failed.
"""
try:
async with async_session() as session:
rows = (await session.execute(
text(
"""
SELECT id::text, name, geojson, active
FROM geofences
WHERE active = 1
AND ST_Intersects(
geom,
ST_SetSRID(ST_MakePoint(:lon, :lat), 4326)
)
"""
),
{"lon": lon, "lat": lat},
)).mappings().all()
return [
{
"id": r["id"],
"name": r["name"],
"geojson": r["geojson"] if isinstance(r["geojson"], dict)
else json.loads(r["geojson"] or "{}"),
"active": True,
}
for r in rows
]
except Exception:
return matching_geofences(lon, lat, cached_fences())
async def record_and_notify(
*,
source_kind: str,
entity_id: str,
lat: float,
lon: float,
payload: dict[str, Any] | None = None,
) -> int:
"""Insert a geofence_alerts row per hit and WS-push to viewport clients.
PostGIS ST_Intersects is the source of truth. The in-process GeoJSON
cache is not a reject filter — the FIRMS ingester never fills it.
"""
hits = await st_intersects(lon, lat)
if not hits:
return 0
from ws_manager import manager
sent = 0
now = datetime.now(timezone.utc)
fresh = []
for fence in hits:
key = (str(fence["id"]), str(entity_id))
prev = _recent_hits.get(key)
if prev is not None and now - prev < _HIT_COOLDOWN:
continue
_recent_hits[key] = now
fresh.append(fence)
if not fresh:
return 0
hits = fresh
async with async_session() as session:
for fence in hits:
aid = str(uuid4())
body = {
"id": aid,
"geofence_id": fence["id"],
"geofence_name": fence.get("name"),
"source_kind": source_kind,
"entity_id": str(entity_id),
"lat": lat,
"lon": lon,
"payload": payload or {},
"created_at": now.isoformat(),
}
try:
await session.execute(
text(
"""
INSERT INTO geofence_alerts
(id, geofence_id, source_kind, entity_id, lat, lon, payload)
VALUES (
CAST(:id AS uuid), CAST(:geofence_id AS uuid),
:source_kind, :entity_id, :lat, :lon, CAST(:payload AS jsonb)
)
"""
),
{
"id": aid,
"geofence_id": fence["id"],
"source_kind": source_kind,
"entity_id": str(entity_id),
"lat": lat,
"lon": lon,
"payload": json.dumps(payload or {}),
},
)
except Exception:
pass
sent += await manager.publish_point(
"geofence_alert", body, lat=lat, lon=lon,
)
try:
await session.commit()
except Exception:
pass
return sent
async def list_alerts(
*,
geofence_id: str | None = None,
since: datetime | None = None,
until: datetime | None = None,
source_kind: str | None = None,
limit: int = 100,
) -> list[dict]:
"""Filterable hit log. Empty list if the DB is down — never raises."""
where = ["TRUE"]
params: dict[str, Any] = {"limit": int(limit)}
if geofence_id:
where.append("geofence_id = CAST(:geofence_id AS uuid)")
params["geofence_id"] = geofence_id
if since is not None:
where.append("created_at >= :since")
params["since"] = since
if until is not None:
where.append("created_at <= :until")
params["until"] = until
if source_kind:
where.append("source_kind = :source_kind")
params["source_kind"] = source_kind
sql = f"""
SELECT id::text, geofence_id::text, source_kind, entity_id,
lat, lon, payload, created_at
FROM geofence_alerts
WHERE {' AND '.join(where)}
ORDER BY created_at DESC
LIMIT :limit
"""
try:
async with async_session() as session:
rows = (await session.execute(text(sql), params)).mappings().all()
out = []
for r in rows:
item = dict(r)
if item.get("created_at") is not None:
item["created_at"] = item["created_at"].isoformat()
out.append(item)
return out
except Exception:
return []
async def get_geofence(gid: str) -> dict | None:
current = next((f for f in _cache if f["id"] == gid), None)
if current is not None:
return current
try:
await refresh_cache()
except Exception:
return None
return next((f for f in _cache if f["id"] == gid), None)
def _marker_from_track(row) -> dict:
from live_layers import to_marker
extra = {"bucket": row["bucket"].isoformat() if row.get("bucket") else None, "dvr": True}
return to_marker(
row["id"], row["lat"], row["lon"],
heading=row.get("heading"), speed=row.get("speed"),
label=row.get("label") or row["id"],
extra=extra,
)
async def _cagg_inside(gid: str, kind: str, bucket: datetime, limit: int = 2000) -> list[dict]:
table = "aircraft_tracks_1min" if kind == "aircraft" else "vessel_tracks_1min"
id_col = "hex" if kind == "aircraft" else "mmsi"
sql = f"""
SELECT {id_col} AS id, lat, lon, heading, speed, label, bucket
FROM {table}
WHERE bucket = :bucket
AND ST_Intersects(
(SELECT geom FROM geofences WHERE id = CAST(:gid AS uuid)),
ST_SetSRID(ST_MakePoint(lon, lat), 4326)
)
LIMIT :limit
"""
try:
async with async_session() as session:
rows = (await session.execute(
text(sql), {"bucket": bucket, "gid": gid, "limit": limit},
)).mappings().all()
return [
_marker_from_track(r)
for r in rows
if r["lat"] is not None and r["lon"] is not None
]
except Exception:
return []
async def _fires_inside(gid: str, ts: datetime, limit: int = 2000) -> list[dict]:
from tracks import minute_bucket
bucket = minute_bucket(ts)
t1 = bucket + timedelta(minutes=1)
sql = """
SELECT latitude, longitude, brightness, confidence, acq_time, satellite,
instrument, bright_ti5, frp, daynight
FROM fires
WHERE acq_time >= :t0 AND acq_time < :t1
AND ST_Intersects(
(SELECT geom FROM geofences WHERE id = CAST(:gid AS uuid)),
ST_SetSRID(ST_MakePoint(longitude, latitude), 4326)
)
LIMIT :limit
"""
try:
async with async_session() as session:
rows = (await session.execute(
text(sql),
{"t0": bucket, "t1": t1, "gid": gid, "limit": limit},
)).mappings().all()
out = []
for r in rows:
item = dict(r)
if item.get("acq_time") is not None:
item["acq_time"] = item["acq_time"].isoformat()
out.append(item)
return out
except Exception:
return []
async def snapshot_at(gid: str, ts: datetime) -> dict | None:
"""Positions inside the fence at time T. None if the fence is missing.
Does not persist or notify. Empty lists if track/fire queries fail.
"""
fence = await get_geofence(gid)
if fence is None:
return None
from tracks import minute_bucket
bucket = minute_bucket(ts)
aircraft = await _cagg_inside(gid, "aircraft", bucket)
vessels = await _cagg_inside(gid, "vessel", bucket)
fires = await _fires_inside(gid, ts)
return {
"geofence_id": gid,
"timestamp": ts.isoformat(),
"aircraft": aircraft,
"vessels": vessels,
"fires": fires,
}