From 9b6c64fbbf170823d129bd20dc4b228391b37c2c Mon Sep 17 00:00:00 2001 From: Logan Cusano Date: Sun, 27 Sep 2026 14:08:07 -0400 Subject: [PATCH 1/2] Secondary SDR priority: every spare SDR runs the next decoder in an ordered list 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 --- drb-edge-node/app/internal/mqtt_manager.py | 19 +- .../app/internal/secondary_priority.py | 29 +++ .../app/internal/secondary_sdr_client.py | 14 +- .../app/internal/telemetry_uplink.py | 9 +- drb-edge-node/app/main.py | 26 ++- drb-edge-node/app/models.py | 31 ++- drb-edge-node/app/routers/api.py | 28 ++- drb-edge-node/app/templates/index.html | 133 +++++++++++ .../tests/test_secondary_priority.py | 58 +++++ .../app/internal/decoder_control.py | 218 +++++++++++------- .../app/routers/secondary_controller.py | 24 +- 11 files changed, 472 insertions(+), 117 deletions(-) create mode 100644 drb-edge-node/app/internal/secondary_priority.py create mode 100644 drb-edge-node/tests/test_secondary_priority.py diff --git a/drb-edge-node/app/internal/mqtt_manager.py b/drb-edge-node/app/internal/mqtt_manager.py index 3e12b7a..cb311b9 100644 --- a/drb-edge-node/app/internal/mqtt_manager.py +++ b/drb-edge-node/app/internal/mqtt_manager.py @@ -210,10 +210,14 @@ class MQTTManager: logger.info("No API key on disk — requesting re-delivery from C2 server.") self._publish(self._t_key_request, {}, qos=1) + async def publish_checkin(self): + await self._publish_checkin() + async def _publish_checkin(self): from app.internal.discord_radio import radio_bot from app.internal.config_manager import load_node_config from app.internal.op25_client import op25_client + from app.internal.secondary_sdr_client import secondary_sdr_client config = load_node_config() devices = await op25_client.devices() payload = { @@ -228,11 +232,18 @@ class MQTTManager: "override_system_id": config.override_system_id, "enforce_override_timeout": config.enforce_override_timeout, "secondary_sdr_mode": config.secondary_sdr_mode, + "secondary_sdr_priority": config.secondary_sdr_priority, } - # Best-effort — the first node-initiated hardware-report field. Omit - # rather than guess if op25's control API is unreachable. - if devices is not None: - payload["sdr_count"] = devices.get("count") + secondary = await secondary_sdr_client.status() + if secondary is not None: + payload["secondary_sdr_running"] = [r["mode"] for r in secondary.get("running", [])] + # Best-effort hardware report — omit rather than guess. op25's :stable + # image has no lsusb, so fall back to the secondary-sdr container's count. + count = devices.get("count") if devices else None + if count is None and secondary is not None: + count = secondary.get("sdr_count") + if count is not None: + payload["sdr_count"] = count self._publish(self._t_checkin, payload, qos=1) def _publish(self, topic: str, payload: dict, qos: int = 0, retain: bool = False): diff --git a/drb-edge-node/app/internal/secondary_priority.py b/drb-edge-node/app/internal/secondary_priority.py new file mode 100644 index 0000000..4df8a5f --- /dev/null +++ b/drb-edge-node/app/internal/secondary_priority.py @@ -0,0 +1,29 @@ +import asyncio +from typing import Iterable, List, Optional + +from app.internal.config_manager import load_node_config, save_node_config +from app.internal.logger import logger +from app.internal.secondary_sdr_client import secondary_sdr_client +from app.models import normalize_secondary_priority + + +async def set_secondary_priority(priority: Iterable[str]) -> Optional[List[str]]: + """Persist and apply the node's secondary SDR priority — one path for the + local dashboard, a C2 command and a config push alike. Never touches op25: + changing what the spare dongles do must not interrupt P25 recording. + + Returns the decoders now running (None if the container is unreachable). + """ + ordered = normalize_secondary_priority(priority) + cfg = load_node_config() + cfg.secondary_sdr_priority = ordered + cfg.secondary_sdr_mode = ordered[0] if ordered else "none" + save_node_config(cfg) + + running = await secondary_sdr_client.apply(ordered) + logger.info(f"Secondary SDR priority set to {ordered!r}; running {running!r}") + + # Report straight away so C2's view doesn't wait for the next heartbeat. + from app.internal.mqtt_manager import mqtt_manager + asyncio.create_task(mqtt_manager.publish_checkin()) + return running diff --git a/drb-edge-node/app/internal/secondary_sdr_client.py b/drb-edge-node/app/internal/secondary_sdr_client.py index 402ce3b..5e6b06e 100644 --- a/drb-edge-node/app/internal/secondary_sdr_client.py +++ b/drb-edge-node/app/internal/secondary_sdr_client.py @@ -1,5 +1,5 @@ import httpx -from typing import Any, Dict, Optional +from typing import Any, Dict, List, Optional from app.config import settings from app.internal.logger import logger @@ -15,6 +15,18 @@ class SecondarySdrClient: def __init__(self): self.api_url = settings.secondary_sdr_api_url + async def apply(self, priority: List[str]) -> Optional[List[str]]: + """Run decoders down the priority list until SDRs run out. Returns the + modes actually running, or None if the container is unreachable.""" + try: + async with httpx.AsyncClient(timeout=30) as client: + r = await client.post(f"{self.api_url}/secondary/apply", json={"priority": priority}) + r.raise_for_status() + return r.json().get("running", []) + except Exception as e: + logger.error(f"Secondary SDR apply (priority={priority!r}) failed: {e}") + return None + async def start(self, mode: str) -> bool: try: async with httpx.AsyncClient(timeout=10) as client: diff --git a/drb-edge-node/app/internal/telemetry_uplink.py b/drb-edge-node/app/internal/telemetry_uplink.py index 4aa7d97..50063a7 100644 --- a/drb-edge-node/app/internal/telemetry_uplink.py +++ b/drb-edge-node/app/internal/telemetry_uplink.py @@ -32,13 +32,14 @@ async def _post_snapshot(path: str, body: dict) -> None: async def telemetry_uplink_loop(): while True: await asyncio.sleep(UPLINK_INTERVAL_SECONDS) - config = load_node_config() - if config.secondary_sdr_mode not in ("adsb", "ais"): + if not load_node_config().secondary_sdr_priority: continue snapshot = await secondary_sdr_client.data() if not snapshot: continue - if config.secondary_sdr_mode == "adsb" and snapshot.get("aircraft"): + # Several decoders can run at once (one per spare SDR), so forward + # whatever each produced rather than keying off a single mode. + if snapshot.get("aircraft"): await _post_snapshot("/telemetry/adsb", {"aircraft": snapshot["aircraft"]}) - elif config.secondary_sdr_mode == "ais" and snapshot.get("vessels"): + if snapshot.get("vessels"): await _post_snapshot("/telemetry/ais", {"vessels": snapshot["vessels"]}) diff --git a/drb-edge-node/app/main.py b/drb-edge-node/app/main.py index e3837c8..58c9bf1 100644 --- a/drb-edge-node/app/main.py +++ b/drb-edge-node/app/main.py @@ -158,6 +158,9 @@ async def on_command(payload: dict): ) elif action == "discord_leave": await radio_bot.leave() + elif action == "set_secondary_priority": + from app.internal.secondary_priority import set_secondary_priority + await set_secondary_priority(payload.get("priority") or []) elif action == "op25_restart": from app.internal.op25_client import op25_client await op25_client.stop() @@ -224,7 +227,8 @@ async def on_config_push(payload: dict): hardware_preset = payload.pop("hardware_preset", None) ppm_override = payload.pop("ppm_override", None) node_type = payload.pop("node_type", None) - secondary_sdr_mode = payload.pop("secondary_sdr_mode", None) + secondary_sdr_priority = payload.pop("secondary_sdr_priority", None) + secondary_sdr_mode = payload.pop("secondary_sdr_mode", None) # legacy single-mode C2 enforce_override_timeout = payload.pop("enforce_override_timeout", None) try: config = SystemConfig(**payload) @@ -244,8 +248,6 @@ async def on_config_push(payload: dict): node_cfg.ppm_override = float(ppm_override) if node_type is not None: node_cfg.node_type = node_type - if secondary_sdr_mode is not None: - node_cfg.secondary_sdr_mode = secondary_sdr_mode if enforce_override_timeout is not None: node_cfg.enforce_override_timeout = bool(enforce_override_timeout) save_node_config(node_cfg) @@ -260,12 +262,11 @@ async def on_config_push(payload: dict): await op25_client.start() logger.info(f"Config push applied: {config.name}") - if secondary_sdr_mode is not None: - from app.internal.secondary_sdr_client import secondary_sdr_client - if secondary_sdr_mode in ("adsb", "ais"): - await secondary_sdr_client.start(secondary_sdr_mode) - else: - await secondary_sdr_client.stop() + if secondary_sdr_priority is None and secondary_sdr_mode is not None: + secondary_sdr_priority = [secondary_sdr_mode] + if secondary_sdr_priority is not None: + from app.internal.secondary_priority import set_secondary_priority + await set_secondary_priority(secondary_sdr_priority) # --------------------------------------------------------------------------- @@ -322,10 +323,11 @@ async def lifespan(app: FastAPI): logger.warning(f"OP25 not ready yet (attempt {attempt + 1}/10), retrying in 3s…") await asyncio.sleep(3) - if node_cfg.secondary_sdr_mode in ("adsb", "ais"): + # After op25 has claimed its dongle, so the decoders only get the spares. + if node_cfg.secondary_sdr_priority: from app.internal.secondary_sdr_client import secondary_sdr_client - logger.info(f"Resuming secondary SDR (mode={node_cfg.secondary_sdr_mode!r}) after restart.") - await secondary_sdr_client.start(node_cfg.secondary_sdr_mode) + logger.info(f"Resuming secondary SDRs (priority={node_cfg.secondary_sdr_priority!r}) after restart.") + await secondary_sdr_client.apply(node_cfg.secondary_sdr_priority) heartbeat_task = asyncio.create_task(mqtt_manager.heartbeat_loop()) from app.internal.telemetry_uplink import telemetry_uplink_loop diff --git a/drb-edge-node/app/models.py b/drb-edge-node/app/models.py index 8cebf61..703edea 100644 --- a/drb-edge-node/app/models.py +++ b/drb-edge-node/app/models.py @@ -1,5 +1,5 @@ -from pydantic import BaseModel -from typing import Optional, Dict, Any +from pydantic import BaseModel, model_validator +from typing import Optional, Dict, Any, Iterable, List from enum import Enum from datetime import datetime @@ -23,6 +23,21 @@ class SystemConfig(BaseModel): config: Dict[str, Any] # OP25-compatible config blob passed through to op25-container +# Decoders a node can run on the SDRs beyond op25's (node-26#9). op25 always +# keeps its own dongle; each further dongle runs the next entry of the node's +# secondary_sdr_priority, so a 3-SDR node with ["adsb", "ais"] runs both. +SECONDARY_SDR_MODES = ("adsb", "ais") + + +def normalize_secondary_priority(items: Iterable[str]) -> List[str]: + """Known modes only, first occurrence wins, order preserved.""" + out: List[str] = [] + for m in items or []: + if m in SECONDARY_SDR_MODES and m not in out: + out.append(m) + return out + + class NodeConfig(BaseModel): node_id: str node_name: str @@ -34,12 +49,22 @@ class NodeConfig(BaseModel): hardware_preset: str = "rtl-sdr-v3" ppm_override: Optional[float] = None node_type: str = "fixed" # fixed or portable - secondary_sdr_mode: str = "none" # none | adsb | ais | op25_2 — requires a second physical SDR + secondary_sdr_priority: List[str] = [] # ordered; SDRs beyond op25's run these top-down + # Legacy single-mode field (pre-priority). Still written as priority[0] so + # an older C2 reading checkins sees something sensible; read only to + # migrate a node_config.json saved before priority existed. + secondary_sdr_mode: str = "none" enforce_override_timeout: bool = True override_system_id: Optional[str] = None override_config: Optional[SystemConfig] = None offline_call_buffer_size: int = 35 # max call_end events to buffer while MQTT is offline + @model_validator(mode="after") + def _migrate_secondary_mode(self) -> "NodeConfig": + if not self.secondary_sdr_priority and self.secondary_sdr_mode in SECONDARY_SDR_MODES: + self.secondary_sdr_priority = [self.secondary_sdr_mode] + return self + class CallEvent(BaseModel): call_id: str diff --git a/drb-edge-node/app/routers/api.py b/drb-edge-node/app/routers/api.py index ec96c83..50ec6d5 100644 --- a/drb-edge-node/app/routers/api.py +++ b/drb-edge-node/app/routers/api.py @@ -1,9 +1,9 @@ from fastapi import APIRouter, Depends, HTTPException, Body -from typing import Optional +from typing import List, Optional import asyncio import httpx from app.config import settings -from app.models import SystemConfig +from app.models import SystemConfig, SECONDARY_SDR_MODES from app.internal.op25_client import op25_client from app.internal.config_manager import load_node_config, save_node_config, apply_system_config from app.internal.call_recorder import call_recorder @@ -204,6 +204,30 @@ async def ack_override(timeout_minutes: int = Body(1440)): raise HTTPException(500, f"Failed to contact C2: {e}") +@router.get("/secondary") +async def get_secondary(): + """Secondary SDR priority plus what's actually running (node-26#9).""" + from app.internal.secondary_sdr_client import secondary_sdr_client + status = await secondary_sdr_client.status() + return { + "priority": load_node_config().secondary_sdr_priority, + "modes": list(SECONDARY_SDR_MODES), + "running": [r["mode"] for r in status.get("running", [])] if status else None, + "sdr_count": status.get("sdr_count") if status else None, + } + + +@router.post("/secondary/priority") +async def set_secondary_priority(priority: List[str] = Body(..., embed=True)): + """Set the ordered list the spare SDRs work through. Never restarts op25.""" + from app.internal.secondary_priority import set_secondary_priority as apply_priority + unknown = [m for m in priority if m not in SECONDARY_SDR_MODES] + if unknown: + raise HTTPException(400, f"Unknown secondary SDR mode(s): {unknown}") + running = await apply_priority(priority) + return {"ok": True, "priority": load_node_config().secondary_sdr_priority, "running": running} + + @router.post("/discord/join") async def discord_join(guild_id: int, channel_id: int): ok = await radio_bot.join(guild_id, channel_id) diff --git a/drb-edge-node/app/templates/index.html b/drb-edge-node/app/templates/index.html index 98c1795..2391461 100644 --- a/drb-edge-node/app/templates/index.html +++ b/drb-edge-node/app/templates/index.html @@ -105,6 +105,30 @@ background: rgba(255, 255, 255, 0.1); } + .sdr-row { + display: flex; + align-items: center; + gap: 0.6rem; + padding: 0.55rem 0; + border-bottom: 1px solid var(--glass-border); + } + .sdr-row:last-child { border-bottom: none; } + .sdr-rank { width: 1.2rem; color: var(--text-muted); font-size: 0.8rem; text-align: right; } + .sdr-name { flex: 1; } + .sdr-name small { display: block; color: var(--text-muted); font-size: 0.75rem; } + .sdr-move { + background: rgba(255, 255, 255, 0.05); + border: 1px solid var(--glass-border); + color: var(--text-main); + border-radius: 6px; + width: 1.9rem; + height: 1.9rem; + cursor: pointer; + } + .sdr-move:disabled { opacity: 0.3; cursor: default; } + .sdr-state { font-size: 0.75rem; min-width: 6.5rem; text-align: right; color: var(--text-muted); } + .sdr-state.on { color: var(--success); } + .grid { display: grid; grid-template-columns: repeat(auto-fit, minmax(320px, 1fr)); @@ -360,6 +384,22 @@ + +
+
+
Secondary SDRs
+
+

