Secondary SDR priority: every spare SDR runs the next decoder in an ordered list
CI / lint (push) Successful in 6s
CI / test (push) Successful in 39s

Replaces the single secondary_sdr_mode with secondary_sdr_priority, e.g.
["adsb", "ais"]. OP25 always keeps its own dongle; the secondary-sdr
container starts decoders top-down until it runs out of free SDRs, so a
3-SDR node runs ADS-B and AIS at once and a 2-SDR node runs the top pick.

- secondary-sdr: one decoder per mode, POST /secondary/apply(priority)
  (no-op when the right prefix is already running), orphaned decoders
  from a uvicorn reload are reaped on start, /status reports sdr_count
  via lsusb (op25's :stable image has none).
- edge-node: one apply path (set_secondary_priority) for the local
  dashboard, a new C2 'set_secondary_priority' MQTT command, and config
  pushes; it never restarts op25. Legacy mode migrates on load. Checkin
  reports priority, what's running, and sdr_count. Uplink forwards
  aircraft and vessels whenever either is present.
- Local dashboard: 'Secondary SDRs' card to enable/reorder/save.

Verified: edge-node pytest 190 passed, flake8 clean.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
Logan Cusano
2026-09-27 14:08:07 -04:00
co-authored by Claude Opus 5.5
parent 1c74edffea
commit 9b6c64fbbf
11 changed files with 472 additions and 117 deletions
@@ -10,64 +10,59 @@ 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()
# One decoder per mode, each on its own SDR (node-26#9). The node's secondary
# SDR *priority* (e.g. ["adsb", "ais"]) is applied by apply(): decoders start
# down the list until no free dongle is left, so every SDR the node has gets
# used and the ones beyond the list's reach simply aren't started. OP25 always
# keeps its own dongle — it is started first and never part of this list.
MODES = ("adsb", "ais")
# 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.
# that happens, so _start_one() 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"
_STATE_DIR = Path("/tmp/secondary_sdr")
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
# Live decoder handles. poll() is the only reliable liveness check: a decoder
# that dies on startup stays an unreaped zombie, and killpg(pgid, 0) still
# succeeds on a zombie.
_procs: Dict[str, subprocess.Popen] = {}
_indices: Dict[str, int] = {}
_lock = threading.Lock()
# AIS-catcher streams one JSON object per received message on stdout rather
# than writing a periodic snapshot file (readsb's approach) — so the
# current-vessel snapshot lives in memory, keyed by mmsi, kept warm by a
# background reader thread for as long as the decoder is alive.
_ais_vessels: Dict[str, Dict[str, Any]] = {}
_ais_lock = threading.Lock()
def _save_state(pgid: int, mode: str) -> None:
Path(_PGID_FILE).write_text(str(pgid))
Path(_MODE_FILE).write_text(mode)
def _pgid_file(mode: str) -> Path:
return _STATE_DIR / f"{mode}.pgid"
def _read_pgid() -> Optional[int]:
try:
return int(Path(_PGID_FILE).read_text().strip())
except Exception:
return None
def _reap_orphans() -> None:
"""Kill decoders left behind by a previous API process (uvicorn --reload
restarts this process, but decoders run in their own session and would
otherwise keep holding their SDRs)."""
if not _STATE_DIR.exists():
return
for f in _STATE_DIR.glob("*.pgid"):
try:
os.killpg(int(f.read_text().strip()), signal.SIGTERM)
LOGGER.info(f"Stopped orphaned secondary decoder from {f.name}")
except Exception:
pass
f.unlink(missing_ok=True)
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
_reap_orphans()
def _adsb_command(index: int) -> List[str]:
@@ -90,6 +85,9 @@ def _ais_command(index: int) -> List[str]:
]
_COMMANDS = {"adsb": _adsb_command, "ais": _ais_command}
def _ais_reader(proc: subprocess.Popen) -> None:
"""
Consume AIS-catcher's stdout, one JSON message per line, and keep the
@@ -134,25 +132,32 @@ def _ais_reader(proc: subprocess.Popen) -> None:
_ais_vessels[str(mmsi)] = existing
def start(mode: str) -> bool:
global _proc
if is_running():
stop()
def is_running(mode: str) -> bool:
proc = _procs.get(mode)
return proc is not None and proc.poll() is None
if mode == "adsb":
build = _adsb_command
elif mode == "ais":
build = _ais_command
def running() -> List[str]:
return [m for m in MODES if is_running(m)]
def _start_one(mode: str) -> bool:
"""Start one decoder on the first free SDR. False when none is free."""
if mode not in _COMMANDS:
raise ValueError(f"Unknown secondary SDR mode: {mode!r}")
if is_running(mode):
return True
if mode == "ais":
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):
if index in {_indices[m] for m in running() if m in _indices}:
continue
try:
proc = subprocess.Popen(
build(index),
_COMMANDS[mode](index),
preexec_fn=os.setsid,
stdout=subprocess.PIPE if needs_stdout else None,
text=True if needs_stdout else None,
@@ -169,46 +174,84 @@ def start(mode: str) -> bool:
pass
if needs_stdout:
threading.Thread(target=_ais_reader, args=(proc,), daemon=True).start()
_proc = proc
_save_state(proc.pid, mode)
_procs[mode] = proc
_indices[mode] = index
_STATE_DIR.mkdir(parents=True, exist_ok=True)
_pgid_file(mode).write_text(str(proc.pid))
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?)")
LOGGER.info(f"Secondary SDR decoder mode={mode!r}: no free SDR left")
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:
def _stop_one(mode: str) -> None:
proc = _procs.pop(mode, None)
_indices.pop(mode, None)
if proc is not None:
try:
_proc.wait(timeout=5)
os.killpg(proc.pid, signal.SIGTERM)
except OSError:
pass
try:
proc.wait(timeout=5)
except subprocess.TimeoutExpired:
pass
_proc = None
_pgid_file(mode).unlink(missing_ok=True)
def start(mode: str) -> bool:
with _lock:
return _start_one(mode)
def stop(mode: Optional[str] = None) -> None:
with _lock:
for m in [mode] if mode else list(_procs):
_stop_one(m)
def apply(priority: List[str]) -> List[str]:
"""Run decoders in priority order until SDRs run out; stop everything else.
Unchanged when the running set is already the right prefix of the list, so
re-applying the same priority (every restart/config push) is a no-op
rather than a decoder restart.
"""
for m in priority:
if m not in _COMMANDS:
raise ValueError(f"Unknown secondary SDR mode: {m!r}")
with _lock:
live = set(running())
if live and live == set(priority[:len(live)]):
# Already running the head of the list; only try to extend it.
for m in priority[len(live):]:
if not _start_one(m):
break
return running()
for m in list(_procs):
_stop_one(m)
for m in priority:
if not _start_one(m):
break
return running()
def sdr_count() -> Optional[int]:
"""RTL-SDR dongles on USB, same heuristic as install.sh. This container has
usbutils; op25's :stable image doesn't, so its /op25/devices can't answer."""
try:
os.remove(_PGID_FILE)
except OSError:
pass
try:
os.remove(_MODE_FILE)
except OSError:
pass
return True
out = subprocess.run(["lsusb"], capture_output=True, text=True, timeout=5).stdout
except Exception:
return None
return sum(1 for line in out.splitlines() if "0bda:2838" in line or "0bda:2832" in line)
def status() -> Dict[str, Any]:
running = is_running()
live = running()
return {
"status": "running" if running else "stopped",
"mode": _read_mode() if running else None,
"running": [{"mode": m, "sdr_index": _indices.get(m)} for m in live],
"sdr_count": sdr_count(),
}
@@ -250,11 +293,10 @@ def _read_adsb_snapshot() -> List[Dict[str, Any]]:
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": []}
with _ais_lock:
vessels = list(_ais_vessels.values()) if is_running("ais") else []
return {
"running": running(),
"aircraft": _read_adsb_snapshot() if is_running("adsb") else [],
"vessels": vessels,
}