Files
node-26/secondary-sdr-container/app/internal/decoder_control.py
T
Logan CusanoandClaude Opus 5.5 9b50ba8114
CI / lint (push) Successful in 6s
CI / test (push) Successful in 41s
Pin each SDR service to a dongle by serial; OP25 always opens its SDR by serial
Fixes node-26#11. OP25's generated config said "rtl" (= whichever dongle
enumerates first), so on 2-SDR nodes a decoder could take OP25's dongle
and stop recording. Now:

- sdr_pins: {op25|adsb|ais: serial}, absent = automatic. Every OP25
  config generation (op25_client.generate_config) rewrites the device to
  rtl=<serial>: the pin, else the first dongle's serial, which is what
  "rtl" always opened. Left as "rtl" only for unknown/shared serials.
- secondary-sdr: /secondary/devices lists dongles + serials via librtlsdr
  (works while claimed). apply(priority, pins, reserved) never touches
  OP25's dongle, gives a pinned service only its own dongle, lets a
  higher-priority service take a spare from a lower one, and still runs a
  pinned lower-priority service when the top pick has no dongle.
- sdr_settings.py replaces secondary_priority.py: one apply path for the
  local dashboard, the new set_sdr_config C2 command (set_secondary_priority
  kept as an alias) and config pushes. OP25 restarts only when its own
  dongle changes. Checkin reports sdr_devices, sdr_pins, op25_sdr_serial.
- Local dashboard: 'SDRs' card with an OP25 SDR dropdown and a per-service
  dongle dropdown, duplicate-serial and double-pin warnings.

Verified: edge-node pytest 194 passed; secondary-sdr tests 7 passed;
flake8 clean; page JS passes node --check.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:35:22 -04:00

354 lines
13 KiB
Python

import ctypes
import ctypes.util
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__)
# 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_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
_STATE_DIR = Path("/tmp/secondary_sdr")
ADSB_JSON_DIR = Path("/tmp/adsb")
# 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 _pgid_file(mode: str) -> Path:
return _STATE_DIR / f"{mode}.pgid"
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)
_reap_orphans()
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", "5", # JSON Full: decoded fields (4 = sparse, "JSON" is rejected)
]
_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
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 is_running(mode: str) -> bool:
proc = _procs.get(mode)
return proc is not None and proc.poll() is None
def running() -> List[str]:
return [m for m in MODES if is_running(m)]
def _start_one(mode: str, candidates: Optional[List[int]] = None) -> bool:
"""Start one decoder on the first free SDR among `candidates` (default:
every index). 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()
needs_stdout = mode == "ais"
for index in candidates if candidates is not None else range(MAX_SDR_INDEX):
if index in {_indices[m] for m in running() if m in _indices}:
continue
try:
proc = subprocess.Popen(
_COMMANDS[mode](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()
_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.info(f"Secondary SDR decoder mode={mode!r}: no free SDR left")
return False
def _stop_one(mode: str) -> None:
proc = _procs.pop(mode, None)
_indices.pop(mode, None)
if proc is not None:
try:
os.killpg(proc.pid, signal.SIGTERM)
except OSError:
pass
try:
proc.wait(timeout=5)
except subprocess.TimeoutExpired:
pass
_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 _candidates(mode: str, pins: Dict[str, str], reserved: List[str], devs: List[Dict[str, Any]]) -> List[int]:
"""SDR indices `mode` may use. A pinned mode gets exactly its dongle; an
unpinned one gets any dongle that isn't op25's (reserved) or pinned to
another service. Without enumeration, fall back to probing every index."""
if not devs:
return [] if pins.get(mode) else list(range(MAX_SDR_INDEX))
if pins.get(mode):
idx = _index_of(pins[mode], devs)
return [] if idx is None else [idx]
taken = set(reserved) | {s for m, s in pins.items() if m != mode and s}
return [d["index"] for d in devs if d["serial"] not in taken]
def apply(priority: List[str], pins: Optional[Dict[str, str]] = None,
reserved: Optional[List[str]] = None) -> List[str]:
"""Run decoders in priority order until SDRs run out; stop everything else.
`pins` maps a mode to the serial of the dongle carrying its antenna;
`reserved` lists serials no decoder may touch (op25's). Pins only bind
enabled modes — a disabled service's dongle is free for the others.
Unchanged when the right decoders already run on allowed dongles, so
re-applying the same settings is a no-op rather than a restart.
"""
for m in priority:
if m not in _COMMANDS:
raise ValueError(f"Unknown secondary SDR mode: {m!r}")
pins = {m: s for m, s in (pins or {}).items() if m in priority and s}
reserved = [s for s in (reserved or []) if s]
devs = devices()
with _lock:
for m in running():
if m not in priority or _indices.get(m) not in _candidates(m, pins, reserved, devs):
_stop_one(m)
for i, m in enumerate(priority):
if m in running():
continue
cands = _candidates(m, pins, reserved, devs)
if _start_one(m, cands):
continue
# Out of free dongles: take one from the lowest-priority decoder
# holding a dongle this mode may use, which then gets its own turn.
lower = [x for x in priority[i + 1:] if x in running() and _indices.get(x) in cands]
if lower:
_stop_one(lower[-1])
_start_one(m, cands)
return running()
def devices() -> List[Dict[str, Any]]:
"""Every RTL-SDR on the node with its USB serial — readable even while a
dongle is claimed (op25's included), since it doesn't open the device.
Cheap dongles often ship with the same serial (00000001); those are
flagged, because pinning a service to a shared serial is ambiguous."""
try:
lib = ctypes.CDLL(ctypes.util.find_library("rtlsdr") or "librtlsdr.so.0")
lib.rtlsdr_get_device_name.restype = ctypes.c_char_p
out = []
for i in range(lib.rtlsdr_get_device_count()):
manufact, product, serial = (ctypes.create_string_buffer(256) for _ in range(3))
lib.rtlsdr_get_device_usb_strings(i, manufact, product, serial)
out.append({
"index": i,
"serial": serial.value.decode(errors="replace") or None,
"name": (lib.rtlsdr_get_device_name(i) or b"").decode(errors="replace"),
})
except Exception as e:
LOGGER.warning(f"SDR enumeration failed: {e}")
return []
serials = [d["serial"] for d in out]
for d in out:
d["duplicate_serial"] = d["serial"] is not None and serials.count(d["serial"]) > 1
return out
def _index_of(serial: str, devs: List[Dict[str, Any]]) -> Optional[int]:
"""Index of a uniquely-identified serial; None if absent or ambiguous."""
matches = [d["index"] for d in devs if d["serial"] == serial]
return matches[0] if len(matches) == 1 else None
def status() -> Dict[str, Any]:
devs = devices()
by_index = {d["index"]: d["serial"] for d in devs}
return {
"running": [
{"mode": m, "sdr_index": _indices.get(m), "serial": by_index.get(_indices.get(m))} for m in running()
],
"sdr_count": len(devs) if devs else None,
"devices": devs,
}
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]:
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,
}