Files
node-26/secondary-sdr-container/app/internal/decoder_control.py
T
Logan CusanoandClaude Opus 5.5 b54624e176
CI / lint (push) Successful in 6s
CI / lint (pull_request) Successful in 7s
CI / test (push) Successful in 39s
CI / test (pull_request) Successful in 40s
secondary-sdr: fix ADS-B on real hardware (tested on radio-box)
First hardware test of node-26#9/#10 found three blockers:
- antirez/dump1090 has no --write-json, so ADS-B mode exited on start.
  Swapped to wiedehopf/readsb; map alt_baro/gs (old names as fallback).
- op25 is not always on RTL-SDR index 0 (radio-box: op25 on 1, 0 free).
  start() now tries each index and keeps the first decoder that stays up.
  AIS-catcher index flag fixed to -d:N ("-d N" selects by serial).
- status() reported "running" for a decoder that died on startup: the
  zombie still answered killpg(pgid, 0). Liveness now via Popen.poll().

install.sh blacklists dvb_usb_rtl28xxu, which claimed the second dongle.

Verified on radio-box: 831 msgs/min, 6 aircraft (3 with position) on a
9cm whip; op25 unaffected. AIS still untested.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 12:34:18 -04:00

261 lines
8.1 KiB
Python

import json
import os
import signal
import subprocess
import threading
from pathlib import Path
from typing import Any, Dict, List, Optional
from internal.logger import create_logger
LOGGER = create_logger(__name__)
# AIS-catcher streams one JSON object per received message on stdout rather
# than writing a periodic snapshot file (dump1090's approach above) — so the
# current-vessel snapshot lives in memory, keyed by mmsi, kept warm by a
# background reader thread for as long as the subprocess we started is
# alive. Unlike the pgid-file state, this does NOT survive this API
# process restarting independently of its subprocess — acceptable since
# nothing here does that today.
_ais_vessels: Dict[str, Dict[str, Any]] = {}
_ais_lock = threading.Lock()
# Which RTL-SDR index op25 holds is NOT fixed — on radio-box op25 had index 1
# and index 0 was free, so "op25 is always 0" was wrong. rtlsdr can't open a
# dongle another process has claimed, and the decoders exit within ~50ms when
# that happens, so start() tries each index and keeps the first that stays up.
# op25 is never disturbed: a failed claim doesn't touch its dongle.
MAX_SDR_INDEX = 4
_STARTUP_GRACE_S = 2.0
_PGID_FILE = "/tmp/secondary_sdr.pgid"
_MODE_FILE = "/tmp/secondary_sdr.mode"
ADSB_JSON_DIR = Path("/tmp/adsb")
# The live decoder handle. poll() is the only reliable liveness check: a
# decoder that dies on startup (e.g. SDR busy) stays an unreaped zombie, and
# killpg(pgid, 0) still succeeds on a zombie — status reported "running".
_proc: Optional[subprocess.Popen] = None
def _save_state(pgid: int, mode: str) -> None:
Path(_PGID_FILE).write_text(str(pgid))
Path(_MODE_FILE).write_text(mode)
def _read_pgid() -> Optional[int]:
try:
return int(Path(_PGID_FILE).read_text().strip())
except Exception:
return None
def _read_mode() -> Optional[str]:
try:
return Path(_MODE_FILE).read_text().strip()
except Exception:
return None
def is_running() -> bool:
if _proc is not None:
return _proc.poll() is None
pgid = _read_pgid()
if pgid is None:
return False
try:
os.killpg(pgid, 0)
return True
except OSError:
return False
def _adsb_command(index: int) -> List[str]:
ADSB_JSON_DIR.mkdir(parents=True, exist_ok=True)
return [
"/opt/readsb/readsb",
"--net",
"--device-type", "rtlsdr",
"--device", str(index),
"--write-json", str(ADSB_JSON_DIR),
"--write-json-every", "1",
]
def _ais_command(index: int) -> List[str]:
return [
"/opt/AIS-catcher/build/AIS-catcher",
f"-d:{index}", # "-d <x>" would select by serial, not index
"-o", "JSON",
]
def _ais_reader(proc: subprocess.Popen) -> None:
"""
Consume AIS-catcher's stdout, one JSON message per line, and keep the
latest report per mmsi. Field names (mmsi/lat/lon/speed/course or
heading/shipname or name) are believed correct from AIS-catcher's
published JSON output docs but UNVERIFIED against a real capture in
this session — same caveat as dump1090's aircraft.json mapping.
Malformed/partial lines (e.g. static-data-only messages with no
position) are skipped rather than raising, since dropping one line must
never kill the reader thread.
"""
if not proc.stdout:
return
for line in proc.stdout:
try:
msg = json.loads(line)
except Exception:
continue
mmsi = msg.get("mmsi")
if not mmsi:
continue
# AIS-catcher emits separate message TYPES per mmsi — static data
# (name, no position) and position reports (lat/lon, no name) arrive
# as distinct lines. Merge onto the existing entry, only overwriting
# a field the new message actually carries, so a position-only
# report doesn't blank out a name learned from an earlier message.
name = (msg.get("shipname") or msg.get("name") or "").strip() or None
heading = msg.get("heading") if msg.get("heading") is not None else msg.get("course")
updates = {
"mmsi": str(mmsi),
"name": name,
"lat": msg.get("lat"),
"lon": msg.get("lon"),
"speed_kt": msg.get("speed"),
"heading_deg": heading,
}
with _ais_lock:
existing = _ais_vessels.get(str(mmsi), {})
for key, value in updates.items():
if value is not None:
existing[key] = value
_ais_vessels[str(mmsi)] = existing
def start(mode: str) -> bool:
global _proc
if is_running():
stop()
if mode == "adsb":
build = _adsb_command
elif mode == "ais":
build = _ais_command
with _ais_lock:
_ais_vessels.clear()
else:
raise ValueError(f"Unknown secondary SDR mode: {mode!r}")
needs_stdout = mode == "ais"
for index in range(MAX_SDR_INDEX):
try:
proc = subprocess.Popen(
build(index),
preexec_fn=os.setsid,
stdout=subprocess.PIPE if needs_stdout else None,
text=True if needs_stdout else None,
bufsize=1 if needs_stdout else -1,
)
except Exception as e:
LOGGER.error(f"Failed to start secondary SDR decoder mode={mode!r}: {e}")
return False
try:
proc.wait(timeout=_STARTUP_GRACE_S)
LOGGER.info(f"Secondary SDR decoder mode={mode!r} could not use SDR index {index}, trying next")
continue
except subprocess.TimeoutExpired:
pass
if needs_stdout:
threading.Thread(target=_ais_reader, args=(proc,), daemon=True).start()
_proc = proc
_save_state(proc.pid, mode)
LOGGER.info(f"Started secondary SDR decoder mode={mode!r} on SDR index {index} pid={proc.pid}")
return True
LOGGER.error(f"Secondary SDR decoder mode={mode!r}: no free SDR found (op25 holds one; is a second plugged in?)")
return False
def stop() -> bool:
global _proc
pgid = _read_pgid()
if pgid is None:
return True
try:
os.killpg(pgid, signal.SIGTERM)
except OSError:
pass
if _proc is not None:
try:
_proc.wait(timeout=5)
except subprocess.TimeoutExpired:
pass
_proc = None
try:
os.remove(_PGID_FILE)
except OSError:
pass
try:
os.remove(_MODE_FILE)
except OSError:
pass
return True
def status() -> Dict[str, Any]:
running = is_running()
return {
"status": "running" if running else "stopped",
"mode": _read_mode() if running else None,
}
def _altitude(a: Dict[str, Any]) -> Optional[int]:
alt = a.get("alt_baro", a.get("altitude"))
if alt == "ground":
return 0
return alt if isinstance(alt, (int, float)) else None
def _read_adsb_snapshot() -> List[Dict[str, Any]]:
"""
Map readsb's aircraft.json (--write-json output) to the server's
telemetry schema. readsb uses the dump1090-fa field names: alt_baro (int,
or the string "ground"), gs, track. Older dump1090 forks used
altitude/speed, kept as a fallback.
"""
path = ADSB_JSON_DIR / "aircraft.json"
try:
raw = json.loads(path.read_text())
except Exception:
return []
out = []
for a in raw.get("aircraft", []):
icao = a.get("hex")
if not icao:
continue
out.append({
"icao": icao.upper(),
"callsign": (a.get("flight") or "").strip() or None,
"lat": a.get("lat"),
"lon": a.get("lon"),
"altitude_ft": _altitude(a),
"ground_speed_kt": a.get("gs", a.get("speed")),
"track_deg": a.get("track"),
})
return out
def data() -> Dict[str, Any]:
mode = _read_mode()
if mode == "adsb":
return {"mode": mode, "aircraft": _read_adsb_snapshot()}
if mode == "ais":
with _ais_lock:
vessels = list(_ais_vessels.values())
return {"mode": mode, "vessels": vessels}
return {"mode": mode, "aircraft": [], "vessels": []}