+ OP25 always keeps its own SDR. Every other SDR runs the next enabled item below, top first. + +

+
+
+ + +
+
+
@@ -492,6 +532,99 @@ } } + // ── Secondary SDR priority ──────────────────────────────────────────── + const SDR_LABELS = { + adsb: ['ADS-B', 'Aircraft · 1090 MHz'], + ais: ['AIS', 'Vessels · 162 MHz'], + }; + let sdrRows = []; // [{mode, enabled}] in display order + let sdrRunning = []; + let sdrDirty = false; + + function renderSdr() { + const list = document.getElementById('sdr-list'); + list.innerHTML = ''; + let rank = 0; + sdrRows.forEach((row, i) => { + const [name, hint] = SDR_LABELS[row.mode] || [row.mode, '']; + const running = sdrRunning.includes(row.mode); + const state = !row.enabled ? 'Off' : running ? 'Running' : (sdrDirty ? 'Unsaved' : 'Waiting for SDR'); + const el = document.createElement('div'); + el.className = 'sdr-row'; + el.innerHTML = ` + ${row.enabled ? ++rank : ''} + + ${name}${hint} + + + ${state}`; + const [box] = el.getElementsByTagName('input'); + const [up, down] = el.getElementsByTagName('button'); + box.onchange = () => { row.enabled = box.checked; markSdrDirty(); }; + up.onclick = () => { [sdrRows[i - 1], sdrRows[i]] = [sdrRows[i], sdrRows[i - 1]]; markSdrDirty(); }; + down.onclick = () => { [sdrRows[i + 1], sdrRows[i]] = [sdrRows[i], sdrRows[i + 1]]; markSdrDirty(); }; + list.appendChild(el); + }); + document.getElementById('sdr-save').disabled = !sdrDirty; + } + + function markSdrDirty() { + sdrDirty = true; + document.getElementById('sdr-msg').textContent = ''; + renderSdr(); + } + + async function loadSdr() { + if (sdrDirty) return; // never clobber unsaved edits with a poll + try { + const r = await fetch('/api/secondary'); + if (!r.ok) return; + const d = await r.json(); + sdrRunning = d.running || []; + sdrRows = [ + ...d.priority.map(mode => ({ mode, enabled: true })), + ...d.modes.filter(m => !d.priority.includes(m)).map(mode => ({ mode, enabled: false })), + ]; + const summary = document.getElementById('sdr-summary'); + if (d.running === null) { + summary.textContent = 'Secondary SDR service is not responding.'; + } else if (d.sdr_count != null) { + const spare = Math.max(d.sdr_count - 1, 0); + summary.textContent = `This node has ${d.sdr_count} SDR${d.sdr_count === 1 ? '' : 's'}, so ${spare} spare.`; + } + renderSdr(); + } catch (e) { + console.error('Secondary SDR load failed:', e); + } + } + + async function saveSdrPriority() { + const priority = sdrRows.filter(r => r.enabled).map(r => r.mode); + const btn = document.getElementById('sdr-save'); + const msg = document.getElementById('sdr-msg'); + btn.disabled = true; + msg.textContent = 'Applying…'; + try { + const r = await fetch('/api/secondary/priority', { + method: 'POST', + headers: { 'Content-Type': 'application/json' }, + body: JSON.stringify({ priority }), + }); + if (!r.ok) throw new Error(await r.text()); + const d = await r.json(); + sdrDirty = false; + msg.textContent = d.running === null ? 'Saved, but the secondary SDR service did not respond.' : 'Saved.'; + await loadSdr(); + } catch (e) { + console.error('Secondary SDR save failed:', e); + msg.textContent = 'Save failed.'; + btn.disabled = false; + } + } + + loadSdr(); + setInterval(loadSdr, 10000); + refresh(); setInterval(refresh, 2000); // Polling every 2 seconds diff --git a/drb-edge-node/tests/test_secondary_priority.py b/drb-edge-node/tests/test_secondary_priority.py new file mode 100644 index 0000000..fc70782 --- /dev/null +++ b/drb-edge-node/tests/test_secondary_priority.py @@ -0,0 +1,58 @@ +""" +node-26#9 — secondary SDR priority: an ordered list the SDRs beyond op25's work +through. Normalisation, migration from the legacy single-mode field, and the one +apply path shared by the dashboard, C2 commands and config pushes. +""" +import asyncio +from unittest.mock import AsyncMock, patch + +from app.models import NodeConfig, normalize_secondary_priority + + +def _cfg(**kw) -> NodeConfig: + return NodeConfig(node_id="n1", node_name="N1", lat=0.0, lon=0.0, **kw) + + +def test_normalize_keeps_order_drops_unknown_and_duplicates(): + assert normalize_secondary_priority(["ais", "bogus", "adsb", "ais"]) == ["ais", "adsb"] + assert normalize_secondary_priority([]) == [] + + +def test_legacy_single_mode_migrates_to_priority(): + assert _cfg(secondary_sdr_mode="adsb").secondary_sdr_priority == ["adsb"] + assert _cfg(secondary_sdr_mode="none").secondary_sdr_priority == [] + + +def test_explicit_priority_wins_over_legacy_mode(): + assert _cfg(secondary_sdr_mode="adsb", secondary_sdr_priority=["ais"]).secondary_sdr_priority == ["ais"] + + +def test_set_priority_persists_applies_and_never_touches_op25(tmp_path): + config_file = tmp_path / "node_config.json" + import app.internal.config_manager as cm + with patch.object(cm, "_CONFIG_FILE", config_file): + cm.save_node_config(_cfg(secondary_sdr_mode="adsb")) + from app.internal import secondary_priority as sp + with patch.object(sp.secondary_sdr_client, "apply", AsyncMock(return_value=["ais"])) as apply, \ + patch("app.internal.mqtt_manager.mqtt_manager.publish_checkin", AsyncMock()), \ + patch("app.internal.op25_client.op25_client.stop", AsyncMock()) as op25_stop: + running = asyncio.run(sp.set_secondary_priority(["ais", "adsb", "ais"])) + saved = cm.load_node_config() + + assert running == ["ais"] + apply.assert_awaited_once_with(["ais", "adsb"]) + op25_stop.assert_not_awaited() + assert saved.secondary_sdr_priority == ["ais", "adsb"] + assert saved.secondary_sdr_mode == "ais" + + +def test_clearing_priority_is_not_undone_by_legacy_migration(tmp_path): + config_file = tmp_path / "node_config.json" + import app.internal.config_manager as cm + with patch.object(cm, "_CONFIG_FILE", config_file): + cm.save_node_config(_cfg(secondary_sdr_mode="adsb")) + from app.internal import secondary_priority as sp + with patch.object(sp.secondary_sdr_client, "apply", AsyncMock(return_value=[])), \ + patch("app.internal.mqtt_manager.mqtt_manager.publish_checkin", AsyncMock()): + asyncio.run(sp.set_secondary_priority([])) + assert cm.load_node_config().secondary_sdr_priority == [] diff --git a/secondary-sdr-container/app/internal/decoder_control.py b/secondary-sdr-container/app/internal/decoder_control.py index 7c34418..2063c29 100644 --- a/secondary-sdr-container/app/internal/decoder_control.py +++ b/secondary-sdr-container/app/internal/decoder_control.py @@ -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, + } diff --git a/secondary-sdr-container/app/routers/secondary_controller.py b/secondary-sdr-container/app/routers/secondary_controller.py index 31492f6..f94117c 100644 --- a/secondary-sdr-container/app/routers/secondary_controller.py +++ b/secondary-sdr-container/app/routers/secondary_controller.py @@ -1,3 +1,5 @@ +from typing import List, Optional + from fastapi import APIRouter, HTTPException from pydantic import BaseModel @@ -11,9 +13,25 @@ class StartBody(BaseModel): mode: str # adsb | ais +class StopBody(BaseModel): + mode: Optional[str] = None # omit to stop every decoder + + +class ApplyBody(BaseModel): + priority: List[str] # ordered, e.g. ["adsb", "ais"] + + def create_secondary_router(): router = APIRouter() + @router.post("/apply") + async def apply(body: ApplyBody): + try: + live = decoder_control.apply(body.priority) + except ValueError as e: + raise HTTPException(status_code=400, detail=str(e)) + return {"running": live} + @router.post("/start") async def start(body: StartBody): try: @@ -21,12 +39,12 @@ def create_secondary_router(): except ValueError as e: raise HTTPException(status_code=400, detail=str(e)) if not ok: - raise HTTPException(status_code=500, detail="Failed to start secondary SDR decoder") + raise HTTPException(status_code=409, detail="No free SDR for this decoder") return {"status": f"secondary SDR started ({body.mode})"} @router.post("/stop") - async def stop(): - decoder_control.stop() + async def stop(body: Optional[StopBody] = None): + decoder_control.stop(body.mode if body else None) return {"status": "secondary SDR stopped"} @router.get("/status") From c82fc6191075a1b25d1e3ae2f20cdd20b38967a4 Mon Sep 17 00:00:00 2001 From: Logan Cusano Date: Sun, 27 Sep 2026 14:14:49 -0400 Subject: [PATCH 2/2] Local Secondary SDRs card: 'Unsaved' outranks 'Running'; unknown when unreachable QA blocker: a reordered but unsaved list still showed the old order's decoder as 'Running'. Also shows 'Unknown' rather than 'Waiting for SDR' when the secondary-sdr service doesn't answer (server-26#187). Co-Authored-By: Claude Opus 5.5 --- drb-edge-node/app/templates/index.html | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/drb-edge-node/app/templates/index.html b/drb-edge-node/app/templates/index.html index 2391461..ed8f6b2 100644 --- a/drb-edge-node/app/templates/index.html +++ b/drb-edge-node/app/templates/index.html @@ -547,8 +547,12 @@ let rank = 0; sdrRows.forEach((row, i) => { const [name, hint] = SDR_LABELS[row.mode] || [row.mode, '']; - const running = sdrRunning.includes(row.mode); - const state = !row.enabled ? 'Off' : running ? 'Running' : (sdrDirty ? 'Unsaved' : 'Waiting for SDR'); + const running = (sdrRunning || []).includes(row.mode); + // Unsaved first: a reorder must not keep claiming the old order is live. + const state = !row.enabled ? 'Off' + : sdrDirty ? 'Unsaved' + : sdrRunning === null ? 'Unknown' + : running ? 'Running' : 'Waiting for SDR'; const el = document.createElement('div'); el.className = 'sdr-row'; el.innerHTML = ` @@ -557,7 +561,7 @@ ${name}${hint} - ${state}`; + ${state}`; const [box] = el.getElementsByTagName('input'); const [up, down] = el.getElementsByTagName('button'); box.onchange = () => { row.enabled = box.checked; markSdrDirty(); }; @@ -580,7 +584,7 @@ const r = await fetch('/api/secondary'); if (!r.ok) return; const d = await r.json(); - sdrRunning = d.running || []; + sdrRunning = d.running; // null = secondary SDR service unreachable sdrRows = [ ...d.priority.map(mode => ({ mode, enabled: true })), ...d.modes.filter(m => !d.priority.includes(m)).map(mode => ({ mode, enabled: false })),