Compare commits
21
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
9b6c64fbbf | ||
|
|
1c74edffea | ||
|
|
c0a05a8ad6 | ||
|
|
52d11c1546 | ||
|
|
b54624e176 | ||
|
|
52f31bbcc0 | ||
|
|
b2e3804dc3 | ||
|
|
8f5fca8757 | ||
|
|
a53000a092 | ||
|
|
bee9e173e6 | ||
|
|
31e1176c45 | ||
|
|
4f942bd770 | ||
|
|
94c3e2a952 | ||
|
|
5b0048e8db | ||
|
|
0c08275482 | ||
|
|
fb13bb8ae3 | ||
|
|
e6aab7589a | ||
|
|
28266b4441 | ||
|
|
7f1d09c753 | ||
|
|
9d86304b8a | ||
|
|
5ff1089551 |
+11
-2
@@ -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
|
||||
|
||||
@@ -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
|
||||
@@ -0,0 +1,52 @@
|
||||
name: Build secondary-sdr
|
||||
|
||||
on:
|
||||
workflow_dispatch:
|
||||
push:
|
||||
branches: [main, master]
|
||||
paths:
|
||||
- "secondary-sdr-container/**"
|
||||
|
||||
jobs:
|
||||
build:
|
||||
runs-on: ubuntu-latest
|
||||
permissions:
|
||||
contents: read
|
||||
packages: write
|
||||
env:
|
||||
CONTAINER_NAME: secondary-sdr
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
|
||||
- uses: docker/setup-qemu-action@v3
|
||||
|
||||
- uses: docker/setup-buildx-action@v3
|
||||
with:
|
||||
config-inline: |
|
||||
[registry."git.vpn.cusano.net"]
|
||||
http = false
|
||||
insecure = false
|
||||
|
||||
- uses: docker/login-action@v3
|
||||
with:
|
||||
registry: git.vpn.cusano.net
|
||||
username: ${{ gitea.actor }}
|
||||
password: ${{ secrets.BUILD_TOKEN }}
|
||||
|
||||
- name: Get version
|
||||
id: meta
|
||||
run: |
|
||||
echo "REPO_NAME=$(echo ${GITHUB_REPOSITORY} | awk -F'/' '{print $2}')" >> $GITHUB_OUTPUT
|
||||
echo "VERSION=$(git describe --tags --always | sed 's/^v//')" >> $GITHUB_OUTPUT
|
||||
|
||||
- uses: docker/build-push-action@v6
|
||||
with:
|
||||
context: ./secondary-sdr-container
|
||||
file: ./secondary-sdr-container/Dockerfile
|
||||
platforms: linux/arm64
|
||||
push: true
|
||||
tags: |
|
||||
git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:${{ steps.meta.outputs.VERSION }}
|
||||
git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:latest
|
||||
cache-from: type=registry,ref=git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:buildcache
|
||||
cache-to: type=registry,ref=git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ env.CONTAINER_NAME }}:buildcache,mode=max
|
||||
@@ -0,0 +1,103 @@
|
||||
name: Publish public images
|
||||
|
||||
# Push the three node images to a PUBLIC registry on a version tag, so a fresh
|
||||
# Pi can `docker pull` them without a Gitea login (git.vpn.cusano.net is now
|
||||
# behind REQUIRE_SIGNIN_VIEW — see INCIDENT-2026-09-06). Source + the private
|
||||
# registry stay walled; only the built node images go public.
|
||||
#
|
||||
# ── DISABLED ────────────────────────────────────────────────────────────────
|
||||
# The job is gated on `vars.NODE_PUBLIC_PUBLISH == 'true'`. Until that repo
|
||||
# variable is set the workflow triggers on tags but the job is skipped, so
|
||||
# this file is wired and inert. We're still building the core; flip it on when
|
||||
# self-serve node install is actually needed.
|
||||
#
|
||||
# To enable:
|
||||
# 1. Repo → Settings → Actions → Variables:
|
||||
# NODE_PUBLIC_PUBLISH = true
|
||||
# PUBLIC_REGISTRY = ghcr.io (or docker.io)
|
||||
# PUBLIC_NAMESPACE = <org-or-user> (images land at <ns>/drb-<name>)
|
||||
# 2. Repo → Settings → Actions → Secrets:
|
||||
# PUBLIC_REGISTRY_USER = <push user>
|
||||
# PUBLIC_REGISTRY_TOKEN = <push token / PAT with write:packages>
|
||||
# 3. Re-push a tag (or run this workflow via workflow_dispatch).
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
on:
|
||||
workflow_dispatch:
|
||||
push:
|
||||
tags:
|
||||
- "v*"
|
||||
|
||||
concurrency:
|
||||
group: publish-public-${{ github.ref }}
|
||||
cancel-in-progress: false
|
||||
|
||||
jobs:
|
||||
publish:
|
||||
# Inert until the repo variable is set. Do NOT convert this to `if: false`
|
||||
# — the variable is the switch, no code change needed to go live.
|
||||
if: ${{ vars.NODE_PUBLIC_PUBLISH == 'true' }}
|
||||
runs-on: ubuntu-latest
|
||||
permissions:
|
||||
contents: read
|
||||
packages: write
|
||||
strategy:
|
||||
fail-fast: false
|
||||
matrix:
|
||||
include:
|
||||
- name: edge-node
|
||||
context: ./drb-edge-node
|
||||
file: ./drb-edge-node/Dockerfile
|
||||
cache_name: edge-node
|
||||
- name: icecast
|
||||
context: ./icecast
|
||||
file: ./icecast/Dockerfile
|
||||
cache_name: icecast
|
||||
- name: op25-client
|
||||
context: ./op25-container
|
||||
file: ./op25-container/Dockerfile
|
||||
cache_name: op25-client
|
||||
steps:
|
||||
- uses: actions/checkout@v4
|
||||
with:
|
||||
fetch-depth: 0 # need tags for `git describe`
|
||||
|
||||
- uses: docker/setup-qemu-action@v3
|
||||
- uses: docker/setup-buildx-action@v3
|
||||
with:
|
||||
config-inline: |
|
||||
[registry."git.vpn.cusano.net"]
|
||||
http = false
|
||||
insecure = false
|
||||
|
||||
# Private Gitea registry — read only, to reuse the existing build cache
|
||||
# (keeps the op25 image off a ~1h from-scratch compile).
|
||||
- uses: docker/login-action@v3
|
||||
with:
|
||||
registry: git.vpn.cusano.net
|
||||
username: ${{ gitea.actor }}
|
||||
password: ${{ secrets.BUILD_TOKEN }}
|
||||
|
||||
# Public registry — where the images are pushed.
|
||||
- uses: docker/login-action@v3
|
||||
with:
|
||||
registry: ${{ vars.PUBLIC_REGISTRY }}
|
||||
username: ${{ secrets.PUBLIC_REGISTRY_USER }}
|
||||
password: ${{ secrets.PUBLIC_REGISTRY_TOKEN }}
|
||||
|
||||
- name: Version
|
||||
id: meta
|
||||
run: |
|
||||
echo "REPO_NAME=$(echo ${GITHUB_REPOSITORY} | awk -F'/' '{print $2}')" >> $GITHUB_OUTPUT
|
||||
echo "VERSION=$(git describe --tags --always | sed 's/^v//')" >> $GITHUB_OUTPUT
|
||||
|
||||
- uses: docker/build-push-action@v6
|
||||
with:
|
||||
context: ${{ matrix.context }}
|
||||
file: ${{ matrix.file }}
|
||||
platforms: linux/arm64
|
||||
push: true
|
||||
tags: |
|
||||
${{ vars.PUBLIC_REGISTRY }}/${{ vars.PUBLIC_NAMESPACE }}/drb-${{ matrix.name }}:${{ steps.meta.outputs.VERSION }}
|
||||
${{ vars.PUBLIC_REGISTRY }}/${{ vars.PUBLIC_NAMESPACE }}/drb-${{ matrix.name }}:latest
|
||||
cache-from: type=registry,ref=git.vpn.cusano.net/${{ vars.DOCKER_ORG }}/${{ steps.meta.outputs.REPO_NAME }}/${{ matrix.cache_name }}:buildcache
|
||||
@@ -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).
|
||||
|
||||
@@ -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
@@ -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
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
|
||||
@@ -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,19 @@ 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,
|
||||
}
|
||||
secondary = await secondary_sdr_client.status()
|
||||
if secondary is not None:
|
||||
payload["secondary_sdr_running"] = [r["mode"] for r in secondary.get("running", [])]
|
||||
# Best-effort hardware report — omit rather than guess. op25's :stable
|
||||
# image has no lsusb, so fall back to the secondary-sdr container's count.
|
||||
count = devices.get("count") if devices else None
|
||||
if count is None and secondary is not None:
|
||||
count = secondary.get("sdr_count")
|
||||
if count is not None:
|
||||
payload["sdr_count"] = count
|
||||
self._publish(self._t_checkin, payload, qos=1)
|
||||
|
||||
def _publish(self, topic: str, payload: dict, qos: int = 0, retain: bool = False):
|
||||
|
||||
@@ -65,6 +65,16 @@ 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:
|
||||
|
||||
@@ -0,0 +1,29 @@
|
||||
import asyncio
|
||||
from typing import Iterable, List, Optional
|
||||
|
||||
from app.internal.config_manager import load_node_config, save_node_config
|
||||
from app.internal.logger import logger
|
||||
from app.internal.secondary_sdr_client import secondary_sdr_client
|
||||
from app.models import normalize_secondary_priority
|
||||
|
||||
|
||||
async def set_secondary_priority(priority: Iterable[str]) -> Optional[List[str]]:
|
||||
"""Persist and apply the node's secondary SDR priority — one path for the
|
||||
local dashboard, a C2 command and a config push alike. Never touches op25:
|
||||
changing what the spare dongles do must not interrupt P25 recording.
|
||||
|
||||
Returns the decoders now running (None if the container is unreachable).
|
||||
"""
|
||||
ordered = normalize_secondary_priority(priority)
|
||||
cfg = load_node_config()
|
||||
cfg.secondary_sdr_priority = ordered
|
||||
cfg.secondary_sdr_mode = ordered[0] if ordered else "none"
|
||||
save_node_config(cfg)
|
||||
|
||||
running = await secondary_sdr_client.apply(ordered)
|
||||
logger.info(f"Secondary SDR priority set to {ordered!r}; running {running!r}")
|
||||
|
||||
# Report straight away so C2's view doesn't wait for the next heartbeat.
|
||||
from app.internal.mqtt_manager import mqtt_manager
|
||||
asyncio.create_task(mqtt_manager.publish_checkin())
|
||||
return running
|
||||
@@ -0,0 +1,71 @@
|
||||
import httpx
|
||||
from typing import Any, Dict, List, Optional
|
||||
from app.config import settings
|
||||
from app.internal.logger import logger
|
||||
|
||||
|
||||
class SecondarySdrClient:
|
||||
"""Talks to the secondary-sdr-container (node-26#9) over its control API.
|
||||
|
||||
Mirrors op25_client.py's shape on purpose — same failure handling (log
|
||||
and return None/False rather than raise), since this container is
|
||||
optional and its absence must never break the primary op25 radio path.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
self.api_url = settings.secondary_sdr_api_url
|
||||
|
||||
async def apply(self, priority: List[str]) -> Optional[List[str]]:
|
||||
"""Run decoders down the priority list until SDRs run out. Returns the
|
||||
modes actually running, or None if the container is unreachable."""
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=30) as client:
|
||||
r = await client.post(f"{self.api_url}/secondary/apply", json={"priority": priority})
|
||||
r.raise_for_status()
|
||||
return r.json().get("running", [])
|
||||
except Exception as e:
|
||||
logger.error(f"Secondary SDR apply (priority={priority!r}) failed: {e}")
|
||||
return None
|
||||
|
||||
async def start(self, mode: str) -> bool:
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10) as client:
|
||||
r = await client.post(f"{self.api_url}/secondary/start", json={"mode": mode})
|
||||
r.raise_for_status()
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"Secondary SDR start (mode={mode!r}) failed: {e}")
|
||||
return False
|
||||
|
||||
async def stop(self) -> bool:
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10) as client:
|
||||
r = await client.post(f"{self.api_url}/secondary/stop")
|
||||
r.raise_for_status()
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"Secondary SDR stop failed: {e}")
|
||||
return False
|
||||
|
||||
async def status(self) -> Optional[Dict[str, Any]]:
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=5) as client:
|
||||
r = await client.get(f"{self.api_url}/secondary/status")
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
except Exception as e:
|
||||
logger.error(f"Secondary SDR status failed: {e}")
|
||||
return None
|
||||
|
||||
async def data(self) -> Optional[Dict[str, Any]]:
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=5) as client:
|
||||
r = await client.get(f"{self.api_url}/secondary/data")
|
||||
r.raise_for_status()
|
||||
return r.json()
|
||||
except Exception as e:
|
||||
logger.error(f"Secondary SDR data fetch failed: {e}")
|
||||
return None
|
||||
|
||||
|
||||
secondary_sdr_client = SecondarySdrClient()
|
||||
@@ -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:
|
||||
@@ -23,7 +24,7 @@ async def fetch_and_cache_systems() -> bool:
|
||||
r = await client.get(url, headers=headers)
|
||||
r.raise_for_status()
|
||||
systems = r.json()
|
||||
|
||||
|
||||
_CACHE_FILE.parent.mkdir(parents=True, exist_ok=True)
|
||||
_CACHE_FILE.write_text(json.dumps(systems, indent=2))
|
||||
logger.info(f"Cached {len(systems)} systems from C2.")
|
||||
@@ -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"]})
|
||||
@@ -158,6 +158,9 @@ async def on_command(payload: dict):
|
||||
)
|
||||
elif action == "discord_leave":
|
||||
await radio_bot.leave()
|
||||
elif action == "set_secondary_priority":
|
||||
from app.internal.secondary_priority import set_secondary_priority
|
||||
await set_secondary_priority(payload.get("priority") or [])
|
||||
elif action == "op25_restart":
|
||||
from app.internal.op25_client import op25_client
|
||||
await op25_client.stop()
|
||||
@@ -224,6 +227,8 @@ async def on_config_push(payload: dict):
|
||||
hardware_preset = payload.pop("hardware_preset", None)
|
||||
ppm_override = payload.pop("ppm_override", None)
|
||||
node_type = payload.pop("node_type", None)
|
||||
secondary_sdr_priority = payload.pop("secondary_sdr_priority", None)
|
||||
secondary_sdr_mode = payload.pop("secondary_sdr_mode", None) # legacy single-mode C2
|
||||
enforce_override_timeout = payload.pop("enforce_override_timeout", None)
|
||||
try:
|
||||
config = SystemConfig(**payload)
|
||||
@@ -257,6 +262,12 @@ async def on_config_push(payload: dict):
|
||||
await op25_client.start()
|
||||
logger.info(f"Config push applied: {config.name}")
|
||||
|
||||
if secondary_sdr_priority is None and secondary_sdr_mode is not None:
|
||||
secondary_sdr_priority = [secondary_sdr_mode]
|
||||
if secondary_sdr_priority is not None:
|
||||
from app.internal.secondary_priority import set_secondary_priority
|
||||
await set_secondary_priority(secondary_sdr_priority)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# App lifecycle
|
||||
@@ -297,7 +308,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 +323,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.secondary_sdr_client import secondary_sdr_client
|
||||
logger.info(f"Resuming secondary SDRs (priority={node_cfg.secondary_sdr_priority!r}) after restart.")
|
||||
await secondary_sdr_client.apply(node_cfg.secondary_sdr_priority)
|
||||
|
||||
heartbeat_task = asyncio.create_task(mqtt_manager.heartbeat_loop())
|
||||
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()
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
from pydantic import BaseModel
|
||||
from typing import Optional, Dict, Any
|
||||
from pydantic import BaseModel, model_validator
|
||||
from typing import Optional, Dict, Any, Iterable, List
|
||||
from enum import Enum
|
||||
from datetime import datetime
|
||||
|
||||
@@ -23,6 +23,21 @@ class SystemConfig(BaseModel):
|
||||
config: Dict[str, Any] # OP25-compatible config blob passed through to op25-container
|
||||
|
||||
|
||||
# Decoders a node can run on the SDRs beyond op25's (node-26#9). op25 always
|
||||
# keeps its own dongle; each further dongle runs the next entry of the node's
|
||||
# secondary_sdr_priority, so a 3-SDR node with ["adsb", "ais"] runs both.
|
||||
SECONDARY_SDR_MODES = ("adsb", "ais")
|
||||
|
||||
|
||||
def normalize_secondary_priority(items: Iterable[str]) -> List[str]:
|
||||
"""Known modes only, first occurrence wins, order preserved."""
|
||||
out: List[str] = []
|
||||
for m in items or []:
|
||||
if m in SECONDARY_SDR_MODES and m not in out:
|
||||
out.append(m)
|
||||
return out
|
||||
|
||||
|
||||
class NodeConfig(BaseModel):
|
||||
node_id: str
|
||||
node_name: str
|
||||
@@ -34,11 +49,22 @@ class NodeConfig(BaseModel):
|
||||
hardware_preset: str = "rtl-sdr-v3"
|
||||
ppm_override: Optional[float] = None
|
||||
node_type: str = "fixed" # fixed or portable
|
||||
secondary_sdr_priority: List[str] = [] # ordered; SDRs beyond op25's run these top-down
|
||||
# Legacy single-mode field (pre-priority). Still written as priority[0] so
|
||||
# an older C2 reading checkins sees something sensible; read only to
|
||||
# migrate a node_config.json saved before priority existed.
|
||||
secondary_sdr_mode: str = "none"
|
||||
enforce_override_timeout: bool = True
|
||||
override_system_id: Optional[str] = None
|
||||
override_config: Optional[SystemConfig] = None
|
||||
offline_call_buffer_size: int = 35 # max call_end events to buffer while MQTT is offline
|
||||
|
||||
@model_validator(mode="after")
|
||||
def _migrate_secondary_mode(self) -> "NodeConfig":
|
||||
if not self.secondary_sdr_priority and self.secondary_sdr_mode in SECONDARY_SDR_MODES:
|
||||
self.secondary_sdr_priority = [self.secondary_sdr_mode]
|
||||
return self
|
||||
|
||||
|
||||
class CallEvent(BaseModel):
|
||||
call_id: str
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
from fastapi import APIRouter, Depends, HTTPException, Body
|
||||
from typing import Optional
|
||||
from typing import List, Optional
|
||||
import asyncio
|
||||
import httpx
|
||||
from app.config import settings
|
||||
from app.models import SystemConfig
|
||||
from app.models import SystemConfig, SECONDARY_SDR_MODES
|
||||
from app.internal.op25_client import op25_client
|
||||
from app.internal.config_manager import load_node_config, save_node_config, apply_system_config
|
||||
from app.internal.call_recorder import call_recorder
|
||||
@@ -25,12 +25,16 @@ router = APIRouter(prefix="/api", tags=["api"], dependencies=[Depends(auth.requi
|
||||
async def get_status():
|
||||
node_cfg = load_node_config()
|
||||
op25_status = await op25_client.status()
|
||||
|
||||
|
||||
active_tgid = metadata_watcher.current_tgid
|
||||
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:
|
||||
@@ -126,7 +130,7 @@ async def set_override(
|
||||
):
|
||||
node_cfg = load_node_config()
|
||||
config = None
|
||||
|
||||
|
||||
if system_id:
|
||||
from app.internal.system_cacher import get_cached_system
|
||||
cached = get_cached_system(system_id)
|
||||
@@ -139,19 +143,19 @@ async def set_override(
|
||||
config = SystemConfig(**system_config)
|
||||
else:
|
||||
raise HTTPException(400, "Must specify system_id or system_config.")
|
||||
|
||||
|
||||
node_cfg.override_system_id = config.system_id
|
||||
node_cfg.override_config = config
|
||||
save_node_config(node_cfg)
|
||||
|
||||
|
||||
from app.main import _generate_op25_config
|
||||
if not await _generate_op25_config(config):
|
||||
raise HTTPException(500, f"Failed to generate OP25 config for override: {config.name}")
|
||||
|
||||
|
||||
await op25_client.stop()
|
||||
await asyncio.sleep(2)
|
||||
await op25_client.start()
|
||||
|
||||
|
||||
await mqtt_manager._publish_checkin()
|
||||
return {"ok": True}
|
||||
|
||||
@@ -161,20 +165,20 @@ async def revert_config():
|
||||
node_cfg = load_node_config()
|
||||
if not node_cfg.override_system_id:
|
||||
return {"ok": True, "message": "No override active."}
|
||||
|
||||
|
||||
node_cfg.override_system_id = None
|
||||
node_cfg.override_config = None
|
||||
save_node_config(node_cfg)
|
||||
|
||||
|
||||
if node_cfg.system_config:
|
||||
from app.main import _generate_op25_config
|
||||
if not await _generate_op25_config(node_cfg.system_config):
|
||||
raise HTTPException(500, "Failed to regenerate original OP25 config.")
|
||||
|
||||
|
||||
await op25_client.stop()
|
||||
await asyncio.sleep(2)
|
||||
await op25_client.start()
|
||||
|
||||
|
||||
await mqtt_manager._publish_checkin()
|
||||
return {"ok": True}
|
||||
|
||||
@@ -183,10 +187,10 @@ async def revert_config():
|
||||
async def ack_override(timeout_minutes: int = Body(1440)):
|
||||
if not settings.c2_url:
|
||||
raise HTTPException(400, "C2_URL not configured.")
|
||||
|
||||
|
||||
api_key = credentials.get_api_key()
|
||||
headers = {"Authorization": f"Bearer {api_key}"} if api_key else {}
|
||||
|
||||
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=10) as client:
|
||||
r = await client.post(
|
||||
@@ -200,6 +204,30 @@ async def ack_override(timeout_minutes: int = Body(1440)):
|
||||
raise HTTPException(500, f"Failed to contact C2: {e}")
|
||||
|
||||
|
||||
@router.get("/secondary")
|
||||
async def get_secondary():
|
||||
"""Secondary SDR priority plus what's actually running (node-26#9)."""
|
||||
from app.internal.secondary_sdr_client import secondary_sdr_client
|
||||
status = await secondary_sdr_client.status()
|
||||
return {
|
||||
"priority": load_node_config().secondary_sdr_priority,
|
||||
"modes": list(SECONDARY_SDR_MODES),
|
||||
"running": [r["mode"] for r in status.get("running", [])] if status else None,
|
||||
"sdr_count": status.get("sdr_count") if status else None,
|
||||
}
|
||||
|
||||
|
||||
@router.post("/secondary/priority")
|
||||
async def set_secondary_priority(priority: List[str] = Body(..., embed=True)):
|
||||
"""Set the ordered list the spare SDRs work through. Never restarts op25."""
|
||||
from app.internal.secondary_priority import set_secondary_priority as apply_priority
|
||||
unknown = [m for m in priority if m not in SECONDARY_SDR_MODES]
|
||||
if unknown:
|
||||
raise HTTPException(400, f"Unknown secondary SDR mode(s): {unknown}")
|
||||
running = await apply_priority(priority)
|
||||
return {"ok": True, "priority": load_node_config().secondary_sdr_priority, "running": running}
|
||||
|
||||
|
||||
@router.post("/discord/join")
|
||||
async def discord_join(guild_id: int, channel_id: int):
|
||||
ok = await radio_bot.join(guild_id, channel_id)
|
||||
|
||||
@@ -105,6 +105,30 @@
|
||||
background: rgba(255, 255, 255, 0.1);
|
||||
}
|
||||
|
||||
.sdr-row {
|
||||
display: flex;
|
||||
align-items: center;
|
||||
gap: 0.6rem;
|
||||
padding: 0.55rem 0;
|
||||
border-bottom: 1px solid var(--glass-border);
|
||||
}
|
||||
.sdr-row:last-child { border-bottom: none; }
|
||||
.sdr-rank { width: 1.2rem; color: var(--text-muted); font-size: 0.8rem; text-align: right; }
|
||||
.sdr-name { flex: 1; }
|
||||
.sdr-name small { display: block; color: var(--text-muted); font-size: 0.75rem; }
|
||||
.sdr-move {
|
||||
background: rgba(255, 255, 255, 0.05);
|
||||
border: 1px solid var(--glass-border);
|
||||
color: var(--text-main);
|
||||
border-radius: 6px;
|
||||
width: 1.9rem;
|
||||
height: 1.9rem;
|
||||
cursor: pointer;
|
||||
}
|
||||
.sdr-move:disabled { opacity: 0.3; cursor: default; }
|
||||
.sdr-state { font-size: 0.75rem; min-width: 6.5rem; text-align: right; color: var(--text-muted); }
|
||||
.sdr-state.on { color: var(--success); }
|
||||
|
||||
.grid {
|
||||
display: grid;
|
||||
grid-template-columns: repeat(auto-fit, minmax(320px, 1fr));
|
||||
@@ -360,6 +384,22 @@
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<!-- Secondary SDRs Card (node-26#9) -->
|
||||
<div class="glass-card">
|
||||
<div class="card-header">
|
||||
<div class="card-title">Secondary SDRs</div>
|
||||
</div>
|
||||
<p style="color: var(--text-muted); font-size: 0.8rem; margin: 0 0 0.5rem;">
|
||||
OP25 always keeps its own SDR. Every other SDR runs the next enabled item below, top first.
|
||||
<span id="sdr-summary"></span>
|
||||
</p>
|
||||
<div id="sdr-list"></div>
|
||||
<div style="display:flex; gap:0.75rem; align-items:center; margin-top:0.75rem;">
|
||||
<button id="sdr-save" class="btn btn-primary" style="padding:0.5rem 1rem;" onclick="saveSdrPriority()" disabled>Save priority</button>
|
||||
<span id="sdr-msg" style="color: var(--text-muted); font-size: 0.8rem;"></span>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<!-- Local Listening Card -->
|
||||
<div class="glass-card player-card">
|
||||
<div class="card-header" style="margin-bottom:0.5rem; border:none;">
|
||||
@@ -492,6 +532,99 @@
|
||||
}
|
||||
}
|
||||
|
||||
// ── Secondary SDR priority ────────────────────────────────────────────
|
||||
const SDR_LABELS = {
|
||||
adsb: ['ADS-B', 'Aircraft · 1090 MHz'],
|
||||
ais: ['AIS', 'Vessels · 162 MHz'],
|
||||
};
|
||||
let sdrRows = []; // [{mode, enabled}] in display order
|
||||
let sdrRunning = [];
|
||||
let sdrDirty = false;
|
||||
|
||||
function renderSdr() {
|
||||
const list = document.getElementById('sdr-list');
|
||||
list.innerHTML = '';
|
||||
let rank = 0;
|
||||
sdrRows.forEach((row, i) => {
|
||||
const [name, hint] = SDR_LABELS[row.mode] || [row.mode, ''];
|
||||
const running = sdrRunning.includes(row.mode);
|
||||
const state = !row.enabled ? 'Off' : running ? 'Running' : (sdrDirty ? 'Unsaved' : 'Waiting for SDR');
|
||||
const el = document.createElement('div');
|
||||
el.className = 'sdr-row';
|
||||
el.innerHTML = `
|
||||
<span class="sdr-rank">${row.enabled ? ++rank : ''}</span>
|
||||
<input type="checkbox" ${row.enabled ? 'checked' : ''} aria-label="Enable ${name}">
|
||||
<span class="sdr-name">${name}<small>${hint}</small></span>
|
||||
<button class="sdr-move" aria-label="Move ${name} up" ${i === 0 ? 'disabled' : ''}>▲</button>
|
||||
<button class="sdr-move" aria-label="Move ${name} down" ${i === sdrRows.length - 1 ? 'disabled' : ''}>▼</button>
|
||||
<span class="sdr-state ${running && row.enabled ? 'on' : ''}">${state}</span>`;
|
||||
const [box] = el.getElementsByTagName('input');
|
||||
const [up, down] = el.getElementsByTagName('button');
|
||||
box.onchange = () => { row.enabled = box.checked; markSdrDirty(); };
|
||||
up.onclick = () => { [sdrRows[i - 1], sdrRows[i]] = [sdrRows[i], sdrRows[i - 1]]; markSdrDirty(); };
|
||||
down.onclick = () => { [sdrRows[i + 1], sdrRows[i]] = [sdrRows[i], sdrRows[i + 1]]; markSdrDirty(); };
|
||||
list.appendChild(el);
|
||||
});
|
||||
document.getElementById('sdr-save').disabled = !sdrDirty;
|
||||
}
|
||||
|
||||
function markSdrDirty() {
|
||||
sdrDirty = true;
|
||||
document.getElementById('sdr-msg').textContent = '';
|
||||
renderSdr();
|
||||
}
|
||||
|
||||
async function loadSdr() {
|
||||
if (sdrDirty) return; // never clobber unsaved edits with a poll
|
||||
try {
|
||||
const r = await fetch('/api/secondary');
|
||||
if (!r.ok) return;
|
||||
const d = await r.json();
|
||||
sdrRunning = d.running || [];
|
||||
sdrRows = [
|
||||
...d.priority.map(mode => ({ mode, enabled: true })),
|
||||
...d.modes.filter(m => !d.priority.includes(m)).map(mode => ({ mode, enabled: false })),
|
||||
];
|
||||
const summary = document.getElementById('sdr-summary');
|
||||
if (d.running === null) {
|
||||
summary.textContent = 'Secondary SDR service is not responding.';
|
||||
} else if (d.sdr_count != null) {
|
||||
const spare = Math.max(d.sdr_count - 1, 0);
|
||||
summary.textContent = `This node has ${d.sdr_count} SDR${d.sdr_count === 1 ? '' : 's'}, so ${spare} spare.`;
|
||||
}
|
||||
renderSdr();
|
||||
} catch (e) {
|
||||
console.error('Secondary SDR load failed:', e);
|
||||
}
|
||||
}
|
||||
|
||||
async function saveSdrPriority() {
|
||||
const priority = sdrRows.filter(r => r.enabled).map(r => r.mode);
|
||||
const btn = document.getElementById('sdr-save');
|
||||
const msg = document.getElementById('sdr-msg');
|
||||
btn.disabled = true;
|
||||
msg.textContent = 'Applying…';
|
||||
try {
|
||||
const r = await fetch('/api/secondary/priority', {
|
||||
method: 'POST',
|
||||
headers: { 'Content-Type': 'application/json' },
|
||||
body: JSON.stringify({ priority }),
|
||||
});
|
||||
if (!r.ok) throw new Error(await r.text());
|
||||
const d = await r.json();
|
||||
sdrDirty = false;
|
||||
msg.textContent = d.running === null ? 'Saved, but the secondary SDR service did not respond.' : 'Saved.';
|
||||
await loadSdr();
|
||||
} catch (e) {
|
||||
console.error('Secondary SDR save failed:', e);
|
||||
msg.textContent = 'Save failed.';
|
||||
btn.disabled = false;
|
||||
}
|
||||
}
|
||||
|
||||
loadSdr();
|
||||
setInterval(loadSdr, 10000);
|
||||
|
||||
refresh();
|
||||
setInterval(refresh, 2000); // Polling every 2 seconds
|
||||
</script>
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -0,0 +1,58 @@
|
||||
"""
|
||||
node-26#9 — secondary SDR priority: an ordered list the SDRs beyond op25's work
|
||||
through. Normalisation, migration from the legacy single-mode field, and the one
|
||||
apply path shared by the dashboard, C2 commands and config pushes.
|
||||
"""
|
||||
import asyncio
|
||||
from unittest.mock import AsyncMock, patch
|
||||
|
||||
from app.models import NodeConfig, normalize_secondary_priority
|
||||
|
||||
|
||||
def _cfg(**kw) -> NodeConfig:
|
||||
return NodeConfig(node_id="n1", node_name="N1", lat=0.0, lon=0.0, **kw)
|
||||
|
||||
|
||||
def test_normalize_keeps_order_drops_unknown_and_duplicates():
|
||||
assert normalize_secondary_priority(["ais", "bogus", "adsb", "ais"]) == ["ais", "adsb"]
|
||||
assert normalize_secondary_priority([]) == []
|
||||
|
||||
|
||||
def test_legacy_single_mode_migrates_to_priority():
|
||||
assert _cfg(secondary_sdr_mode="adsb").secondary_sdr_priority == ["adsb"]
|
||||
assert _cfg(secondary_sdr_mode="none").secondary_sdr_priority == []
|
||||
|
||||
|
||||
def test_explicit_priority_wins_over_legacy_mode():
|
||||
assert _cfg(secondary_sdr_mode="adsb", secondary_sdr_priority=["ais"]).secondary_sdr_priority == ["ais"]
|
||||
|
||||
|
||||
def test_set_priority_persists_applies_and_never_touches_op25(tmp_path):
|
||||
config_file = tmp_path / "node_config.json"
|
||||
import app.internal.config_manager as cm
|
||||
with patch.object(cm, "_CONFIG_FILE", config_file):
|
||||
cm.save_node_config(_cfg(secondary_sdr_mode="adsb"))
|
||||
from app.internal import secondary_priority as sp
|
||||
with patch.object(sp.secondary_sdr_client, "apply", AsyncMock(return_value=["ais"])) as apply, \
|
||||
patch("app.internal.mqtt_manager.mqtt_manager.publish_checkin", AsyncMock()), \
|
||||
patch("app.internal.op25_client.op25_client.stop", AsyncMock()) as op25_stop:
|
||||
running = asyncio.run(sp.set_secondary_priority(["ais", "adsb", "ais"]))
|
||||
saved = cm.load_node_config()
|
||||
|
||||
assert running == ["ais"]
|
||||
apply.assert_awaited_once_with(["ais", "adsb"])
|
||||
op25_stop.assert_not_awaited()
|
||||
assert saved.secondary_sdr_priority == ["ais", "adsb"]
|
||||
assert saved.secondary_sdr_mode == "ais"
|
||||
|
||||
|
||||
def test_clearing_priority_is_not_undone_by_legacy_migration(tmp_path):
|
||||
config_file = tmp_path / "node_config.json"
|
||||
import app.internal.config_manager as cm
|
||||
with patch.object(cm, "_CONFIG_FILE", config_file):
|
||||
cm.save_node_config(_cfg(secondary_sdr_mode="adsb"))
|
||||
from app.internal import secondary_priority as sp
|
||||
with patch.object(sp.secondary_sdr_client, "apply", AsyncMock(return_value=[])), \
|
||||
patch("app.internal.mqtt_manager.mqtt_manager.publish_checkin", AsyncMock()):
|
||||
asyncio.run(sp.set_secondary_priority([]))
|
||||
assert cm.load_node_config().secondary_sdr_priority == []
|
||||
+13
-2
@@ -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
@@ -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 }}
|
||||
@@ -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 .
|
||||
@@ -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
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
# Secondary-SDR Container — node-26#9
|
||||
#
|
||||
# Claims the node's SECOND physical SDR (the first is always op25's). Mode is
|
||||
# chosen at runtime via the control API, not baked in: dump1090 for ADS-B,
|
||||
# AIS-catcher for AIS (op25_2 mode is not handled here yet — see node-26#9).
|
||||
#
|
||||
# Device claiming is by RTL-SDR index, not serial (op25's DeviceConfig.args
|
||||
# has no serial concept either — see op25-container/app/models.py). Index 0
|
||||
# is reserved for op25; this container always addresses index 1. That's a
|
||||
# real limitation once serial-stable device binding matters (hot-unplug /
|
||||
# replug can swap indices) — tracked in node-26#9, not fixed here.
|
||||
#
|
||||
# UNVERIFIED: this image has not been built or run against real hardware in
|
||||
# this session (sandboxed authoring machine, no docker). dump1090 and
|
||||
# AIS-catcher's exact CLI flags below are believed correct from their
|
||||
# published docs but not confirmed against a real capture — the CTO/QA
|
||||
# review before this ships to a real node should build and smoke-test it.
|
||||
FROM python:3.14-slim
|
||||
|
||||
ENV DEBIAN_FRONTEND=noninteractive
|
||||
|
||||
RUN apt-get update && \
|
||||
apt-get upgrade -y && \
|
||||
apt-get install -y --no-install-recommends \
|
||||
git build-essential cmake pkg-config \
|
||||
librtlsdr-dev libusb-1.0-0-dev libssl-dev zlib1g-dev libzstd-dev libncurses-dev usbutils
|
||||
|
||||
# readsb (wiedehopf) — ADS-B decoder. antirez/dump1090 was used here first but
|
||||
# has no --write-json at all (it only serves /data.json over --net), so the
|
||||
# decoder exited on an unknown flag. readsb writes the dump1090-fa style
|
||||
# aircraft.json that _read_adsb_snapshot() parses.
|
||||
RUN git clone --depth 1 https://github.com/wiedehopf/readsb /opt/readsb && \
|
||||
cd /opt/readsb && make RTLSDR=yes
|
||||
|
||||
# AIS-catcher — AIS decoder.
|
||||
RUN git clone https://github.com/jvde-github/AIS-catcher /opt/AIS-catcher && \
|
||||
cd /opt/AIS-catcher && mkdir build && cd build && cmake .. && make
|
||||
|
||||
EXPOSE 8002
|
||||
|
||||
VOLUME ["/configs"]
|
||||
|
||||
WORKDIR /app
|
||||
COPY ./app /app
|
||||
|
||||
COPY docker-entrypoint.sh /usr/local/bin/
|
||||
RUN sed -i 's/\r$//' /usr/local/bin/docker-entrypoint.sh && \
|
||||
chmod +x /usr/local/bin/docker-entrypoint.sh
|
||||
|
||||
COPY requirements.txt /tmp/requirements.txt
|
||||
RUN pip3 install --no-cache-dir -r /tmp/requirements.txt
|
||||
|
||||
ENTRYPOINT ["/usr/local/bin/docker-entrypoint.sh"]
|
||||
CMD ["python", "main.py"]
|
||||
@@ -0,0 +1,20 @@
|
||||
from pydantic_settings import BaseSettings
|
||||
|
||||
|
||||
class Settings(BaseSettings):
|
||||
# Same rationale as op25-container's OP25_DEBUG_EXPOSE: both containers
|
||||
# share the host network namespace (network_mode: host), so edge-node
|
||||
# reaches this control API over localhost regardless of this flag. False
|
||||
# (default) binds 127.0.0.1; true exposes unauthenticated start/stop to
|
||||
# the node's LAN and should only ever be set for local development.
|
||||
secondary_sdr_debug_expose: bool = False
|
||||
|
||||
class Config:
|
||||
env_file = ".env"
|
||||
|
||||
|
||||
settings = Settings()
|
||||
|
||||
|
||||
def bind_host() -> str:
|
||||
return "0.0.0.0" if settings.secondary_sdr_debug_expose else "127.0.0.1"
|
||||
@@ -0,0 +1,302 @@
|
||||
import json
|
||||
import os
|
||||
import signal
|
||||
import subprocess
|
||||
import threading
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from internal.logger import create_logger
|
||||
|
||||
LOGGER = create_logger(__name__)
|
||||
|
||||
# One decoder per mode, each on its own SDR (node-26#9). The node's secondary
|
||||
# SDR *priority* (e.g. ["adsb", "ais"]) is applied by apply(): decoders start
|
||||
# down the list until no free dongle is left, so every SDR the node has gets
|
||||
# used and the ones beyond the list's reach simply aren't started. OP25 always
|
||||
# keeps its own dongle — it is started first and never part of this list.
|
||||
MODES = ("adsb", "ais")
|
||||
|
||||
# Which RTL-SDR index op25 holds is NOT fixed — on radio-box op25 had index 1
|
||||
# and index 0 was free, so "op25 is always 0" was wrong. rtlsdr can't open a
|
||||
# dongle another process has claimed, and the decoders exit within ~50ms when
|
||||
# that happens, so _start_one() tries each index and keeps the first that
|
||||
# stays up. op25 is never disturbed: a failed claim doesn't touch its dongle.
|
||||
MAX_SDR_INDEX = 4
|
||||
_STARTUP_GRACE_S = 2.0
|
||||
|
||||
_STATE_DIR = Path("/tmp/secondary_sdr")
|
||||
ADSB_JSON_DIR = Path("/tmp/adsb")
|
||||
|
||||
# Live decoder handles. poll() is the only reliable liveness check: a decoder
|
||||
# that dies on startup stays an unreaped zombie, and killpg(pgid, 0) still
|
||||
# succeeds on a zombie.
|
||||
_procs: Dict[str, subprocess.Popen] = {}
|
||||
_indices: Dict[str, int] = {}
|
||||
_lock = threading.Lock()
|
||||
|
||||
# AIS-catcher streams one JSON object per received message on stdout rather
|
||||
# than writing a periodic snapshot file (readsb's approach) — so the
|
||||
# current-vessel snapshot lives in memory, keyed by mmsi, kept warm by a
|
||||
# background reader thread for as long as the decoder is alive.
|
||||
_ais_vessels: Dict[str, Dict[str, Any]] = {}
|
||||
_ais_lock = threading.Lock()
|
||||
|
||||
|
||||
def _pgid_file(mode: str) -> Path:
|
||||
return _STATE_DIR / f"{mode}.pgid"
|
||||
|
||||
|
||||
def _reap_orphans() -> None:
|
||||
"""Kill decoders left behind by a previous API process (uvicorn --reload
|
||||
restarts this process, but decoders run in their own session and would
|
||||
otherwise keep holding their SDRs)."""
|
||||
if not _STATE_DIR.exists():
|
||||
return
|
||||
for f in _STATE_DIR.glob("*.pgid"):
|
||||
try:
|
||||
os.killpg(int(f.read_text().strip()), signal.SIGTERM)
|
||||
LOGGER.info(f"Stopped orphaned secondary decoder from {f.name}")
|
||||
except Exception:
|
||||
pass
|
||||
f.unlink(missing_ok=True)
|
||||
|
||||
|
||||
_reap_orphans()
|
||||
|
||||
|
||||
def _adsb_command(index: int) -> List[str]:
|
||||
ADSB_JSON_DIR.mkdir(parents=True, exist_ok=True)
|
||||
return [
|
||||
"/opt/readsb/readsb",
|
||||
"--net",
|
||||
"--device-type", "rtlsdr",
|
||||
"--device", str(index),
|
||||
"--write-json", str(ADSB_JSON_DIR),
|
||||
"--write-json-every", "1",
|
||||
]
|
||||
|
||||
|
||||
def _ais_command(index: int) -> List[str]:
|
||||
return [
|
||||
"/opt/AIS-catcher/build/AIS-catcher",
|
||||
f"-d:{index}", # "-d <x>" would select by serial, not index
|
||||
"-o", "5", # JSON Full: decoded fields (4 = sparse, "JSON" is rejected)
|
||||
]
|
||||
|
||||
|
||||
_COMMANDS = {"adsb": _adsb_command, "ais": _ais_command}
|
||||
|
||||
|
||||
def _ais_reader(proc: subprocess.Popen) -> None:
|
||||
"""
|
||||
Consume AIS-catcher's stdout, one JSON message per line, and keep the
|
||||
latest report per mmsi. Field names (mmsi/lat/lon/speed/course or
|
||||
heading/shipname or name) are believed correct from AIS-catcher's
|
||||
published JSON output docs but UNVERIFIED against a real capture in
|
||||
this session — same caveat as dump1090's aircraft.json mapping.
|
||||
Malformed/partial lines (e.g. static-data-only messages with no
|
||||
position) are skipped rather than raising, since dropping one line must
|
||||
never kill the reader thread.
|
||||
"""
|
||||
if not proc.stdout:
|
||||
return
|
||||
for line in proc.stdout:
|
||||
try:
|
||||
msg = json.loads(line)
|
||||
except Exception:
|
||||
continue
|
||||
mmsi = msg.get("mmsi")
|
||||
if not mmsi:
|
||||
continue
|
||||
# AIS-catcher emits separate message TYPES per mmsi — static data
|
||||
# (name, no position) and position reports (lat/lon, no name) arrive
|
||||
# as distinct lines. Merge onto the existing entry, only overwriting
|
||||
# a field the new message actually carries, so a position-only
|
||||
# report doesn't blank out a name learned from an earlier message.
|
||||
name = (msg.get("shipname") or msg.get("name") or "").strip() or None
|
||||
heading = msg.get("heading") if msg.get("heading") is not None else msg.get("course")
|
||||
updates = {
|
||||
"mmsi": str(mmsi),
|
||||
"name": name,
|
||||
"lat": msg.get("lat"),
|
||||
"lon": msg.get("lon"),
|
||||
"speed_kt": msg.get("speed"),
|
||||
"heading_deg": heading,
|
||||
}
|
||||
with _ais_lock:
|
||||
existing = _ais_vessels.get(str(mmsi), {})
|
||||
for key, value in updates.items():
|
||||
if value is not None:
|
||||
existing[key] = value
|
||||
_ais_vessels[str(mmsi)] = existing
|
||||
|
||||
|
||||
def is_running(mode: str) -> bool:
|
||||
proc = _procs.get(mode)
|
||||
return proc is not None and proc.poll() is None
|
||||
|
||||
|
||||
def running() -> List[str]:
|
||||
return [m for m in MODES if is_running(m)]
|
||||
|
||||
|
||||
def _start_one(mode: str) -> bool:
|
||||
"""Start one decoder on the first free SDR. False when none is free."""
|
||||
if mode not in _COMMANDS:
|
||||
raise ValueError(f"Unknown secondary SDR mode: {mode!r}")
|
||||
if is_running(mode):
|
||||
return True
|
||||
if mode == "ais":
|
||||
with _ais_lock:
|
||||
_ais_vessels.clear()
|
||||
|
||||
needs_stdout = mode == "ais"
|
||||
for index in range(MAX_SDR_INDEX):
|
||||
if index in {_indices[m] for m in running() if m in _indices}:
|
||||
continue
|
||||
try:
|
||||
proc = subprocess.Popen(
|
||||
_COMMANDS[mode](index),
|
||||
preexec_fn=os.setsid,
|
||||
stdout=subprocess.PIPE if needs_stdout else None,
|
||||
text=True if needs_stdout else None,
|
||||
bufsize=1 if needs_stdout else -1,
|
||||
)
|
||||
except Exception as e:
|
||||
LOGGER.error(f"Failed to start secondary SDR decoder mode={mode!r}: {e}")
|
||||
return False
|
||||
try:
|
||||
proc.wait(timeout=_STARTUP_GRACE_S)
|
||||
LOGGER.info(f"Secondary SDR decoder mode={mode!r} could not use SDR index {index}, trying next")
|
||||
continue
|
||||
except subprocess.TimeoutExpired:
|
||||
pass
|
||||
if needs_stdout:
|
||||
threading.Thread(target=_ais_reader, args=(proc,), daemon=True).start()
|
||||
_procs[mode] = proc
|
||||
_indices[mode] = index
|
||||
_STATE_DIR.mkdir(parents=True, exist_ok=True)
|
||||
_pgid_file(mode).write_text(str(proc.pid))
|
||||
LOGGER.info(f"Started secondary SDR decoder mode={mode!r} on SDR index {index} pid={proc.pid}")
|
||||
return True
|
||||
|
||||
LOGGER.info(f"Secondary SDR decoder mode={mode!r}: no free SDR left")
|
||||
return False
|
||||
|
||||
|
||||
def _stop_one(mode: str) -> None:
|
||||
proc = _procs.pop(mode, None)
|
||||
_indices.pop(mode, None)
|
||||
if proc is not None:
|
||||
try:
|
||||
os.killpg(proc.pid, signal.SIGTERM)
|
||||
except OSError:
|
||||
pass
|
||||
try:
|
||||
proc.wait(timeout=5)
|
||||
except subprocess.TimeoutExpired:
|
||||
pass
|
||||
_pgid_file(mode).unlink(missing_ok=True)
|
||||
|
||||
|
||||
def start(mode: str) -> bool:
|
||||
with _lock:
|
||||
return _start_one(mode)
|
||||
|
||||
|
||||
def stop(mode: Optional[str] = None) -> None:
|
||||
with _lock:
|
||||
for m in [mode] if mode else list(_procs):
|
||||
_stop_one(m)
|
||||
|
||||
|
||||
def apply(priority: List[str]) -> List[str]:
|
||||
"""Run decoders in priority order until SDRs run out; stop everything else.
|
||||
|
||||
Unchanged when the running set is already the right prefix of the list, so
|
||||
re-applying the same priority (every restart/config push) is a no-op
|
||||
rather than a decoder restart.
|
||||
"""
|
||||
for m in priority:
|
||||
if m not in _COMMANDS:
|
||||
raise ValueError(f"Unknown secondary SDR mode: {m!r}")
|
||||
with _lock:
|
||||
live = set(running())
|
||||
if live and live == set(priority[:len(live)]):
|
||||
# Already running the head of the list; only try to extend it.
|
||||
for m in priority[len(live):]:
|
||||
if not _start_one(m):
|
||||
break
|
||||
return running()
|
||||
for m in list(_procs):
|
||||
_stop_one(m)
|
||||
for m in priority:
|
||||
if not _start_one(m):
|
||||
break
|
||||
return running()
|
||||
|
||||
|
||||
def sdr_count() -> Optional[int]:
|
||||
"""RTL-SDR dongles on USB, same heuristic as install.sh. This container has
|
||||
usbutils; op25's :stable image doesn't, so its /op25/devices can't answer."""
|
||||
try:
|
||||
out = subprocess.run(["lsusb"], capture_output=True, text=True, timeout=5).stdout
|
||||
except Exception:
|
||||
return None
|
||||
return sum(1 for line in out.splitlines() if "0bda:2838" in line or "0bda:2832" in line)
|
||||
|
||||
|
||||
def status() -> Dict[str, Any]:
|
||||
live = running()
|
||||
return {
|
||||
"running": [{"mode": m, "sdr_index": _indices.get(m)} for m in live],
|
||||
"sdr_count": sdr_count(),
|
||||
}
|
||||
|
||||
|
||||
def _altitude(a: Dict[str, Any]) -> Optional[int]:
|
||||
alt = a.get("alt_baro", a.get("altitude"))
|
||||
if alt == "ground":
|
||||
return 0
|
||||
return alt if isinstance(alt, (int, float)) else None
|
||||
|
||||
|
||||
def _read_adsb_snapshot() -> List[Dict[str, Any]]:
|
||||
"""
|
||||
Map readsb's aircraft.json (--write-json output) to the server's
|
||||
telemetry schema. readsb uses the dump1090-fa field names: alt_baro (int,
|
||||
or the string "ground"), gs, track. Older dump1090 forks used
|
||||
altitude/speed, kept as a fallback.
|
||||
"""
|
||||
path = ADSB_JSON_DIR / "aircraft.json"
|
||||
try:
|
||||
raw = json.loads(path.read_text())
|
||||
except Exception:
|
||||
return []
|
||||
|
||||
out = []
|
||||
for a in raw.get("aircraft", []):
|
||||
icao = a.get("hex")
|
||||
if not icao:
|
||||
continue
|
||||
out.append({
|
||||
"icao": icao.upper(),
|
||||
"callsign": (a.get("flight") or "").strip() or None,
|
||||
"lat": a.get("lat"),
|
||||
"lon": a.get("lon"),
|
||||
"altitude_ft": _altitude(a),
|
||||
"ground_speed_kt": a.get("gs", a.get("speed")),
|
||||
"track_deg": a.get("track"),
|
||||
})
|
||||
return out
|
||||
|
||||
|
||||
def data() -> Dict[str, Any]:
|
||||
with _ais_lock:
|
||||
vessels = list(_ais_vessels.values()) if is_running("ais") else []
|
||||
return {
|
||||
"running": running(),
|
||||
"aircraft": _read_adsb_snapshot() if is_running("adsb") else [],
|
||||
"vessels": vessels,
|
||||
}
|
||||
@@ -0,0 +1,31 @@
|
||||
import logging
|
||||
from logging.handlers import RotatingFileHandler
|
||||
|
||||
|
||||
def create_logger(name, level=logging.DEBUG, max_bytes=10485760, backup_count=2):
|
||||
debug_log_file = "./secondary-sdr.debug.log"
|
||||
info_log_file = "./secondary-sdr.log"
|
||||
|
||||
logger = logging.getLogger(name)
|
||||
logger.setLevel(level)
|
||||
|
||||
if not logger.hasHandlers():
|
||||
console_handler = logging.StreamHandler()
|
||||
console_handler.setLevel(level)
|
||||
|
||||
debug_file_handler = RotatingFileHandler(debug_log_file, maxBytes=max_bytes, backupCount=backup_count)
|
||||
debug_file_handler.setLevel(logging.DEBUG)
|
||||
|
||||
info_file_handler = RotatingFileHandler(info_log_file, maxBytes=max_bytes, backupCount=backup_count)
|
||||
info_file_handler.setLevel(logging.INFO)
|
||||
|
||||
formatter = logging.Formatter('%(asctime)s - %(name)s - %(levelname)s - %(message)s')
|
||||
console_handler.setFormatter(formatter)
|
||||
debug_file_handler.setFormatter(formatter)
|
||||
info_file_handler.setFormatter(formatter)
|
||||
|
||||
logger.addHandler(console_handler)
|
||||
logger.addHandler(debug_file_handler)
|
||||
logger.addHandler(info_file_handler)
|
||||
|
||||
return logger
|
||||
@@ -0,0 +1,13 @@
|
||||
from fastapi import FastAPI
|
||||
import routers.secondary_controller as secondary_controller
|
||||
from config import bind_host
|
||||
|
||||
app = FastAPI()
|
||||
|
||||
app.include_router(secondary_controller.create_secondary_router(), prefix="/secondary")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
import uvicorn
|
||||
|
||||
uvicorn.run("main:app", host=bind_host(), port=8002, reload=True)
|
||||
@@ -0,0 +1,58 @@
|
||||
from typing import List, Optional
|
||||
|
||||
from fastapi import APIRouter, HTTPException
|
||||
from pydantic import BaseModel
|
||||
|
||||
from internal import decoder_control
|
||||
from internal.logger import create_logger
|
||||
|
||||
LOGGER = create_logger(__name__)
|
||||
|
||||
|
||||
class StartBody(BaseModel):
|
||||
mode: str # adsb | ais
|
||||
|
||||
|
||||
class StopBody(BaseModel):
|
||||
mode: Optional[str] = None # omit to stop every decoder
|
||||
|
||||
|
||||
class ApplyBody(BaseModel):
|
||||
priority: List[str] # ordered, e.g. ["adsb", "ais"]
|
||||
|
||||
|
||||
def create_secondary_router():
|
||||
router = APIRouter()
|
||||
|
||||
@router.post("/apply")
|
||||
async def apply(body: ApplyBody):
|
||||
try:
|
||||
live = decoder_control.apply(body.priority)
|
||||
except ValueError as e:
|
||||
raise HTTPException(status_code=400, detail=str(e))
|
||||
return {"running": live}
|
||||
|
||||
@router.post("/start")
|
||||
async def start(body: StartBody):
|
||||
try:
|
||||
ok = decoder_control.start(body.mode)
|
||||
except ValueError as e:
|
||||
raise HTTPException(status_code=400, detail=str(e))
|
||||
if not ok:
|
||||
raise HTTPException(status_code=409, detail="No free SDR for this decoder")
|
||||
return {"status": f"secondary SDR started ({body.mode})"}
|
||||
|
||||
@router.post("/stop")
|
||||
async def stop(body: Optional[StopBody] = None):
|
||||
decoder_control.stop(body.mode if body else None)
|
||||
return {"status": "secondary SDR stopped"}
|
||||
|
||||
@router.get("/status")
|
||||
async def get_status():
|
||||
return decoder_control.status()
|
||||
|
||||
@router.get("/data")
|
||||
async def get_data():
|
||||
return decoder_control.data()
|
||||
|
||||
return router
|
||||
@@ -0,0 +1,3 @@
|
||||
#!/bin/bash
|
||||
mkdir -p /tmp/adsb
|
||||
exec "$@"
|
||||
@@ -0,0 +1,3 @@
|
||||
uvicorn
|
||||
fastapi
|
||||
pydantic-settings
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user