diff --git a/drb-edge-node/app/internal/mqtt_manager.py b/drb-edge-node/app/internal/mqtt_manager.py index 81efc70..73aba28 100644 --- a/drb-edge-node/app/internal/mqtt_manager.py +++ b/drb-edge-node/app/internal/mqtt_manager.py @@ -234,9 +234,17 @@ class MQTTManager: "secondary_sdr_mode": config.secondary_sdr_mode, "secondary_sdr_priority": config.secondary_sdr_priority, } + payload["sdr_pins"] = config.sdr_pins secondary = await secondary_sdr_client.status() if secondary is not None: payload["secondary_sdr_running"] = [r["mode"] for r in secondary.get("running", [])] + if secondary.get("devices") is not None: + payload["sdr_devices"] = [ + {k: d.get(k) for k in ("index", "serial", "name", "duplicate_serial")} + for d in secondary["devices"] + ] + from app.internal.sdr_settings import op25_serial + payload["op25_sdr_serial"] = op25_serial(config, secondary["devices"]) # Best-effort hardware report — omit rather than guess. Prefer the # secondary-sdr container's count: op25's :stable image has no lsusb and # its /op25/devices answers 0 rather than "unknown" when that fails. A diff --git a/drb-edge-node/app/internal/op25_client.py b/drb-edge-node/app/internal/op25_client.py index a1cc0c7..c7956d7 100644 --- a/drb-edge-node/app/internal/op25_client.py +++ b/drb-edge-node/app/internal/op25_client.py @@ -80,10 +80,14 @@ class OP25Client: async with httpx.AsyncClient(timeout=10) as client: r = await client.post(f"{self.api_url}/op25/generate-config", json=config) r.raise_for_status() - return True except Exception as e: logger.error(f"OP25 generate-config failed: {e}") return False + # Every generated config opens its dongle by serial, never "first + # found" (node-26#11) — here so no generation path can skip it. + from app.internal.sdr_settings import pin_op25_device + await pin_op25_device() + return True async def poll_terminal(self) -> Optional[TerminalUpdate]: """ diff --git a/drb-edge-node/app/internal/sdr_settings.py b/drb-edge-node/app/internal/sdr_settings.py new file mode 100644 index 0000000..8d58369 --- /dev/null +++ b/drb-edge-node/app/internal/sdr_settings.py @@ -0,0 +1,102 @@ +""" +Which SDR does what on this node (node-26#9, node-26#11). + +- OP25 always has exactly one dongle. `sdr_pins["op25"]` names it by serial; + unpinned, OP25 keeps its historical "first dongle" — but that dongle is now + named by serial too, so the decoders can never take it out from under OP25. +- Every other dongle runs the next enabled service in `secondary_sdr_priority`. + `sdr_pins[mode]` optionally binds a service to the dongle carrying its + antenna; unpinned services take any spare. + +One apply path for the local dashboard, C2 commands and config pushes. +""" +import asyncio +import json +from pathlib import Path +from typing import Any, Dict, Iterable, List, Optional + +from app.config import settings +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 NodeConfig, normalize_sdr_pins, normalize_secondary_priority + +_OP25_CONFIG = Path(settings.config_path) / "active.cfg.json" + + +def op25_serial(cfg: NodeConfig, devs: Optional[List[Dict[str, Any]]]) -> Optional[str]: + """The serial of OP25's dongle: the pin, else the first detected dongle + (what OP25's plain "rtl" device string has always opened).""" + if cfg.sdr_pins.get("op25"): + return cfg.sdr_pins["op25"] + return devs[0]["serial"] if devs else None + + +async def pin_op25_device() -> Optional[str]: + """Rewrite OP25's generated config to open its dongle by serial. + + Runs after every generate-config (op25_client.generate_config). A plain + "rtl" means "device 0", and which dongle is device 0 depends on who opened + what first — that is how OP25 lost its SDR to readsb on 2026-09-27. + Left as "rtl" only when the serial is unknown or shared by two dongles. + """ + devs = await secondary_sdr_client.devices() + serial = op25_serial(load_node_config(), devs) + if not serial or (devs and sum(d["serial"] == serial for d in devs) > 1): + return None + try: + cfg = json.loads(_OP25_CONFIG.read_text()) + for dev in cfg.get("devices", []): + dev["args"] = f"rtl={serial}" + _OP25_CONFIG.write_text(json.dumps(cfg, indent=2)) + except Exception as e: + logger.error(f"Could not pin OP25 to SDR {serial}: {e}") + return None + return serial + + +async def apply_secondaries() -> Optional[List[str]]: + """Start/stop decoders to match the saved priority and pins, never on OP25's dongle.""" + cfg = load_node_config() + if not cfg.secondary_sdr_priority: + return await secondary_sdr_client.apply([], {}, []) + devs = await secondary_sdr_client.devices() + reserved = [s for s in [op25_serial(cfg, devs)] if s] + pins = {m: s for m, s in cfg.sdr_pins.items() if m != "op25"} + return await secondary_sdr_client.apply(cfg.secondary_sdr_priority, pins, reserved) + + +async def set_sdr_settings(priority: Optional[Iterable[str]] = None, + pins: Optional[Dict[str, Optional[str]]] = None) -> Optional[List[str]]: + """Persist and apply priority and/or pins. OP25 restarts only if its own + dongle changed; reordering the spare dongles never interrupts P25. + + Returns the decoders now running (None if the container is unreachable). + """ + cfg = load_node_config() + old_op25 = cfg.sdr_pins.get("op25") + if priority is not None: + ordered = normalize_secondary_priority(priority) + cfg.secondary_sdr_priority = ordered + cfg.secondary_sdr_mode = ordered[0] if ordered else "none" + if pins is not None: + cfg.sdr_pins = normalize_sdr_pins(pins) + save_node_config(cfg) + + if cfg.sdr_pins.get("op25") != old_op25: + from app.internal.op25_client import op25_client + # Free the new dongle first if a decoder holds it, then move OP25. + await secondary_sdr_client.apply([], {}, []) + serial = await pin_op25_device() + await op25_client.stop() + await asyncio.sleep(2) + await op25_client.start() + logger.info(f"OP25 moved to SDR {serial or 'rtl (first dongle)'}") + + running = await apply_secondaries() + logger.info(f"SDR settings: priority={cfg.secondary_sdr_priority!r} pins={cfg.sdr_pins!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_priority.py b/drb-edge-node/app/internal/secondary_priority.py deleted file mode 100644 index 4df8a5f..0000000 --- a/drb-edge-node/app/internal/secondary_priority.py +++ /dev/null @@ -1,29 +0,0 @@ -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 5e6b06e..3fab60e 100644 --- a/drb-edge-node/app/internal/secondary_sdr_client.py +++ b/drb-edge-node/app/internal/secondary_sdr_client.py @@ -15,12 +15,25 @@ 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.""" + async def devices(self) -> Optional[List[Dict[str, Any]]]: + """Every RTL-SDR on the node with its serial, or None if unreachable.""" + try: + async with httpx.AsyncClient(timeout=5) as client: + r = await client.get(f"{self.api_url}/secondary/devices") + r.raise_for_status() + return r.json().get("devices", []) + except Exception as e: + logger.error(f"Secondary SDR device list failed: {e}") + return None + + async def apply(self, priority: List[str], pins: Dict[str, str], reserved: List[str]) -> Optional[List[str]]: + """Run decoders down the priority list until SDRs run out, honouring + pins and never touching `reserved` (op25's) dongles. Returns the modes + actually running, or None if the container is unreachable.""" + body = {"priority": priority, "pins": pins, "reserved": reserved} try: async with httpx.AsyncClient(timeout=30) as client: - r = await client.post(f"{self.api_url}/secondary/apply", json={"priority": priority}) + r = await client.post(f"{self.api_url}/secondary/apply", json=body) r.raise_for_status() return r.json().get("running", []) except Exception as e: diff --git a/drb-edge-node/app/main.py b/drb-edge-node/app/main.py index 58c9bf1..eadb50e 100644 --- a/drb-edge-node/app/main.py +++ b/drb-edge-node/app/main.py @@ -5,7 +5,7 @@ from typing import Optional from fastapi import FastAPI from app.config import settings -from app.models import SystemConfig +from app.models import SystemConfig, normalize_sdr_pins, normalize_secondary_priority from app.internal.logger import logger from app.internal.mqtt_manager import mqtt_manager from app.internal import credentials @@ -158,9 +158,12 @@ 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 in ("set_sdr_config", "set_secondary_priority"): + from app.internal.sdr_settings import set_sdr_settings + try: + await set_sdr_settings(payload.get("priority"), payload.get("pins")) + except ValueError as e: + logger.error(f"Rejected SDR settings from C2: {e}") elif action == "op25_restart": from app.internal.op25_client import op25_client await op25_client.stop() @@ -228,6 +231,7 @@ async def on_config_push(payload: dict): ppm_override = payload.pop("ppm_override", None) node_type = payload.pop("node_type", None) secondary_sdr_priority = payload.pop("secondary_sdr_priority", None) + sdr_pins = payload.pop("sdr_pins", None) secondary_sdr_mode = payload.pop("secondary_sdr_mode", None) # legacy single-mode C2 enforce_override_timeout = payload.pop("enforce_override_timeout", None) try: @@ -250,6 +254,16 @@ async def on_config_push(payload: dict): node_cfg.node_type = node_type if enforce_override_timeout is not None: node_cfg.enforce_override_timeout = bool(enforce_override_timeout) + 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: + node_cfg.secondary_sdr_priority = normalize_secondary_priority(secondary_sdr_priority) + node_cfg.secondary_sdr_mode = (node_cfg.secondary_sdr_priority or ["none"])[0] + if sdr_pins is not None: + try: + node_cfg.sdr_pins = normalize_sdr_pins(sdr_pins) + except ValueError as e: + logger.error(f"Ignoring invalid sdr_pins in config push: {e}") save_node_config(node_cfg) from app.internal.op25_client import op25_client @@ -257,16 +271,16 @@ async def on_config_push(payload: dict): logger.error(f"Failed to generate OP25 config for {config.name}") return + # OP25's (pinned) dongle may be one a decoder holds right now: free the + # spares, restart OP25 on its own dongle, then hand the rest back out. + from app.internal.secondary_sdr_client import secondary_sdr_client + from app.internal.sdr_settings import apply_secondaries + await secondary_sdr_client.apply([], {}, []) await op25_client.stop() await asyncio.sleep(2) await op25_client.start() logger.info(f"Config push applied: {config.name}") - - 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) + await apply_secondaries() # --------------------------------------------------------------------------- @@ -325,9 +339,9 @@ async def lifespan(app: FastAPI): # 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 + from app.internal.sdr_settings import apply_secondaries 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) + await apply_secondaries() 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 703edea..bd06fb2 100644 --- a/drb-edge-node/app/models.py +++ b/drb-edge-node/app/models.py @@ -38,6 +38,19 @@ def normalize_secondary_priority(items: Iterable[str]) -> List[str]: return out +# Services that can be bound to a specific dongle by its USB serial. +SDR_PIN_KEYS = ("op25",) + SECONDARY_SDR_MODES + + +def normalize_sdr_pins(pins: Dict[str, Optional[str]]) -> Dict[str, str]: + """Known services only, blank = automatic (dropped). Two services can't + share a dongle, so a serial claimed twice raises.""" + out = {k: str(v).strip() for k, v in (pins or {}).items() if k in SDR_PIN_KEYS and v and str(v).strip()} + if len(set(out.values())) != len(out): + raise ValueError("Two services can't be pinned to the same SDR.") + return out + + class NodeConfig(BaseModel): node_id: str node_name: str @@ -50,6 +63,7 @@ class NodeConfig(BaseModel): ppm_override: Optional[float] = None node_type: str = "fixed" # fixed or portable secondary_sdr_priority: List[str] = [] # ordered; SDRs beyond op25's run these top-down + sdr_pins: Dict[str, str] = {} # service (op25/adsb/ais) -> dongle serial; absent = automatic # 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. diff --git a/drb-edge-node/app/routers/api.py b/drb-edge-node/app/routers/api.py index 50ec6d5..7cbf696 100644 --- a/drb-edge-node/app/routers/api.py +++ b/drb-edge-node/app/routers/api.py @@ -1,5 +1,6 @@ from fastapi import APIRouter, Depends, HTTPException, Body -from typing import List, Optional +from pydantic import BaseModel +from typing import Dict, List, Optional import asyncio import httpx from app.config import settings @@ -204,28 +205,42 @@ 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).""" +@router.get("/sdr") +async def get_sdr(): + """Which SDR does what: detected dongles, pins, priority, what's running.""" from app.internal.secondary_sdr_client import secondary_sdr_client + from app.internal.sdr_settings import op25_serial + cfg = load_node_config() + devs = await secondary_sdr_client.devices() status = await secondary_sdr_client.status() return { - "priority": load_node_config().secondary_sdr_priority, + "devices": devs, # None = secondary-sdr service unreachable + "op25_serial": op25_serial(cfg, devs), + "pins": cfg.sdr_pins, + "priority": cfg.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} +class SdrSettingsBody(BaseModel): + priority: Optional[List[str]] = None + pins: Optional[Dict[str, Optional[str]]] = None # service -> serial; null/"" = automatic + + +@router.post("/sdr") +async def set_sdr(body: SdrSettingsBody): + """Set priority and/or pins. Restarts OP25 only if OP25's own SDR changes.""" + from app.internal.sdr_settings import set_sdr_settings + if body.priority is not None: + unknown = [m for m in body.priority if m not in SECONDARY_SDR_MODES] + if unknown: + raise HTTPException(400, f"Unknown secondary SDR mode(s): {unknown}") + try: + running = await set_sdr_settings(body.priority, body.pins) + except ValueError as e: + raise HTTPException(400, str(e)) + return {"ok": True, "running": running} @router.post("/discord/join") diff --git a/drb-edge-node/app/templates/index.html b/drb-edge-node/app/templates/index.html index ed8f6b2..de7e04e 100644 --- a/drb-edge-node/app/templates/index.html +++ b/drb-edge-node/app/templates/index.html @@ -128,6 +128,16 @@ .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); } + .sdr-select { + background: rgba(255, 255, 255, 0.05); + border: 1px solid var(--glass-border); + color: var(--text-main); + border-radius: 6px; + padding: 0.3rem 0.4rem; + font-size: 0.8rem; + max-width: 13rem; + } + .sdr-select option { background: #111827; } .grid { display: grid; @@ -384,18 +394,24 @@ - +
-
Secondary SDRs
+
SDRs
+
+ OP25 SDR + +
+

