27 Commits
Author SHA1 Message Date
Logan CusanoandClaude Opus 5.5 6bf2c9dd93 Merge feat/sdr-pins: pin each SDR service to a dongle by serial; OP25 always by serial (closes #11)
CI / lint (push) Successful in 10s
CI / test (push) Successful in 51s
Build edge-node / build (push) Successful in 58s
Build secondary-sdr / build (push) Successful in 2m39s
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:41:28 -04:00
Logan CusanoandClaude Opus 5.5 575a99feb5 checkin: report an unreachable secondary-sdr as null, never keep stale 'Running'
QA blocker: when the secondary-sdr container was down, the checkin
omitted secondary_sdr_running / sdr_devices / op25_sdr_serial, and C2
only overwrites keys that are present, so the dashboard kept the last
'Running' forever. Send explicit nulls. The local card also says 'not
reported' rather than 'not plugged in' for a pin it can't check.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:41:27 -04:00
Logan CusanoandClaude Opus 5.5 9b50ba8114 Pin each SDR service to a dongle by serial; OP25 always opens its SDR by serial
CI / lint (push) Successful in 6s
CI / test (push) Successful in 41s
Fixes node-26#11. OP25's generated config said "rtl" (= whichever dongle
enumerates first), so on 2-SDR nodes a decoder could take OP25's dongle
and stop recording. Now:

- sdr_pins: {op25|adsb|ais: serial}, absent = automatic. Every OP25
  config generation (op25_client.generate_config) rewrites the device to
  rtl=<serial>: the pin, else the first dongle's serial, which is what
  "rtl" always opened. Left as "rtl" only for unknown/shared serials.
- secondary-sdr: /secondary/devices lists dongles + serials via librtlsdr
  (works while claimed). apply(priority, pins, reserved) never touches
  OP25's dongle, gives a pinned service only its own dongle, lets a
  higher-priority service take a spare from a lower one, and still runs a
  pinned lower-priority service when the top pick has no dongle.
- sdr_settings.py replaces secondary_priority.py: one apply path for the
  local dashboard, the new set_sdr_config C2 command (set_secondary_priority
  kept as an alias) and config pushes. OP25 restarts only when its own
  dongle changes. Checkin reports sdr_devices, sdr_pins, op25_sdr_serial.
- Local dashboard: 'SDRs' card with an OP25 SDR dropdown and a per-service
  dongle dropdown, duplicate-serial and double-pin warnings.

