Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6bf2c9dd93 | ||
|
|
575a99feb5 | ||
|
|
9b50ba8114 | ||
|
|
3ecc7eff1a | ||
|
|
358766d88b | ||
|
|
c82fc61910 |
@@ -234,14 +234,30 @@ 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()
|
||||
# Unreachable secondary-sdr: say so with explicit nulls. C2 only
|
||||
# overwrites keys present in the checkin, so omitting them would leave
|
||||
# a dead decoder showing "Running" forever.
|
||||
payload["secondary_sdr_running"] = None
|
||||
payload["sdr_devices"] = None
|
||||
payload["op25_sdr_serial"] = None
|
||||
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 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
|
||||
# node running op25 has at least one SDR, so 0 is never a real reading.
|
||||
count = secondary.get("sdr_count") if secondary is not None else None
|
||||
if not count and devices:
|
||||
count = devices.get("count") or None
|
||||
if count is not None:
|
||||
payload["sdr_count"] = count
|
||||
self._publish(self._t_checkin, payload, qos=1)
|
||||
|
||||
@@ -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]:
|
||||
"""
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -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:
|
||||
|
||||
+26
-12
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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 @@
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<!-- Secondary SDRs Card (node-26#9) -->
|
||||
<!-- SDRs Card (node-26#9, #11) -->
|
||||
<div class="glass-card">
|
||||
<div class="card-header">
|
||||
<div class="card-title">Secondary SDRs</div>
|
||||
<div class="card-title">SDRs</div>
|
||||
</div>
|
||||
<div class="data-row" style="align-items:center;">
|
||||
<span class="data-label">OP25 SDR</span>
|
||||
<select id="sdr-op25" class="sdr-select" aria-label="OP25 SDR"></select>
|
||||
</div>
|
||||
<p id="sdr-op25-note" style="color: var(--text-muted); font-size: 0.75rem; margin: 0.25rem 0 0.75rem;"></p>
|
||||
<p style="color: var(--text-muted); font-size: 0.8rem; margin: 0 0 0.5rem;">
|
||||
OP25 always keeps its own SDR. Every other SDR runs the next enabled item below, top first.
|
||||
<span id="sdr-summary"></span>
|
||||
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".
|
||||
</p>
|
||||
<div id="sdr-list"></div>
|
||||
<p id="sdr-warning" style="color: var(--danger); font-size: 0.8rem; margin: 0.5rem 0 0; display:none;"></p>
|
||||
<div style="display:flex; gap:0.75rem; align-items:center; margin-top:0.75rem;">
|
||||
<button id="sdr-save" class="btn btn-primary" style="padding:0.5rem 1rem;" onclick="saveSdrPriority()" disabled>Save priority</button>
|
||||
<button id="sdr-save" class="btn btn-primary" style="padding:0.5rem 1rem;" onclick="saveSdr()" disabled>Save</button>
|
||||
<span id="sdr-msg" style="color: var(--text-muted); font-size: 0.8rem;"></span>
|
||||
</div>
|
||||
</div>
|
||||
@@ -532,40 +548,81 @@
|
||||
}
|
||||
}
|
||||
|
||||
// ── 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 `<option value="">${esc(autoLabel)}</option>` +
|
||||
devs.map(d => `<option value="${esc(d.serial)}" ${d.serial === selected ? 'selected' : ''}>` +
|
||||
`SDR ${d.index + 1} · serial ${esc(d.serial)}${d.duplicate_serial ? ' (shared serial!)' : ''}</option>`).join('') +
|
||||
(selected && !known ? `<option value="${esc(selected)}" selected>serial ${esc(selected)} ` +
|
||||
`(${sdr.devices === null ? 'not reported' : 'not plugged in'})</option>` : '');
|
||||
}
|
||||
|
||||
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);
|
||||
const state = !row.enabled ? 'Off' : running ? 'Running' : (sdrDirty ? 'Unsaved' : 'Waiting for SDR');
|
||||
const pinMissing = row.pin && sdr.devices !== null && !sdr.devices.some(d => d.serial === row.pin);
|
||||
const state = !row.enabled ? 'Off'
|
||||
: sdrDirty ? 'Unsaved'
|
||||
: 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 = `
|
||||
<span class="sdr-rank">${row.enabled ? ++rank : ''}</span>
|
||||
<input type="checkbox" ${row.enabled ? 'checked' : ''} aria-label="Enable ${name}">
|
||||
<span class="sdr-name">${name}<small>${hint}</small></span>
|
||||
<select class="sdr-select" aria-label="${name} SDR">${deviceOptions(row.pin, 'Any spare SDR')}</select>
|
||||
<button class="sdr-move" aria-label="Move ${name} up" ${i === 0 ? 'disabled' : ''}>▲</button>
|
||||
<button class="sdr-move" aria-label="Move ${name} down" ${i === sdrRows.length - 1 ? 'disabled' : ''}>▼</button>
|
||||
<span class="sdr-state ${running && row.enabled ? 'on' : ''}">${state}</span>`;
|
||||
<span class="sdr-state ${state === 'Running' ? 'on' : ''}">${state}</span>`;
|
||||
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() {
|
||||
@@ -577,47 +634,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 || [];
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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"}
|
||||
@@ -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 == []
|
||||
@@ -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,
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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}
|
||||
Reference in New Issue
Block a user