- OP25 always keeps its own SDR. Every other SDR runs the next enabled item below, top first. - + Every other SDR runs the next enabled service, top first. Pin a service to the SDR + that has its antenna, or leave it on "Any spare SDR".

+
- +
@@ -532,44 +548,80 @@ } } - // ── Secondary SDR priority ──────────────────────────────────────────── + // ── SDRs: OP25's dongle, secondary priority and per-service pins ──────── const SDR_LABELS = { adsb: ['ADS-B', 'Aircraft · 1090 MHz'], ais: ['AIS', 'Vessels · 162 MHz'], }; - let sdrRows = []; // [{mode, enabled}] in display order - let sdrRunning = []; + let sdr = null; // last GET /api/sdr + let sdrRows = []; // [{mode, enabled, pin}] in display order + let sdrOp25Pin = ''; // '' = automatic let sdrDirty = false; + function esc(v) { + return String(v ?? '').replace(/[&<>"']/g, c => ({'&': '&', '<': '<', '>': '>', '"': '"', "'": '''}[c])); + } + + function deviceOptions(selected, autoLabel) { + const devs = sdr.devices || []; + const known = devs.some(d => d.serial === selected); + return `` + + devs.map(d => ``).join('') + + (selected && !known ? `` : ''); + } + function renderSdr() { + const op25 = document.getElementById('sdr-op25'); + op25.innerHTML = deviceOptions(sdrOp25Pin, 'Automatic (first SDR)'); + op25.onchange = () => { sdrOp25Pin = op25.value; markSdrDirty(); }; + document.getElementById('sdr-op25-note').textContent = sdrOp25Pin !== (sdr.pins.op25 || '') + ? 'Saving will restart OP25 on the selected SDR.' + : (sdr.op25_serial ? `OP25 is using serial ${sdr.op25_serial}.` : ''); + const list = document.getElementById('sdr-list'); list.innerHTML = ''; let rank = 0; + const running = sdr.running || []; sdrRows.forEach((row, i) => { const [name, hint] = SDR_LABELS[row.mode] || [row.mode, '']; - const running = (sdrRunning || []).includes(row.mode); - // Unsaved first: a reorder must not keep claiming the old order is live. + const pinMissing = row.pin && !(sdr.devices || []).some(d => d.serial === row.pin); const state = !row.enabled ? 'Off' : sdrDirty ? 'Unsaved' - : sdrRunning === null ? 'Unknown' - : running ? 'Running' : 'Waiting for SDR'; + : sdr.running === null ? 'Unknown' + : running.includes(row.mode) ? 'Running' + : pinMissing ? 'Pinned SDR missing' : '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 [pick] = el.getElementsByTagName('select'); const [up, down] = el.getElementsByTagName('button'); box.onchange = () => { row.enabled = box.checked; markSdrDirty(); }; + pick.onchange = () => { row.pin = pick.value; 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; + + const warn = document.getElementById('sdr-warning'); + const pinned = [sdrOp25Pin, ...sdrRows.map(r => r.pin)].filter(Boolean); + const problems = []; + if (sdr.devices === null) problems.push('The secondary SDR service is not responding, so SDRs can\'t be listed.'); + if ((sdr.devices || []).some(d => d.duplicate_serial)) { + problems.push('Two SDRs share a serial number, so they can\'t be told apart. Give each a unique serial (rtl_eeprom -s) before pinning.'); + } + if (new Set(pinned).size !== pinned.length) problems.push('Two services are pinned to the same SDR.'); + warn.textContent = problems.join(' '); + warn.style.display = problems.length ? 'block' : 'none'; + document.getElementById('sdr-save').disabled = !sdrDirty || new Set(pinned).size !== pinned.length; } function markSdrDirty() { @@ -581,47 +633,41 @@ async function loadSdr() { if (sdrDirty) return; // never clobber unsaved edits with a poll try { - const r = await fetch('/api/secondary'); + const r = await fetch('/api/sdr'); if (!r.ok) return; - const d = await r.json(); - sdrRunning = d.running; // null = secondary SDR service unreachable + sdr = await r.json(); + sdrOp25Pin = sdr.pins.op25 || ''; sdrRows = [ - ...d.priority.map(mode => ({ mode, enabled: true })), - ...d.modes.filter(m => !d.priority.includes(m)).map(mode => ({ mode, enabled: false })), + ...sdr.priority.map(mode => ({ mode, enabled: true, pin: sdr.pins[mode] || '' })), + ...sdr.modes.filter(m => !sdr.priority.includes(m)).map(mode => ({ mode, enabled: false, pin: sdr.pins[mode] || '' })), ]; - 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); + console.error('SDR load failed:', e); } } - async function saveSdrPriority() { - const priority = sdrRows.filter(r => r.enabled).map(r => r.mode); + async function saveSdr() { + const pins = { op25: sdrOp25Pin || null }; + sdrRows.forEach(r => { pins[r.mode] = r.pin || null; }); 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', { + const r = await fetch('/api/sdr', { method: 'POST', headers: { 'Content-Type': 'application/json' }, - body: JSON.stringify({ priority }), + body: JSON.stringify({ priority: sdrRows.filter(r => r.enabled).map(r => r.mode), pins }), }); - if (!r.ok) throw new Error(await r.text()); + if (!r.ok) throw new Error((await r.json().catch(() => ({}))).detail || r.statusText); 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.'; + console.error('SDR save failed:', e); + msg.textContent = `Save failed: ${e.message}`; btn.disabled = false; } } diff --git a/drb-edge-node/tests/test_sdr_settings.py b/drb-edge-node/tests/test_sdr_settings.py new file mode 100644 index 0000000..4eabc16 --- /dev/null +++ b/drb-edge-node/tests/test_sdr_settings.py @@ -0,0 +1,91 @@ +""" +node-26#9 / #11 — which SDR does what: secondary priority, per-service pins, +and OP25 always opening its dongle by serial. +""" +import asyncio +import json +from unittest.mock import AsyncMock, patch + +import pytest + +from app.models import NodeConfig, normalize_sdr_pins, normalize_secondary_priority + +DEVS = [ + {"index": 0, "serial": "69420", "name": "RTL", "duplicate_serial": False}, + {"index": 1, "serial": "00000001", "name": "RTL", "duplicate_serial": False}, +] + + +def _cfg(**kw) -> NodeConfig: + return NodeConfig(node_id="n1", node_name="N1", lat=0.0, lon=0.0, **kw) + + +def test_normalize_priority_keeps_order_drops_unknown_and_duplicates(): + assert normalize_secondary_priority(["ais", "bogus", "adsb", "ais"]) == ["ais", "adsb"] + + +def test_legacy_single_mode_migrates_to_priority(): + assert _cfg(secondary_sdr_mode="adsb").secondary_sdr_priority == ["adsb"] + assert _cfg(secondary_sdr_mode="adsb", secondary_sdr_priority=["ais"]).secondary_sdr_priority == ["ais"] + + +def test_pins_drop_blanks_and_unknown_services(): + assert normalize_sdr_pins({"op25": "00000001", "adsb": "", "ais": None, "x": "1"}) == {"op25": "00000001"} + + +def test_two_services_cannot_share_a_dongle(): + with pytest.raises(ValueError): + normalize_sdr_pins({"op25": "69420", "adsb": "69420"}) + + +@pytest.fixture +def node(tmp_path): + """Isolated node_config.json + OP25 active.cfg.json, mocked decoders/op25/mqtt.""" + import app.internal.config_manager as cm + from app.internal import sdr_settings as ss + op25_cfg = tmp_path / "active.cfg.json" + op25_cfg.write_text(json.dumps({"devices": [{"args": "rtl", "name": "sdr"}]})) + with patch.object(cm, "_CONFIG_FILE", tmp_path / "node_config.json"), \ + patch.object(ss, "_OP25_CONFIG", op25_cfg), \ + patch.object(ss.secondary_sdr_client, "devices", AsyncMock(return_value=DEVS)), \ + patch.object(ss.secondary_sdr_client, "apply", AsyncMock(return_value=["adsb"])) as apply, \ + patch("app.internal.op25_client.op25_client.stop", AsyncMock()) as op25_stop, \ + patch("app.internal.op25_client.op25_client.start", AsyncMock()), \ + patch("app.internal.mqtt_manager.mqtt_manager.publish_checkin", AsyncMock()), \ + patch("asyncio.sleep", AsyncMock()): + cm.save_node_config(_cfg()) + yield {"cm": cm, "ss": ss, "apply": apply, "op25_stop": op25_stop, + "op25_args": lambda: json.loads(op25_cfg.read_text())["devices"][0]["args"]} + + +def test_unpinned_op25_is_still_opened_by_serial_of_the_first_dongle(node): + assert asyncio.run(node["ss"].pin_op25_device()) == "69420" + assert node["op25_args"]() == "rtl=69420" + + +def test_pinned_op25_opens_its_pinned_dongle(node): + node["cm"].save_node_config(_cfg(sdr_pins={"op25": "00000001"})) + asyncio.run(node["ss"].pin_op25_device()) + assert node["op25_args"]() == "rtl=00000001" + + +def test_shared_serial_leaves_op25_on_plain_rtl(node): + dup = [dict(d, serial="00000001") for d in DEVS] + with patch.object(node["ss"].secondary_sdr_client, "devices", AsyncMock(return_value=dup)): + assert asyncio.run(node["ss"].pin_op25_device()) is None + assert node["op25_args"]() == "rtl" + + +def test_priority_change_never_restarts_op25_and_reserves_its_dongle(node): + node["cm"].save_node_config(_cfg(sdr_pins={"op25": "00000001"})) + asyncio.run(node["ss"].set_sdr_settings(priority=["adsb", "ais"])) + node["op25_stop"].assert_not_awaited() + node["apply"].assert_awaited_with(["adsb", "ais"], {}, ["00000001"]) + + +def test_moving_op25_restarts_it_on_the_new_dongle(node): + asyncio.run(node["ss"].set_sdr_settings(priority=["adsb"], pins={"op25": "00000001", "adsb": "69420"})) + node["op25_stop"].assert_awaited_once() + assert node["op25_args"]() == "rtl=00000001" + node["apply"].assert_awaited_with(["adsb"], {"adsb": "69420"}, ["00000001"]) + assert node["cm"].load_node_config().sdr_pins == {"op25": "00000001", "adsb": "69420"} diff --git a/drb-edge-node/tests/test_secondary_priority.py b/drb-edge-node/tests/test_secondary_priority.py deleted file mode 100644 index fc70782..0000000 --- a/drb-edge-node/tests/test_secondary_priority.py +++ /dev/null @@ -1,58 +0,0 @@ -""" -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 2063c29..3d4a568 100644 --- a/secondary-sdr-container/app/internal/decoder_control.py +++ b/secondary-sdr-container/app/internal/decoder_control.py @@ -1,3 +1,5 @@ +import ctypes +import ctypes.util import json import os import signal @@ -141,8 +143,9 @@ 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.""" +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): @@ -152,7 +155,7 @@ def _start_one(mode: str) -> bool: _ais_vessels.clear() needs_stdout = mode == "ais" - for index in range(MAX_SDR_INDEX): + 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: @@ -211,47 +214,95 @@ def stop(mode: Optional[str] = None) -> None: _stop_one(m) -def apply(priority: List[str]) -> List[str]: +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. - 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. + `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: - 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 + 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 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.""" +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: - 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) + 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]: - live = running() + devs = devices() + by_index = {d["index"]: d["serial"] for d in devs} return { - "running": [{"mode": m, "sdr_index": _indices.get(m)} for m in live], - "sdr_count": sdr_count(), + "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, } diff --git a/secondary-sdr-container/app/routers/secondary_controller.py b/secondary-sdr-container/app/routers/secondary_controller.py index f94117c..23b98db 100644 --- a/secondary-sdr-container/app/routers/secondary_controller.py +++ b/secondary-sdr-container/app/routers/secondary_controller.py @@ -1,4 +1,4 @@ -from typing import List, Optional +from typing import Dict, List, Optional from fastapi import APIRouter, HTTPException from pydantic import BaseModel @@ -19,6 +19,8 @@ class StopBody(BaseModel): class ApplyBody(BaseModel): priority: List[str] # ordered, e.g. ["adsb", "ais"] + pins: Dict[str, str] = {} # mode -> serial of the dongle with its antenna + reserved: List[str] = [] # serials no decoder may touch (op25's) def create_secondary_router(): @@ -27,7 +29,7 @@ def create_secondary_router(): @router.post("/apply") async def apply(body: ApplyBody): try: - live = decoder_control.apply(body.priority) + live = decoder_control.apply(body.priority, body.pins, body.reserved) except ValueError as e: raise HTTPException(status_code=400, detail=str(e)) return {"running": live} @@ -51,6 +53,10 @@ def create_secondary_router(): async def get_status(): return decoder_control.status() + @router.get("/devices") + async def get_devices(): + return {"devices": decoder_control.devices()} + @router.get("/data") async def get_data(): return decoder_control.data() diff --git a/secondary-sdr-container/tests/test_decoder_control.py b/secondary-sdr-container/tests/test_decoder_control.py new file mode 100644 index 0000000..a2ae3a6 --- /dev/null +++ b/secondary-sdr-container/tests/test_decoder_control.py @@ -0,0 +1,84 @@ +""" +node-26#11 — which dongle each secondary decoder may use. + +Run from secondary-sdr-container/: PYTHONPATH=app python -m pytest -q tests +""" +from unittest.mock import patch + +import pytest + +from internal import decoder_control as dc + +# radio-box's real pair: op25's dongle and the one with the 1090 antenna. +DEVS = [ + {"index": 0, "serial": "69420", "name": "RTL", "duplicate_serial": False}, + {"index": 1, "serial": "00000001", "name": "RTL", "duplicate_serial": False}, +] + + +def test_unpinned_mode_never_gets_op25s_dongle(): + assert dc._candidates("adsb", {}, ["00000001"], DEVS) == [0] + + +def test_pinned_mode_gets_exactly_its_dongle(): + assert dc._candidates("ais", {"ais": "69420"}, ["00000001"], DEVS) == [0] + + +def test_unpinned_mode_skips_a_dongle_pinned_to_another_service(): + assert dc._candidates("ais", {"adsb": "69420"}, ["00000001"], DEVS) == [] + + +def test_missing_or_ambiguous_pin_gets_nothing(): + assert dc._candidates("adsb", {"adsb": "nope"}, [], DEVS) == [] + dup = [dict(d, serial="00000001") for d in DEVS] + assert dc._candidates("adsb", {"adsb": "00000001"}, [], dup) == [] + + +class FakeDecoders: + """Stands in for real processes: one decoder per free index.""" + + def __init__(self): + self.live = {} + + def start(self, mode, candidates=None): + free = [i for i in candidates if i not in self.live.values()] + if not free: + return False + self.live[mode] = free[0] + return True + + def stop(self, mode): + self.live.pop(mode, None) + + +@pytest.fixture +def fake(): + f = FakeDecoders() + with patch.object(dc, "devices", return_value=DEVS), \ + patch.object(dc, "_start_one", side_effect=f.start), \ + patch.object(dc, "_stop_one", side_effect=f.stop), \ + patch.object(dc, "running", side_effect=lambda: [m for m in dc.MODES if m in f.live]), \ + patch.dict(dc._indices, clear=True): + dc._indices.update(f.live) + yield f + + +def _apply(fake, priority, pins=None): + dc.apply(priority, pins, ["00000001"]) + dc._indices.clear() + dc._indices.update(fake.live) + return fake.live + + +def test_one_spare_goes_to_the_top_pick(fake): + assert _apply(fake, ["ais", "adsb"]) == {"ais": 0} + + +def test_reordering_hands_the_spare_to_the_new_top_pick(fake): + _apply(fake, ["adsb", "ais"]) + assert _apply(fake, ["ais", "adsb"]) == {"ais": 0} + + +def test_pinned_lower_priority_still_runs_when_top_pick_has_no_dongle(fake): + # AIS ranked first but its only candidate is ADS-B's pinned antenna dongle. + assert _apply(fake, ["ais", "adsb"], {"adsb": "69420"}) == {"adsb": 0}