Verified: edge-node pytest 194 passed; secondary-sdr tests 7 passed;
flake8 clean; page JS passes node --check.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:35:22 -04:00
Logan CusanoandClaude Opus 5.5 3ecc7eff1a checkin: take sdr_count from secondary-sdr first; op25's 0 is never real
CI / lint (push) Successful in 7s
Build edge-node / build (push) Successful in 44s
CI / test (push) Successful in 48s
op25's :stable image has no lsusb, and /op25/devices answers count 0
instead of unknown, so the dashboard said radio-box 'reports 0 SDRs'
with 2 plugged in. Prefer the secondary-sdr container's lsusb count;
treat 0 from op25 as unknown.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:23:30 -04:00
Logan CusanoandClaude Opus 5.5 358766d88b Merge feat/secondary-sdr-priority: every spare SDR runs the next decoder in priority order
CI / lint (push) Successful in 11s
CI / test (push) Successful in 55s
Build edge-node / build (push) Successful in 1m12s
Build secondary-sdr / build (push) Successful in 3m29s
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:14:53 -04:00
Logan CusanoandClaude Opus 5.5 c82fc61910 Local Secondary SDRs card: 'Unsaved' outranks 'Running'; unknown when unreachable
QA blocker: a reordered but unsaved list still showed the old order's
decoder as 'Running'. Also shows 'Unknown' rather than 'Waiting for SDR'
when the secondary-sdr service doesn't answer (server-26#187).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-27 14:14:49 -04:00
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
logan 31e1176c45 install.sh: re-run collects the key via pickup_secret, never re-enrolls (#7)
CI / test (push) Successful in 39s
CI / lint (push) Successful in 5s
2026-09-06 22:32:19 -04:00
Logan CusanoandClaude Sonnet 5 4f942bd770 install.sh: re-run collects the key via pickup_secret, never re-enrolls
CI / lint (push) Successful in 12s
CI / test (push) Successful in 51s
CI / test (pull_request) Failing after 11m12s
CI / lint (pull_request) Failing after 11m21s
Confirmed prod bug: §5 keyed idempotency on configs/credentials.json only.
Re-running the installer on a pending-then-approved node (the exact flow the
script's own output tells you to do) found no credentials.json and fell
through to a fresh POST /nodes/enroll — which enrollment.py's CRITICAL GUARD
answers with 403 for an already-approved node_id, rendered as "use Reissue
key". The working path (GET /nodes/{id}/credentials with the saved
pickup_secret — no already-approved guard on that endpoint) was only ever
tried within a single run.

§5 rewritten as a decision tree that runs before any POST /nodes/enroll:
  - credentials.json has api_key            -> skip (unchanged)
  - configs/pickup_secret exists            -> GET /credentials with it:
      200 + api_key   -> write credentials.json, done
      200, no key     -> say "approve it, re-run"; clean exit, start stack
      401 (rotated)   -> re-enroll iff --token, else specific die
      404 (deleted)   -> re-enroll iff --token, else specific die
      000             -> connectivity die
  - no creds, no pickup_secret              -> fresh enroll (do_fresh_enroll)

Also: fresh-enroll 403/401/429 handlers are now specific and point at
pickup-secret recovery, not just "Reissue key"; --wait-approval polling now
applies on the re-run path; §7's "re-run the installer" banner is
conditional on pickup_secret existing.

Response shapes verified against enrollment.py @ v1. approve_node mints
node_keys/{id}.api_key synchronously — no second bug; assigning a system is
independent and not required for a key. bash -n clean.

Recovery for a node stuck by the old behaviour (approved, no credentials.json,
pickup_secret on disk): re-run the patched install.sh (no --token needed), or
  curl -fsS $C2_URL/nodes/$NODE_ID/credentials \
    -H "X-Pickup-Secret: $(cat configs/pickup_secret)" | jq '{api_key}' > configs/credentials.json

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-06 21:37:14 -04:00
logan 94c3e2a952 Merge pull request 'One-shot install.sh for a fresh Pi; retire setup.sh (node-26#4)' (#5) from feat/one-shot-install into main
CI / lint (push) Successful in 7s
CI / test (push) Successful in 42s
Reviewed-on: #5
2026-09-06 19:29:30 -04:00
Logan CusanoandClaude Sonnet 5 5b0048e8db node: one-shot install.sh for a fresh Pi; retire setup.sh (node-26#4)
CI / lint (pull_request) Successful in 10s
CI / lint (push) Successful in 11s
CI / test (pull_request) Successful in 46s
CI / test (push) Successful in 46s
install.sh is the `curl -fsSL <url> | sudo bash -s -- --token ...` bootstrap:
preflight (root/arch/apt/SDR) → install docker + compose + git + jq → clone
node-26 at a pinned ref (default `v1`, `--track-main` opt-in) → write .env
non-interactively from flags/env with a /dev/tty interactive fallback →
enroll with C2 (POST /nodes/enroll, poll GET /nodes/{id}/credentials, matches
drb-c2-core/app/routers/enrollment.py exactly) and write configs/credentials.json
→ docker compose pull && up -d (prebuilt; --build opts into the ~1h op25 build)
→ print the admin-approval step. Idempotent: re-run picks up the api_key after
approval; existing .env is preserved.

- setup.sh deleted — two scripts writing .env drift. install.sh owns it now.
- Makefile `setup:` no longer calls the removed script (cp .env.example fallback).
- README Setup section rewritten around the one-liner; `make setup`/`make up`
  kept as the local-dev path.

Notes carried in the script header: git.vpn.cusano.net is public (D1, settled);
`v1` must be re-cut at this change's merge commit so the tag actually contains
install.sh. Standing hazard D3: docker-compose.yml bind-mounts the app source
over the image, so a pinned ref and the pulled image tags must not diverge.

Client-side enrollment still belongs in the edge-node app (mqtt_manager.py:73-85);
install.sh doing it is the interim. Tracked for follow-up.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
2026-09-06 19:14:40 -04:00
Logan CusanoandClaude Opus 5 0c08275482 Stop throwing away the audio Whisper has to read
Build edge-node / build (push) Successful in 35s
CI / lint (push) Successful in 7s
CI / test (push) Successful in 42s
The single encode at save time was %mp3(bitrate=16), and the comment said why:
it matched what Liquidsoap pushes to Icecast. That was the wrong thing to
match. Icecast is the LISTENING path and 16 kbps is a bandwidth budget for a
live stream; this file is the ACCURACY path -- it is what Whisper transcribes,
and CLAUDE.md is explicit that everything downstream is hostage to it. P25 has
already been through a vocoder, so 16 kbps MP3 stacked a second lossy stage on
the one copy that had to stay faithful.

FLAC instead. Lossless, so the bytes Whisper receives are the bytes PulseAudio
captured. ~1.3 MB/min against 120 KB/min, which keeps a 600 s call (the time
cap) around 13 MB -- inside Whisper's 25 MB request cap and well inside
upload_max_bytes. Icecast's own 16 kbps stream is untouched; nothing about
live listening changes.

Capture, buffering, silence detection and the byte-offset trim are all
unchanged: they operate on raw PCM and never saw the encode. The sample rate
stays pinned to pcm.SAMPLE_RATE so the encode remains a straight pass -- the
trim arithmetic depends on that, and Whisper resamples to 16 kHz itself.

encode_mp3 is now encode_recording, the upload sends audio/flac, and the test
that pinned the old contract now pins losslessness instead, including an
assertion that no bitrate constant comes back.

This is the before/after boundary for STT quality. Last night's window is the
16 kbps baseline.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-23 12:48:32 -04:00
Logan Cusano fb13bb8ae3 Pin *.sh to LF so Windows checkouts cannot ship a CRLF shebang
CI / lint (push) Successful in 10s
CI / test (push) Successful in 57s
2026-08-20 03:13:10 -04:00
Logan Cusano e6aab7589a Stop shipping "hackme" as the Icecast password
Build edge-node / build (push) Failing after 6s
CI / lint (push) Successful in 6s
CI / test (push) Successful in 40s
Build op25 / build (push) Successful in 1m48s
Build icecast / build (push) Successful in 3m16s
The source password had a `hackme` fallback in five places -- entrypoint.sh,
docker-compose.yml, setup.sh's prompt default, .env.example, edge-node's
config.py -- plus op25-container's os.getenv default and two README rows. Any
node whose operator pressed Enter through setup.sh is running a credential
that is written down in this repo.

That matters more than the usual default-password case because Icecast binds
all interfaces and the SOURCE password is write access: it does not just let a
LAN neighbour listen, it lets them PUSH audio into the stream the frontend and
mobile clients play as live radio. Injecting fake traffic into a public-safety
feed is the failure worth preventing.

Approach: remove every fallback rather than change them to a better default.
- icecast/entrypoint.sh refuses to start if either password is empty, and says
  how to generate one. This is the single hard gate; everything else is
  defence in depth behind it.
- docker-compose.yml uses ${VAR:?message} so a missing value stops the stack
  at compose time with a readable error instead of becoming an empty string.
- setup.sh GENERATES a random password when the operator presses Enter, via
  openssl rand -base64 24 with a /dev/urandom fallback. Pressing Enter now
  gives you a random password rather than a known one, which is the actual
  behaviour change -- a prompt default nobody types over is not a default, it
  is the value.
- .env.example ships the keys empty with the generation command in a comment,
  and README.md now marks both as required with no default.

Client suite: 185 passed.

Note this does NOT rotate anything already deployed. node-002's .env still has
whatever it was set up with; that is an operational step, tracked in the issue.

Closes logan/node-26#3
2026-08-20 03:12:44 -04:00
Logan Cusano 28266b4441 Pin the op25 container to a Python major version
CI / lint (push) Successful in 6s
Build op25 / build (push) Successful in 21s
CI / test (push) Successful in 38s
`python:slim-trixie` carried no version at all, so a rebuild could move the
interpreter across a major release without anything in the repo changing. That
is not hypothetical here: app/models.py referenced IcecastConfig about 85 lines
before its definition and ran only because trixie currently ships Python 3.14,
where PEP 649 defers annotation evaluation. On 3.13 it was a hard NameError.
The ordering was fixed on 2026-08-16; the unpinned base outlived it.

Pinned to 3.14-slim rather than 3.14-slim-trixie so it matches drb-edge-node,
which was already on 3.14-slim. Patch releases still float, which is what we
want for security updates -- only the major version is nailed down.

Every other Dockerfile in both repos already pinned a major version
(python:3.12-slim, python:3.14-slim, node:20-slim, debian:bookworm-slim), so
this was the only genuinely unpinned base image, despite server-26#11 claiming
none of them were pinned.

Closes logan/server-26#11 (filed against the wrong repo -- the file lives in
the client repo).
2026-08-20 03:01:38 -04:00
Logan CusanoandClaude Opus 5 7f1d09c753 Satisfy flake8 so the lint job stops failing
CI / lint (push) Successful in 6s
CI / test (push) Successful in 39s
Build edge-node / build (push) Successful in 7m44s
Pure formatting, no behaviour change: strip trailing whitespace from blank
lines, give top-level defs in system_cacher.py their two blank lines, wrap
the long discord_radio.join signature, and split the duplicated
active_config ternary in main.py and routers/api.py across lines.

Verified clean with flake8 --max-line-length=120, matching CI.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-16 13:56:22 -04:00
Logan CusanoandClaude Opus 5 9d86304b8a Revert the CI full-fetch workaround and drop unreachable workflows
CI / lint (push) Failing after 17s
CI / test (push) Successful in 49s
Build op25 / build (push) Successful in 1h29m25s
Shallow clones were never a Gitea packing bug. An intruder had set
uploadpack.packObjectsHook in Gitea's HOME gitconfig, pointing at a
non-executable dropper, so every upload-pack died mid-pack. That hook is
gone and --depth=1 clones are verified working, so fetch-depth: 0 buys
nothing but slower CI. See INCIDENT-2026-08-11.md.

op25-container/.gitea/workflows/* never ran: Gitea only executes
workflows under .gitea/workflows at the repo root, and op25-container is
a subdirectory of this repo, not a repo of its own. The live OP25 image
build is .gitea/workflows/build-op25.yml.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-16 13:00:02 -04:00
Logan CusanoandClaude Opus 5 5ff1089551 Use a full fetch in CI: Gitea fails to pack a shallow clone
CI / test (push) Failing after 30s
CI / lint (push) Failing after 39s
actions/checkout defaults to depth=1, and Gitea aborted generating that pack
with a bad pack header protocol error on every retry, failing the run before
any image was built. Same fix already applied on the server repo; applied here
to all four workflows so the edge-node, op25 and icecast images can build.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-16 10:09:48 -04:00
40 changed files with 2201 additions and 374 deletions
+11 -2
View File
@@ -38,8 +38,10 @@ C2_URL=http://localhost:8888
# Icecast (local container — usually no need to change)
# Live listening only. Call recording and Discord voice use PulseAudio instead.
ICECAST_SOURCE_PASSWORD=hackme
ICECAST_ADMIN_PASSWORD=admin
# REQUIRED, no default — the container refuses to start without them.
# Generate with: openssl rand -base64 24
ICECAST_SOURCE_PASSWORD=
ICECAST_ADMIN_PASSWORD=
ICECAST_HOST=localhost
ICECAST_PORT=8000
ICECAST_MOUNT=/radio
@@ -109,6 +111,13 @@ OP25_TERMINAL_URL=http://localhost:8081
# for local development off a real node; leave false everywhere else.
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 ---------------------------------------------
# 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
+4
View File
@@ -0,0 +1,4 @@
# Shell scripts run inside Linux containers. A CRLF shebang there fails as
# "bad interpreter: /bin/sh^M", which surfaces only as a container that will
# not start. Windows checkouts have core.autocrlf=true, so pin these to LF.
*.sh text eol=lf
+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
+6 -1
View File
@@ -1,7 +1,12 @@
.PHONY: setup test up up-prebuilt pull down logs
# Local-dev only: seed .env from the example so `make up` has something to read.
# A real edge node is provisioned with install.sh (one-shot bootstrap for a
# clean Pi — clones, enrolls with C2, pulls prebuilt images):
# curl -fsSL https://git.vpn.cusano.net/logan/node-26/raw/tag/v1/install.sh | sudo bash -s -- --help
setup:
@bash setup.sh
@test -f .env || cp .env.example .env
@echo ".env ready — edit it, then 'make up' (local build) or 'make up-prebuilt'"
# Run pytest inside the running edge-node container.
# Requires: docker compose up (or at least the edge-node image built).
+25 -14
View File
@@ -135,21 +135,32 @@ Client/
## Setup
### Provision a real node — `install.sh`
One-shot bootstrap for a clean Raspberry Pi OS (arm64). Installs Docker, clones
this repo at the `v1` tag, enrols with C2, pulls the prebuilt images and starts:
```bash
# 1. Copy env template
cp .env.example .env
# 2. Fill in at minimum: NODE_ID and MQTT_BROKER
nano .env
# 3. Build all images (op25 takes ~10-15 minutes first time)
docker compose build
# 4. Start
docker compose up -d
curl -fsSL https://git.vpn.cusano.net/logan/node-26/raw/tag/v1/install.sh \
| sudo bash -s -- --token DRB-xxxx --node-id node-003 \
--c2-url https://api.<domain> --mqtt-broker mqtt.<domain>
```
The node will appear as **pending** in the server admin dashboard. An admin must approve it before it becomes operational. After approval, assign a radio system in the dashboard and the node will start decoding automatically.
Mint the `--token` at **Settings → Nodes** in the web app (the panel prints the
whole command). Run `install.sh --help` for every flag; each also has a
`DRB_*` env var. `--build` compiles op25 on the Pi (~1h) instead of pulling.
### Local dev / manual
```bash
make setup # seeds .env from .env.example
nano .env # at minimum: NODE_ID, MQTT_BROKER, C2_URL
make up # build locally (op25 ~10-15 min first time)
# or: make up-prebuilt # pull images, no local build
```
The node appears as **pending** in the admin dashboard. An admin approves it,
then assigns a radio system, and the node starts decoding automatically.
## Environment Variables (`.env`)
@@ -167,8 +178,8 @@ The node will appear as **pending** in the server admin dashboard. An admin must
| `ICECAST_HOST` | No | `localhost` | Icecast hostname (leave as localhost — host network mode) |
| `ICECAST_PORT` | No | `8000` | Icecast HTTP port |
| `ICECAST_MOUNT` | No | `/radio` | Icecast mount point |
| `ICECAST_SOURCE_PASSWORD` | No | `hackme` | Icecast source password — change this |
| `ICECAST_ADMIN_PASSWORD` | No | `hackme` | Icecast admin password — change this |
| `ICECAST_SOURCE_PASSWORD` | **Yes** | none | Icecast source password. No default — the container refuses to start without it. `install.sh` generates one; otherwise `openssl rand -base64 24` |
| `ICECAST_ADMIN_PASSWORD` | **Yes** | none | Icecast admin password. Same rules |
| `OP25_API_URL` | No | `http://localhost:8001` | OP25 container HTTP API |
| `OP25_TERMINAL_URL` | No | `http://localhost:8081` | OP25 HTTP terminal (live talkgroup metadata) |
+19 -2
View File
@@ -5,8 +5,10 @@ services:
restart: unless-stopped
network_mode: host
environment:
ICECAST_SOURCE_PASSWORD: ${ICECAST_SOURCE_PASSWORD:-hackme}
ICECAST_ADMIN_PASSWORD: ${ICECAST_ADMIN_PASSWORD:-admin}
# :? not :- — a missing password must stop the stack, not silently
# become a credential that is published in this file.
ICECAST_SOURCE_PASSWORD: ${ICECAST_SOURCE_PASSWORD:?set ICECAST_SOURCE_PASSWORD in .env (run setup.sh, or openssl rand -base64 24)}
ICECAST_ADMIN_PASSWORD: ${ICECAST_ADMIN_PASSWORD:?set ICECAST_ADMIN_PASSWORD in .env (run setup.sh, or openssl rand -base64 24)}
# No `ports:` here — network_mode: host makes it a no-op either way. The
# control API (:8001) and OP25's HTTP terminal (:8081) are unauthenticated,
@@ -30,6 +32,21 @@ services:
depends_on:
- 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:
image: ${IMAGE_REGISTRY:-git.vpn.cusano.net}/${DOCKER_ORG:-logan}/${DOCKER_REPO:-node-26}/edge-node:latest
build: ./drb-edge-node
+5 -1
View File
@@ -42,7 +42,8 @@ class Settings(BaseSettings):
icecast_host: str = "localhost"
icecast_port: int = 8000
icecast_mount: str = "/radio"
icecast_source_password: str = "hackme"
# No default: see icecast/entrypoint.sh, which refuses to start without one.
icecast_source_password: str = ""
# PulseAudio — the low-latency path used for call recording and Discord voice.
# Liquidsoap (op25 container) writes into the `drb_sink` null sink; we capture
@@ -135,6 +136,9 @@ class Settings(BaseSettings):
op25_api_url: str = "http://localhost:8001"
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)
config_path: str = "/configs"
recordings_path: str = "/recordings"
+46 -23
View File
@@ -7,14 +7,16 @@ A persistent capture process runs for the lifetime of the node. Spawning FFmpeg
per call used to lose the first 1-2 s to process startup, which meant short
transmissions produced empty files, so capture never stops.
RAW PCM, NOT MP3 — this is the change everything else hangs off. FFmpeg is
asked for s16le/22050/mono on stdout instead of an MP3 stream, so:
RAW PCM, NOT A COMPRESSED STREAM — this is the change everything else hangs
off. FFmpeg is asked for s16le/22050/mono on stdout instead of an encoded
stream, so:
* silence detection is integer arithmetic over each chunk as it arrives, with
no decode, which is what makes AUDIO-DRIVEN call boundaries possible;
* trimming is a byte-offset slice, not a second FFmpeg pass;
* MP3 encoding happens exactly ONCE, at save time, so uploads are no longer
double-encoded.
* encoding happens exactly ONCE, at save time, so uploads are no longer
double-encoded. That encode is now FLAC (lossless) rather than 16 kbps MP3
— see the AUDIO_* constants below for why.
TWO BUFFERS, TWO JOBS — this split is load-bearing:
@@ -94,14 +96,32 @@ RING_BUFFER_SECONDS = 30
# under PRE_ROLL_SECONDS and well under the shortest utterance we care about.
READ_CHUNK_BYTES = 2048
# Encoder settings for the single encode at save time, matched on purpose to
# what Liquidsoap already pushes to Icecast — %mp3(bitrate=16, samplerate=22050,
# stereo=false) — so the C2 /upload endpoint keeps receiving exactly the kind of
# MP3 it has always received (multipart "audio/mpeg", stored to GCS as .mp3,
# then fed to Whisper). MP3_SAMPLE_RATE MUST equal pcm.SAMPLE_RATE: the encode
# is a straight pass with no resampling.
MP3_BITRATE = "16k"
MP3_SAMPLE_RATE = str(pcm.SAMPLE_RATE)
# Encoder settings for the single encode at save time.
#
# This used to be %mp3(bitrate=16) chosen to match what Liquidsoap pushes to
# Icecast. That was the wrong thing to match: Icecast is the LISTENING path and
# 16 kbps is a bandwidth budget for a live stream, while this file is the
# ACCURACY path — it is what Whisper transcribes, and the transcript is what
# every downstream stage is hostage to. P25 audio has already been through a
# vocoder; 16 kbps MP3 stacked a second lossy stage on top of that, on the one
# copy that had to stay faithful.
#
# FLAC instead: lossless, so the bytes Whisper receives are the bytes PulseAudio
# captured. Roughly 1.3 MB/min against 120 KB/min for 16k MP3 — larger, but a
# 600 s call is still ~13 MB, inside both Whisper's 25 MB request cap and the
# C2 upload_max_bytes (100 MB). Icecast's own 16 kbps stream is untouched;
# nothing about the listening path changes.
#
# AUDIO_SAMPLE_RATE MUST equal pcm.SAMPLE_RATE: the encode is a straight pass
# with no resampling. Whisper resamples to 16 kHz itself, so handing it 22050
# unresampled keeps the one resample in the pipeline inside the model.
AUDIO_SAMPLE_RATE = str(pcm.SAMPLE_RATE)
AUDIO_FORMAT = "flac"
AUDIO_SUFFIX = ".flac"
AUDIO_MIME = "audio/flac"
# -compression_level 5 is ffmpeg's default: near-best ratio, and the encode is
# off the hot path anyway (once per call, at save).
FLAC_COMPRESSION_LEVEL = "5"
# Bounded so a wedged encoder can never stall the upload path.
ENCODE_TIMEOUT_SECONDS = 60.0
@@ -218,9 +238,12 @@ class Recording:
all_silence: bool = False
async def encode_mp3(audio: bytes, path: Path) -> bool:
async def encode_recording(audio: bytes, path: Path) -> bool:
"""
The one and only encode in the pipeline: raw PCM in, MP3 file out.
The one and only encode in the pipeline: raw PCM in, FLAC file out.
Lossless on purpose — see the AUDIO_* constants above. This file is what
Whisper transcribes, so the encode must not throw anything away.
Module-level rather than a method so tests can substitute it without
needing FFmpeg, and so the "exactly one encode per call" property is
@@ -233,13 +256,13 @@ async def encode_mp3(audio: bytes, path: Path) -> bool:
"-hide_banner", "-nostdin", "-nostats",
"-loglevel", "warning", "-y",
"-f", "s16le",
"-ar", MP3_SAMPLE_RATE,
"-ar", AUDIO_SAMPLE_RATE,
"-ac", str(pcm.CHANNELS),
"-i", "pipe:0",
"-ar", MP3_SAMPLE_RATE,
"-ar", AUDIO_SAMPLE_RATE,
"-ac", str(pcm.CHANNELS),
"-b:a", MP3_BITRATE,
"-f", "mp3", str(path),
"-compression_level", FLAC_COMPRESSION_LEVEL,
"-f", AUDIO_FORMAT, str(path),
]
try:
proc = await asyncio.create_subprocess_exec(
@@ -327,7 +350,7 @@ class CallRecorder:
"-loglevel", "warning",
"-f", "pulse", "-i", settings.pulse_source,
"-ac", str(pcm.CHANNELS),
"-ar", MP3_SAMPLE_RATE,
"-ar", AUDIO_SAMPLE_RATE,
# Raw PCM on stdout. No muxer, so no -flush_packets games: s16le is
# a bare byte stream and every byte FFmpeg produces is immediately
# readable, which is what keeps arrival timestamps honest.
@@ -360,7 +383,7 @@ class CallRecorder:
async def _run_capture(self) -> None:
cmd = self._ffmpeg_command()
logger.info(f"Starting capture: ffmpeg -f pulse -i {settings.pulse_source} (s16le/{MP3_SAMPLE_RATE}/mono)")
logger.info(f"Starting capture: ffmpeg -f pulse -i {settings.pulse_source} (s16le/{AUDIO_SAMPLE_RATE}/mono)")
self._last_stderr_lines.clear()
proc = await asyncio.create_subprocess_exec(
*cmd,
@@ -674,9 +697,9 @@ class CallRecorder:
self._recordings_dir.mkdir(parents=True, exist_ok=True)
ts_str = datetime.now(timezone.utc).strftime("%Y%m%d_%H%M%S")
output_path = self._recordings_dir / f"{ts_str}_{recording.call_id}.mp3"
output_path = self._recordings_dir / f"{ts_str}_{recording.call_id}{AUDIO_SUFFIX}"
if not await encode_mp3(audio, output_path):
if not await encode_recording(audio, output_path):
output_path.unlink(missing_ok=True)
return None
@@ -766,7 +789,7 @@ class CallRecorder:
with open(file_path, "rb") as f:
r = await client.post(
upload_url,
files={"file": (file_path.name, f, "audio/mpeg")},
files={"file": (file_path.name, f, AUDIO_MIME)},
data=form,
headers=headers,
)
+2 -1
View File
@@ -27,7 +27,8 @@ class RadioBot:
self._channel_id: Optional[int] = None
self._was_streaming: bool = False
async def join(self, guild_id: int, channel_id: int, token: str, call_active: bool = False, system_name: str = None) -> bool:
async def join(self, guild_id: int, channel_id: int, token: str,
call_active: bool = False, system_name: str = None) -> bool:
# (Re)start the bot if the token changed or the bot isn't running
if self._current_token != token or not self._is_bot_running():
if not await self._start_bot(token):
@@ -210,10 +210,16 @@ class MQTTManager:
logger.info("No API key on disk — requesting re-delivery from C2 server.")
self._publish(self._t_key_request, {}, qos=1)
async def publish_checkin(self):
await self._publish_checkin()
async def _publish_checkin(self):
from app.internal.discord_radio import radio_bot
from app.internal.config_manager import load_node_config
from app.internal.op25_client import op25_client
from app.internal.secondary_sdr_client import secondary_sdr_client
config = load_node_config()
devices = await op25_client.devices()
payload = {
"node_id": settings.node_id,
"name": settings.node_name,
@@ -225,7 +231,35 @@ class MQTTManager:
"is_overridden": config.override_system_id is not None and config.node_type != "portable",
"override_system_id": config.override_system_id,
"enforce_override_timeout": config.enforce_override_timeout,
"secondary_sdr_mode": config.secondary_sdr_mode,
"secondary_sdr_priority": config.secondary_sdr_priority,
}
payload["sdr_pins"] = config.sdr_pins
secondary = await secondary_sdr_client.status()
# Unreachable secondary-sdr: say so with explicit nulls. C2 only
# overwrites keys present in the checkin, so omitting them would leave
# a dead decoder showing "Running" forever.
payload["secondary_sdr_running"] = None
payload["sdr_devices"] = None
payload["op25_sdr_serial"] = None
if secondary is not None:
payload["secondary_sdr_running"] = [r["mode"] for r in secondary.get("running", [])]
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)
def _publish(self, topic: str, payload: dict, qos: int = 0, retain: bool = False):
+15 -1
View File
@@ -65,15 +65,29 @@ class OP25Client:
logger.error(f"OP25 status failed: {e}")
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:
try:
async with httpx.AsyncClient(timeout=10) as client:
r = await client.post(f"{self.api_url}/op25/generate-config", json=config)
r.raise_for_status()
return True
except Exception as e:
logger.error(f"OP25 generate-config failed: {e}")
return False
# Every generated config opens its dongle by serial, never "first
# found" (node-26#11) — here so no generation path can skip it.
from app.internal.sdr_settings import pin_op25_device
await pin_op25_device()
return True
async def poll_terminal(self) -> Optional[TerminalUpdate]:
"""
+102
View File
@@ -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()
@@ -8,6 +8,7 @@ from app.internal import credentials
_CACHE_FILE = Path(settings.config_path) / "systems_cache.json"
async def fetch_and_cache_systems() -> bool:
"""Fetch all systems from the C2 server and cache them locally."""
if not settings.c2_url:
@@ -32,6 +33,7 @@ async def fetch_and_cache_systems() -> bool:
logger.warning(f"Failed to fetch systems from C2: {e}. Offline cache will be used.")
return False
def load_cached_systems() -> List[Dict[str, Any]]:
"""Load cached systems from disk."""
if _CACHE_FILE.exists():
@@ -41,6 +43,7 @@ def load_cached_systems() -> List[Dict[str, Any]]:
logger.error(f"Failed to read systems cache: {e}")
return []
def get_cached_system(system_id: str) -> Optional[Dict[str, Any]]:
"""Retrieve a single system config from the cache."""
systems = load_cached_systems()
@@ -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"]})
+40 -2
View File
@@ -5,7 +5,7 @@ from typing import Optional
from fastapi import FastAPI
from app.config import settings
from app.models import SystemConfig
from app.models import SystemConfig, normalize_sdr_pins, normalize_secondary_priority
from app.internal.logger import logger
from app.internal.mqtt_manager import mqtt_manager
from app.internal import credentials
@@ -158,6 +158,12 @@ async def on_command(payload: dict):
)
elif action == "discord_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":
from app.internal.op25_client import op25_client
await op25_client.stop()
@@ -224,6 +230,9 @@ async def on_config_push(payload: dict):
hardware_preset = payload.pop("hardware_preset", None)
ppm_override = payload.pop("ppm_override", None)
node_type = payload.pop("node_type", None)
secondary_sdr_priority = payload.pop("secondary_sdr_priority", None)
sdr_pins = payload.pop("sdr_pins", None)
secondary_sdr_mode = payload.pop("secondary_sdr_mode", None) # legacy single-mode C2
enforce_override_timeout = payload.pop("enforce_override_timeout", None)
try:
config = SystemConfig(**payload)
@@ -245,6 +254,16 @@ async def on_config_push(payload: dict):
node_cfg.node_type = node_type
if enforce_override_timeout is not None:
node_cfg.enforce_override_timeout = bool(enforce_override_timeout)
if secondary_sdr_priority is None and secondary_sdr_mode is not None:
secondary_sdr_priority = [secondary_sdr_mode]
if secondary_sdr_priority is not None:
node_cfg.secondary_sdr_priority = normalize_secondary_priority(secondary_sdr_priority)
node_cfg.secondary_sdr_mode = (node_cfg.secondary_sdr_priority or ["none"])[0]
if sdr_pins is not None:
try:
node_cfg.sdr_pins = normalize_sdr_pins(sdr_pins)
except ValueError as e:
logger.error(f"Ignoring invalid sdr_pins in config push: {e}")
save_node_config(node_cfg)
from app.internal.op25_client import op25_client
@@ -252,10 +271,16 @@ async def on_config_push(payload: dict):
logger.error(f"Failed to generate OP25 config for {config.name}")
return
# OP25's (pinned) dongle may be one a decoder holds right now: free the
# spares, restart OP25 on its own dongle, then hand the rest back out.
from app.internal.secondary_sdr_client import secondary_sdr_client
from app.internal.sdr_settings import apply_secondaries
await secondary_sdr_client.apply([], {}, [])
await op25_client.stop()
await asyncio.sleep(2)
await op25_client.start()
logger.info(f"Config push applied: {config.name}")
await apply_secondaries()
# ---------------------------------------------------------------------------
@@ -297,7 +322,11 @@ async def lifespan(app: FastAPI):
initial_status = "online" if node_cfg.configured else "unconfigured"
await mqtt_manager.publish_status(initial_status)
active_config = node_cfg.override_config if (node_cfg.override_system_id and node_cfg.override_config) else node_cfg.system_config
active_config = (
node_cfg.override_config
if (node_cfg.override_system_id and node_cfg.override_config)
else node_cfg.system_config
)
if node_cfg.configured and active_config:
from app.internal.op25_client import op25_client
logger.info("Node is configured — waiting for OP25 API then generating config.")
@@ -308,12 +337,21 @@ async def lifespan(app: FastAPI):
logger.warning(f"OP25 not ready yet (attempt {attempt + 1}/10), retrying in 3s…")
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())
from app.internal.telemetry_uplink import telemetry_uplink_loop
telemetry_task = asyncio.create_task(telemetry_uplink_loop())
yield # --- app running ---
logger.info("Edge node shutting down.")
heartbeat_task.cancel()
telemetry_task.cancel()
await metadata_watcher.stop()
await call_recorder.stop()
await radio_bot.stop()
+42 -2
View File
@@ -1,5 +1,5 @@
from pydantic import BaseModel
from typing import Optional, Dict, Any
from pydantic import BaseModel, model_validator
from typing import Optional, Dict, Any, Iterable, List
from enum import Enum
from datetime import datetime
@@ -23,6 +23,34 @@ class SystemConfig(BaseModel):
config: Dict[str, Any] # OP25-compatible config blob passed through to op25-container
# Decoders a node can run on the SDRs beyond op25's (node-26#9). op25 always
# keeps its own dongle; each further dongle runs the next entry of the node's
# secondary_sdr_priority, so a 3-SDR node with ["adsb", "ais"] runs both.
SECONDARY_SDR_MODES = ("adsb", "ais")
def normalize_secondary_priority(items: Iterable[str]) -> List[str]:
"""Known modes only, first occurrence wins, order preserved."""
out: List[str] = []
for m in items or []:
if m in SECONDARY_SDR_MODES and m not in out:
out.append(m)
return out
# Services that can be bound to a specific dongle by its USB serial.
SDR_PIN_KEYS = ("op25",) + SECONDARY_SDR_MODES
def normalize_sdr_pins(pins: Dict[str, Optional[str]]) -> Dict[str, str]:
"""Known services only, blank = automatic (dropped). Two services can't
share a dongle, so a serial claimed twice raises."""
out = {k: str(v).strip() for k, v in (pins or {}).items() if k in SDR_PIN_KEYS and v and str(v).strip()}
if len(set(out.values())) != len(out):
raise ValueError("Two services can't be pinned to the same SDR.")
return out
class NodeConfig(BaseModel):
node_id: str
node_name: str
@@ -34,11 +62,23 @@ class NodeConfig(BaseModel):
hardware_preset: str = "rtl-sdr-v3"
ppm_override: Optional[float] = None
node_type: str = "fixed" # fixed or portable
secondary_sdr_priority: List[str] = [] # ordered; SDRs beyond op25's run these top-down
sdr_pins: Dict[str, str] = {} # service (op25/adsb/ais) -> dongle serial; absent = automatic
# Legacy single-mode field (pre-priority). Still written as priority[0] so
# an older C2 reading checkins sees something sensible; read only to
# migrate a node_config.json saved before priority existed.
secondary_sdr_mode: str = "none"
enforce_override_timeout: bool = True
override_system_id: Optional[str] = None
override_config: Optional[SystemConfig] = None
offline_call_buffer_size: int = 35 # max call_end events to buffer while MQTT is offline
@model_validator(mode="after")
def _migrate_secondary_mode(self) -> "NodeConfig":
if not self.secondary_sdr_priority and self.secondary_sdr_mode in SECONDARY_SDR_MODES:
self.secondary_sdr_priority = [self.secondary_sdr_mode]
return self
class CallEvent(BaseModel):
call_id: str
+46 -3
View File
@@ -1,9 +1,10 @@
from fastapi import APIRouter, Depends, HTTPException, Body
from typing import Optional
from pydantic import BaseModel
from typing import Dict, List, Optional
import asyncio
import httpx
from app.config import settings
from app.models import SystemConfig
from app.models import SystemConfig, SECONDARY_SDR_MODES
from app.internal.op25_client import op25_client
from app.internal.config_manager import load_node_config, save_node_config, apply_system_config
from app.internal.call_recorder import call_recorder
@@ -30,7 +31,11 @@ async def get_status():
active_tgid_name = metadata_watcher.current_tgid_name
system_name = None
active_config = node_cfg.override_config if (node_cfg.override_system_id and node_cfg.override_config) else node_cfg.system_config
active_config = (
node_cfg.override_config
if (node_cfg.override_system_id and node_cfg.override_config)
else node_cfg.system_config
)
if active_config:
system_name = active_config.name
if active_tgid:
@@ -200,6 +205,44 @@ async def ack_override(timeout_minutes: int = Body(1440)):
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")
async def discord_join(guild_id: int, channel_id: int):
ok = await radio_bot.join(guild_id, channel_id)
+184
View File
@@ -105,6 +105,40 @@
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 {
display: grid;
grid-template-columns: repeat(auto-fit, minmax(320px, 1fr));
@@ -360,6 +394,28 @@
</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 -->
<div class="glass-card player-card">
<div class="card-header" style="margin-bottom:0.5rem; border:none;">
@@ -492,6 +548,134 @@
}
}
// ── 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 => ({'&': '&amp;', '<': '&lt;', '>': '&gt;', '"': '&quot;', "'": '&#39;'}[c]));
}
function deviceOptions(selected, autoLabel) {
const devs = sdr.devices || [];
const known = devs.some(d => d.serial === selected);
return `<option value="">${esc(autoLabel)}</option>` +
devs.map(d => `<option value="${esc(d.serial)}" ${d.serial === selected ? 'selected' : ''}>` +
`SDR ${d.index + 1} · serial ${esc(d.serial)}${d.duplicate_serial ? ' (shared serial!)' : ''}</option>`).join('') +
(selected && !known ? `<option value="${esc(selected)}" selected>serial ${esc(selected)} ` +
`(${sdr.devices === null ? 'not reported' : 'not plugged in'})</option>` : '');
}
function renderSdr() {
const op25 = document.getElementById('sdr-op25');
op25.innerHTML = deviceOptions(sdrOp25Pin, 'Automatic (first SDR)');
op25.onchange = () => { sdrOp25Pin = op25.value; markSdrDirty(); };
document.getElementById('sdr-op25-note').textContent = sdrOp25Pin !== (sdr.pins.op25 || '')
? 'Saving will restart OP25 on the selected SDR.'
: (sdr.op25_serial ? `OP25 is using serial ${sdr.op25_serial}.` : '');
const list = document.getElementById('sdr-list');
list.innerHTML = '';
let rank = 0;
const running = sdr.running || [];
sdrRows.forEach((row, i) => {
const [name, hint] = SDR_LABELS[row.mode] || [row.mode, ''];
const pinMissing = row.pin && sdr.devices !== null && !sdr.devices.some(d => d.serial === row.pin);
const state = !row.enabled ? 'Off'
: sdrDirty ? 'Unsaved'
: sdr.running === null ? 'Unknown'
: running.includes(row.mode) ? 'Running'
: pinMissing ? 'Pinned SDR missing' : 'Waiting for SDR';
const el = document.createElement('div');
el.className = 'sdr-row';
el.innerHTML = `
<span class="sdr-rank">${row.enabled ? ++rank : ''}</span>
<input type="checkbox" ${row.enabled ? 'checked' : ''} aria-label="Enable ${name}">
<span class="sdr-name">${name}<small>${hint}</small></span>
<select class="sdr-select" aria-label="${name} SDR">${deviceOptions(row.pin, 'Any spare SDR')}</select>
<button class="sdr-move" aria-label="Move ${name} up" ${i === 0 ? 'disabled' : ''}>▲</button>
<button class="sdr-move" aria-label="Move ${name} down" ${i === sdrRows.length - 1 ? 'disabled' : ''}>▼</button>
<span class="sdr-state ${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();
setInterval(refresh, 2000); // Polling every 2 seconds
</script>
+18 -8
View File
@@ -59,7 +59,7 @@ def encodes(monkeypatch):
path.write_bytes(audio)
return True
monkeypatch.setattr(recorder_mod, "encode_mp3", _encode)
monkeypatch.setattr(recorder_mod, "encode_recording", _encode)
return calls
@@ -412,21 +412,31 @@ async def test_a_failed_encode_leaves_no_file_and_no_recording(recorder, monkeyp
async def _fail(audio, path):
return False
monkeypatch.setattr(recorder_mod, "encode_mp3", _fail)
monkeypatch.setattr(recorder_mod, "encode_recording", _fail)
ingest(recorder, T0, T0 + 5.0)
await recorder.start_recording("call-1", start_epoch=T0 + 1.0)
assert await recorder.stop_recording(end_epoch=T0 + 3.0) is None
assert list(recorder._recordings_dir.glob("*.mp3")) == []
assert list(recorder._recordings_dir.glob("*.flac")) == []
def test_encoder_command_contract_matches_what_c2_expects():
"""
/upload has always received mono MP3 at 22050 Hz / 16 kbps, and Whisper
consumes it downstream. The single encode must not quietly change that.
The saved file is what Whisper transcribes, so the encode must stay
LOSSLESS and must not resample. It was 16 kbps MP3 — a bitrate copied from
Icecast's live stream, i.e. the listening path's budget applied to the
accuracy path — which put a second lossy stage on top of the P25 vocoder.
The sample rate must equal pcm.SAMPLE_RATE or the encode stops being a
straight pass and the byte-offset trim arithmetic no longer lines up.
"""
assert recorder_mod.MP3_SAMPLE_RATE == str(pcm.SAMPLE_RATE) == "22050"
assert recorder_mod.MP3_BITRATE == "16k"
assert recorder_mod.AUDIO_SAMPLE_RATE == str(pcm.SAMPLE_RATE) == "22050"
assert recorder_mod.AUDIO_FORMAT == "flac"
assert recorder_mod.AUDIO_SUFFIX == ".flac"
assert recorder_mod.AUDIO_MIME == "audio/flac"
# No bitrate constant should exist: a bitrate on a lossless codec would mean
# someone reintroduced lossy encoding.
assert not hasattr(recorder_mod, "MP3_BITRATE")
# ---------------------------------------------------------------------------
@@ -523,7 +533,7 @@ async def test_discard_drops_the_audio_without_writing_anything(recorder, encode
assert not recorder.is_recording
assert encodes == []
assert list(recorder._recordings_dir.glob("*.mp3")) == []
assert list(recorder._recordings_dir.glob("*.flac")) == []
# ...and the recorder is immediately reusable.
assert await recorder.start_recording("call-next", start_epoch=T0 + 2.0) is True
+91
View File
@@ -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"}
+13 -2
View File
@@ -1,8 +1,19 @@
#!/bin/sh
set -e
ICECAST_SOURCE_PASSWORD="${ICECAST_SOURCE_PASSWORD:-hackme}"
ICECAST_ADMIN_PASSWORD="${ICECAST_ADMIN_PASSWORD:-admin}"
# No defaults here on purpose. This container binds all interfaces, so a
# fallback password is a published credential on every node that ever accepted
# it -- and the source password is what lets a caller PUSH audio into the
# stream the frontend plays as live radio. Refuse to start instead.
for var in ICECAST_SOURCE_PASSWORD ICECAST_ADMIN_PASSWORD; do
eval "value=\${$var}"
if [ -z "$value" ]; then
echo "icecast: $var is not set." >&2
echo "icecast: set it in the node's .env -- 'bash setup.sh' generates a random one," >&2
echo "icecast: or run: openssl rand -base64 24" >&2
exit 1
fi
done
export ICECAST_SOURCE_PASSWORD ICECAST_ADMIN_PASSWORD
+545
View File
@@ -0,0 +1,545 @@
#!/usr/bin/env bash
#
# DRB edge node — one-shot bootstrap for a clean Raspberry Pi OS (arm64).
#
# curl -fsSL https://git.vpn.cusano.net/logan/node-26/raw/tag/v1/install.sh \
# | sudo bash -s -- --token DRB-xxxx --node-id node-003 \
# --c2-url https://api.<domain> --mqtt-broker mqtt.<domain>
#
# (fetch install.sh from the same tag it installs; use raw/branch/main with
# --track-main for the owner's own always-latest nodes)
#
# Hosting, the pinned ref, and setup.sh's retirement are settled (owner
# decisions, 2026-09-06). One standing hazard remains — D3 below.
#
# What it does, in order:
# 1. Preflight (root, arch, apt, SDR present)
# 2. Install Docker + compose plugin + git + curl + jq
# 3. Clone logan/node-26 at a PINNED ref into $INSTALL_DIR
# 4. Write .env (non-interactive from flags/env, interactive fallback)
# 5. Enroll with C2 (POST /nodes/enroll) and poll for the api_key
# 6. docker compose pull && up -d (prebuilt; --build opts into the ~1h build)
# 7. Print the approval step
#
# It does NOT set up WireGuard. WireGuard-per-node was evaluated and rejected
# (Server/MQTT-PUBLIC-AUTH-PLAN.md, Server/infra/main.tf:67-73). A field node
# reaches production over the public internet only:
# https://api.<domain> enrollment + /upload (Caddy, real TLS)
# mqtt.<domain>:8883 MQTT over TLS, username=NODE_ID password=api_key
# The two stale "nodes reach it via WireGuard" comments in
# Server/docker-compose.prod.yml:6 and Server/infra/main.tf:69 are leftovers.
#
# ---------------------------------------------------------------------------
# SETTLED (owner decisions, 2026-09-06 — node-26#4)
#
# D1 HOSTING. git.vpn.cusano.net is PUBLIC — it is a CNAME to
# cusano-net.duckdns.org (71.117.95.129) with a real Let's Encrypt cert,
# resolvable from any public resolver. install.sh, `git clone` and
# `docker compose pull` all work from a customer's Pi with no VPN.
# Caveat, not a blocker: that IP is a dynamic-DNS record on the owner's
# home uplink, so it is a single point of failure and a bandwidth limit
# — fine for the beachhead, revisit before scaling node count.
#
# D2 PIN. node-26 gets a `v1` tag (owner cuts it — see the git command in
# the mint panel / node-26#4). DEFAULT_REF below is `v1`. `--track-main`
# stays as an opt-in for the owner's own nodes.
#
# STANDING HAZARD
#
# D3 SOURCE-OVERLAY. docker-compose.yml bind-mounts ./drb-edge-node/app and
# ./op25-container/app OVER the image's /app. So even in prebuilt mode
# the running Python is the CLONED REF's code against the pulled image's
# dependencies. `v1` == the commit CI built the current :latest/:stable
# from, so today they match — but the moment `v1` and the image tags
# diverge this silently mixes them. Fix: drop those two mounts from the
# prod compose path, or always retag images from the same ref as `v1`.
# ---------------------------------------------------------------------------
set -euo pipefail
# ── Defaults ────────────────────────────────────────────────────────────────
DEFAULT_REF="v1" # D2
DEFAULT_REPO_URL="https://git.vpn.cusano.net/logan/node-26.git" # D1 (public)
DEFAULT_INSTALL_DIR="/opt/drb/node-26"
REPO_URL="${DRB_REPO_URL:-$DEFAULT_REPO_URL}"
REF="${DRB_REF:-$DEFAULT_REF}"
INSTALL_DIR="${DRB_INSTALL_DIR:-$DEFAULT_INSTALL_DIR}"
NODE_ID="${DRB_NODE_ID:-}"
NODE_NAME="${DRB_NODE_NAME:-}"
NODE_LAT="${DRB_NODE_LAT:-}"
NODE_LON="${DRB_NODE_LON:-}"
C2_URL="${DRB_C2_URL:-}"
MQTT_BROKER="${DRB_MQTT_BROKER:-}"
MQTT_PORT="${DRB_MQTT_PORT:-8883}"
MQTT_TLS="${DRB_MQTT_TLS:-true}"
ENROLLMENT_TOKEN="${DRB_ENROLLMENT_TOKEN:-}"
DASHBOARD_USER="${DRB_DASHBOARD_USERNAME:-admin}"
DASHBOARD_PASS="${DRB_DASHBOARD_PASSWORD:-}"
REGISTRY="${DRB_IMAGE_REGISTRY:-git.vpn.cusano.net}" # D1 (public host)
DOCKER_ORG="${DRB_DOCKER_ORG:-logan}"
DOCKER_REPO="${DRB_DOCKER_REPO:-node-26}"
REGISTRY_USER="${DRB_REGISTRY_USER:-}"
REGISTRY_PASS="${DRB_REGISTRY_PASS:-}"
DO_BUILD=0 # 0 = pull prebuilt images (default), 1 = build on the Pi (~1h for op25)
DO_START=1
ASSUME_YES=0
ENROLL_WAIT="${DRB_ENROLL_WAIT:-0}" # seconds to block waiting for admin approval; 0 = don't block
C='\033[0;36m'; G='\033[0;32m'; Y='\033[1;33m'; R='\033[0;31m'; N='\033[0m'
say() { printf "${C}==>${N} %s\n" "$*"; }
ok() { printf "${G} ok${N} %s\n" "$*"; }
warn() { printf "${Y} !!${N} %s\n" "$*" >&2; }
die() { printf "${R}error:${N} %s\n" "$*" >&2; exit 1; }
usage() {
cat <<'USAGE'
Usage: install.sh [options]
--token TOKEN Enrollment token (Settings -> Nodes -> New token)
--node-id ID Unique node id, e.g. node-003
--name NAME Display name (default: node id)
--lat N --lon N Decimal degrees for the map
--c2-url URL e.g. https://api.drb.example.net
--mqtt-broker HOST e.g. mqtt.drb.example.net
--mqtt-port N default 8883
--no-tls plaintext MQTT (LAN/dev brokers only)
--dashboard-pass PW local dashboard password (generated if omitted)
--ref REF git ref to install (default: v1)
--track-main install main HEAD instead of the v1 tag
--dir PATH install location (default /opt/drb/node-26)
--build build images locally instead of pulling (~1h for op25)
--no-start configure and enroll, but do not start containers
--wait-approval SEC block up to SEC seconds polling for admin approval
-y, --yes never prompt; fail instead of asking
Every option also has an env var: DRB_NODE_ID, DRB_C2_URL, DRB_ENROLLMENT_TOKEN,
DRB_MQTT_BROKER, DRB_REF, DRB_INSTALL_DIR, DRB_REGISTRY_USER/PASS, ...
Secrets are read from the environment or prompted on the TTY, never from a pipe.
USAGE
}
while [ $# -gt 0 ]; do
case "$1" in
--token) ENROLLMENT_TOKEN="$2"; shift 2 ;;
--node-id) NODE_ID="$2"; shift 2 ;;
--name) NODE_NAME="$2"; shift 2 ;;
--lat) NODE_LAT="$2"; shift 2 ;;
--lon) NODE_LON="$2"; shift 2 ;;
--c2-url) C2_URL="$2"; shift 2 ;;
--mqtt-broker) MQTT_BROKER="$2"; shift 2 ;;
--mqtt-port) MQTT_PORT="$2"; shift 2 ;;
--no-tls) MQTT_TLS=false; [ "$MQTT_PORT" = 8883 ] && MQTT_PORT=1883; shift ;;
--dashboard-pass) DASHBOARD_PASS="$2"; shift 2 ;;
--ref) REF="$2"; shift 2 ;;
--track-main) REF="main"; shift ;;
--dir) INSTALL_DIR="$2"; shift 2 ;;
--build) DO_BUILD=1; shift ;;
--no-start) DO_START=0; shift ;;
--wait-approval) ENROLL_WAIT="$2"; shift 2 ;;
-y|--yes) ASSUME_YES=1; shift ;;
-h|--help) usage; exit 0 ;;
*) die "unknown option: $1 (try --help)" ;;
esac
done
# Prompts must come from the terminal, not from the `curl |` pipe on stdin.
ask() { # ask VAR "prompt" "default"
local __v="$1" __p="$2" __d="${3:-}" __r=""
if [ "$ASSUME_YES" = 1 ] || [ ! -r /dev/tty ]; then
[ -n "$__d" ] || die "$__p is required (non-interactive: pass the flag or env var)"
printf -v "$__v" '%s' "$__d"; return
fi
read -rp "$__p${__d:+ [$__d]}: " __r </dev/tty
printf -v "$__v" '%s' "${__r:-$__d}"
}
ask_secret() {
local __v="$1" __p="$2" __r=""
if [ "$ASSUME_YES" = 1 ] || [ ! -r /dev/tty ]; then printf -v "$__v" '%s' ''; return; fi
read -rsp "$__p: " __r </dev/tty; echo >/dev/tty
printf -v "$__v" '%s' "$__r"
}
genpw() { head -c 24 /dev/urandom | base64 | tr -d '\n=' ; }
# ── 1. Preflight ────────────────────────────────────────────────────────────
say "Preflight"
[ "$(id -u)" -eq 0 ] || die "run as root: curl -fsSL <url> | sudo bash -s -- ..."
command -v apt-get >/dev/null || die "apt-get not found — this script targets Raspberry Pi OS / Debian"
ARCH="$(dpkg --print-architecture)"
case "$ARCH" in
arm64|aarch64) ok "arch $ARCH" ;;
*) warn "arch is $ARCH — CI only builds linux/arm64 images. Prebuilt pull will fail; use --build." ;;
esac
ok "running as root"
# Not fatal: the dongle can be plugged in after install.
if command -v lsusb >/dev/null 2>&1 && lsusb | grep -qiE 'rtl2838|realtek.*283[28]|sdr'; then
ok "SDR dongle detected on USB"
else
warn "no RTL-SDR dongle detected on USB — plug one in before expecting audio"
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 ─────────────────────────────────────────────────────────
say "Installing dependencies"
export DEBIAN_FRONTEND=noninteractive
apt-get update -qq
apt-get install -y -qq git curl ca-certificates jq usbutils openssl >/dev/null
ok "git curl jq openssl"
if ! command -v docker >/dev/null 2>&1; then
say "Installing Docker (get.docker.com)"
curl -fsSL https://get.docker.com | sh
fi
docker --version >/dev/null || die "docker install failed"
ok "$(docker --version)"
if ! docker compose version >/dev/null 2>&1; then
apt-get install -y -qq docker-compose-plugin >/dev/null
fi
docker compose version >/dev/null 2>&1 || die "docker compose plugin missing"
ok "$(docker compose version --short 2>/dev/null || echo 'compose plugin')"
systemctl enable --now docker >/dev/null 2>&1 || true
# The docker group only matters for the human who logs in LATER — this script
# is already root, so nothing below needs a re-login. setup.sh's bug (add the
# group, then immediately run compose in the same unprivileged shell) does not
# apply here.
TARGET_USER="${SUDO_USER:-}"
if [ -n "$TARGET_USER" ] && [ "$TARGET_USER" != root ]; then
usermod -aG docker "$TARGET_USER" || true
ok "added '$TARGET_USER' to the docker group (takes effect at their next login)"
fi
# ── 3. Fetch the repo at a pinned ref ───────────────────────────────────────
say "Fetching node-26 @ ${REF}"
mkdir -p "$(dirname "$INSTALL_DIR")"
if [ -d "$INSTALL_DIR/.git" ]; then
ok "existing install at $INSTALL_DIR — updating in place (.env is preserved)"
git -C "$INSTALL_DIR" remote set-url origin "$REPO_URL"
git -C "$INSTALL_DIR" fetch --tags --prune origin
else
git clone --no-checkout "$REPO_URL" "$INSTALL_DIR"
fi
git -C "$INSTALL_DIR" -c advice.detachedHead=false checkout --force "$REF"
RESOLVED_SHA="$(git -C "$INSTALL_DIR" rev-parse HEAD)"
ok "checked out $REF ($RESOLVED_SHA)"
[ "$REF" = main ] && warn "tracking main — two nodes installed on different days run different software"
cd "$INSTALL_DIR"
mkdir -p configs recordings
chmod 700 configs
# ── 4. Configuration ────────────────────────────────────────────────────────
say "Configuring"
if [ -f .env ]; then
ok ".env already present — keeping it (delete it to reconfigure)"
# Re-read the three values section 5 needs. Not `source .env` — that would
# execute whatever is in the file.
envget() { grep -E "^$1=" .env | head -1 | cut -d= -f2- | tr -d '"'; }
NODE_ID="$(envget NODE_ID)"
C2_URL="${C2_URL:-$(envget C2_URL)}"; C2_URL="${C2_URL%/}"
NODE_NAME="${NODE_NAME:-$(envget NODE_NAME)}"
NODE_LAT="${NODE_LAT:-$(envget NODE_LAT)}"
NODE_LON="${NODE_LON:-$(envget NODE_LON)}"
DASHBOARD_USER="$(envget DASHBOARD_USERNAME)"
[ -n "$NODE_ID" ] || die ".env exists but has no NODE_ID — fix or delete it"
else
[ -n "$NODE_ID" ] || ask NODE_ID "Node ID (e.g. node-003)"
[[ "$NODE_ID" =~ ^[A-Za-z0-9_-]+$ ]] || die "NODE_ID must be letters/numbers/dash/underscore only"
[ -n "$NODE_NAME" ] || ask NODE_NAME "Display name" "$NODE_ID"
[ -n "$NODE_LAT" ] || ask NODE_LAT "Latitude" "0.0"
[ -n "$NODE_LON" ] || ask NODE_LON "Longitude" "0.0"
[ -n "$C2_URL" ] || ask C2_URL "C2 API base URL (https://api.<domain>)"
C2_URL="${C2_URL%/}"
[ -n "$MQTT_BROKER" ] || ask MQTT_BROKER "MQTT broker host (mqtt.<domain>)"
if [ -z "$DASHBOARD_PASS" ]; then
ask_secret DASHBOARD_PASS "Local dashboard password (Enter to generate)"
[ -n "$DASHBOARD_PASS" ] || { DASHBOARD_PASS="$(genpw)"; GENERATED_DASH=1; }
fi
ICE_SRC="$(genpw)"; ICE_ADM="$(genpw)"
umask 077
cat > .env <<EOF
# Written by install.sh on $(date -Is) from ref ${RESOLVED_SHA}
NODE_ID=${NODE_ID}
NODE_NAME="${NODE_NAME}"
NODE_LAT=${NODE_LAT}
NODE_LON=${NODE_LON}
# MQTT — post-cutover auth. There is NO shared node login: the node
# authenticates as username=NODE_ID, password=<its C2-issued api_key>, which
# section 5 below fetches into configs/credentials.json. Deliberately no
# MQTT_USER/MQTT_PASS here; a dynsec broker rejects them.
MQTT_BROKER=${MQTT_BROKER}
MQTT_PORT=${MQTT_PORT}
MQTT_TLS=${MQTT_TLS}
C2_URL=${C2_URL}
ICECAST_SOURCE_PASSWORD=${ICE_SRC}
ICECAST_ADMIN_PASSWORD=${ICE_ADM}
ICECAST_HOST=localhost
ICECAST_PORT=8000
ICECAST_MOUNT=/radio
DASHBOARD_USERNAME=${DASHBOARD_USER}
DASHBOARD_PASSWORD=${DASHBOARD_PASS}
PULSE_SOURCE=drb_sink.monitor
OP25_API_URL=http://localhost:8001
OP25_TERMINAL_URL=http://localhost:8081
OP25_DEBUG_EXPOSE=false
IMAGE_REGISTRY=${REGISTRY}
DOCKER_ORG=${DOCKER_ORG}
DOCKER_REPO=${DOCKER_REPO}
EOF
umask 022
chmod 600 .env
[ -n "$TARGET_USER" ] && chown "$TARGET_USER" .env 2>/dev/null || true
ok ".env written for '$NODE_ID'"
fi
# ── 5. Enrollment ───────────────────────────────────────────────────────────
# Client half of Server/drb-c2-core/app/routers/enrollment.py. It does NOT
# exist in the edge-node app today (mqtt_manager.py:73-85 says so explicitly),
# so without this section a fresh node can never obtain an api_key against a
# dynsec broker: MQTT needs the key, and the legacy key-over-MQTT delivery
# needs MQTT. Doing it here breaks that loop.
#
# Two server endpoints, and the exact response shapes verified against
# enrollment.py @ v1:
#
# POST /nodes/enroll (X-Enrollment-Token)
# 200 -> {node_id, pickup_secret, approval_status}
# 403 -> node_id is ALREADY APPROVED. The CRITICAL GUARD in enrollment.py
# refuses to mint a fresh pickup_secret off the shared fleet token.
# So we must only ever POST this for a node we have not enrolled
# from this machine before — i.e. when configs/pickup_secret is
# absent. Re-running the installer must NOT re-POST here.
# 401 bad/revoked token · 400 missing node_id · 429 rate limited
#
# GET /nodes/{id}/credentials (X-Pickup-Secret)
# Always HTTP 200 with {approval_status, api_key} unless the secret or
# node is bad. api_key is null until an admin approves the node in the UI
# (nodes.py approve_node() mints node_keys/{id}.api_key synchronously in
# the same call — approve is enough; assigning a system is independent and
# NOT required for a key). This endpoint has NO already-approved guard, so
# it is the correct — and only working — re-run path after approval.
# 401 -> missing/invalid/rotated pickup secret
# 404 -> node unknown to C2 (deleted server-side, or never enrolled)
CREDS="$INSTALL_DIR/configs/credentials.json"
PICKUP_FILE="$INSTALL_DIR/configs/pickup_secret"
# GET /nodes/{id}/credentials. Sets CRED_HTTP + CRED_BODY (no -f: we need the
# body and status on a 4xx). One implementation so first-run and re-run agree.
creds_pickup() { # creds_pickup PICKUP_SECRET
local _tmp; _tmp="$(mktemp)"
CRED_HTTP="$(curl -sS -o "$_tmp" -w '%{http_code}' \
"$C2_URL/nodes/$NODE_ID/credentials" -H "X-Pickup-Secret: $1" 2>/dev/null || echo 000)"
CRED_BODY="$(cat "$_tmp" 2>/dev/null || true)"
rm -f "$_tmp"
}
cred_field() { printf '%s' "${CRED_BODY:-}" | jq -r "$1 // empty" 2>/dev/null || true; }
write_creds() { # write_creds API_KEY
umask 077; jq -n --arg k "$1" '{api_key:$k}' > "$CREDS"; umask 022
ok "api_key received and written to configs/credentials.json"
}
say_pending() { # say_pending APPROVAL_STATUS — not an error: node is enrolled, key not minted yet
warn "not approved yet — C2 reports approval_status=${1:-pending}, no api_key minted."
warn "an admin must, at <app-url>/settings/nodes : Approve '$NODE_ID' (assigning a system is separate)."
warn "then re-run this installer — it reuses configs/pickup_secret — or fetch it directly:"
warn " curl -fsS $C2_URL/nodes/$NODE_ID/credentials -H \"X-Pickup-Secret: \$(cat $PICKUP_FILE)\" | jq -r .api_key"
}
# Poll creds_pickup for up to ENROLL_WAIT seconds while still pending. Result
# left in CRED_HTTP/CRED_BODY. --wait-approval is what would have avoided the
# original prod bug; the default (0) does not block, so the re-run path below
# must stand on its own.
wait_for_key() { # wait_for_key PICKUP_SECRET
[ "${ENROLL_WAIT:-0}" -gt 0 ] || return 0
local _end; _end=$(( $(date +%s) + ENROLL_WAIT ))
say "Waiting up to ${ENROLL_WAIT}s for an admin to approve '$NODE_ID'"
while [ "$(date +%s)" -lt "$_end" ]; do
sleep 10
creds_pickup "$1"
[ "$CRED_HTTP" = 200 ] || return 0
[ -z "$(cred_field '.api_key')" ] || return 0
done
}
do_fresh_enroll() {
say "Enrolling '$NODE_ID' with $C2_URL"
[ -n "$ENROLLMENT_TOKEN" ] || ask_secret ENROLLMENT_TOKEN "Enrollment token"
[ -n "$ENROLLMENT_TOKEN" ] || die "no enrollment token — mint one at Settings -> Nodes, then re-run with --token"
local BODY RESP PICKUP STATUS KEY
BODY="$(jq -nc --arg id "$NODE_ID" --arg n "${NODE_NAME:-$NODE_ID}" \
--argjson lat "${NODE_LAT:-0}" --argjson lon "${NODE_LON:-0}" \
'{node_id:$id,name:$n,lat:$lat,lon:$lon}')"
# Token goes in a header from a shell variable — never on a process command
# line, never echoed.
if ! RESP="$(curl -fsS -X POST "$C2_URL/nodes/enroll" \
-H "Content-Type: application/json" \
-H "X-Enrollment-Token: $ENROLLMENT_TOKEN" \
--data "$BODY" 2>&1)"; then
case "$RESP" in
*403*) die "enroll refused (403): node '$NODE_ID' is already approved on C2, and
this machine has no configs/pickup_secret to collect its key with. A token
alone cannot re-issue an approved node's key (enrollment.py CRITICAL GUARD).
Recover by restoring this node's original configs/pickup_secret and re-running,
or have an admin Reissue key (Settings -> Nodes) and write the key into
$CREDS by hand (see node-26#6)." ;;
*401*) die "enroll refused (401): enrollment token missing/invalid/revoked —
mint a fresh one at Settings -> Nodes and re-run with --token." ;;
*429*) die "enroll refused (429): rate limited. Wait ~1 minute, then re-run." ;;
*) die "enroll failed: $RESP" ;;
esac
fi
PICKUP="$(printf '%s' "$RESP" | jq -r '.pickup_secret')"
STATUS="$(printf '%s' "$RESP" | jq -r '.approval_status')"
[ -n "$PICKUP" ] && [ "$PICKUP" != null ] || die "enroll returned no pickup_secret: $RESP"
umask 077; printf '%s' "$PICKUP" > "$PICKUP_FILE"; umask 022
ok "enrolled — approval_status=$STATUS (pickup secret saved to configs/pickup_secret)"
creds_pickup "$PICKUP"
[ "$CRED_HTTP" = 200 ] || die "post-enroll credential fetch failed (HTTP $CRED_HTTP): ${CRED_BODY:-<no body>}"
KEY="$(cred_field '.api_key')"
if [ -z "$KEY" ]; then
wait_for_key "$PICKUP" || true
KEY="$(cred_field '.api_key')"
fi
if [ -n "$KEY" ]; then
write_creds "$KEY"
else
say_pending "$(cred_field '.approval_status')"
fi
}
if [ -s "$CREDS" ] && jq -e '.api_key // empty' "$CREDS" >/dev/null 2>&1; then
say "Enrollment"
ok "api_key already on disk — skipping enrollment"
elif [ -z "${C2_URL:-}" ]; then
say "Enrollment"
warn "no C2_URL — skipping enrollment"
elif [ -s "$PICKUP_FILE" ]; then
# RE-RUN. This node already enrolled from this machine. Do NOT POST
# /nodes/enroll again — an approved node_id gets 403 there, and the v1
# installer's own "re-run to pick up the key" advice then dead-ends on
# "use Reissue key". The pickup endpoint has no such guard: use it.
say "Enrollment — collecting credentials for '$NODE_ID' (pickup secret from a previous run)"
PICKUP="$(cat "$PICKUP_FILE")"
creds_pickup "$PICKUP"
case "$CRED_HTTP" in
200)
API_KEY="$(cred_field '.api_key')"
if [ -z "$API_KEY" ]; then
wait_for_key "$PICKUP" || true
API_KEY="$(cred_field '.api_key')"
fi
if [ -n "$API_KEY" ]; then
write_creds "$API_KEY"
else
# Still pending. Enrolled and idempotent — next run collects the key.
# Clean exit, fall through to start the stack. NOT a failure.
say_pending "$(cred_field '.approval_status')"
fi
;;
401)
if [ -n "$ENROLLMENT_TOKEN" ]; then
warn "saved pickup secret rejected (401) — likely rotated by a re-enroll elsewhere. Re-enrolling with --token."
rm -f "$PICKUP_FILE"
do_fresh_enroll
else
die "saved pickup secret is stale (401) and no --token was given. Re-run with
--token DRB-… to re-enroll (only works while the node is still pending), or
have an admin Reissue key for an approved node and write $CREDS by hand."
fi
;;
404)
if [ -n "$ENROLLMENT_TOKEN" ]; then
warn "C2 does not know node '$NODE_ID' (404) — deleted server-side or never fully enrolled. Re-enrolling with --token."
rm -f "$PICKUP_FILE"
do_fresh_enroll
else
die "C2 does not know node '$NODE_ID' (404) and no --token was given.
Re-run with --token DRB-… to enroll it again."
fi
;;
000)
die "could not reach $C2_URL/nodes/$NODE_ID/credentials — check --c2-url and connectivity." ;;
*)
die "credential pickup failed (HTTP $CRED_HTTP): ${CRED_BODY:-<no body>}" ;;
esac
else
say "Enrollment"
do_fresh_enroll
fi
# ── 6. Images + start ───────────────────────────────────────────────────────
if [ "$DO_START" = 1 ]; then
if [ -n "$REGISTRY_USER" ] && [ -n "$REGISTRY_PASS" ]; then
printf '%s' "$REGISTRY_PASS" | docker login "$REGISTRY" -u "$REGISTRY_USER" --password-stdin >/dev/null
ok "logged in to $REGISTRY"
fi
if [ "$DO_BUILD" = 1 ]; then
say "Building images locally — op25 takes roughly an hour on a Pi"
docker compose build
docker compose up -d
else
say "Pulling prebuilt images from $REGISTRY/$DOCKER_ORG/$DOCKER_REPO"
if ! docker compose pull; then
die "pull failed. $REGISTRY is public, so this is most likely a login
requirement or a transient network error: re-run with DRB_REGISTRY_USER /
DRB_REGISTRY_PASS set, or with --build to compile on the Pi (~1h for op25)."
fi
docker compose up --no-build -d
fi
ok "stack started"
else
say "Skipping start (--no-start). Run: cd $INSTALL_DIR && make up-prebuilt"
fi
# ── 7. What the operator does next ──────────────────────────────────────────
IP="$(hostname -I 2>/dev/null | awk '{print $1}')"
cat <<EOF
$(printf "${G}Node '%s' installed at %s${N}" "$NODE_ID" "$INSTALL_DIR")
ref ${RESOLVED_SHA}
images $([ "$DO_BUILD" = 1 ] && echo "built locally" || echo "pulled from $REGISTRY")
dashboard http://${IP:-<node-ip>}/ (user: ${DASHBOARD_USER})
logs cd $INSTALL_DIR && docker compose logs -f edge-node
NEXT — an admin must approve this node before it can do anything:
1. Open <app-url>/settings/nodes
2. Approve "$NODE_ID"
3. Assign it a radio system
EOF
if [ "${GENERATED_DASH:-0}" = 1 ]; then
printf "${Y}Generated dashboard password (shown once): %s${N}\n\n" "$DASHBOARD_PASS"
fi
if [ ! -s "$CREDS" ]; then
if [ -s "$PICKUP_FILE" ]; then
printf "${Y}This node has no api_key yet. After an admin approves it, re-run the same\ninstall command — it reuses configs/pickup_secret and will collect the key\n(no --token needed for the re-run).${N}\n\n"
else
printf "${Y}This node has no api_key and no saved pickup secret, so a plain re-run cannot\nfix it. Re-run with --token DRB-… to enroll; or, if the node is already\napproved, have an admin Reissue key and write it into\n%s by hand.${N}\n\n" "$CREDS"
fi
fi
@@ -1,57 +0,0 @@
name: release-tag
on:
push:
branches:
- dev
jobs:
release-image:
runs-on: ubuntu-latest
env:
DOCKER_LATEST: stable
CONTAINER_NAME: drb-client-discord-bot
steps:
- name: Checkout
uses: actions/checkout@v4
- name: Set up QEMU
uses: docker/setup-qemu-action@v3
- name: Set up Docker BuildX
uses: docker/setup-buildx-action@v3
with: # replace it with your local IP
config-inline: |
[registry."git.vpn.cusano.net"]
http = false
insecure = false
- name: Login to DockerHub
uses: docker/login-action@v3
with:
registry: git.vpn.cusano.net # replace it with your local IP
username: ${{ secrets.GIT_REPO_USERNAME }}
password: ${{ secrets.GIT_REPO_PASSWORD }}
- name: Get Meta
id: meta
run: |
echo REPO_NAME=$(echo ${GITHUB_REPOSITORY} | awk -F"/" '{print $2}') >> $GITHUB_OUTPUT
echo REPO_VERSION=$(git describe --tags --always | sed 's/^v//') >> $GITHUB_OUTPUT
- name: Validate build configuration
uses: docker/build-push-action@v6
with:
call: check
- name: Build and push
uses: docker/build-push-action@v6
with:
context: .
file: ./Dockerfile
platforms: |
linux/arm64
push: true
tags: | # replace it with your local IP and tags
git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:${{ steps.meta.outputs.REPO_VERSION }}
git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:${{ env.DOCKER_LATEST }}
@@ -1,60 +0,0 @@
name: release-tag
on:
push:
branches:
- master
jobs:
release-image:
runs-on: ubuntu-latest
permissions:
contents: read
packages: write
env:
DOCKER_LATEST: stable
CONTAINER_NAME: op25-client
steps:
- name: Checkout
uses: actions/checkout@v5
- name: Set up QEMU
uses: docker/setup-qemu-action@v3
- name: Set up Docker BuildX
uses: docker/setup-buildx-action@v3
with:
config-inline: |
[registry."git.vpn.cusano.net"]
http = false
insecure = false
- name: Login to Gitea Container Registry
uses: docker/login-action@v3
with:
registry: git.vpn.cusano.net
username: ${{ gitea.actor }} # Uses the user or bot that triggered the workflow
password: ${{ secrets.GITHUB_COM_TOKEN }} # The built-in, temporary token
- name: Get Meta
id: meta
run: |
echo REPO_NAME=$(echo ${GITHUB_REPOSITORY} | awk -F"/" '{print $2}') >> $GITHUB_OUTPUT
echo REPO_VERSION=$(git describe --tags --always | sed 's/^v//') >> $GITHUB_OUTPUT
- name: Validate build configuration
uses: docker/build-push-action@v6
with:
call: check
- name: Build and push
uses: docker/build-push-action@v6
with:
context: .
file: ./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.REPO_VERSION }}
git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:${{ env.DOCKER_LATEST }}
-30
View File
@@ -1,30 +0,0 @@
name: Lint
on:
push:
branches:
- master
pull_request:
branches:
- "*"
jobs:
lint:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: Set up Python
uses: actions/setup-python@v5
with:
python-version: '3.13'
- name: Install dependencies
run: |
python -m pip install --upgrade pip
pip install flake8
- name: Run Lint
run: |
flake8 --max-line-length=88 --ignore=E203,E302,E501 .
+7 -2
View File
@@ -1,5 +1,10 @@
# OP25 Core Container
FROM python:slim-trixie
# Pinned to a Python major version deliberately. The bare `slim-trixie` tag
# carries no version at all, so a rebuild could move the interpreter across a
# major release -- which this repo has already been bitten by once, when
# app/models.py only ran because trixie happened to ship 3.14 and PEP 649
# defers annotation evaluation. Matches drb-edge-node, which is already 3.14.
FROM python:3.14-slim
# Set environment variables
ENV DEBIAN_FRONTEND=noninteractive
@@ -7,7 +12,7 @@ ENV DEBIAN_FRONTEND=noninteractive
# Install system dependencies
RUN apt-get update && \
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)
COPY system.pa /etc/pulse/system.pa
+1 -1
View File
@@ -24,7 +24,7 @@ async def lifespan(app: FastAPI):
icecast_host=os.getenv("ICECAST_HOST", "localhost"),
icecast_port=int(os.getenv("ICECAST_PORT", "8000")),
icecast_mountpoint=os.getenv("ICECAST_MOUNT", "/radio"),
icecast_password=os.getenv("ICECAST_SOURCE_PASSWORD", "hackme"),
icecast_password=os.getenv("ICECAST_SOURCE_PASSWORD", ""),
)
generate_liquid_script(config)
LOGGER.info("op25.liq generated from environment variables.")
@@ -1,6 +1,7 @@
from fastapi import HTTPException, APIRouter
import subprocess
import os
import re
import signal
import json
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():
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")
async def generate_config(generator: ConfigGenerator):
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,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
+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,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 "$@"
+3
View File
@@ -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}
-148
View File
@@ -1,148 +0,0 @@
#!/usr/bin/env bash
# Interactive first-time setup for a DRB edge node.
# Installs system dependencies (Docker, make, curl) then writes .env
# and optionally builds + starts the stack.
set -e
GREEN='\033[0;32m'; YELLOW='\033[1;33m'; CYAN='\033[0;36m'; RED='\033[0;31m'; NC='\033[0m'
cd "$(dirname "$0")"
echo -e "${CYAN}DRB Edge Node Setup${NC}"
echo "-------------------"
# ── Dependency installation ──────────────────────────────────────────────────
install_deps() {
if ! command -v apt-get &>/dev/null; then
echo -e "${YELLOW}⚠ apt-get not found — skipping auto-install. Ensure docker, make, and curl are installed.${NC}"
return
fi
echo ""
echo -e "${CYAN}Installing system dependencies…${NC}"
sudo apt-get update -qq
local pkgs=()
command -v make &>/dev/null || pkgs+=(make)
command -v curl &>/dev/null || pkgs+=(curl)
command -v git &>/dev/null || pkgs+=(git)
if [ ${#pkgs[@]} -gt 0 ]; then
echo " Installing: ${pkgs[*]}"
sudo apt-get install -y -qq "${pkgs[@]}"
fi
# Docker — use get.docker.com if not present
if ! command -v docker &>/dev/null; then
echo " Installing Docker via get.docker.com…"
curl -fsSL https://get.docker.com | sudo sh
# Allow current user to run docker without sudo
sudo usermod -aG docker "$USER"
echo -e "${YELLOW} ⚠ Docker group added. You may need to log out and back in for it to take effect.${NC}"
echo -e "${YELLOW} If 'docker compose' fails below, run: newgrp docker${NC}"
else
echo -e "${GREEN} ✓ docker$(docker --version | grep -oP ' \d+\.\d+\.\d+' | head -1)${NC}"
fi
# Docker Compose plugin check (comes with Docker Engine ≥ 20.10)
if ! docker compose version &>/dev/null 2>&1; then
echo -e "${RED} docker compose plugin not found. Installing…${NC}"
sudo apt-get install -y -qq docker-compose-plugin
else
echo -e "${GREEN} ✓ docker compose $(docker compose version --short 2>/dev/null || true)${NC}"
fi
echo -e "${GREEN}✓ Dependencies ready${NC}"
}
install_deps
if [ -f .env ]; then
echo -e "${YELLOW}Warning: .env already exists.${NC}"
read -rp "Overwrite? [y/N] " yn
[[ "$yn" =~ ^[Yy]$ ]] || { echo "Aborted."; exit 0; }
fi
# --- Node identity ---
echo ""
echo "Unique node ID — no spaces (e.g. node-ossining, node-002)"
read -rp "NODE_ID: " NODE_ID
while [[ ! "$NODE_ID" =~ ^[a-zA-Z0-9_-]+$ ]]; do
echo " Use letters, numbers, dashes, underscores only."
read -rp "NODE_ID: " NODE_ID
done
echo ""
read -rp "Node display name [$NODE_ID]: " NODE_NAME
NODE_NAME="${NODE_NAME:-$NODE_ID}"
# --- GPS ---
echo ""
echo "GPS coordinates (decimal degrees — used for the map)"
read -rp "Latitude [0.0]: " NODE_LAT; NODE_LAT="${NODE_LAT:-0.0}"
read -rp "Longitude [0.0]: " NODE_LON; NODE_LON="${NODE_LON:-0.0}"
# --- C2 server ---
echo ""
echo "C2 server — hostname or IP of the machine running the server stack"
read -rp "C2 server host: " C2_HOST; C2_HOST="${C2_HOST:-localhost}"
read -rp "C2 API port [8888]: " C2_PORT; C2_PORT="${C2_PORT:-8888}"
# --- MQTT ---
echo ""
echo "MQTT credentials (must match MQTT_NODE_USER/PASS in the server .env)"
read -rp "MQTT port [1883]: " MQTT_PORT; MQTT_PORT="${MQTT_PORT:-1883}"
read -rp "MQTT username [drb-node]: " MQTT_USER; MQTT_USER="${MQTT_USER:-drb-node}"
read -rsp "MQTT password: " MQTT_PASS; echo ""; MQTT_PASS="${MQTT_PASS:-change-me-node}"
# --- Icecast ---
echo ""
echo "Icecast passwords (local container)"
read -rsp "Source password [hackme]: " ICECAST_SOURCE; echo ""; ICECAST_SOURCE="${ICECAST_SOURCE:-hackme}"
read -rsp "Admin password [admin]: " ICECAST_ADMIN; echo ""; ICECAST_ADMIN="${ICECAST_ADMIN:-admin}"
# --- Write .env ---
cat > .env <<EOF
# Node Identity
NODE_ID=${NODE_ID}
NODE_NAME="${NODE_NAME}"
NODE_LAT=${NODE_LAT}
NODE_LON=${NODE_LON}
# MQTT — point to your C2 server
MQTT_BROKER=${C2_HOST}
MQTT_PORT=${MQTT_PORT}
MQTT_USER=${MQTT_USER}
MQTT_PASS=${MQTT_PASS}
# C2 server for audio upload
C2_URL=http://${C2_HOST}:${C2_PORT}
# API key is provisioned automatically via MQTT after admin approves the node
# Icecast (local container — usually no need to change)
ICECAST_SOURCE_PASSWORD=${ICECAST_SOURCE}
ICECAST_ADMIN_PASSWORD=${ICECAST_ADMIN}
ICECAST_HOST=localhost
ICECAST_PORT=8000
ICECAST_MOUNT=/radio
# OP25 container (usually no need to change)
OP25_API_URL=http://localhost:8001
OP25_TERMINAL_URL=http://localhost:8081
EOF
echo ""
echo -e "${GREEN}✓ .env written for node '${NODE_ID}'${NC}"
echo ""
read -rp "Build and start now? [Y/n] " start
if [[ ! "$start" =~ ^[Nn]$ ]]; then
echo ""
echo "Building images (op25 takes ~10 min on first run)…"
docker compose build
docker compose up -d
echo ""
echo -e "${GREEN}✓ Node '${NODE_ID}' started.${NC}"
echo " → Check the dashboard — it will appear as pending approval."
else
echo "Run 'make up' when ready."
fi