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..ed8f6b2 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 @@ + +
+ OP25 always keeps its own SDR. Every other SDR runs the next enabled item below, top first. + +
+ +