Author SHA1 Message Date
Logan CusanoandClaude Opus 5.5 9b6c64fbbf Secondary SDR priority: every spare SDR runs the next decoder in an ordered list
CI / lint (push) Successful in 6s
CI / test (push) Successful in 39s
Replaces the single secondary_sdr_mode with secondary_sdr_priority, e.g.
["adsb", "ais"]. OP25 always keeps its own dongle; the secondary-sdr
container starts decoders top-down until it runs out of free SDRs, so a
3-SDR node runs ADS-B and AIS at once and a 2-SDR node runs the top pick.

- secondary-sdr: one decoder per mode, POST /secondary/apply(priority)
  (no-op when the right prefix is already running), orphaned decoders
  from a uvicorn reload are reaped on start, /status reports sdr_count
  via lsusb (op25's :stable image has none).
- edge-node: one apply path (set_secondary_priority) for the local
  dashboard, a new C2 'set_secondary_priority' MQTT command, and config
  pushes; it never restarts op25. Legacy mode migrates on load. Checkin
  reports priority, what's running, and sdr_count. Uplink forwards
  aircraft and vessels whenever either is present.
- Local dashboard: 'Secondary SDRs' card to enable/reorder/save.

Verified: edge-node pytest 190 passed, flake8 clean.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:08:07 -04:00
Logan CusanoandClaude Opus 5.5 1c74edffea secondary-sdr: AIS-catcher output flag is -o 5, not -o JSON
CI / lint (push) Successful in 7s
CI / test (push) Successful in 48s
Build secondary-sdr / build (push) Successful in 47m56s
AIS-catcher rejects '-o JSON' (Unknown message format) and exited on
every index. -o 5 is JSON Full, whose field names the reader already
expects. Verified on radio-box: 16 AIS transmitters (Hudson AtoN buoys)
decoded on a 9cm whip. (node-26#9)

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 13:32:23 -04:00
Logan CusanoandClaude Opus 5.5 c0a05a8ad6 Merge feat/second-sdr-adsb-ais: second SDR (ADS-B verified on hardware, AIS untested) (node-26#9)
CI / lint (push) Successful in 7s
Build op25 / build (push) Failing after 31s
CI / test (push) Successful in 48s
Build edge-node / build (push) Successful in 8m40s
Build secondary-sdr / build (push) Successful in 47m9s
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 13:17:00 -04:00
Logan CusanoandClaude Opus 5.5 52d11c1546 ci: build and publish the secondary-sdr image
Without it, compose references secondary-sdr:latest that no registry
has, and 'make pull' on every node fails once this branch lands.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 13:16:58 -04:00
Logan CusanoandClaude Opus 5.5 b54624e176 secondary-sdr: fix ADS-B on real hardware (tested on radio-box)
CI / lint (push) Successful in 6s
CI / lint (pull_request) Successful in 7s
CI / test (push) Successful in 39s
CI / test (pull_request) Successful in 40s
First hardware test of node-26#9/#10 found three blockers:
- antirez/dump1090 has no --write-json, so ADS-B mode exited on start.
  Swapped to wiedehopf/readsb; map alt_baro/gs (old names as fallback).
- op25 is not always on RTL-SDR index 0 (radio-box: op25 on 1, 0 free).
  start() now tries each index and keeps the first decoder that stays up.
  AIS-catcher index flag fixed to -d:N ("-d N" selects by serial).
- status() reported "running" for a decoder that died on startup: the
  zombie still answered killpg(pgid, 0). Liveness now via Popen.poll().

install.sh blacklists dvb_usb_rtl28xxu, which claimed the second dongle.

Verified on radio-box: 831 msgs/min, 6 aircraft (3 with position) on a
9cm whip; op25 unaffected. AIS still untested.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 12:34:18 -04:00
Logan CusanoandClaude Sonnet 5 52f31bbcc0 Wire AIS end to end: AIS-catcher support in secondary-sdr container
CI / lint (push) Successful in 5s
CI / lint (pull_request) Successful in 5s
CI / test (push) Successful in 38s
CI / test (pull_request) Successful in 38s
node-26#9. decoder_control.py's AIS mode now actually launches AIS-catcher
instead of 400ing: a background thread reads its stdout JSON stream
(one message per line) and keeps a live in-memory snapshot keyed by mmsi,
since AIS-catcher streams rather than writing a periodic file the way
dump1090's --write-json does. Messages are merged onto the existing entry
per mmsi rather than replacing it, since AIS-catcher emits static data
(name) and position reports (lat/lon) as separate message types — a naive
overwrite would blank the name back to null on every position-only update
(caught by a standalone unit check against a fake stdout stream before
this fix, not by pytest — this container has no test suite yet).

edge-node's telemetry_uplink_loop now posts non-empty vessel snapshots to
the new /telemetry/ais the same way it already does for aircraft.

UNVERIFIED against real hardware/binary in this session, same caveat as
the ADS-B commit — AIS-catcher's JSON field names are believed correct
from its docs, not confirmed against a real capture.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-20 19:10:30 -04:00
Logan CusanoandClaude Sonnet 5 b2e3804dc3 Wire ADS-B end to end: secondary-sdr container running dump1090
node-26#9. New secondary-sdr-container claims the node's SECOND physical
SDR (RTL-SDR index 1 — op25 always claims index 0; no serial-based binding
yet, same gap op25 itself has). Its control API (start/stop/status/data,
mirroring op25_controller.py) launches dump1090 in adsb mode and exposes
the decoded aircraft.json snapshot; AIS mode 400s until it's wired next.

edge-node: on_config_push starts/stops it when secondary_sdr_mode changes,
lifespan resumes it after a restart if already configured, and a new
telemetry_uplink_loop polls its /secondary/data every 10s and POSTs
non-empty snapshots to C2's new /telemetry/adsb (same bearer-key pattern
call_recorder.py already uses for audio upload).

UNVERIFIED: this container has not been built or run against real hardware
in this session (sandboxed authoring machine, no docker) — dump1090's
--write-json field names are believed correct from its docs but not
confirmed against a real capture. Build + hardware smoke test before this
reaches a real node.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-20 19:10:24 -04:00
Logan CusanoandClaude Sonnet 5 8f5fca8757 Add second-SDR plumbing: secondary_sdr_mode config + hardware reporting
Adds a per-node secondary_sdr_mode config field (none|adsb|ais|op25_2),
applied the same way as hardware_preset/ppm_override so it survives system
reassignment. op25-container gets a GET /devices endpoint that counts
connected SDRs via lsusb; the edge-node checkin now reports sdr_count and
secondary_sdr_mode up to the server, the first node-initiated hardware
report (everything else was C2 pushing config down). Tracked as node-26#9.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-20 19:10:17 -04:00
logan a53000a092 ci: publish node images to a public registry on a version tag, DISABLED (#8)
CI / lint (push) Successful in 6s
CI / test (push) Successful in 37s
2026-09-06 23:49:02 -04:00
Logan CusanoandClaude Sonnet 5 bee9e173e6 ci: workflow to publish node images to a public registry on a version tag (DISABLED)
CI / lint (pull_request) Successful in 6s
CI / lint (push) Successful in 7s
CI / test (pull_request) Successful in 37s
CI / test (push) Successful in 39s
git.vpn.cusano.net is now behind REQUIRE_SIGNIN_VIEW (INCIDENT-2026-09-06),
so a fresh Pi can no longer anonymously `docker pull` the node images or
fetch install.sh. Plan (owner, option 3): on a `v*` tag, build the three
node images (edge-node, icecast, op25-client, arm64) and push them to a
PUBLIC registry; source and the private Gitea registry stay walled.

This workflow is wired but INERT — the job is gated on
`vars.NODE_PUBLIC_PUBLISH == 'true'`, which is unset. It triggers on tags
and skips. Enabling is three repo variables + two secrets, documented in the
file header; no code change. Reuses the existing Gitea buildcache refs so
op25 doesn't recompile from scratch.

Not turning it on now — still building the core.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-06 22:39:42 -04:00
26 changed files with 1129 additions and 5 deletions
+7
View File
@@ -111,6 +111,13 @@ OP25_TERMINAL_URL=http://localhost:8081
# for local development off a real node; leave false everywhere else. # for local development off a real node; leave false everywhere else.
OP25_DEBUG_EXPOSE=false OP25_DEBUG_EXPOSE=false
# Secondary SDR container (node-26#9) — only matters if a second physical SDR
# is present and secondary_sdr_mode is set to adsb|ais via the edge dashboard
# or C2. Usually no need to change.
SECONDARY_SDR_API_URL=http://localhost:8002
# Same caveat as OP25_DEBUG_EXPOSE — debugging aid only, leave false.
SECONDARY_SDR_DEBUG_EXPOSE=false
# --- Local dashboard / API login --------------------------------------------- # --- Local dashboard / API login ---------------------------------------------
# Protects the node's local dashboard (port 80) and JSON API. The node is # Protects the node's local dashboard (port 80) and JSON API. The node is
# reachable by anyone on whatever site's LAN it's deployed to, so this MUST be # reachable by anyone on whatever site's LAN it's deployed to, so this MUST be
+52
View File
@@ -0,0 +1,52 @@
name: Build secondary-sdr
on:
workflow_dispatch:
push:
branches: [main, master]
paths:
- "secondary-sdr-container/**"
jobs:
build:
runs-on: ubuntu-latest
permissions:
contents: read
packages: write
env:
CONTAINER_NAME: secondary-sdr
steps:
- uses: actions/checkout@v4
- uses: docker/setup-qemu-action@v3
- uses: docker/setup-buildx-action@v3
with:
config-inline: |
[registry."git.vpn.cusano.net"]
http = false
insecure = false
- uses: docker/login-action@v3
with:
registry: git.vpn.cusano.net
username: ${{ gitea.actor }}
password: ${{ secrets.BUILD_TOKEN }}
- name: Get version
id: meta
run: |
echo "REPO_NAME=$(echo ${GITHUB_REPOSITORY} | awk -F'/' '{print $2}')" >> $GITHUB_OUTPUT
echo "VERSION=$(git describe --tags --always | sed 's/^v//')" >> $GITHUB_OUTPUT
- uses: docker/build-push-action@v6
with:
context: ./secondary-sdr-container
file: ./secondary-sdr-container/Dockerfile
platforms: linux/arm64
push: true
tags: |
git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:${{ steps.meta.outputs.VERSION }}
git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:latest
cache-from: type=registry,ref=git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:buildcache
cache-to: type=registry,ref=git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:buildcache,mode=max
+103
View File
@@ -0,0 +1,103 @@
name: Publish public images
# Push the three node images to a PUBLIC registry on a version tag, so a fresh
# Pi can `docker pull` them without a Gitea login (git.vpn.cusano.net is now
# behind REQUIRE_SIGNIN_VIEW — see INCIDENT-2026-09-06). Source + the private
# registry stay walled; only the built node images go public.
#
# ── DISABLED ────────────────────────────────────────────────────────────────
# The job is gated on `vars.NODE_PUBLIC_PUBLISH == 'true'`. Until that repo
# variable is set the workflow triggers on tags but the job is skipped, so
# this file is wired and inert. We're still building the core; flip it on when
# self-serve node install is actually needed.
#
# To enable:
# 1. Repo → Settings → Actions → Variables:
# NODE_PUBLIC_PUBLISH = true
# PUBLIC_REGISTRY = ghcr.io (or docker.io)
# PUBLIC_NAMESPACE = <org-or-user> (images land at <ns>/drb-<name>)
# 2. Repo → Settings → Actions → Secrets:
# PUBLIC_REGISTRY_USER = <push user>
# PUBLIC_REGISTRY_TOKEN = <push token / PAT with write:packages>
# 3. Re-push a tag (or run this workflow via workflow_dispatch).
# ---------------------------------------------------------------------------
on:
workflow_dispatch:
push:
tags:
- "v*"
concurrency:
group: publish-public-${{ github.ref }}
cancel-in-progress: false
jobs:
publish:
# Inert until the repo variable is set. Do NOT convert this to `if: false`
# — the variable is the switch, no code change needed to go live.
if: ${{ vars.NODE_PUBLIC_PUBLISH == 'true' }}
runs-on: ubuntu-latest
permissions:
contents: read
packages: write
strategy:
fail-fast: false
matrix:
include:
- name: edge-node
context: ./drb-edge-node
file: ./drb-edge-node/Dockerfile
cache_name: edge-node
- name: icecast
context: ./icecast
file: ./icecast/Dockerfile
cache_name: icecast
- name: op25-client
context: ./op25-container
file: ./op25-container/Dockerfile
cache_name: op25-client
steps:
- uses: actions/checkout@v4
with:
fetch-depth: 0 # need tags for `git describe`
- uses: docker/setup-qemu-action@v3
- uses: docker/setup-buildx-action@v3
with:
config-inline: |
[registry."git.vpn.cusano.net"]
http = false
insecure = false
# Private Gitea registry — read only, to reuse the existing build cache
# (keeps the op25 image off a ~1h from-scratch compile).
- uses: docker/login-action@v3
with:
registry: git.vpn.cusano.net
username: ${{ gitea.actor }}
password: ${{ secrets.BUILD_TOKEN }}
# Public registry — where the images are pushed.
- uses: docker/login-action@v3
with:
registry: ${{ vars.PUBLIC_REGISTRY }}
username: ${{ secrets.PUBLIC_REGISTRY_USER }}
password: ${{ secrets.PUBLIC_REGISTRY_TOKEN }}
- name: Version
id: meta
run: |
echo "REPO_NAME=$(echo ${GITHUB_REPOSITORY} | awk -F'/' '{print $2}')" >> $GITHUB_OUTPUT
echo "VERSION=$(git describe --tags --always | sed 's/^v//')" >> $GITHUB_OUTPUT
- uses: docker/build-push-action@v6
with:
context: ${{ matrix.context }}
file: ${{ matrix.file }}
platforms: linux/arm64
push: true
tags: |
${{ vars.PUBLIC_REGISTRY }}/${{ vars.PUBLIC_NAMESPACE }}/drb-${{ matrix.name }}:${{ steps.meta.outputs.VERSION }}
${{ vars.PUBLIC_REGISTRY }}/${{ vars.PUBLIC_NAMESPACE }}/drb-${{ matrix.name }}:latest
cache-from: type=registry,ref=git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ matrix.cache_name }}:buildcache
+15
View File
@@ -32,6 +32,21 @@ services:
depends_on: depends_on:
- icecast - icecast
# Claims the node's SECOND physical SDR (op25 always claims the first).
# Only useful if secondary_sdr_mode is set to adsb|ais via the edge-node
# config; otherwise it just sits idle answering /secondary/status. See
# node-26#9. Same network/device access as op25 for the same reason: it
# needs the raw USB device, not a virtualized one.
secondary-sdr:
image: ${IMAGE_REGISTRY:-git.vpn.cusano.net}/${DOCKER_ORG:-logan}/${DOCKER_REPO:-node-26}/secondary-sdr:latest
build: ./secondary-sdr-container
restart: unless-stopped
privileged: true
network_mode: host
env_file: .env
volumes:
- /dev:/dev
edge-node: edge-node:
image: ${IMAGE_REGISTRY:-git.vpn.cusano.net}/${DOCKER_ORG:-logan}/${DOCKER_REPO:-node-26}/edge-node:latest image: ${IMAGE_REGISTRY:-git.vpn.cusano.net}/${DOCKER_ORG:-logan}/${DOCKER_REPO:-node-26}/edge-node:latest
build: ./drb-edge-node build: ./drb-edge-node
+3
View File
@@ -136,6 +136,9 @@ class Settings(BaseSettings):
op25_api_url: str = "http://localhost:8001" op25_api_url: str = "http://localhost:8001"
op25_terminal_url: str = "http://localhost:8081" op25_terminal_url: str = "http://localhost:8081"
# Secondary SDR container (node-26#9) — ADS-B / AIS on a second SDR
secondary_sdr_api_url: str = "http://localhost:8002"
# Paths (volume mounts) # Paths (volume mounts)
config_path: str = "/configs" config_path: str = "/configs"
recordings_path: str = "/recordings" recordings_path: str = "/recordings"
@@ -210,10 +210,16 @@ class MQTTManager:
logger.info("No API key on disk — requesting re-delivery from C2 server.") logger.info("No API key on disk — requesting re-delivery from C2 server.")
self._publish(self._t_key_request, {}, qos=1) self._publish(self._t_key_request, {}, qos=1)
async def publish_checkin(self):
await self._publish_checkin()
async def _publish_checkin(self): async def _publish_checkin(self):
from app.internal.discord_radio import radio_bot from app.internal.discord_radio import radio_bot
from app.internal.config_manager import load_node_config 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() config = load_node_config()
devices = await op25_client.devices()
payload = { payload = {
"node_id": settings.node_id, "node_id": settings.node_id,
"name": settings.node_name, "name": settings.node_name,
@@ -225,7 +231,19 @@ class MQTTManager:
"is_overridden": config.override_system_id is not None and config.node_type != "portable", "is_overridden": config.override_system_id is not None and config.node_type != "portable",
"override_system_id": config.override_system_id, "override_system_id": config.override_system_id,
"enforce_override_timeout": config.enforce_override_timeout, "enforce_override_timeout": config.enforce_override_timeout,
"secondary_sdr_mode": config.secondary_sdr_mode,
"secondary_sdr_priority": config.secondary_sdr_priority,
} }
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) self._publish(self._t_checkin, payload, qos=1)
def _publish(self, topic: str, payload: dict, qos: int = 0, retain: bool = False): def _publish(self, topic: str, payload: dict, qos: int = 0, retain: bool = False):
+10
View File
@@ -65,6 +65,16 @@ class OP25Client:
logger.error(f"OP25 status failed: {e}") logger.error(f"OP25 status failed: {e}")
return None return None
async def devices(self) -> Optional[Dict[str, Any]]:
try:
async with httpx.AsyncClient(timeout=5) as client:
r = await client.get(f"{self.api_url}/op25/devices")
r.raise_for_status()
return r.json()
except Exception as e:
logger.error(f"OP25 device enumeration failed: {e}")
return None
async def generate_config(self, config: Dict[str, Any]) -> bool: async def generate_config(self, config: Dict[str, Any]) -> bool:
try: try:
async with httpx.AsyncClient(timeout=10) as client: async with httpx.AsyncClient(timeout=10) as client:
@@ -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
@@ -0,0 +1,71 @@
import httpx
from typing import Any, Dict, List, Optional
from app.config import settings
from app.internal.logger import logger
class SecondarySdrClient:
"""Talks to the secondary-sdr-container (node-26#9) over its control API.
Mirrors op25_client.py's shape on purpose — same failure handling (log
and return None/False rather than raise), since this container is
optional and its absence must never break the primary op25 radio path.
"""
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:
r = await client.post(f"{self.api_url}/secondary/start", json={"mode": mode})
r.raise_for_status()
return True
except Exception as e:
logger.error(f"Secondary SDR start (mode={mode!r}) failed: {e}")
return False
async def stop(self) -> bool:
try:
async with httpx.AsyncClient(timeout=10) as client:
r = await client.post(f"{self.api_url}/secondary/stop")
r.raise_for_status()
return True
except Exception as e:
logger.error(f"Secondary SDR stop failed: {e}")
return False
async def status(self) -> Optional[Dict[str, Any]]:
try:
async with httpx.AsyncClient(timeout=5) as client:
r = await client.get(f"{self.api_url}/secondary/status")
r.raise_for_status()
return r.json()
except Exception as e:
logger.error(f"Secondary SDR status failed: {e}")
return None
async def data(self) -> Optional[Dict[str, Any]]:
try:
async with httpx.AsyncClient(timeout=5) as client:
r = await client.get(f"{self.api_url}/secondary/data")
r.raise_for_status()
return r.json()
except Exception as e:
logger.error(f"Secondary SDR data fetch failed: {e}")
return None
secondary_sdr_client = SecondarySdrClient()
@@ -0,0 +1,45 @@
import asyncio
import httpx
from app.config import settings
from app.internal import credentials
from app.internal.config_manager import load_node_config
from app.internal.logger import logger
from app.internal.secondary_sdr_client import secondary_sdr_client
# How often the second-SDR decoder's current snapshot is forwarded to C2
# (node-26#9). This is a live-map overlay, not a flight/vessel history, so
# there is no backlog/retry on a missed tick — the next one supersedes it.
UPLINK_INTERVAL_SECONDS = 10
async def _post_snapshot(path: str, body: dict) -> None:
if not settings.c2_url:
return
api_key = credentials.get_api_key()
if not api_key:
return
headers = {"Authorization": f"Bearer {api_key}"}
try:
async with httpx.AsyncClient(timeout=10) as client:
r = await client.post(f"{settings.c2_url}{path}", json=body, headers=headers)
r.raise_for_status()
except Exception as e:
logger.debug(f"Telemetry uplink to {path} failed: {e}")
async def telemetry_uplink_loop():
while True:
await asyncio.sleep(UPLINK_INTERVAL_SECONDS)
if not load_node_config().secondary_sdr_priority:
continue
snapshot = await secondary_sdr_client.data()
if not snapshot:
continue
# 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"]})
if snapshot.get("vessels"):
await _post_snapshot("/telemetry/ais", {"vessels": snapshot["vessels"]})
+20
View File
@@ -158,6 +158,9 @@ async def on_command(payload: dict):
) )
elif action == "discord_leave": elif action == "discord_leave":
await radio_bot.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": elif action == "op25_restart":
from app.internal.op25_client import op25_client from app.internal.op25_client import op25_client
await op25_client.stop() await op25_client.stop()
@@ -224,6 +227,8 @@ async def on_config_push(payload: dict):
hardware_preset = payload.pop("hardware_preset", None) hardware_preset = payload.pop("hardware_preset", None)
ppm_override = payload.pop("ppm_override", None) ppm_override = payload.pop("ppm_override", None)
node_type = payload.pop("node_type", None) node_type = payload.pop("node_type", 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) enforce_override_timeout = payload.pop("enforce_override_timeout", None)
try: try:
config = SystemConfig(**payload) config = SystemConfig(**payload)
@@ -257,6 +262,12 @@ async def on_config_push(payload: dict):
await op25_client.start() await op25_client.start()
logger.info(f"Config push applied: {config.name}") 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)
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
# App lifecycle # App lifecycle
@@ -312,12 +323,21 @@ async def lifespan(app: FastAPI):
logger.warning(f"OP25 not ready yet (attempt {attempt + 1}/10), retrying in 3s…") logger.warning(f"OP25 not ready yet (attempt {attempt + 1}/10), retrying in 3s…")
await asyncio.sleep(3) await asyncio.sleep(3)
# 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 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()) heartbeat_task = asyncio.create_task(mqtt_manager.heartbeat_loop())
from app.internal.telemetry_uplink import telemetry_uplink_loop
telemetry_task = asyncio.create_task(telemetry_uplink_loop())
yield # --- app running --- yield # --- app running ---
logger.info("Edge node shutting down.") logger.info("Edge node shutting down.")
heartbeat_task.cancel() heartbeat_task.cancel()
telemetry_task.cancel()
await metadata_watcher.stop() await metadata_watcher.stop()
await call_recorder.stop() await call_recorder.stop()
await radio_bot.stop() await radio_bot.stop()
+28 -2
View File
@@ -1,5 +1,5 @@
from pydantic import BaseModel from pydantic import BaseModel, model_validator
from typing import Optional, Dict, Any from typing import Optional, Dict, Any, Iterable, List
from enum import Enum from enum import Enum
from datetime import datetime from datetime import datetime
@@ -23,6 +23,21 @@ class SystemConfig(BaseModel):
config: Dict[str, Any] # OP25-compatible config blob passed through to op25-container 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): class NodeConfig(BaseModel):
node_id: str node_id: str
node_name: str node_name: str
@@ -34,11 +49,22 @@ class NodeConfig(BaseModel):
hardware_preset: str = "rtl-sdr-v3" hardware_preset: str = "rtl-sdr-v3"
ppm_override: Optional[float] = None ppm_override: Optional[float] = None
node_type: str = "fixed" # fixed or portable node_type: str = "fixed" # fixed or portable
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 enforce_override_timeout: bool = True
override_system_id: Optional[str] = None override_system_id: Optional[str] = None
override_config: Optional[SystemConfig] = None override_config: Optional[SystemConfig] = None
offline_call_buffer_size: int = 35 # max call_end events to buffer while MQTT is offline 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): class CallEvent(BaseModel):
call_id: str call_id: str
+26 -2
View File
@@ -1,9 +1,9 @@
from fastapi import APIRouter, Depends, HTTPException, Body from fastapi import APIRouter, Depends, HTTPException, Body
from typing import Optional from typing import List, Optional
import asyncio import asyncio
import httpx import httpx
from app.config import settings 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.op25_client import op25_client
from app.internal.config_manager import load_node_config, save_node_config, apply_system_config from app.internal.config_manager import load_node_config, save_node_config, apply_system_config
from app.internal.call_recorder import call_recorder 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}") 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") @router.post("/discord/join")
async def discord_join(guild_id: int, channel_id: int): async def discord_join(guild_id: int, channel_id: int):
ok = await radio_bot.join(guild_id, channel_id) ok = await radio_bot.join(guild_id, channel_id)
+133
View File
@@ -105,6 +105,30 @@
background: rgba(255, 255, 255, 0.1); 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 { .grid {
display: grid; display: grid;
grid-template-columns: repeat(auto-fit, minmax(320px, 1fr)); grid-template-columns: repeat(auto-fit, minmax(320px, 1fr));
@@ -360,6 +384,22 @@
</div> </div>
</div> </div>
<!-- Secondary SDRs Card (node-26#9) -->
<div class="glass-card">
<div class="card-header">
<div class="card-title">Secondary SDRs</div>
</div>
<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>
</p>
<div id="sdr-list"></div>
<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>
<span id="sdr-msg" style="color: var(--text-muted); font-size: 0.8rem;"></span>
</div>
</div>
<!-- Local Listening Card --> <!-- Local Listening Card -->
<div class="glass-card player-card"> <div class="glass-card player-card">
<div class="card-header" style="margin-bottom:0.5rem; border:none;"> <div class="card-header" style="margin-bottom:0.5rem; border:none;">
@@ -492,6 +532,99 @@
} }
} }
// ── Secondary SDR priority ────────────────────────────────────────────
const SDR_LABELS = {
adsb: ['ADS-B', 'Aircraft · 1090 MHz'],
ais: ['AIS', 'Vessels · 162 MHz'],
};
let sdrRows = []; // [{mode, enabled}] in display order
let sdrRunning = [];
let sdrDirty = false;
function renderSdr() {
const list = document.getElementById('sdr-list');
list.innerHTML = '';
let rank = 0;
sdrRows.forEach((row, i) => {
const [name, hint] = SDR_LABELS[row.mode] || [row.mode, ''];
const running = sdrRunning.includes(row.mode);
const state = !row.enabled ? 'Off' : running ? 'Running' : (sdrDirty ? 'Unsaved' : 'Waiting for SDR');
const el = document.createElement('div');
el.className = 'sdr-row';
el.innerHTML = `
<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>
<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>`;
const [box] = el.getElementsByTagName('input');
const [up, down] = el.getElementsByTagName('button');
box.onchange = () => { row.enabled = box.checked; markSdrDirty(); };
up.onclick = () => { [sdrRows[i - 1], sdrRows[i]] = [sdrRows[i], sdrRows[i - 1]]; markSdrDirty(); };
down.onclick = () => { [sdrRows[i + 1], sdrRows[i]] = [sdrRows[i], sdrRows[i + 1]]; markSdrDirty(); };
list.appendChild(el);
});
document.getElementById('sdr-save').disabled = !sdrDirty;
}
function markSdrDirty() {
sdrDirty = true;
document.getElementById('sdr-msg').textContent = '';
renderSdr();
}
async function loadSdr() {
if (sdrDirty) return; // never clobber unsaved edits with a poll
try {
const r = await fetch('/api/secondary');
if (!r.ok) return;
const d = await r.json();
sdrRunning = d.running || [];
sdrRows = [
...d.priority.map(mode => ({ mode, enabled: true })),
...d.modes.filter(m => !d.priority.includes(m)).map(mode => ({ mode, enabled: false })),
];
const summary = document.getElementById('sdr-summary');
if (d.running === null) {
summary.textContent = 'Secondary SDR service is not responding.';
} else if (d.sdr_count != null) {
const spare = Math.max(d.sdr_count - 1, 0);
summary.textContent = `This node has ${d.sdr_count} SDR${d.sdr_count === 1 ? '' : 's'}, so ${spare} spare.`;
}
renderSdr();
} catch (e) {
console.error('Secondary SDR load failed:', e);
}
}
async function saveSdrPriority() {
const priority = sdrRows.filter(r => r.enabled).map(r => r.mode);
const btn = document.getElementById('sdr-save');
const msg = document.getElementById('sdr-msg');
btn.disabled = true;
msg.textContent = 'Applying…';
try {
const r = await fetch('/api/secondary/priority', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ priority }),
});
if (!r.ok) throw new Error(await r.text());
const d = await r.json();
sdrDirty = false;
msg.textContent = d.running === null ? 'Saved, but the secondary SDR service did not respond.' : 'Saved.';
await loadSdr();
} catch (e) {
console.error('Secondary SDR save failed:', e);
msg.textContent = 'Save failed.';
btn.disabled = false;
}
}
loadSdr();
setInterval(loadSdr, 10000);
refresh(); refresh();
setInterval(refresh, 2000); // Polling every 2 seconds setInterval(refresh, 2000); // Polling every 2 seconds
</script> </script>
@@ -0,0 +1,58 @@
"""
node-26#9 — secondary SDR priority: an ordered list the SDRs beyond op25's work
through. Normalisation, migration from the legacy single-mode field, and the one
apply path shared by the dashboard, C2 commands and config pushes.
"""
import asyncio
from unittest.mock import AsyncMock, patch
from app.models import NodeConfig, normalize_secondary_priority
def _cfg(**kw) -> NodeConfig:
return NodeConfig(node_id="n1", node_name="N1", lat=0.0, lon=0.0, **kw)
def test_normalize_keeps_order_drops_unknown_and_duplicates():
assert normalize_secondary_priority(["ais", "bogus", "adsb", "ais"]) == ["ais", "adsb"]
assert normalize_secondary_priority([]) == []
def test_legacy_single_mode_migrates_to_priority():
assert _cfg(secondary_sdr_mode="adsb").secondary_sdr_priority == ["adsb"]
assert _cfg(secondary_sdr_mode="none").secondary_sdr_priority == []
def test_explicit_priority_wins_over_legacy_mode():
assert _cfg(secondary_sdr_mode="adsb", secondary_sdr_priority=["ais"]).secondary_sdr_priority == ["ais"]
def test_set_priority_persists_applies_and_never_touches_op25(tmp_path):
config_file = tmp_path / "node_config.json"
import app.internal.config_manager as cm
with patch.object(cm, "_CONFIG_FILE", config_file):
cm.save_node_config(_cfg(secondary_sdr_mode="adsb"))
from app.internal import secondary_priority as sp
with patch.object(sp.secondary_sdr_client, "apply", AsyncMock(return_value=["ais"])) as apply, \
patch("app.internal.mqtt_manager.mqtt_manager.publish_checkin", AsyncMock()), \
patch("app.internal.op25_client.op25_client.stop", AsyncMock()) as op25_stop:
running = asyncio.run(sp.set_secondary_priority(["ais", "adsb", "ais"]))
saved = cm.load_node_config()
assert running == ["ais"]
apply.assert_awaited_once_with(["ais", "adsb"])
op25_stop.assert_not_awaited()
assert saved.secondary_sdr_priority == ["ais", "adsb"]
assert saved.secondary_sdr_mode == "ais"
def test_clearing_priority_is_not_undone_by_legacy_migration(tmp_path):
config_file = tmp_path / "node_config.json"
import app.internal.config_manager as cm
with patch.object(cm, "_CONFIG_FILE", config_file):
cm.save_node_config(_cfg(secondary_sdr_mode="adsb"))
from app.internal import secondary_priority as sp
with patch.object(sp.secondary_sdr_client, "apply", AsyncMock(return_value=[])), \
patch("app.internal.mqtt_manager.mqtt_manager.publish_checkin", AsyncMock()):
asyncio.run(sp.set_secondary_priority([]))
assert cm.load_node_config().secondary_sdr_priority == []
+7
View File
@@ -181,6 +181,13 @@ else
warn "no RTL-SDR dongle detected on USB — plug one in before expecting audio" warn "no RTL-SDR dongle detected on USB — plug one in before expecting audio"
fi fi
# The kernel's DVB-TV driver auto-binds RTL2838 dongles. op25 detaches it from
# the one it opens, but any other dongle stays claimed and the secondary-sdr
# decoder fails with "usb_claim_interface error -6" (node-26#9).
echo "blacklist dvb_usb_rtl28xxu" > /etc/modprobe.d/blacklist-rtl-sdr.conf
modprobe -r rtl2832_sdr dvb_usb_rtl28xxu 2>/dev/null || true
ok "DVB-TV kernel driver blacklisted for RTL-SDR"
# ── 2. Dependencies ───────────────────────────────────────────────────────── # ── 2. Dependencies ─────────────────────────────────────────────────────────
say "Installing dependencies" say "Installing dependencies"
export DEBIAN_FRONTEND=noninteractive export DEBIAN_FRONTEND=noninteractive
+1 -1
View File
@@ -12,7 +12,7 @@ ENV DEBIAN_FRONTEND=noninteractive
# Install system dependencies # Install system dependencies
RUN apt-get update && \ RUN apt-get update && \
apt-get upgrade -y && \ apt-get upgrade -y && \
apt-get install git pulseaudio pulseaudio-utils liquidsoap -y apt-get install git pulseaudio pulseaudio-utils liquidsoap usbutils -y
# Install custom PulseAudio system config (enables anonymous access for edge-node) # Install custom PulseAudio system config (enables anonymous access for edge-node)
COPY system.pa /etc/pulse/system.pa COPY system.pa /etc/pulse/system.pa
@@ -1,6 +1,7 @@
from fastapi import HTTPException, APIRouter from fastapi import HTTPException, APIRouter
import subprocess import subprocess
import os import os
import re
import signal import signal
import json import json
from models import ConfigGenerator, DecodeMode, ChannelConfig, DeviceConfig, TrunkingConfig, TrunkingChannelConfig, TerminalConfig, MetadataConfig, MetadataStreamConfig, HARDWARE_PRESETS from models import ConfigGenerator, DecodeMode, ChannelConfig, DeviceConfig, TrunkingConfig, TrunkingChannelConfig, TerminalConfig, MetadataConfig, MetadataStreamConfig, HARDWARE_PRESETS
@@ -69,6 +70,24 @@ def create_op25_router():
async def get_status(): async def get_status():
return {"status": "running" if _is_running() else "stopped"} return {"status": "running" if _is_running() else "stopped"}
@router.get("/devices")
async def list_sdr_devices():
"""Enumerate connected SDR-looking USB devices via lsusb.
Same match heuristic as install.sh's one-shot host check — good enough
to answer "is a second SDR plugged in", not a serial-level device
binding (op25's DeviceConfig.args has no serial concept yet either).
"""
devices = []
try:
out = subprocess.run(["lsusb"], capture_output=True, text=True, timeout=5).stdout
for line in out.splitlines():
if re.search(r"rtl2838|realtek.*283[28]|sdr", line, re.IGNORECASE):
devices.append(line.strip())
except Exception as e:
LOGGER.warning(f"SDR device enumeration failed: {e}")
return {"count": len(devices), "devices": devices}
@router.post("/generate-config") @router.post("/generate-config")
async def generate_config(generator: ConfigGenerator): async def generate_config(generator: ConfigGenerator):
try: try:
+54
View File
@@ -0,0 +1,54 @@
# Secondary-SDR Container — node-26#9
#
# Claims the node's SECOND physical SDR (the first is always op25's). Mode is
# chosen at runtime via the control API, not baked in: dump1090 for ADS-B,
# AIS-catcher for AIS (op25_2 mode is not handled here yet — see node-26#9).
#
# Device claiming is by RTL-SDR index, not serial (op25's DeviceConfig.args
# has no serial concept either — see op25-container/app/models.py). Index 0
# is reserved for op25; this container always addresses index 1. That's a
# real limitation once serial-stable device binding matters (hot-unplug /
# replug can swap indices) — tracked in node-26#9, not fixed here.
#
# UNVERIFIED: this image has not been built or run against real hardware in
# this session (sandboxed authoring machine, no docker). dump1090 and
# AIS-catcher's exact CLI flags below are believed correct from their
# published docs but not confirmed against a real capture — the CTO/QA
# review before this ships to a real node should build and smoke-test it.
FROM python:3.14-slim
ENV DEBIAN_FRONTEND=noninteractive
RUN apt-get update && \
apt-get upgrade -y && \
apt-get install -y --no-install-recommends \
git build-essential cmake pkg-config \
librtlsdr-dev libusb-1.0-0-dev libssl-dev zlib1g-dev libzstd-dev libncurses-dev usbutils
# readsb (wiedehopf) — ADS-B decoder. antirez/dump1090 was used here first but
# has no --write-json at all (it only serves /data.json over --net), so the
# decoder exited on an unknown flag. readsb writes the dump1090-fa style
# aircraft.json that _read_adsb_snapshot() parses.
RUN git clone --depth 1 https://github.com/wiedehopf/readsb /opt/readsb && \
cd /opt/readsb && make RTLSDR=yes
# AIS-catcher — AIS decoder.
RUN git clone https://github.com/jvde-github/AIS-catcher /opt/AIS-catcher && \
cd /opt/AIS-catcher && mkdir build && cd build && cmake .. && make
EXPOSE 8002
VOLUME ["/configs"]
WORKDIR /app
COPY ./app /app
COPY docker-entrypoint.sh /usr/local/bin/
RUN sed -i 's/\r$//' /usr/local/bin/docker-entrypoint.sh && \
chmod +x /usr/local/bin/docker-entrypoint.sh
COPY requirements.txt /tmp/requirements.txt
RUN pip3 install --no-cache-dir -r /tmp/requirements.txt
ENTRYPOINT ["/usr/local/bin/docker-entrypoint.sh"]
CMD ["python", "main.py"]
+20
View File
@@ -0,0 +1,20 @@
from pydantic_settings import BaseSettings
class Settings(BaseSettings):
# Same rationale as op25-container's OP25_DEBUG_EXPOSE: both containers
# share the host network namespace (network_mode: host), so edge-node
# reaches this control API over localhost regardless of this flag. False
# (default) binds 127.0.0.1; true exposes unauthenticated start/stop to
# the node's LAN and should only ever be set for local development.
secondary_sdr_debug_expose: bool = False
class Config:
env_file = ".env"
settings = Settings()
def bind_host() -> str:
return "0.0.0.0" if settings.secondary_sdr_debug_expose else "127.0.0.1"
@@ -0,0 +1,302 @@
import json
import os
import signal
import subprocess
import threading
from pathlib import Path
from typing import Any, Dict, List, Optional
from internal.logger import create_logger
LOGGER = create_logger(__name__)
# One decoder per mode, each on its own SDR (node-26#9). The node's secondary
# SDR *priority* (e.g. ["adsb", "ais"]) is applied by apply(): decoders start
# down the list until no free dongle is left, so every SDR the node has gets
# used and the ones beyond the list's reach simply aren't started. OP25 always
# keeps its own dongle — it is started first and never part of this list.
MODES = ("adsb", "ais")
# Which RTL-SDR index op25 holds is NOT fixed — on radio-box op25 had index 1
# and index 0 was free, so "op25 is always 0" was wrong. rtlsdr can't open a
# dongle another process has claimed, and the decoders exit within ~50ms when
# that happens, so _start_one() tries each index and keeps the first that
# stays up. op25 is never disturbed: a failed claim doesn't touch its dongle.
MAX_SDR_INDEX = 4
_STARTUP_GRACE_S = 2.0
_STATE_DIR = Path("/tmp/secondary_sdr")
ADSB_JSON_DIR = Path("/tmp/adsb")
# Live decoder handles. poll() is the only reliable liveness check: a decoder
# that dies on startup stays an unreaped zombie, and killpg(pgid, 0) still
# succeeds on a zombie.
_procs: Dict[str, subprocess.Popen] = {}
_indices: Dict[str, int] = {}
_lock = threading.Lock()
# AIS-catcher streams one JSON object per received message on stdout rather
# than writing a periodic snapshot file (readsb's approach) — so the
# current-vessel snapshot lives in memory, keyed by mmsi, kept warm by a
# background reader thread for as long as the decoder is alive.
_ais_vessels: Dict[str, Dict[str, Any]] = {}
_ais_lock = threading.Lock()
def _pgid_file(mode: str) -> Path:
return _STATE_DIR / f"{mode}.pgid"
def _reap_orphans() -> None:
"""Kill decoders left behind by a previous API process (uvicorn --reload
restarts this process, but decoders run in their own session and would
otherwise keep holding their SDRs)."""
if not _STATE_DIR.exists():
return
for f in _STATE_DIR.glob("*.pgid"):
try:
os.killpg(int(f.read_text().strip()), signal.SIGTERM)
LOGGER.info(f"Stopped orphaned secondary decoder from {f.name}")
except Exception:
pass
f.unlink(missing_ok=True)
_reap_orphans()
def _adsb_command(index: int) -> List[str]:
ADSB_JSON_DIR.mkdir(parents=True, exist_ok=True)
return [
"/opt/readsb/readsb",
"--net",
"--device-type", "rtlsdr",
"--device", str(index),
"--write-json", str(ADSB_JSON_DIR),
"--write-json-every", "1",
]
def _ais_command(index: int) -> List[str]:
return [
"/opt/AIS-catcher/build/AIS-catcher",
f"-d:{index}", # "-d <x>" would select by serial, not index
"-o", "5", # JSON Full: decoded fields (4 = sparse, "JSON" is rejected)
]
_COMMANDS = {"adsb": _adsb_command, "ais": _ais_command}
def _ais_reader(proc: subprocess.Popen) -> None:
"""
Consume AIS-catcher's stdout, one JSON message per line, and keep the
latest report per mmsi. Field names (mmsi/lat/lon/speed/course or
heading/shipname or name) are believed correct from AIS-catcher's
published JSON output docs but UNVERIFIED against a real capture in
this session — same caveat as dump1090's aircraft.json mapping.
Malformed/partial lines (e.g. static-data-only messages with no
position) are skipped rather than raising, since dropping one line must
never kill the reader thread.
"""
if not proc.stdout:
return
for line in proc.stdout:
try:
msg = json.loads(line)
except Exception:
continue
mmsi = msg.get("mmsi")
if not mmsi:
continue
# AIS-catcher emits separate message TYPES per mmsi — static data
# (name, no position) and position reports (lat/lon, no name) arrive
# as distinct lines. Merge onto the existing entry, only overwriting
# a field the new message actually carries, so a position-only
# report doesn't blank out a name learned from an earlier message.
name = (msg.get("shipname") or msg.get("name") or "").strip() or None
heading = msg.get("heading") if msg.get("heading") is not None else msg.get("course")
updates = {
"mmsi": str(mmsi),
"name": name,
"lat": msg.get("lat"),
"lon": msg.get("lon"),
"speed_kt": msg.get("speed"),
"heading_deg": heading,
}
with _ais_lock:
existing = _ais_vessels.get(str(mmsi), {})
for key, value in updates.items():
if value is not None:
existing[key] = value
_ais_vessels[str(mmsi)] = existing
def is_running(mode: str) -> bool:
proc = _procs.get(mode)
return proc is not None and proc.poll() is None
def running() -> List[str]:
return [m for m in MODES if is_running(m)]
def _start_one(mode: str) -> bool:
"""Start one decoder on the first free SDR. False when none is free."""
if mode not in _COMMANDS:
raise ValueError(f"Unknown secondary SDR mode: {mode!r}")
if is_running(mode):
return True
if mode == "ais":
with _ais_lock:
_ais_vessels.clear()
needs_stdout = mode == "ais"
for index in range(MAX_SDR_INDEX):
if index in {_indices[m] for m in running() if m in _indices}:
continue
try:
proc = subprocess.Popen(
_COMMANDS[mode](index),
preexec_fn=os.setsid,
stdout=subprocess.PIPE if needs_stdout else None,
text=True if needs_stdout else None,
bufsize=1 if needs_stdout else -1,
)
except Exception as e:
LOGGER.error(f"Failed to start secondary SDR decoder mode={mode!r}: {e}")
return False
try:
proc.wait(timeout=_STARTUP_GRACE_S)
LOGGER.info(f"Secondary SDR decoder mode={mode!r} could not use SDR index {index}, trying next")
continue
except subprocess.TimeoutExpired:
pass
if needs_stdout:
threading.Thread(target=_ais_reader, args=(proc,), daemon=True).start()
_procs[mode] = proc
_indices[mode] = index
_STATE_DIR.mkdir(parents=True, exist_ok=True)
_pgid_file(mode).write_text(str(proc.pid))
LOGGER.info(f"Started secondary SDR decoder mode={mode!r} on SDR index {index} pid={proc.pid}")
return True
LOGGER.info(f"Secondary SDR decoder mode={mode!r}: no free SDR left")
return False
def _stop_one(mode: str) -> None:
proc = _procs.pop(mode, None)
_indices.pop(mode, None)
if proc is not None:
try:
os.killpg(proc.pid, signal.SIGTERM)
except OSError:
pass
try:
proc.wait(timeout=5)
except subprocess.TimeoutExpired:
pass
_pgid_file(mode).unlink(missing_ok=True)
def start(mode: str) -> bool:
with _lock:
return _start_one(mode)
def stop(mode: Optional[str] = None) -> None:
with _lock:
for m in [mode] if mode else list(_procs):
_stop_one(m)
def apply(priority: List[str]) -> List[str]:
"""Run decoders in priority order until SDRs run out; stop everything else.
Unchanged when the running set is already the right prefix of the list, so
re-applying the same priority (every restart/config push) is a no-op
rather than a decoder restart.
"""
for m in priority:
if m not in _COMMANDS:
raise ValueError(f"Unknown secondary SDR mode: {m!r}")
with _lock:
live = set(running())
if live and live == set(priority[:len(live)]):
# Already running the head of the list; only try to extend it.
for m in priority[len(live):]:
if not _start_one(m):
break
return running()
for m in list(_procs):
_stop_one(m)
for m in priority:
if not _start_one(m):
break
return running()
def sdr_count() -> Optional[int]:
"""RTL-SDR dongles on USB, same heuristic as install.sh. This container has
usbutils; op25's :stable image doesn't, so its /op25/devices can't answer."""
try:
out = subprocess.run(["lsusb"], capture_output=True, text=True, timeout=5).stdout
except Exception:
return None
return sum(1 for line in out.splitlines() if "0bda:2838" in line or "0bda:2832" in line)
def status() -> Dict[str, Any]:
live = running()
return {
"running": [{"mode": m, "sdr_index": _indices.get(m)} for m in live],
"sdr_count": sdr_count(),
}
def _altitude(a: Dict[str, Any]) -> Optional[int]:
alt = a.get("alt_baro", a.get("altitude"))
if alt == "ground":
return 0
return alt if isinstance(alt, (int, float)) else None
def _read_adsb_snapshot() -> List[Dict[str, Any]]:
"""
Map readsb's aircraft.json (--write-json output) to the server's
telemetry schema. readsb uses the dump1090-fa field names: alt_baro (int,
or the string "ground"), gs, track. Older dump1090 forks used
altitude/speed, kept as a fallback.
"""
path = ADSB_JSON_DIR / "aircraft.json"
try:
raw = json.loads(path.read_text())
except Exception:
return []
out = []
for a in raw.get("aircraft", []):
icao = a.get("hex")
if not icao:
continue
out.append({
"icao": icao.upper(),
"callsign": (a.get("flight") or "").strip() or None,
"lat": a.get("lat"),
"lon": a.get("lon"),
"altitude_ft": _altitude(a),
"ground_speed_kt": a.get("gs", a.get("speed")),
"track_deg": a.get("track"),
})
return out
def data() -> Dict[str, Any]:
with _ais_lock:
vessels = list(_ais_vessels.values()) if is_running("ais") else []
return {
"running": running(),
"aircraft": _read_adsb_snapshot() if is_running("adsb") else [],
"vessels": vessels,
}
@@ -0,0 +1,31 @@
import logging
from logging.handlers import RotatingFileHandler
def create_logger(name, level=logging.DEBUG, max_bytes=10485760, backup_count=2):
debug_log_file = "./secondary-sdr.debug.log"
info_log_file = "./secondary-sdr.log"
logger = logging.getLogger(name)
logger.setLevel(level)
if not logger.hasHandlers():
console_handler = logging.StreamHandler()
console_handler.setLevel(level)
debug_file_handler = RotatingFileHandler(debug_log_file, maxBytes=max_bytes, backupCount=backup_count)
debug_file_handler.setLevel(logging.DEBUG)
info_file_handler = RotatingFileHandler(info_log_file, maxBytes=max_bytes, backupCount=backup_count)
info_file_handler.setLevel(logging.INFO)
formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
console_handler.setFormatter(formatter)
debug_file_handler.setFormatter(formatter)
info_file_handler.setFormatter(formatter)
logger.addHandler(console_handler)
logger.addHandler(debug_file_handler)
logger.addHandler(info_file_handler)
return logger
+13
View File
@@ -0,0 +1,13 @@
from fastapi import FastAPI
import routers.secondary_controller as secondary_controller
from config import bind_host
app = FastAPI()
app.include_router(secondary_controller.create_secondary_router(), prefix="/secondary")
if __name__ == "__main__":
import uvicorn
uvicorn.run("main:app", host=bind_host(), port=8002, reload=True)
@@ -0,0 +1,58 @@
from typing import List, Optional
from fastapi import APIRouter, HTTPException
from pydantic import BaseModel
from internal import decoder_control
from internal.logger import create_logger
LOGGER = create_logger(__name__)
class StartBody(BaseModel):
mode: str # adsb | ais
class StopBody(BaseModel):
mode: Optional[str] = None # omit to stop every decoder
class ApplyBody(BaseModel):
priority: List[str] # ordered, e.g. ["adsb", "ais"]
def create_secondary_router():
router = APIRouter()
@router.post("/apply")
async def apply(body: ApplyBody):
try:
live = decoder_control.apply(body.priority)
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
return {"running": live}
@router.post("/start")
async def start(body: StartBody):
try:
ok = decoder_control.start(body.mode)
except ValueError as e:
raise HTTPException(status_code=400, detail=str(e))
if not ok:
raise HTTPException(status_code=409, detail="No free SDR for this decoder")
return {"status": f"secondary SDR started ({body.mode})"}
@router.post("/stop")
async def stop(body: Optional[StopBody] = None):
decoder_control.stop(body.mode if body else None)
return {"status": "secondary SDR stopped"}
@router.get("/status")
async def get_status():
return decoder_control.status()
@router.get("/data")
async def get_data():
return decoder_control.data()
return router
@@ -0,0 +1,3 @@
#!/bin/bash
mkdir -p /tmp/adsb
exec "$@"
+3
View File
@@ -0,0 +1,3 @@
uvicorn
fastapi
pydantic-settings