mirror of
https://github.com/BigBodyCobain/Shadowbroker.git
synced 2026-07-31 16:07:31 +02:00
perf: live-data deltas, payload caps, and map render polish
Cut fast-tier payload cost with zoom-aware sampling, row deltas, CCTV bbox columns, and MapLibre/motion polish; force viewport snapshot refetches so regional pans refill aircraft immediately. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -1446,16 +1446,67 @@ def run_all_ingestors():
|
||||
logger.warning(f"Ingestor {ingestor.__class__.__name__} failed during seed: {e}")
|
||||
|
||||
|
||||
def get_all_cameras() -> List[Dict[str, Any]]:
|
||||
# Columns the map / live-data path actually needs — skip refresh_rate / timestamps.
|
||||
_CAMERA_SELECT_COLS = (
|
||||
"id",
|
||||
"source_agency",
|
||||
"lat",
|
||||
"lon",
|
||||
"direction_facing",
|
||||
"media_url",
|
||||
"media_type",
|
||||
)
|
||||
|
||||
|
||||
def get_camera_count() -> int:
|
||||
"""Cheap total for UI badges without loading the full table."""
|
||||
conn = sqlite3.connect(str(DB_PATH))
|
||||
try:
|
||||
row = conn.execute("SELECT COUNT(*) FROM cameras").fetchone()
|
||||
return int(row[0] if row else 0)
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
|
||||
def get_all_cameras(
|
||||
*,
|
||||
south: float | None = None,
|
||||
west: float | None = None,
|
||||
north: float | None = None,
|
||||
east: float | None = None,
|
||||
) -> List[Dict[str, Any]]:
|
||||
"""Load map-facing camera rows, optionally SQL-scoped to a viewport.
|
||||
|
||||
Without bounds, returns the full catalog (store still holds world list;
|
||||
/api/live-data/fast applies #288 bbox + world-zoom sampling).
|
||||
"""
|
||||
cols = ", ".join(_CAMERA_SELECT_COLS)
|
||||
sql = f"SELECT {cols} FROM cameras"
|
||||
params: list[Any] = []
|
||||
if None not in (south, west, north, east):
|
||||
# Inclusive bbox; antimeridian (west > east) uses OR on lon.
|
||||
if west <= east:
|
||||
sql += " WHERE lat BETWEEN ? AND ? AND lon BETWEEN ? AND ?"
|
||||
params.extend([south, north, west, east])
|
||||
else:
|
||||
sql += (
|
||||
" WHERE lat BETWEEN ? AND ?"
|
||||
" AND (lon >= ? OR lon <= ?)"
|
||||
)
|
||||
params.extend([south, north, west, east])
|
||||
|
||||
conn = sqlite3.connect(str(DB_PATH))
|
||||
conn.row_factory = sqlite3.Row
|
||||
cursor = conn.cursor()
|
||||
cursor.execute("SELECT * FROM cameras")
|
||||
rows = cursor.fetchall()
|
||||
conn.close()
|
||||
try:
|
||||
rows = conn.execute(sql, params).fetchall()
|
||||
finally:
|
||||
conn.close()
|
||||
|
||||
cameras = []
|
||||
for row in rows:
|
||||
cam = dict(row)
|
||||
cam["media_type"] = str(cam.get("media_type") or _detect_media_type(cam.get("media_url", "")) or "image")
|
||||
cam = {k: row[k] for k in _CAMERA_SELECT_COLS}
|
||||
cam["media_type"] = str(
|
||||
cam.get("media_type") or _detect_media_type(cam.get("media_url", "")) or "image"
|
||||
)
|
||||
cameras.append(cam)
|
||||
return cameras
|
||||
|
||||
@@ -146,6 +146,123 @@ source_timestamps = {}
|
||||
source_freshness: dict[str, dict] = {}
|
||||
|
||||
|
||||
# Layers that support row-level live-data deltas (P2).
|
||||
_DELTA_LAYER_KEYS: frozenset[str] = frozenset(
|
||||
{
|
||||
"ships",
|
||||
"commercial_flights",
|
||||
"military_flights",
|
||||
"tracked_flights",
|
||||
"private_flights",
|
||||
"private_jets",
|
||||
}
|
||||
)
|
||||
# Ring of (layer_version, id→item) snapshots for delta computation.
|
||||
_LAYER_ID_RING: dict[str, list[tuple[int, dict[str, Any]]]] = {}
|
||||
_LAYER_ID_RING_MAX = 4
|
||||
|
||||
|
||||
def entity_id_for_layer(layer: str, item: dict) -> str:
|
||||
"""Stable entity id for delta upsert/delete keys."""
|
||||
if not isinstance(item, dict):
|
||||
return ""
|
||||
if layer == "ships":
|
||||
return str(item.get("mmsi") or item.get("id") or "").strip()
|
||||
return str(
|
||||
item.get("icao24") or item.get("icao") or item.get("id") or item.get("hex") or ""
|
||||
).strip().lower()
|
||||
|
||||
|
||||
def _track_fingerprint(item: dict) -> tuple:
|
||||
lng = item.get("lng")
|
||||
if lng is None:
|
||||
lng = item.get("lon")
|
||||
return (
|
||||
item.get("lat"),
|
||||
lng,
|
||||
item.get("heading") or item.get("hdg") or item.get("cog") or item.get("true_track"),
|
||||
item.get("speed_knots") or item.get("sog") or item.get("spd"),
|
||||
item.get("alt") or item.get("altitude") or item.get("alt_km"),
|
||||
)
|
||||
|
||||
|
||||
def _capture_layer_id_map(layer: str, items: list) -> dict[str, Any]:
|
||||
out: dict[str, Any] = {}
|
||||
for item in items or []:
|
||||
if not isinstance(item, dict):
|
||||
continue
|
||||
eid = entity_id_for_layer(layer, item)
|
||||
if eid:
|
||||
out[eid] = item
|
||||
return out
|
||||
|
||||
|
||||
def _record_delta_layer_snapshot_locked(layer: str) -> None:
|
||||
"""Caller must hold ``_data_lock``."""
|
||||
if layer not in _DELTA_LAYER_KEYS:
|
||||
return
|
||||
ver = _layer_versions.get(layer, 0)
|
||||
val = latest_data.get(layer)
|
||||
items = val if isinstance(val, list) else []
|
||||
id_map = _capture_layer_id_map(layer, items)
|
||||
ring = _LAYER_ID_RING.setdefault(layer, [])
|
||||
ring.append((ver, id_map))
|
||||
while len(ring) > _LAYER_ID_RING_MAX:
|
||||
ring.pop(0)
|
||||
|
||||
|
||||
def compute_layer_row_delta(layer: str, since_version: int) -> dict[str, Any] | None:
|
||||
"""Build upsert/delete delta vs a prior layer version.
|
||||
|
||||
Returns ``None`` when the base version is no longer in the ring (client
|
||||
must take a full snapshot). Returns an empty upsert/delete when unchanged.
|
||||
"""
|
||||
if layer not in _DELTA_LAYER_KEYS:
|
||||
return None
|
||||
try:
|
||||
since = int(since_version)
|
||||
except (TypeError, ValueError):
|
||||
return None
|
||||
|
||||
with _data_lock:
|
||||
current_ver = int(_layer_versions.get(layer, 0) or 0)
|
||||
if since == current_ver or since > current_ver:
|
||||
return {
|
||||
"upsert": [],
|
||||
"delete": [],
|
||||
"version": current_ver,
|
||||
"unchanged": True,
|
||||
}
|
||||
ring = list(_LAYER_ID_RING.get(layer) or [])
|
||||
base_map = None
|
||||
for ver, id_map in ring:
|
||||
if ver == since:
|
||||
base_map = id_map
|
||||
break
|
||||
if base_map is None:
|
||||
# Base version aged out of the ring — full snapshot required.
|
||||
return None
|
||||
if ring and ring[-1][0] == current_ver:
|
||||
cur_map = ring[-1][1]
|
||||
else:
|
||||
val = latest_data.get(layer)
|
||||
items = val if isinstance(val, list) else []
|
||||
cur_map = _capture_layer_id_map(layer, items)
|
||||
|
||||
upsert: list[Any] = []
|
||||
for eid, item in cur_map.items():
|
||||
prev = base_map.get(eid)
|
||||
if prev is None or _track_fingerprint(prev) != _track_fingerprint(item):
|
||||
upsert.append(item)
|
||||
delete = [eid for eid in base_map if eid not in cur_map]
|
||||
return {
|
||||
"upsert": upsert,
|
||||
"delete": delete,
|
||||
"version": current_ver,
|
||||
"unchanged": not upsert and not delete,
|
||||
}
|
||||
|
||||
|
||||
def _mark_fresh(*keys):
|
||||
"""Record the current UTC time for one or more data source keys."""
|
||||
now = datetime.utcnow().isoformat()
|
||||
@@ -159,6 +276,7 @@ def _mark_fresh(*keys):
|
||||
val = latest_data.get(k)
|
||||
count = len(val) if isinstance(val, list) else (1 if val is not None else 0)
|
||||
changed.append((k, _layer_versions[k], count))
|
||||
_record_delta_layer_snapshot_locked(k)
|
||||
# Publish partial fetch progress immediately so the frontend can
|
||||
# observe newly available data without waiting for the entire tier.
|
||||
_data_version += 1
|
||||
@@ -273,7 +391,10 @@ def get_latest_data_subset(*keys: str) -> DashboardData:
|
||||
|
||||
|
||||
def get_latest_data_deepcopy_snapshot() -> DashboardData:
|
||||
"""Deep-copy the full dashboard for /api/health and legacy /api/live-data.
|
||||
"""Deep-copy the full dashboard for callers that need an isolated mutate-safe snap.
|
||||
|
||||
Prefer ``get_latest_data_refs_snapshot`` / ``get_latest_data_subset_refs`` on
|
||||
read-only hot paths (legacy ``/api/live-data`` uses refs + orjson).
|
||||
|
||||
The per-value deepcopy runs OUTSIDE ``_data_lock`` so a large clone cannot
|
||||
block fetcher writers (#375). The store contract is replace-don't-mutate,
|
||||
@@ -310,6 +431,17 @@ def get_latest_data_subset_refs(*keys: str) -> DashboardData:
|
||||
return snap
|
||||
|
||||
|
||||
def get_latest_data_refs_snapshot() -> DashboardData:
|
||||
"""Return a shallow dict of all top-level store refs (read-only callers).
|
||||
|
||||
Copies the mapping under the lock without deep-copying values. Safe for
|
||||
orjson serialization and other read-only consumers; callers MUST NOT
|
||||
mutate nested objects.
|
||||
"""
|
||||
with _data_lock:
|
||||
return {key: value for key, value in latest_data.items()}
|
||||
|
||||
|
||||
def get_source_timestamps_snapshot() -> dict[str, str]:
|
||||
"""Return a stable copy of per-source freshness timestamps."""
|
||||
with _data_lock:
|
||||
|
||||
@@ -10,13 +10,14 @@ import logging
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Inline — local DB / static files only.
|
||||
# Inline — tiny local/static loads (ms).
|
||||
_INSTANT_LAYER_KEYS: frozenset[str] = frozenset(
|
||||
{"cctv", "power_plants", "datacenters"}
|
||||
{"power_plants", "datacenters"}
|
||||
)
|
||||
# Background — network-bound; may take seconds.
|
||||
# Background — network-bound OR large local scans (full CCTV SELECT can stall
|
||||
# the single uvicorn worker if run inline on enable).
|
||||
_SLOW_LAYER_KEYS: frozenset[str] = frozenset(
|
||||
{"firms", "psk_reporter", "fishing_activity"}
|
||||
{"cctv", "firms", "psk_reporter", "fishing_activity"}
|
||||
)
|
||||
|
||||
|
||||
@@ -33,12 +34,6 @@ def _was_off_now_on(before: dict[str, bool], key: str) -> bool:
|
||||
|
||||
|
||||
def _instant_fetch(key: str) -> None:
|
||||
if key == "cctv":
|
||||
from services.fetchers.infrastructure import fetch_cctv
|
||||
|
||||
fetch_cctv()
|
||||
logger.info("CCTV loaded (layer enabled)")
|
||||
return
|
||||
if key == "power_plants":
|
||||
from services.fetchers.infrastructure import fetch_power_plants
|
||||
|
||||
@@ -55,6 +50,12 @@ def _instant_fetch(key: str) -> None:
|
||||
|
||||
|
||||
def _slow_fetch(key: str) -> None:
|
||||
if key == "cctv":
|
||||
from services.fetchers.infrastructure import fetch_cctv
|
||||
|
||||
fetch_cctv()
|
||||
logger.info("CCTV loaded (layer enabled)")
|
||||
return
|
||||
if key == "firms":
|
||||
from services.fetchers.earth_observation import (
|
||||
fetch_firms_country_fires,
|
||||
|
||||
@@ -39,6 +39,7 @@ transports.
|
||||
"""
|
||||
|
||||
import json
|
||||
import orjson
|
||||
import os
|
||||
import time
|
||||
import hmac
|
||||
@@ -109,9 +110,30 @@ def _atomic_write_text(target: Path, content: str, encoding: str = "utf-8") -> N
|
||||
pass
|
||||
raise
|
||||
|
||||
|
||||
def _atomic_write_bytes(target: Path, content: bytes) -> None:
|
||||
"""Write binary content atomically via temp file + os.replace()."""
|
||||
parent = target.parent
|
||||
parent.mkdir(parents=True, exist_ok=True)
|
||||
fd, tmp_path = tempfile.mkstemp(dir=str(parent), suffix=".tmp")
|
||||
try:
|
||||
with os.fdopen(fd, "wb") as handle:
|
||||
handle.write(content)
|
||||
handle.flush()
|
||||
os.fsync(handle.fileno())
|
||||
os.replace(tmp_path, str(target))
|
||||
except BaseException:
|
||||
try:
|
||||
os.unlink(tmp_path)
|
||||
except OSError:
|
||||
pass
|
||||
raise
|
||||
|
||||
|
||||
DATA_DIR = Path(__file__).resolve().parents[2] / "data"
|
||||
CHAIN_FILE = DATA_DIR / "infonet.json"
|
||||
WAL_FILE = DATA_DIR / "infonet.wal"
|
||||
CHAIN_COLD_DIR = DATA_DIR / "infonet_cold"
|
||||
GATE_STORE_DIR = DATA_DIR / "gate_messages"
|
||||
GATE_STORAGE_DOMAIN = "gates"
|
||||
|
||||
@@ -1469,6 +1491,7 @@ class Infonet:
|
||||
self.event_index: dict[str, int] = {} # {event_id: index in events list}
|
||||
self.public_key_bindings: dict[str, str] = {} # {public_key: canonical node_id}
|
||||
self.revocations: dict[str, dict] = {}
|
||||
self._cold_segments: list[dict] = [] # Archived prefix metadata (P8)
|
||||
self._replay_filter = ReplayFilter()
|
||||
self._last_validated_index: int = 0 # For incremental validation
|
||||
# Running counters — avoid O(N) scans in get_info()
|
||||
@@ -1549,9 +1572,15 @@ class Infonet:
|
||||
self.head_hash = data.get("head_hash", GENESIS_HASH)
|
||||
self.node_sequences = data.get("node_sequences", {})
|
||||
self.sequence_domains = data.get("sequence_domains", {})
|
||||
cold = data.get("cold_segments") or []
|
||||
self._cold_segments = list(cold) if isinstance(cold, list) else []
|
||||
self._rebuild_state()
|
||||
self._rebuild_revocations()
|
||||
self._rebuild_counters()
|
||||
# Honor MAX_CHAIN_MEMORY after load (archive overflow, keep hot window).
|
||||
if len(self.events) > MAX_CHAIN_MEMORY:
|
||||
self._enforce_memory_cap()
|
||||
self._flush()
|
||||
logger.info(
|
||||
f"Loaded Infonet: {len(self.events)} events, head={self.head_hash[:16]}..."
|
||||
)
|
||||
@@ -1794,6 +1823,7 @@ class Infonet:
|
||||
and schedule a single write after _SAVE_INTERVAL seconds. Multiple
|
||||
rapid calls collapse into one I/O operation.
|
||||
"""
|
||||
self._enforce_memory_cap()
|
||||
self._dirty = True
|
||||
with self._save_lock:
|
||||
if self._save_timer is None or not self._save_timer.is_alive():
|
||||
@@ -1807,6 +1837,8 @@ class Infonet:
|
||||
Sprint 2 / Rec #8: clears the WAL only after the chain file has
|
||||
been durably written. A crash before _flush() succeeds leaves
|
||||
the WAL in place so _replay_wal() can recover on next boot.
|
||||
|
||||
P8: compact orjson (no indent=2) — flush cost scales with chain size.
|
||||
"""
|
||||
if not self._dirty:
|
||||
return
|
||||
@@ -1818,9 +1850,10 @@ class Infonet:
|
||||
"head_hash": self.head_hash,
|
||||
"node_sequences": self.node_sequences,
|
||||
"sequence_domains": self.sequence_domains,
|
||||
"cold_segments": list(getattr(self, "_cold_segments", []) or []),
|
||||
"events": self.events,
|
||||
}
|
||||
_atomic_write_text(CHAIN_FILE, json.dumps(data, indent=2), encoding="utf-8")
|
||||
_atomic_write_bytes(CHAIN_FILE, orjson.dumps(data, option=orjson.OPT_NON_STR_KEYS))
|
||||
self._dirty = False
|
||||
# Chain file is now durable — safe to retire the WAL entry.
|
||||
self._clear_wal()
|
||||
@@ -1837,13 +1870,56 @@ class Infonet:
|
||||
"head_hash": self.head_hash,
|
||||
"node_sequences": self.node_sequences,
|
||||
"sequence_domains": self.sequence_domains,
|
||||
"cold_segments": list(getattr(self, "_cold_segments", []) or []),
|
||||
"events": self.events,
|
||||
}
|
||||
_atomic_write_text(CHAIN_FILE, json.dumps(data, indent=2), encoding="utf-8")
|
||||
_atomic_write_bytes(CHAIN_FILE, orjson.dumps(data, option=orjson.OPT_NON_STR_KEYS))
|
||||
except Exception as e:
|
||||
logger.error(f"Failed to materialize Infonet: {e}")
|
||||
raise
|
||||
|
||||
def _enforce_memory_cap(self) -> None:
|
||||
"""Spill oldest events to cold segments when RAM exceeds MAX_CHAIN_MEMORY.
|
||||
|
||||
History is archived on disk (not deleted). Hot chain + indexes stay
|
||||
coherent for head/sync; deep history remains in cold segments.
|
||||
"""
|
||||
if len(self.events) <= MAX_CHAIN_MEMORY:
|
||||
return
|
||||
overflow = len(self.events) - MAX_CHAIN_MEMORY
|
||||
cold = self.events[:overflow]
|
||||
CHAIN_COLD_DIR.mkdir(parents=True, exist_ok=True)
|
||||
if not hasattr(self, "_cold_segments") or self._cold_segments is None:
|
||||
self._cold_segments = []
|
||||
seg_no = len(self._cold_segments)
|
||||
filename = f"cold_{seg_no:08d}.orjson"
|
||||
path = CHAIN_COLD_DIR / filename
|
||||
_atomic_write_bytes(path, orjson.dumps(cold, option=orjson.OPT_NON_STR_KEYS))
|
||||
self._cold_segments.append(
|
||||
{
|
||||
"filename": filename,
|
||||
"count": len(cold),
|
||||
"first_id": str(cold[0].get("event_id", "") or "") if cold else "",
|
||||
"last_id": str(cold[-1].get("event_id", "") or "") if cold else "",
|
||||
}
|
||||
)
|
||||
self.events = self.events[overflow:]
|
||||
# Rebuild indexes for the hot window only.
|
||||
self.event_index = {
|
||||
str(evt.get("event_id", "") or ""): idx
|
||||
for idx, evt in enumerate(self.events)
|
||||
if evt.get("event_id")
|
||||
}
|
||||
self._rebuild_counters()
|
||||
self._invalidate_merkle_cache()
|
||||
self._dirty = True
|
||||
logger.info(
|
||||
"Infonet memory cap: archived %d events to %s (hot=%d)",
|
||||
overflow,
|
||||
filename,
|
||||
len(self.events),
|
||||
)
|
||||
|
||||
def confirmations_for_event(self, event_id: str) -> int:
|
||||
idx = self.event_index.get(event_id)
|
||||
if idx is None:
|
||||
|
||||
@@ -112,9 +112,9 @@ def _evaluate_watches() -> list[dict[str, Any]]:
|
||||
|
||||
# Load telemetry once for all watches
|
||||
try:
|
||||
from services.telemetry import get_cached_telemetry, get_cached_slow_telemetry
|
||||
fast = get_cached_telemetry() or {}
|
||||
slow = get_cached_slow_telemetry() or {}
|
||||
from services.telemetry import get_cached_telemetry_refs, get_cached_slow_telemetry_refs
|
||||
fast = get_cached_telemetry_refs() or {}
|
||||
slow = get_cached_slow_telemetry_refs() or {}
|
||||
except Exception:
|
||||
return []
|
||||
|
||||
|
||||
Reference in New Issue
Block a user