Compare commits
14
Commits
v1
...
feat/sdr-pins
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9b50ba8114 | ||
|
|
3ecc7eff1a | ||
|
|
358766d88b | ||
|
|
c82fc61910 | ||
|
|
9b6c64fbbf | ||
|
|
1c74edffea | ||
|
|
c0a05a8ad6 | ||
|
|
52d11c1546 | ||
|
|
b54624e176 | ||
|
|
52f31bbcc0 | ||
|
|
b2e3804dc3 | ||
|
|
8f5fca8757 | ||
|
|
a53000a092 | ||
|
|
bee9e173e6 |
@@ -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
|
||||||
|
|||||||
@@ -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
|
||||||
@@ -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
|
||||||
@@ -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
|
||||||
|
|||||||
@@ -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,29 @@ 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,
|
||||||
}
|
}
|
||||||
|
payload["sdr_pins"] = config.sdr_pins
|
||||||
|
secondary = await secondary_sdr_client.status()
|
||||||
|
if secondary is not None:
|
||||||
|
payload["secondary_sdr_running"] = [r["mode"] for r in secondary.get("running", [])]
|
||||||
|
if secondary.get("devices") is not None:
|
||||||
|
payload["sdr_devices"] = [
|
||||||
|
{k: d.get(k) for k in ("index", "serial", "name", "duplicate_serial")}
|
||||||
|
for d in secondary["devices"]
|
||||||
|
]
|
||||||
|
from app.internal.sdr_settings import op25_serial
|
||||||
|
payload["op25_sdr_serial"] = op25_serial(config, secondary["devices"])
|
||||||
|
# Best-effort hardware report — omit rather than guess. Prefer the
|
||||||
|
# secondary-sdr container's count: op25's :stable image has no lsusb and
|
||||||
|
# its /op25/devices answers 0 rather than "unknown" when that fails. A
|
||||||
|
# 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)
|
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):
|
||||||
|
|||||||
@@ -65,15 +65,29 @@ 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:
|
||||||
r = await client.post(f"{self.api_url}/op25/generate-config", json=config)
|
r = await client.post(f"{self.api_url}/op25/generate-config", json=config)
|
||||||
r.raise_for_status()
|
r.raise_for_status()
|
||||||
return True
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logger.error(f"OP25 generate-config failed: {e}")
|
logger.error(f"OP25 generate-config failed: {e}")
|
||||||
return False
|
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]:
|
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
|
||||||
@@ -0,0 +1,84 @@
|
|||||||
|
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 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=body)
|
||||||
|
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"]})
|
||||||
@@ -5,7 +5,7 @@ from typing import Optional
|
|||||||
|
|
||||||
from fastapi import FastAPI
|
from fastapi import FastAPI
|
||||||
from app.config import settings
|
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.logger import logger
|
||||||
from app.internal.mqtt_manager import mqtt_manager
|
from app.internal.mqtt_manager import mqtt_manager
|
||||||
from app.internal import credentials
|
from app.internal import credentials
|
||||||
@@ -158,6 +158,12 @@ async def on_command(payload: dict):
|
|||||||
)
|
)
|
||||||
elif action == "discord_leave":
|
elif action == "discord_leave":
|
||||||
await radio_bot.leave()
|
await radio_bot.leave()
|
||||||
|
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":
|
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 +230,9 @@ 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)
|
||||||
|
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)
|
enforce_override_timeout = payload.pop("enforce_override_timeout", None)
|
||||||
try:
|
try:
|
||||||
config = SystemConfig(**payload)
|
config = SystemConfig(**payload)
|
||||||
@@ -245,6 +254,16 @@ async def on_config_push(payload: dict):
|
|||||||
node_cfg.node_type = node_type
|
node_cfg.node_type = node_type
|
||||||
if enforce_override_timeout is not None:
|
if enforce_override_timeout is not None:
|
||||||
node_cfg.enforce_override_timeout = bool(enforce_override_timeout)
|
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)
|
save_node_config(node_cfg)
|
||||||
|
|
||||||
from app.internal.op25_client import op25_client
|
from app.internal.op25_client import op25_client
|
||||||
@@ -252,10 +271,16 @@ async def on_config_push(payload: dict):
|
|||||||
logger.error(f"Failed to generate OP25 config for {config.name}")
|
logger.error(f"Failed to generate OP25 config for {config.name}")
|
||||||
return
|
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 op25_client.stop()
|
||||||
await asyncio.sleep(2)
|
await asyncio.sleep(2)
|
||||||
await op25_client.start()
|
await op25_client.start()
|
||||||
logger.info(f"Config push applied: {config.name}")
|
logger.info(f"Config push applied: {config.name}")
|
||||||
|
await apply_secondaries()
|
||||||
|
|
||||||
|
|
||||||
# ---------------------------------------------------------------------------
|
# ---------------------------------------------------------------------------
|
||||||
@@ -312,12 +337,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.sdr_settings import apply_secondaries
|
||||||
|
logger.info(f"Resuming secondary SDRs (priority={node_cfg.secondary_sdr_priority!r}) after restart.")
|
||||||
|
await apply_secondaries()
|
||||||
|
|
||||||
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()
|
||||||
|
|||||||
@@ -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,34 @@ 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
|
||||||
|
|
||||||
|
|
||||||
|
# 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):
|
class NodeConfig(BaseModel):
|
||||||
node_id: str
|
node_id: str
|
||||||
node_name: str
|
node_name: str
|
||||||
@@ -34,11 +62,23 @@ 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
|
||||||
|
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.
|
||||||
|
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
|
||||||
|
|||||||
@@ -1,9 +1,10 @@
|
|||||||
from fastapi import APIRouter, Depends, HTTPException, Body
|
from fastapi import APIRouter, Depends, HTTPException, Body
|
||||||
from typing import Optional
|
from pydantic import BaseModel
|
||||||
|
from typing import Dict, 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 +205,44 @@ 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("/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 {
|
||||||
|
"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,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
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")
|
@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)
|
||||||
|
|||||||
@@ -105,6 +105,40 @@
|
|||||||
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); }
|
||||||
|
.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 {
|
.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 +394,28 @@
|
|||||||
</div>
|
</div>
|
||||||
</div>
|
</div>
|
||||||
|
|
||||||
|
<!-- SDRs Card (node-26#9, #11) -->
|
||||||
|
<div class="glass-card">
|
||||||
|
<div class="card-header">
|
||||||
|
<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;">
|
||||||
|
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="saveSdr()" disabled>Save</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 +548,133 @@
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ── 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 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)} (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 pinMissing = row.pin && !(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 ${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);
|
||||||
|
});
|
||||||
|
|
||||||
|
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() {
|
||||||
|
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/sdr');
|
||||||
|
if (!r.ok) return;
|
||||||
|
sdr = await r.json();
|
||||||
|
sdrOp25Pin = sdr.pins.op25 || '';
|
||||||
|
sdrRows = [
|
||||||
|
...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] || '' })),
|
||||||
|
];
|
||||||
|
renderSdr();
|
||||||
|
} catch (e) {
|
||||||
|
console.error('SDR load failed:', e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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/sdr', {
|
||||||
|
method: 'POST',
|
||||||
|
headers: { 'Content-Type': 'application/json' },
|
||||||
|
body: JSON.stringify({ priority: sdrRows.filter(r => r.enabled).map(r => r.mode), pins }),
|
||||||
|
});
|
||||||
|
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('SDR save failed:', e);
|
||||||
|
msg.textContent = `Save failed: ${e.message}`;
|
||||||
|
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,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"}
|
||||||
@@ -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
|
||||||
|
|||||||
@@ -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:
|
||||||
|
|||||||
@@ -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"]
|
||||||
@@ -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,353 @@
|
|||||||
|
import ctypes
|
||||||
|
import ctypes.util
|
||||||
|
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, 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):
|
||||||
|
return True
|
||||||
|
if mode == "ais":
|
||||||
|
with _ais_lock:
|
||||||
|
_ais_vessels.clear()
|
||||||
|
|
||||||
|
needs_stdout = mode == "ais"
|
||||||
|
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:
|
||||||
|
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 _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.
|
||||||
|
|
||||||
|
`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:
|
||||||
|
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 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:
|
||||||
|
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]:
|
||||||
|
devs = devices()
|
||||||
|
by_index = {d["index"]: d["serial"] for d in devs}
|
||||||
|
return {
|
||||||
|
"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,
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
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
|
||||||
@@ -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,64 @@
|
|||||||
|
from typing import Dict, 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"]
|
||||||
|
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():
|
||||||
|
router = APIRouter()
|
||||||
|
|
||||||
|
@router.post("/apply")
|
||||||
|
async def apply(body: ApplyBody):
|
||||||
|
try:
|
||||||
|
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}
|
||||||
|
|
||||||
|
@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("/devices")
|
||||||
|
async def get_devices():
|
||||||
|
return {"devices": decoder_control.devices()}
|
||||||
|
|
||||||
|
@router.get("/data")
|
||||||
|
async def get_data():
|
||||||
|
return decoder_control.data()
|
||||||
|
|
||||||
|
return router
|
||||||
@@ -0,0 +1,3 @@
|
|||||||
|
#!/bin/bash
|
||||||
|
mkdir -p /tmp/adsb
|
||||||
|
exec "$@"
|
||||||
@@ -0,0 +1,3 @@
|
|||||||
|
uvicorn
|
||||||
|
fastapi
|
||||||
|
pydantic-settings
|
||||||
@@ -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