Compare commits

..
Author SHA1 Message Date
sjg 39e59dca96 [fix](trx-server): gate test-only history imports
CI / lint (pull_request) Failing after 1s
CI / test (pull_request) Failing after 1s
CI / reuse (pull_request) Failing after 1s
CI / lint (push) Failing after 3s
CI / test (push) Failing after 2s
CI / reuse (push) Failing after 1s
2026-08-01 01:55:33 +02:00
sjg 5654520901 [fix](workspace): clear build and clippy warnings
CI / lint (push) Failing after 1s
CI / test (push) Failing after 1s
CI / reuse (push) Failing after 1s
2026-08-01 01:52:08 +02:00
sjg ff75fcc692 [refactor](trx-server): extract decoder history store
CI / lint (pull_request) Failing after 2s
CI / test (pull_request) Failing after 2s
CI / reuse (pull_request) Failing after 1s
CI / lint (push) Failing after 2s
CI / test (push) Failing after 2s
CI / reuse (push) Failing after 1s
2026-08-01 01:45:19 +02:00
sjg 4728b578ae [refactor](trx-server): extract history policy 2026-08-01 01:44:26 +02:00
sjg 7db7fad9b0 [refactor](trx-frontend): define module service boundaries 2026-08-01 01:44:26 +02:00
sjg 830f7299fe [fix](trx-frontend): vendor Opus decoder 2026-08-01 01:44:26 +02:00
sjg 061738a63b [fix](trx-frontend): serialize plugin loading
CI / lint (push) Failing after 1s
CI / test (push) Failing after 1s
CI / reuse (push) Failing after 1s
2026-08-01 01:43:50 +02:00
sjg e8bd97655f [fix](trx-server): offload history persistence
CI / reuse (push) Failing after 1s
CI / lint (push) Failing after 2s
CI / test (push) Failing after 2s
2026-08-01 01:43:47 +02:00
sjg 7d0b36450d [fix](trx-server): persist all decoder histories
CI / reuse (push) Failing after 1s
CI / lint (pull_request) Failing after 3s
CI / test (pull_request) Failing after 1s
CI / reuse (pull_request) Failing after 2s
CI / lint (push) Failing after 1s
CI / test (push) Failing after 1s
2026-08-01 01:23:36 +02:00
Stan Grams b12c83e8b5 [fix](trx-frontend): preload AIS map icons
CI / lint (pull_request) Failing after 6s
CI / test (pull_request) Failing after 2s
CI / reuse (pull_request) Failing after 1s
CI / lint (push) Failing after 1s
CI / test (push) Failing after 2s
CI / reuse (push) Failing after 0s
2026-08-01 01:14:07 +02:00
30 changed files with 1366 additions and 871 deletions
-22
View File
@@ -1,22 +0,0 @@
{
"name": "trx-rs SDK",
"image": "git.haxx.space/sjg/trx-rs/sdk:latest",
"workspaceFolder": "/work",
"workspaceMount": "source=${localWorkspaceFolder},target=/work,type=bind",
"mounts": [
"source=trx-rs-sccache,target=/sccache,type=volume"
],
"containerEnv": {
"RUSTC_WRAPPER": "sccache",
"CARGO_INCREMENTAL": "0",
"SCCACHE_DIR": "/sccache"
},
"customizations": {
"vscode": {
"extensions": [
"rust-lang.rust-analyzer",
"tamasfe.even-better-toml"
]
}
}
}
+9 -20
View File
@@ -2,10 +2,10 @@
#
# SPDX-License-Identifier: GPL-2.0-or-later
# CI for the Docker-executor runner (VM). The lint/test jobs run inside the
# shared trx-rs SDK image (container/Containerfile), which bakes in the pinned
# Rust toolchain and all build dependencies. The reuse job uses the upstream
# Docker action, which the Docker executor launches as a sibling container.
# CI for the self-hosted, host-executor Podman runners (see container/).
# The runner image bakes in the Rust toolchain and all build dependencies,
# so jobs go straight to cargo — no apt/rustup setup steps (which also
# collided on the dpkg lock when jobs ran concurrently in the same runner).
name: CI
@@ -16,43 +16,32 @@ on:
env:
CARGO_TERM_COLOR: always
# sccache: shared compilation cache persisted on the runner host (see the
# -v mount in runner-config.example.yaml). CARGO_INCREMENTAL=0 because
# sccache cannot cache incremental artifacts.
RUSTC_WRAPPER: sccache
CARGO_INCREMENTAL: "0"
SCCACHE_DIR: /sccache
SCCACHE_CACHE_SIZE: "20G"
jobs:
lint:
runs-on: ubuntu-latest
container: git.haxx.space/sjg/trx-rs/sdk:latest
steps:
- uses: actions/checkout@v4
- name: rustfmt
run: cargo fmt --all -- --check
- name: clippy
run: cargo clippy --workspace --all-targets --all-features -- -D warnings
- name: sccache stats
if: always()
run: sccache --show-stats
test:
runs-on: ubuntu-latest
container: git.haxx.space/sjg/trx-rs/sdk:latest
steps:
- uses: actions/checkout@v4
- name: Build
run: cargo build --workspace --all-targets --locked
- name: Test
run: cargo test --workspace --locked
- name: sccache stats
if: always()
run: sccache --show-stats
reuse:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: fsfe/reuse-action@v5
- name: REUSE compliance
# `reuse` CLI instead of fsfe/reuse-action: the latter is a Docker
# action, which the host-executor runners cannot run. `reuse` is baked
# into the runner image (see container/Containerfile).
run: reuse lint
+21
View File
@@ -0,0 +1,21 @@
MIT License
Copyright (c) <year> <copyright holders>
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in all
copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN THE
SOFTWARE.
+3 -3
View File
@@ -142,6 +142,6 @@ a unified set of frontends.
## License
GPL-2.0-or-later. See [`LICENSES`](LICENSES) for the full license text and
bundled third-party license files. Bundled third-party components (Leaflet and
the Leaflet AIS tracksymbol plugin under `assets/web/vendor/`) retain their
original BSD-2-Clause license.
bundled third-party license files. Bundled third-party components retain their
original licenses: Leaflet is BSD-2-Clause, DSEG is OFL-1.1, and opus-decoder
is MIT.
+8 -2
View File
@@ -12,8 +12,6 @@ path = [
"trx-rs.toml.example",
"docs/**",
"aidocs/**",
"container/**",
".devcontainer/**",
"src/decoders/trx-ftx/README.md",
"src/decoders/trx-wxsat/README.md",
"assets/trx-logo.png",
@@ -44,3 +42,11 @@ SPDX-License-Identifier = "BSD-2-Clause"
path = ["src/trx-client/trx-frontend/trx-frontend-http/assets/web/vendor/dseg14-classic-latin-400-normal.woff2"]
SPDX-FileCopyrightText = "2020 The DSEG Authors (https://github.com/keshikan/DSEG)"
SPDX-License-Identifier = "OFL-1.1"
# Vendored opus-decoder 0.7.11 browser build
# (https://github.com/eshaz/wasm-audio-decoders), MIT.
# SHA-256: fd73ee0a9c8a5e233c0b88234df14f46003ad89bb4ba27435bc8714db2a6dc62
[[annotations]]
path = ["src/trx-client/trx-frontend/trx-frontend-http/assets/web/vendor/opus-decoder-0.7.11.min.js"]
SPDX-FileCopyrightText = "2021-2025 Ethan Halsall"
SPDX-License-Identifier = "MIT"
+38 -32
View File
@@ -2,54 +2,60 @@
#
# SPDX-License-Identifier: GPL-2.0-or-later
# trx-rs SDK / build image.
# Gitea Actions runner image for trx-rs CI (host-executor / "Pattern B").
#
# Single source of truth for the build environment. Used two ways:
# * CI — as the job container for the lint/test jobs (Docker executor).
# * Dev — run locally or via .devcontainer for a reproducible toolchain.
#
# Pinning the Rust version here (and in rust-toolchain.toml) means CI and every
# developer share the exact same rustc/clippy, so "works locally, fails in CI"
# cannot happen.
# All build dependencies, the Rust toolchain, Node.js (for JS actions such as
# actions/checkout and actions/cache) and the `reuse` tool are baked in, so CI
# runs skip the per-run apt/rustup install cost. `sudo` is present so the
# existing workflow's `sudo apt-get ...` / rustup steps remain valid — they
# just become fast no-ops because everything is already installed.
FROM docker.io/library/debian:bookworm-slim
# Keep in sync with rust-toolchain.toml.
ARG RUST_VERSION=1.97.1
ARG ACT_RUNNER_VERSION=0.2.11
ARG NODE_MAJOR=20
ENV DEBIAN_FRONTEND=noninteractive \
RUSTUP_HOME=/usr/local/rustup \
CARGO_HOME=/usr/local/cargo \
PATH=/usr/local/cargo/bin:/usr/local/bin:/usr/bin:/bin
RUSTUP_HOME=/opt/rustup \
CARGO_HOME=/opt/cargo \
PATH=/opt/cargo/bin:/usr/local/bin:/usr/bin:/bin
# Build dependencies (mirror README's manual instructions).
# Base tooling + trx-rs build dependencies (mirrors .gitea/workflows/ci.yml).
RUN apt-get update && apt-get install -y --no-install-recommends \
ca-certificates curl git \
ca-certificates curl xz-utils git sudo pipx \
build-essential pkg-config cmake clang libclang-dev \
libopus-dev libasound2-dev libsoapysdr-dev \
&& rm -rf /var/lib/apt/lists/*
# Node.js JS-based actions (actions/checkout, actions/cache) run *inside*
# the job container under the Docker executor, so node must be present.
# Node.js (JS-based actions need node in PATH under the host executor).
RUN curl -fsSL https://deb.nodesource.com/setup_${NODE_MAJOR}.x | bash - \
&& apt-get install -y --no-install-recommends nodejs \
&& rm -rf /var/lib/apt/lists/*
# Pinned Rust toolchain, installed world-readable so any UID the runner or a
# devcontainer uses can invoke cargo.
# REUSE >= 3 (Debian's packaged reuse is too old for REUSE.toml).
# The [charset-normalizer] extra provides an encoding-detection backend;
# without it (and without libmagic) reuse fails to import at runtime.
RUN PIPX_HOME=/opt/pipx PIPX_BIN_DIR=/usr/local/bin pipx install 'reuse[charset-normalizer]'
# Rust stable with rustfmt + clippy, installed system-wide.
RUN curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs \
| sh -s -- -y --no-modify-path \
--default-toolchain "${RUST_VERSION}" --profile minimal \
--component rustfmt --component clippy \
&& chmod -R a+rwX "$RUSTUP_HOME" "$CARGO_HOME"
| sh -s -- -y --no-modify-path --profile minimal \
--component rustfmt --component clippy \
&& chmod -R a+rwX "$CARGO_HOME" "$RUSTUP_HOME"
# sccache — shared compilation cache. Enabled at build time via
# RUSTC_WRAPPER (see the CI workflow and .devcontainer), not repo-wide, so
# non-SDK builds are unaffected. musl build is static and runs anywhere.
ARG SCCACHE_VERSION=0.8.2
RUN curl -fsSL "https://github.com/mozilla/sccache/releases/download/v${SCCACHE_VERSION}/sccache-v${SCCACHE_VERSION}-x86_64-unknown-linux-musl.tar.gz" \
| tar -xz -C /tmp \
&& install -m755 "/tmp/sccache-v${SCCACHE_VERSION}-x86_64-unknown-linux-musl/sccache" /usr/local/bin/sccache \
&& rm -rf /tmp/sccache-*
# act_runner binary.
RUN arch="$(dpkg --print-architecture)"; \
case "$arch" in amd64) rarch=amd64;; arm64) rarch=arm64;; *) echo "unsupported arch $arch" >&2; exit 1;; esac; \
curl -fsSL -o /usr/local/bin/act_runner \
"https://gitea.com/gitea/act_runner/releases/download/v${ACT_RUNNER_VERSION}/act_runner-${ACT_RUNNER_VERSION}-linux-${rarch}" \
&& chmod +x /usr/local/bin/act_runner
WORKDIR /work
# Default config template (seeded into the /data volume on first boot).
COPY config.yaml /etc/act_runner/config.yaml
COPY entrypoint.sh /usr/local/bin/entrypoint.sh
RUN chmod +x /usr/local/bin/entrypoint.sh
# /data holds the .runner registration, cache and workflow workspaces.
VOLUME /data
WORKDIR /data
ENTRYPOINT ["/usr/local/bin/entrypoint.sh"]
+90 -102
View File
@@ -3,128 +3,116 @@ SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
SPDX-License-Identifier: GPL-2.0-or-later
-->
# trx-rs SDK image
# Podman-based Gitea Actions runners
A single container image that is the canonical build environment for trx-rs,
used **both** by CI and by developers. It bakes in the pinned Rust toolchain
(matching `rust-toolchain.toml`) and every build dependency, so the compiler
and `clippy` are identical everywhere — no "works on my machine".
Run two independent Gitea Actions runners on one host as rootless Podman
containers managed by systemd (Quadlet) — one per project — instead of two
VMs. Uses the **host executor**: workflow steps run directly inside a
purpose-built runner image that already has the Rust toolchain and all build
dependencies baked in, so CI runs skip the per-run install cost and no
Docker/Podman socket is needed.
## Files
| File | Purpose |
|------|---------|
| `Containerfile` | The SDK image (Debian + build deps + pinned Rust + Node + git). |
| `runner-config.example.yaml` | Example act_runner config for the CI VM (Docker executor). |
| `Containerfile` | Runner image: Debian + build deps + clang + Rust + Node + `reuse` + `act_runner`. |
| `entrypoint.sh` | Registers on first boot (if needed), then runs the daemon. |
| `config.yaml` | act_runner config template (seeded into each runner's volume). |
| `trx-rs-runner.container` | Quadlet unit for the trx-rs runner. |
| `project2-runner.container` | Quadlet unit for the second project's runner. |
## Build and publish
## Prerequisites (once per host)
Rootless Podman with cgroups v2 (default on modern distros). As the unprivileged
user that will own the runners:
```bash
# from the repo root
podman build -t git.haxx.space/sjg/trx-rs/sdk:latest container
podman login git.haxx.space
podman push git.haxx.space/sjg/trx-rs/sdk:latest
# Survive logout / start on boot without an interactive session.
loginctl enable-linger "$USER"
```
Tag with the Rust version too (e.g. `:1.97.1`) if you want reproducible pins.
Make the package **public** (Gitea → Packages → the image → Settings) so the CI
runner and developers can pull it without credentials. If you keep it private,
add `credentials:` under the workflow's `container:` and log the runner into the
registry.
No `podman.socket` is required for the host executor.
## Developer use
Reproducible one-off build, no local toolchain needed:
## 1. Build the image
```bash
podman run --rm -it -v "$PWD":/work -w /work \
git.haxx.space/sjg/trx-rs/sdk:latest \
cargo build --release
cd container
podman build -t trx-rs-ci:latest .
```
Or open the repo in the image via VS Code / JetBrains "Reopen in Container"
(`.devcontainer/devcontainer.json` points at the same image).
## 2. Get a registration token
Building outside the container? `rust-toolchain.toml` pins the same rustc, so
`rustup` installs the matching toolchain automatically.
For **each** repo: *Settings → Actions → Runners → Create new Runner* and copy
the token. (Org- or instance-level tokens work too if you prefer wider scope.)
## CI use
`.gitea/workflows/ci.yml` runs the `lint` and `test` jobs *inside* this image
via the `container:` key, so they skip all setup and go straight to `cargo`.
The `reuse` job stays on the upstream `fsfe/reuse-action` (a Docker action the
Docker executor launches as a sibling container) — nothing REUSE-related is
baked into the SDK.
## Compilation cache (sccache)
The SDK image ships [`sccache`](https://github.com/mozilla/sccache). It is
enabled via `RUSTC_WRAPPER=sccache` in CI and the devcontainer (not repo-wide,
so plain `cargo` builds outside the SDK are unaffected).
- **CI** persists the cache on the runner host — create the dir once:
`mkdir -p /var/cache/sccache`. It is bind-mounted into each job container at
`/sccache` (see `runner-config.example.yaml`), so cache survives across runs
and is shared between the lint/test jobs and both projects.
- **Devcontainer** uses a named volume (`trx-rs-sccache`).
- Check effectiveness with `sccache --show-stats` (the CI jobs print it).
`CARGO_INCREMENTAL=0` is set wherever sccache is on, since sccache cannot cache
incremental artifacts.
## CI runner (Alpine / OpenRC)
The runner uses the **Docker executor** (not the host executor): per-job
container isolation and standard `ubuntu-latest` semantics. `act_runner` runs
as an OpenRC service. Files provided:
| File | Purpose |
|------|---------|
| `act_runner.openrc` | OpenRC init script (`supervise-daemon`, depends on docker). |
| `act_runner.confd.example` | Per-instance `conf.d` settings for multi-runner hosts. |
**Cap the thread budget.** In a VM, pin its vCPUs to specific host threads
(libvirt/KVM):
```xml
<vcpu placement='static'>2</vcpu>
<cputune>
<vcpupin vcpu='0' cpuset='4'/>
<vcpupin vcpu='1' cpuset='5'/>
</cputune>
```
On bare metal, the `container.options: "--cpus=2"` and `capacity: 1` in
`runner-config.example.yaml` already bound each runner.
**Set it up:**
## 3. Install and start the runners
```bash
# 1. Docker + a dedicated user with socket access
apk add docker docker-cli
rc-update add docker default && rc-service docker start
adduser -S -D -H -h /var/lib/act_runner act
addgroup act docker
mkdir -p ~/.config/containers/systemd
cp trx-rs-runner.container project2-runner.container ~/.config/containers/systemd/
# 2. act_runner binary (static Go build, works on musl)
curl -fsSL -o /usr/local/bin/act_runner \
https://gitea.com/gitea/act_runner/releases/download/v0.2.11/act_runner-0.2.11-linux-amd64
chmod +x /usr/local/bin/act_runner
# Paste each repo's token for the FIRST boot only:
# Environment=GITEA_RUNNER_REGISTRATION_TOKEN=xxxx…
$EDITOR ~/.config/containers/systemd/trx-rs-runner.container
$EDITOR ~/.config/containers/systemd/project2-runner.container
# 3. Config + register one runner per project (scope keeps their jobs apart)
install -Dm644 container/runner-config.example.yaml /etc/act_runner/trx-rs.yaml
install -d -o act /var/lib/act_runner/trx-rs
su act -s /bin/sh -c 'cd /var/lib/act_runner/trx-rs && \
act_runner register --no-interactive \
--instance https://git.haxx.space --token <TOKEN> \
--name trx-rs-ci \
--labels "ubuntu-latest:docker://catthehacker/ubuntu:act-latest"'
systemctl --user daemon-reload
systemctl --user start trx-rs-runner
systemctl --user start project2-runner
# 4. OpenRC service (repeat the symlink+conf.d for the second project)
install -m755 container/act_runner.openrc /etc/init.d/act_runner
ln -s act_runner /etc/init.d/act_runner.trx-rs
install -m644 container/act_runner.confd.example /etc/conf.d/act_runner.trx-rs
rc-update add act_runner.trx-rs default
rc-service act_runner.trx-rs start
systemctl --user status trx-rs-runner
podman logs -f gitea-runner-trx-rs
```
Check it with `rc-service act_runner.trx-rs status` and
`tail -f /var/log/act_runner.trx-rs.log`.
Once each runner shows **online** in the repo's runner list, blank out the
`GITEA_RUNNER_REGISTRATION_TOKEN` line again (the registration is persisted in
the `…-data` volume) and `systemctl --user daemon-reload`.
## Required workflow change: the `reuse` job
The host executor runs steps directly in the container and therefore **cannot
run Docker-based actions**. The current `reuse` job uses `fsfe/reuse-action@v5`,
which is a Docker action. `reuse` is baked into the image, so replace that job
with a plain command:
```yaml
reuse:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- name: REUSE compliance
run: reuse lint
```
The `lint` and `test` jobs need no changes: their `sudo apt-get …` and rustup
steps still run, but become fast no-ops because the image already has those
packages and the toolchain. (`sudo` is included in the image for exactly this
reason.)
> If you would rather keep Docker-based actions and per-run images, use the
> **Docker executor** instead: drop the `:host` suffix from the label in
> `config.yaml`, enable `systemctl --user --now enable podman.socket`, mount it
> into the container, and set `container.docker_host` to the socket path. That
> trades the baked-in speed for stronger per-job isolation.
## Tuning
- **`capacity`** (in `config.yaml`) — concurrent jobs per runner. Rust builds
are heavy; 12 is sensible when two runners share a host.
- **`PodmanArgs=--cpus/--memory`** (in each `.container`) — hard resource caps
so one project cannot starve the other.
- **SELinux** — the `:Z` volume flag is already set; keep it if SELinux is
enforcing.
## Committing these files
If you add this directory to a REUSE-checked repo, register the markdown in
`REUSE.toml` (the other files carry inline SPDX headers):
```toml
[[annotations]]
path = ["container/**"]
SPDX-FileCopyrightText = "2026 Stan Grams <sjg@haxx.space>"
SPDX-License-Identifier = "GPL-2.0-or-later"
```
-14
View File
@@ -1,14 +0,0 @@
# SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
# SPDX-License-Identifier: GPL-2.0-or-later
#
# Per-instance settings for an act_runner OpenRC service.
# Copy to /etc/conf.d/<service-name>, e.g. /etc/conf.d/act_runner.trx-rs
# (the name must match the /etc/init.d/ symlink).
# User that runs the daemon. Must be a member of the `docker` group.
runner_user="act"
# Per-instance state dir (holds the .runner registration) and config file,
# so two runners on one host stay independent.
runner_dir="/var/lib/act_runner/trx-rs"
runner_config="/etc/act_runner/trx-rs.yaml"
-44
View File
@@ -1,44 +0,0 @@
#!/sbin/openrc-run
# SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
# SPDX-License-Identifier: GPL-2.0-or-later
#
# OpenRC service for a Gitea act_runner (Docker executor) on Alpine.
#
# Install as /etc/init.d/act_runner (chmod +x). Single instance uses
# /etc/act_runner/config.yaml. For one runner per project, symlink this script
# and add a matching conf.d file:
#
# ln -s act_runner /etc/init.d/act_runner.trx-rs
# cp container/act_runner.confd.example /etc/conf.d/act_runner.trx-rs
# $EDITOR /etc/conf.d/act_runner.trx-rs # set runner_dir / runner_config
# rc-update add act_runner.trx-rs default
# rc-service act_runner.trx-rs start
description="Gitea Actions runner"
: "${runner_user:=act}"
: "${runner_dir:=/var/lib/act_runner}"
: "${runner_config:=/etc/act_runner/config.yaml}"
command="/usr/local/bin/act_runner"
command_args="daemon --config ${runner_config}"
# No group given, so supplementary groups (incl. docker) are initialised.
command_user="${runner_user}"
directory="${runner_dir}"
supervisor="supervise-daemon"
respawn_delay=5
respawn_max=0
pidfile="/run/${RC_SVCNAME}.pid"
output_log="/var/log/${RC_SVCNAME}.log"
error_log="/var/log/${RC_SVCNAME}.log"
depend() {
need docker
use net dns
}
start_pre() {
checkpath -d -m 0750 -o "${runner_user}" "${runner_dir}"
checkpath -f -m 0640 -o "${runner_user}" "${output_log}"
}
+32
View File
@@ -0,0 +1,32 @@
# SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
#
# SPDX-License-Identifier: GPL-2.0-or-later
#
# act_runner configuration template. Seeded into /data/config.yaml on first
# boot; edit the copy inside the volume to change settings per runner.
log:
level: info
runner:
# Registration state. Relative to the daemon's working directory (/data).
file: .runner
# Concurrent jobs this runner will pick up. Rust builds are heavy — keep this
# modest, especially if two runners share one host. The trx-rs workflow has
# three parallel jobs (lint, test, reuse); capacity 2 lets two overlap.
capacity: 2
timeout: 3h
# Map the workflow's `runs-on: ubuntu-latest` to the HOST executor, i.e. run
# steps directly inside THIS container (which already has all the toolchain).
# No Docker/Podman socket is required in this mode.
labels:
- "ubuntu-latest:host"
cache:
# Built-in actions cache server (used by actions/cache). Stored in the volume.
enabled: true
dir: "/data/cache"
host:
# Where per-job workspaces are created.
workdir_parent: /data/workflows
+38
View File
@@ -0,0 +1,38 @@
#!/usr/bin/env bash
# SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
#
# SPDX-License-Identifier: GPL-2.0-or-later
#
# Registers the runner on first boot (if no .runner state exists in /data),
# then runs the act_runner daemon. Idempotent: on subsequent boots it reuses
# the stored registration and ignores the token.
set -euo pipefail
CONFIG_FILE="${CONFIG_FILE:-/data/config.yaml}"
cd /data
# Seed the config from the image's template on first boot so it lives in the
# persistent volume and can be edited there.
if [ ! -f "$CONFIG_FILE" ]; then
cp /etc/act_runner/config.yaml "$CONFIG_FILE"
fi
# runner.file in config.yaml is ".runner" (relative to this CWD => /data/.runner).
if [ ! -f /data/.runner ]; then
if [ -z "${GITEA_RUNNER_REGISTRATION_TOKEN:-}" ]; then
echo "ERROR: no /data/.runner registration and GITEA_RUNNER_REGISTRATION_TOKEN is empty." >&2
echo " Grab a token from the repo's Settings -> Actions -> Runners and set it" >&2
echo " in the Quadlet unit for the first boot only." >&2
exit 1
fi
echo "Registering runner '${GITEA_RUNNER_NAME:-podman}' with ${GITEA_INSTANCE_URL} ..."
act_runner register --no-interactive \
--config "$CONFIG_FILE" \
--instance "${GITEA_INSTANCE_URL:?set GITEA_INSTANCE_URL}" \
--token "$GITEA_RUNNER_REGISTRATION_TOKEN" \
--name "${GITEA_RUNNER_NAME:-podman}" \
--labels "${GITEA_RUNNER_LABELS:-ubuntu-latest:host}"
fi
exec act_runner daemon --config "$CONFIG_FILE"
+40
View File
@@ -0,0 +1,40 @@
# SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
#
# SPDX-License-Identifier: GPL-2.0-or-later
#
# Rootless Podman Quadlet for the SECOND project's Gitea Actions runner,
# co-located on the same host as the trx-rs runner.
#
# It has its own name, its own data volume and its own registration token, so
# the two runners are fully independent. They share the `ubuntu-latest` label,
# but registration SCOPE (which repo each token came from) keeps their jobs
# separate — neither will pick up the other's work.
#
# If project 2 needs different build dependencies, build it its own image from
# an adjusted Containerfile and point Image= at that instead of reusing the
# trx-rs image below.
[Unit]
Description=Gitea Actions runner — project 2
After=network-online.target
Wants=network-online.target
[Container]
Image=localhost/gitea-act-runner:latest
ContainerName=gitea-runner-project2
Volume=gitea-runner-project2-data:/data:Z
Environment=CONFIG_FILE=/data/config.yaml
Environment=GITEA_INSTANCE_URL=https://git.haxx.space
Environment=GITEA_RUNNER_NAME=project2-podman
Environment=GITEA_RUNNER_LABELS=ubuntu-latest:host
Environment=GITEA_RUNNER_REGISTRATION_TOKEN=
PodmanArgs=--cpus=4.0 --memory=6g
[Service]
Restart=always
TimeoutStartSec=0
[Install]
WantedBy=default.target
-35
View File
@@ -1,35 +0,0 @@
# SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
#
# SPDX-License-Identifier: GPL-2.0-or-later
#
# Example act_runner config for the Docker-executor runner that lives in the
# CI VM. This is NOT the SDK image — it configures the runner that launches
# per-job containers (including the trx-rs SDK image referenced by the
# workflow's `container:` key). Copy to the VM and pass with
# `act_runner daemon --config`.
log:
level: info
runner:
file: .runner
# One concurrent job. With one runner per project on a 2-vCPU VM this keeps
# total CI usage at ~2 threads.
capacity: 1
timeout: 3h
# Docker executor: no ":host" suffix. Maps runs-on labels to base images
# (the workflow overrides these per job via `container:`).
labels:
- "ubuntu-latest:docker://catthehacker/ubuntu:act-latest"
cache:
enabled: true
container:
# Cap every job container's CPU so CI stays within the 2-thread budget even
# if capacity is raised later. The -v mount persists the sccache cache on the
# host (create it first: `mkdir -p /var/cache/sccache`), matching SCCACHE_DIR
# in the workflow.
options: "--cpus=2 -v /var/cache/sccache:/sccache"
# Reuse the host VM's Docker network for the built-in cache/artifact server.
network: "host"
+41
View File
@@ -0,0 +1,41 @@
# SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
#
# SPDX-License-Identifier: GPL-2.0-or-later
#
# Rootless Podman Quadlet for the trx-rs Gitea Actions runner.
# Install to ~/.config/containers/systemd/trx-rs-runner.container then:
# systemctl --user daemon-reload
# systemctl --user start trx-rs-runner
#
# First boot only: paste a registration token (repo Settings -> Actions ->
# Runners) into GITEA_RUNNER_REGISTRATION_TOKEN. After the runner appears
# online you can blank it again — the registration is persisted in the volume.
[Unit]
Description=Gitea Actions runner — trx-rs
After=network-online.target
Wants=network-online.target
[Container]
Image=localhost/trx-rs-ci:latest
ContainerName=gitea-runner-trx-rs
# Persistent state: .runner registration, cache, workspaces.
Volume=gitea-runner-trx-rs-data:/data:Z
Environment=CONFIG_FILE=/data/config.yaml
Environment=GITEA_INSTANCE_URL=https://git.haxx.space
Environment=GITEA_RUNNER_NAME=trx-rs-podman
Environment=GITEA_RUNNER_LABELS=ubuntu-latest:host
Environment=GITEA_RUNNER_REGISTRATION_TOKEN=
# Resource caps so a heavy Rust build here cannot starve the other project's
# runner on the same host. Tune to your box.
PodmanArgs=--cpus=4.0 --memory=6g
[Service]
Restart=always
# A cold Rust build can be slow; don't let systemd consider startup failed.
TimeoutStartSec=0
[Install]
WantedBy=default.target
-11
View File
@@ -1,11 +0,0 @@
# SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
#
# SPDX-License-Identifier: GPL-2.0-or-later
#
# Pins the Rust toolchain for reproducible builds. Keep in sync with the
# SDK image (container/Containerfile, ARG RUST_VERSION). rustup honours this
# automatically for local builds outside the SDK container.
[toolchain]
channel = "1.97.1"
components = ["rustfmt", "clippy"]
+2 -1
View File
@@ -274,7 +274,8 @@ mod tests {
// Pseudo-random noise vs gradient — correlation should be low.
let noise: Vec<u8> = (0..256)
.map(|i| ((i * 1103515245 + 12345) as u32 >> 8 & 0xff) as u8)
.map(|i| (i as u32).wrapping_mul(1_103_515_245).wrapping_add(12_345))
.map(|value| ((value >> 8) & 0xff) as u8)
.collect();
let r = asm.correlation_with_last(&noise).expect("r");
assert!(
@@ -645,7 +645,7 @@ function flushDeferredDecodeMapSync() {
if (!decodeMapSyncPending || decodeHistoryReplayActive || !window.trx?.map?.aprsMap) return;
decodeMapSyncPending = false;
scheduleUiFrameJob("decode-map-maintenance", () => {
window.trx.map?.pruneMapHistory();
window.trx.modules.map?.pruneMapHistory();
});
}
@@ -1279,7 +1279,7 @@ function applyRigList(activeRigId, rigIds, displayNames) {
if (typeof bmPopulateScopePicker === "function") bmPopulateScopePicker();
if (typeof bmFetch === "function") bmFetch(document.getElementById("bm-category-filter")?.value || "");
}
window.trx.map?.updateMapRigFilter();
window.trx.modules.map?.updateMapRigFilter();
}
@@ -1305,7 +1305,7 @@ async function refreshRigList() {
serverRigs = rigs;
serverActiveRigId = data.active_remote || null;
applyRigList(data.active_remote, rigIds, displayNames);
window.trx.map?.syncAprsReceiverMarker();
window.trx.modules.map?.syncAprsReceiverMarker();
} catch (e) {
// Non-fatal: SSE/status path still drives main UI.
}
@@ -3143,9 +3143,9 @@ function render(update) {
const grid = latLonToMaidenhead(serverLat, serverLon);
locationSubtitle.textContent = `Location: ${grid}`;
locationSubtitle.style.display = "";
window.trx.map?.reverseGeocodeLocation(serverLat, serverLon, grid);
window.trx.modules.map?.reverseGeocodeLocation(serverLat, serverLon, grid);
}
window.trx.map?.syncAprsReceiverMarker();
window.trx.modules.map?.syncAprsReceiverMarker();
if (typeof update.initial_map_zoom === "number" && Number.isFinite(update.initial_map_zoom)) {
initialMapZoom = Math.max(1, Math.round(update.initial_map_zoom));
}
@@ -3894,7 +3894,7 @@ async function switchRigFromSelect(selectEl) {
if (typeof setSchedulerRig === "function") setSchedulerRig(lastActiveRigId);
if (typeof setBackgroundDecodeRig === "function") setBackgroundDecodeRig(lastActiveRigId);
if (typeof bmFetch === "function") bmFetch(document.getElementById("bm-category-filter")?.value || "");
window.trx.map?.syncAprsReceiverMarker();
window.trx.modules.map?.syncAprsReceiverMarker();
// Switch this session's rig and reconnect SSE to the new rig's
// state channel.
try {
@@ -4480,23 +4480,23 @@ function updateTabHistory(name, replaceHistory = false) {
}
// Initialise the Leaflet map, waiting for both Leaflet (L) and map-core.js
// (window.trx.map) if they haven't loaded yet.
// (window.trx.modules.map) if they haven't loaded yet.
let _mapInitTimer = null;
function _initMapWhenReady() {
const loadingEl = document.getElementById("map-loading");
if (window.trx.map && typeof L !== "undefined") {
if (window.trx.modules.map && typeof L !== "undefined") {
if (_mapInitTimer) { clearInterval(_mapInitTimer); _mapInitTimer = null; }
if (loadingEl) loadingEl.classList.add("is-hidden");
window.trx.map.initAprsMap();
window.trx.map.sizeAprsMapToViewport();
window.trx.modules.map.initAprsMap();
window.trx.modules.map.sizeAprsMapToViewport();
// The map panel was just made visible (display:none → ""); the browser
// may not have laid it out yet, so getBoundingClientRect() can return
// stale/zero dimensions. Double-rAF ensures a full layout pass has
// completed before we re-measure and tell Leaflet about its real size.
requestAnimationFrame(() => {
requestAnimationFrame(() => {
window.trx.map.sizeAprsMapToViewport();
if (window.trx.map.aprsMap) window.trx.map.aprsMap.invalidateSize();
window.trx.modules.map.sizeAprsMapToViewport();
if (window.trx.modules.map.aprsMap) window.trx.modules.map.aprsMap.invalidateSize();
});
});
return;
@@ -4544,7 +4544,7 @@ function navigateToTab(name, options = {}) {
_initMapWhenReady();
}
if (name === "statistics") {
window.trx.map?.scheduleStatsRender();
window.trx.modules.map?.scheduleStatsRender();
}
if (name === "recorder") {
refreshRecorderStatus();
@@ -4721,10 +4721,11 @@ if (headerAuthBtn) {
// ── Shared namespace for lazy-loaded modules ────────────────────────────────
// Modules (map-core.js, screenshot.js) access core state and utilities via
// window.trx. Modules register their own APIs as sub-namespaces
// (e.g. window.trx.map, window.trx.screenshot).
window.trx = Object.create(null);
// (e.g. window.trx.modules.map, window.trx.modules.screenshot).
const trxState = Object.create(null);
const trxModules = Object.create(null);
// -- State getters (backed by core-scoped variables) --
Object.defineProperties(window.trx, {
Object.defineProperties(trxState, {
serverLat: { get() { return serverLat; }, set(v) { serverLat = v; } },
serverLon: { get() { return serverLon; }, set(v) { serverLon = v; } },
lastFreqHz: { get() { return lastFreqHz; } },
@@ -4762,7 +4763,7 @@ Object.defineProperties(window.trx, {
signalOverlayGl: { get() { return signalOverlayGl; } },
});
// -- Shared utility functions --
Object.assign(window.trx, {
const trxCore = Object.freeze({
saveSetting, loadSetting, showHint, escapeMapHtml, formatFreq, formatFreqForHumans,
formatWavelength, formatBwLabel, formatUptime, formatSigStrength, formatSignal,
postPath, scheduleUiFrameJob, navigateToTab, rigBadgeColor,
@@ -4772,18 +4773,19 @@ Object.assign(window.trx, {
currentTheme, canvasPalette, currentStyle,
cssColorToRgba, rgbaWithAlpha, isBinsArray, estimateNoiseFloorDb,
spectrumVisibleRange, drawSpectrum,
bandForHz: function(hz) { return window.trx.map?.bandForHz?.(hz); },
bandForHz: function(hz) { return trxModules.map?.bandForHz?.(hz); },
markDecodeMapSyncPending,
decodeHistoryMapRenderingDeferred,
updateDocumentTitle,
activeChannelRds,
});
Object.defineProperties(window.trx, {
Object.defineProperties(trxState, {
decodeHistoryReplayActive: { get() { return decodeHistoryReplayActive; } },
decodeMapSyncPending: { get() { return decodeMapSyncPending; } },
_activeTab: { get() { return _activeTab; } },
locationSubtitle: { get() { return locationSubtitle; } },
});
window.trx = Object.freeze({ state: trxState, core: trxCore, modules: trxModules });
// Load plugin scripts now that window.trx is populated. Dynamic scripts are
// async so they must not be created before the namespace they depend on exists.
@@ -4797,7 +4799,7 @@ window.addEventListener("resize", resizeHeaderSignalCanvas);
// ── Map module (extracted to map-core.js, lazy-loaded) ──────────────────────
// The map, statistics, and geolocation code (~3,450 lines) has been moved to
// map-core.js and is loaded on demand when the Map tab is first activated.
// Core communicates with the map module via window.trx.map.* namespace.
// Core communicates with the map module via window.trx.modules.map.* namespace.
// ── Geo utilities (shared with map-core.js via window.trx) ─────────────────
function haversineKm(lat1, lon1, lat2, lon2) {
@@ -4932,7 +4934,7 @@ document.querySelectorAll(".sub-tab-bar").forEach(_wireSubTabBar);
window.addEventListener("resize", () => {
const mapTab = document.getElementById("tab-map");
if (!mapTab || mapTab.style.display === "none") return;
window.trx.map?.sizeAprsMapToViewport();
window.trx.modules.map?.sizeAprsMapToViewport();
});
// --- Signal measurement ---
@@ -6102,8 +6104,8 @@ function dispatchDecodeMessage(msg, skipStats) {
if (msg.type === "wefax" && window.onServerWefax) window.onServerWefax(msg);
if (msg.type === "wefax_progress" && window.onServerWefaxProgress) window.onServerWefaxProgress(msg);
if (!skipStats && msg.type && msg.type !== "lrpt_image" && msg.type !== "lrpt_progress" && msg.type !== "wefax" && msg.type !== "wefax_progress") {
window.trx.map?.statsRecordDecode(msg.type, msg.rig_id || msg.remote || null);
window.trx.map?.scheduleStatsRender();
window.trx.modules.map?.statsRecordDecode(msg.type, msg.rig_id || msg.remote || null);
window.trx.modules.map?.scheduleStatsRender();
}
}
@@ -6112,10 +6114,10 @@ function dispatchDecodeBatch(batch) {
// Record statistics for every message in the batch regardless of dispatch path.
for (const msg of batch) {
if (msg.type && msg.type !== "lrpt_image" && msg.type !== "lrpt_progress" && msg.type !== "wefax" && msg.type !== "wefax_progress") {
window.trx.map?.statsRecordDecode(msg.type, msg.rig_id || msg.remote || null);
window.trx.modules.map?.statsRecordDecode(msg.type, msg.rig_id || msg.remote || null);
}
}
window.trx.map?.scheduleStatsRender();
window.trx.modules.map?.scheduleStatsRender();
const type = String(batch[0]?.type || "");
const uniformType = batch.every((msg) => String(msg?.type || "") === type);
if (uniformType) {
@@ -6200,9 +6202,9 @@ function restoreDecodeHistoryGroup(kind, messages) {
// Record statistics for restored history messages.
if (kind !== "lrpt_image" && kind !== "lrpt_progress" && kind !== "wefax" && kind !== "wefax_progress") {
for (const msg of messages) {
window.trx.map?.statsRecordDecode(kind, msg.rig_id || msg.remote || null, msg.ts_ms || undefined);
window.trx.modules.map?.statsRecordDecode(kind, msg.rig_id || msg.remote || null, msg.ts_ms || undefined);
}
window.trx.map?.scheduleStatsRender();
window.trx.modules.map?.scheduleStatsRender();
}
if (kind === "ais") {
if (window.restoreAisHistory) { window.restoreAisHistory(messages); }
@@ -8000,12 +8002,12 @@ window.addEventListener("keydown", (event) => {
// S — spectrum screenshot (lazy-loads screenshot.js on first use)
if (key === "s") {
event.preventDefault();
if (window.trx.screenshot) {
void window.trx.screenshot.captureSpectrumScreenshot();
if (window.trx.modules.screenshot) {
void window.trx.modules.screenshot.captureSpectrumScreenshot();
} else {
const s = document.createElement("script");
s.src = "/screenshot.js";
s.onload = () => { void window.trx.screenshot?.captureSpectrumScreenshot(); };
s.onload = () => { void window.trx.modules.screenshot?.captureSpectrumScreenshot(); };
document.body.appendChild(s);
}
return;
@@ -1623,7 +1623,9 @@ SPDX-License-Identifier: GPL-2.0-or-later
<div id="decode-history-overlay-sub" class="decode-history-overlay-sub">Preparing recent decodes for the UI</div>
</div>
</div>
<script defer src="https://cdn.jsdelivr.net/npm/opus-decoder@0.7.11/dist/opus-decoder.min.js" charset="UTF-8"></script>
<script defer src="/vendor/opus-decoder-0.7.11.min.js" charset="UTF-8"></script>
<script defer src="/vendor/leaflet.js"></script>
<script defer src="/leaflet-ais-tracksymbol.js"></script>
<script defer src="/webgl-renderer.js"></script>
<script defer src="/app.js"></script>
<script>
@@ -1632,41 +1634,67 @@ SPDX-License-Identifier: GPL-2.0-or-later
var pluginScripts = {
'digital-modes': ['/ft8.js', '/ft4.js', '/ft2.js', '/wspr.js', '/cw.js', '/background-decode.js', '/sat.js', '/wefax.js'],
'map-data': ['/map-core.js', '/ais.js', '/vdes.js', '/aprs.js', '/hf-aprs.js'],
'map': ['/map-core.js', '/leaflet-ais-tracksymbol.js', '/ais.js', '/vdes.js', '/aprs.js', '/hf-aprs.js', '/sat.js', '/sat-scheduler.js'],
'map': ['/map-core.js', '/ais.js', '/vdes.js', '/aprs.js', '/hf-aprs.js', '/sat.js', '/sat-scheduler.js'],
'statistics': ['/map-core.js'],
'bookmarks': ['/bookmarks.js'],
'recorder': [],
'settings': ['/vchan.js', '/scheduler.js']
};
var loaded = new Set();
function loadPlugins(tab) {
var scripts = pluginScripts[tab];
if (!scripts) return;
scripts.forEach(function(src) {
if (loaded.has(src)) return;
loaded.add(src);
var loading = new Map();
function loadScript(src) {
if (loaded.has(src)) return Promise.resolve();
if (loading.has(src)) return loading.get(src);
var request = new Promise(function(resolve, reject) {
var s = document.createElement('script');
s.src = src;
s.defer = true;
s.onload = function() {
loaded.add(src);
loading.delete(src);
resolve();
};
s.onerror = function() {
loading.delete(src);
reject(new Error('Failed to load plugin script: ' + src));
};
document.body.appendChild(s);
});
loading.set(src, request);
return request;
}
function loadPlugins(tab) {
var scripts = pluginScripts[tab];
if (!scripts) return Promise.resolve();
return scripts.reduce(function(sequence, src) {
return sequence.then(function() { return loadScript(src); });
}, Promise.resolve());
}
function requestPlugins(tab) {
return loadPlugins(tab).catch(function(err) {
console.error(err);
});
}
// Eager plugin loading is triggered by app.js (after window.trx is set up)
// via window.loadEagerPlugins(). Dynamic scripts are effectively async, so
// loading them before app.js would cause map-core.js to crash when
// window.trx is not yet defined.
window.loadEagerPlugins = function() {
['digital-modes', 'map-data', 'bookmarks', 'settings'].forEach(loadPlugins);
return Promise.all(
['digital-modes', 'map-data', 'bookmarks', 'settings'].map(requestPlugins)
);
};
// Load others on tab switch
document.addEventListener('click', function(e) {
var tab = e.target.closest('[data-tab]');
if (tab) loadPlugins(tab.dataset.tab);
if (tab) requestPlugins(tab.dataset.tab);
});
window.loadPluginsForTab = loadPlugins;
})();
</script>
<!-- Template cloning is handled by navigateToTab() in app.js -->
<script defer src="/vendor/leaflet.js"></script>
</body>
</html>
@@ -3,10 +3,10 @@
// SPDX-License-Identifier: GPL-2.0-or-later
// Map, statistics, and geolocation module (lazy-loaded on map tab activation).
// Communicates with app.js core via window.trx namespace.
// Communicates with app.js through explicit state/core/module services.
(function () {
"use strict";
const T = window.trx;
const { state: T, core: C, modules } = window.trx;
// Destructure shared utility functions for convenience
const { saveSetting, loadSetting, showHint, escapeMapHtml, formatFreq, formatFreqForHumans,
@@ -14,7 +14,7 @@
formatUptime, latLonToMaidenhead, locatorToLatLon, haversineKm,
formatDistanceKm, formatTimeAgo, currentDecodeHistoryRetentionMs,
formatWavelength, bookmarkDistanceText, buildBookmarkTooltipText,
nearestBookmarkForHz } = T;
nearestBookmarkForHz } = C;
function updateMapRigFilter() {
const el = document.getElementById("map-rig-filter");
@@ -261,7 +261,7 @@
if (canRenderMap) {
refreshAprsTrack(call, entry);
} else {
T.markDecodeMapSyncPending();
C.markDecodeMapSyncPending();
}
if (!visible) {
if (canRenderMap && selectedAprsTrackCall && String(selectedAprsTrackCall) === String(call)) {
@@ -294,7 +294,7 @@
if (canRenderMap) {
refreshAisTrack(key, entry);
} else {
T.markDecodeMapSyncPending();
C.markDecodeMapSyncPending();
}
if (!visible) {
if (canRenderMap && selectedAisTrackMmsi && String(selectedAisTrackMmsi) === String(key)) {
@@ -337,7 +337,7 @@
entry.stations = new Set();
entry.bandMeta = new Map();
if (canRenderMap) setRetainedMapMarkerVisible(entry.marker, false);
else T.markDecodeMapSyncPending();
else C.markDecodeMapSyncPending();
return false;
}
const nextStations = new Set();
@@ -352,7 +352,7 @@
);
const count = Math.max(nextDetails.size, nextStations.size || 0, 1);
if (!canRenderMap) {
T.markDecodeMapSyncPending();
C.markDecodeMapSyncPending();
return true;
}
ensureDecodeLocatorMarker(entry);
@@ -392,7 +392,7 @@
pruneLocatorEntry(key, entry, cutoffMs);
}
if (!aprsMap || T.decodeHistoryReplayActive) {
T.markDecodeMapSyncPending();
C.markDecodeMapSyncPending();
return;
}
rebuildDecodeContactPaths();
@@ -415,7 +415,7 @@
function locatorFilterColor(type) {
const hues = locatorThemeHues();
const lightTheme = T.currentTheme() === "light";
const lightTheme = C.currentTheme() === "light";
const sat = lightTheme ? 66 : 76;
const light = lightTheme ? 42 : 56;
const hue = type === "bookmark"
@@ -539,7 +539,7 @@
}
function locatorThemeHues() {
const pal = T.canvasPalette();
const pal = C.canvasPalette();
const baseHue = paletteHue(pal?.spectrumLine, 145);
const waveHue = paletteHue(pal?.waveformLine, baseHue + 34);
const peakHue = paletteHue(pal?.waveformPeak, baseHue - 42);
@@ -560,7 +560,7 @@
function locatorBandChipColor(label) {
const hues = locatorThemeHues();
const lightTheme = T.currentTheme() === "light";
const lightTheme = C.currentTheme() === "light";
const hue = wrapHue(hues.bandBase + locatorBandIndex(label) * 137.508);
const sat = lightTheme ? 68 : 78;
const light = lightTheme ? 44 : 58;
@@ -606,7 +606,7 @@
const safeCount = Math.max(1, Number.isFinite(count) ? count : 1);
const intensity = Math.min(1, Math.log2(safeCount + 1) / 5);
const hue = locatorHueForEntry(entry);
const lightTheme = T.currentTheme() === "light";
const lightTheme = C.currentTheme() === "light";
const strokeSat = lightTheme ? 62 : 74;
const fillSat = lightTheme ? 68 : 78;
const strokeLight = lightTheme ? 40 : 56;
@@ -1573,7 +1573,7 @@
stageResizeObserver = new ResizeObserver(() => sizeAprsMapToViewport());
stageResizeObserver.observe(stage);
}
updateMapBaseLayerForTheme(T.currentTheme());
updateMapBaseLayerForTheme(C.currentTheme());
syncAprsReceiverMarker();
// Rebuild popup content on open (keeps age/distance/rig list fresh)
@@ -2307,18 +2307,18 @@
existing.rigIds.add(msgRigId);
}
if (!visible) {
if (!T.decodeHistoryMapRenderingDeferred()) {
if (!C.decodeHistoryMapRenderingDeferred()) {
setRetainedMapMarkerVisible(existing.marker, false);
} else {
T.markDecodeMapSyncPending();
C.markDecodeMapSyncPending();
}
return;
}
if (!T.decodeHistoryMapRenderingDeferred()) {
if (!C.decodeHistoryMapRenderingDeferred()) {
ensureVdesMarker(key, existing);
setRetainedMapMarkerVisible(existing.marker, true);
} else {
T.markDecodeMapSyncPending();
C.markDecodeMapSyncPending();
}
if (aprsMap && existing.marker && !T.decodeHistoryReplayActive) {
existing.marker.setLatLng([msg.lat, msg.lon]);
@@ -2334,11 +2334,11 @@
};
vdesMarkers.set(key, entry);
if (!visible) return;
if (!T.decodeHistoryMapRenderingDeferred()) {
if (!C.decodeHistoryMapRenderingDeferred()) {
ensureVdesMarker(key, entry);
setRetainedMapMarkerVisible(entry.marker, true);
} else {
T.markDecodeMapSyncPending();
C.markDecodeMapSyncPending();
}
if (aprsMap && entry.marker && !T.decodeHistoryReplayActive) {
entry.marker.setPopupContent(popupHtml);
@@ -2365,7 +2365,7 @@
if (T.locationSubtitle) {
T.locationSubtitle.textContent = `Location: ${grid} · ${label}`;
}
T.updateDocumentTitle(T.activeChannelRds());
C.updateDocumentTitle(C.activeChannelRds());
})
.catch(() => {});
}
@@ -2443,8 +2443,8 @@
}
function scheduleDecodeMapMaintenance() {
if (T.decodeHistoryMapRenderingDeferred()) {
T.markDecodeMapSyncPending();
if (C.decodeHistoryMapRenderingDeferred()) {
C.markDecodeMapSyncPending();
return;
}
scheduleUiFrameJob("decode-map-maintenance", () => {
@@ -3461,7 +3461,7 @@
}
// Register module API for core to call
window.trx.map = {
modules.map = {
initAprsMap,
sizeAprsMapToViewport,
syncAprsReceiverMarker,
@@ -257,7 +257,7 @@
}
// Register module API
window.trx.screenshot = {
window.trx.modules.screenshot = {
captureSpectrumScreenshot,
buildSpectrumSnapshotCanvas,
saveCanvasAsPng,
File diff suppressed because one or more lines are too long
@@ -75,6 +75,13 @@ define_gz_cache!(gz_bandplan_json, status::BANDPLAN_JSON, "bandplan.json");
// Vendored DSEG14 Classic font
// (binary woff2 — served directly, not through gz_cache)
// Vendored opus-decoder 0.7.11
define_gz_cache!(
gz_opus_decoder_js,
status::OPUS_DECODER_JS,
"opus-decoder-0.7.11.min.js"
);
// Vendored Leaflet 1.9.4
define_gz_cache!(gz_leaflet_js, status::LEAFLET_JS, "leaflet.js");
define_gz_cache!(gz_leaflet_css, status::LEAFLET_CSS, "leaflet.css");
@@ -341,6 +348,12 @@ pub(crate) async fn dseg14_classic_woff2() -> impl Responder {
.body(status::DSEG14_CLASSIC_WOFF2)
}
#[get("/vendor/opus-decoder-0.7.11.min.js")]
pub(crate) async fn opus_decoder_js(req: HttpRequest) -> impl Responder {
let c = gz_opus_decoder_js();
static_asset_response(&req, "application/javascript; charset=utf-8", c)
}
// ---------------------------------------------------------------------------
// Vendored Leaflet 1.9.4
// ---------------------------------------------------------------------------
@@ -668,6 +668,8 @@ pub fn configure(cfg: &mut web::ServiceConfig) {
.service(assets::bandplan_json)
// Vendored DSEG14 Classic font
.service(assets::dseg14_classic_woff2)
// Vendored opus-decoder 0.7.11
.service(assets::opus_decoder_js)
// Vendored Leaflet 1.9.4
.service(assets::leaflet_js)
.service(assets::leaflet_css)
@@ -40,6 +40,9 @@ pub const BANDPLAN_JSON: &str = include_str!("../assets/web/bandplan.json");
pub const DSEG14_CLASSIC_WOFF2: &[u8] =
include_bytes!("../assets/web/vendor/dseg14-classic-latin-400-normal.woff2");
// Vendored opus-decoder 0.7.11 browser build (WebAssembly embedded in JS)
pub const OPUS_DECODER_JS: &str = include_str!("../assets/web/vendor/opus-decoder-0.7.11.min.js");
// Vendored Leaflet 1.9.4
pub const LEAFLET_JS: &str = include_str!("../assets/web/vendor/leaflet.js");
pub const LEAFLET_CSS: &str = include_str!("../assets/web/vendor/leaflet.css");
+11 -513
View File
@@ -6,13 +6,14 @@
#[cfg(feature = "ft2")]
use std::collections::HashMap;
use std::collections::{HashSet, VecDeque};
use std::collections::HashSet;
#[cfg(test)]
use std::collections::VecDeque;
use std::net::SocketAddr;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use base64::Engine as _;
use bytes::Bytes;
use flate2::write::GzEncoder;
use flate2::Compression;
@@ -34,10 +35,9 @@ use trx_core::audio::{
AUDIO_MSG_VCHAN_MODE, AUDIO_MSG_VCHAN_REMOVE, AUDIO_MSG_VCHAN_SUB, AUDIO_MSG_VCHAN_UNSUB,
AUDIO_MSG_VDES_DECODE, AUDIO_MSG_WEFAX_DECODE, AUDIO_MSG_WEFAX_PROGRESS, AUDIO_MSG_WSPR_DECODE,
};
use trx_core::decode::{
AisMessage, AprsPacket, CwEvent, DecodedMessage, Ft8Message, LrptImage, LrptProgress,
VdesMessage, WefaxMessage, WsprMessage,
};
#[cfg(test)]
use trx_core::decode::{AisMessage, AprsPacket, CwEvent};
use trx_core::decode::{DecodedMessage, Ft8Message, LrptImage, LrptProgress, WsprMessage};
use trx_core::rig::state::{RigMode, RigState};
use trx_core::vchan::SharedVChanManager;
use trx_cw::CwDecoder;
@@ -47,21 +47,11 @@ use trx_wspr::WsprDecoder;
use uuid::Uuid;
use crate::config::AudioConfig;
use crate::history_policy::{current_timestamp_ms, lock_or_recover};
#[cfg(test)]
use crate::history_policy::{enforce_capacity, prune_by_age, MAX_HISTORY_ENTRIES};
use trx_decode_log::DecoderLoggers;
const APRS_HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60);
const HF_APRS_HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60);
const AIS_HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60);
const VDES_HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60);
const CW_HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60);
const FT8_HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60);
const WSPR_HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60);
const LRPT_HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60);
const WEFAX_HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60);
/// Maximum entries per decoder history queue. Prevents unbounded memory growth
/// on busy channels (e.g. AIS near a port). Oldest entries are evicted when
/// the limit is reached, independent of the time-based pruning.
const MAX_HISTORY_ENTRIES: usize = 10_000;
/// Silence timeout before auto-finalising an LRPT pass (30 s without new MCUs).
const LRPT_PASS_SILENCE_TIMEOUT: Duration = Duration::from_secs(30);
const FT8_SAMPLE_RATE: u32 = 12_000;
@@ -75,13 +65,6 @@ const DECODE_AUDIO_GATE_RMS: f32 = 2.5e-4;
const AUDIO_STREAM_ERROR_LOG_INTERVAL: Duration = Duration::from_secs(60);
const AUDIO_STREAM_RECOVERY_DELAY: Duration = Duration::from_secs(1);
fn current_timestamp_ms() -> i64 {
match std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH) {
Ok(dur) => dur.as_millis() as i64,
Err(_) => 0,
}
}
#[cfg(feature = "ft2")]
fn retain_ft2_window(buf: &mut Vec<f32>) {
if buf.len() > FT2_ASYNC_BUFFER_SAMPLES {
@@ -360,492 +343,7 @@ fn classify_stream_error(err: &str) -> &'static str {
}
}
/// Per-rig decoder history store.
///
/// Replaces the previous process-wide `OnceLock` statics so that each rig
/// instance can maintain its own independent history. Pass an
/// `Arc<DecoderHistories>` into every decoder task and into the audio listener.
pub struct DecoderHistories {
pub ais: Mutex<VecDeque<(Instant, AisMessage)>>,
pub vdes: Mutex<VecDeque<(Instant, VdesMessage)>>,
pub aprs: Mutex<VecDeque<(Instant, AprsPacket)>>,
pub hf_aprs: Mutex<VecDeque<(Instant, AprsPacket)>>,
pub cw: Mutex<VecDeque<(Instant, CwEvent)>>,
pub ft8: Mutex<VecDeque<(Instant, Ft8Message)>>,
pub ft4: Mutex<VecDeque<(Instant, Ft8Message)>>,
pub ft2: Mutex<VecDeque<(Instant, Ft8Message)>>,
pub wspr: Mutex<VecDeque<(Instant, WsprMessage)>>,
pub lrpt: Mutex<VecDeque<(Instant, LrptImage)>>,
pub wefax: Mutex<VecDeque<(Instant, WefaxMessage)>>,
/// Approximate total entry count across all decoders, maintained
/// atomically so `estimated_total_count()` avoids 9 lock acquisitions.
total_count: AtomicUsize,
}
/// Acquire a mutex, recovering from poisoning with a warning log.
fn lock_or_recover<'a, T>(mutex: &'a Mutex<T>, label: &str) -> std::sync::MutexGuard<'a, T> {
mutex.lock().unwrap_or_else(|e| {
tracing::warn!(
"Mutex for {} was poisoned (prior panic); recovering with potentially inconsistent data",
label
);
e.into_inner()
})
}
/// Enforce capacity limit on a history deque by evicting oldest entries.
fn enforce_capacity<T>(deque: &mut VecDeque<T>, max: usize) {
while deque.len() > max {
deque.pop_front();
}
}
/// Drop entries older than `retention` from the front of a time-tagged deque.
/// Uses `checked_sub` so an early `now` (before `retention` elapses since the
/// monotonic clock origin) is treated as "no entries are old enough to prune"
/// rather than panicking.
fn prune_by_age<T>(deque: &mut VecDeque<(Instant, T)>, retention: Duration, now: Instant) {
let Some(cutoff) = now.checked_sub(retention) else {
return;
};
while let Some((ts, _)) = deque.front() {
if *ts < cutoff {
deque.pop_front();
} else {
break;
}
}
}
impl DecoderHistories {
pub fn new() -> Arc<Self> {
Arc::new(Self {
ais: Mutex::new(VecDeque::new()),
vdes: Mutex::new(VecDeque::new()),
aprs: Mutex::new(VecDeque::new()),
hf_aprs: Mutex::new(VecDeque::new()),
cw: Mutex::new(VecDeque::new()),
ft8: Mutex::new(VecDeque::new()),
ft4: Mutex::new(VecDeque::new()),
ft2: Mutex::new(VecDeque::new()),
wspr: Mutex::new(VecDeque::new()),
lrpt: Mutex::new(VecDeque::new()),
wefax: Mutex::new(VecDeque::new()),
total_count: AtomicUsize::new(0),
})
}
/// Adjust the atomic total count after a record/prune/clear operation.
///
/// Uses a CAS loop for decrements to prevent underflow wrapping the
/// counter to `usize::MAX` (which would cause a capacity-overflow panic
/// when pre-allocating the history replay blob).
fn adjust_total_count(&self, old_len: usize, new_len: usize) {
if new_len > old_len {
self.total_count
.fetch_add(new_len - old_len, Ordering::Relaxed);
} else if old_len > new_len {
let delta = old_len - new_len;
let mut current = self.total_count.load(Ordering::Relaxed);
loop {
let next = current.saturating_sub(delta);
match self.total_count.compare_exchange_weak(
current,
next,
Ordering::Relaxed,
Ordering::Relaxed,
) {
Ok(_) => break,
Err(actual) => current = actual,
}
}
}
}
// --- AIS ---
fn prune_ais(history: &mut VecDeque<(Instant, AisMessage)>) {
prune_by_age(history, AIS_HISTORY_RETENTION, Instant::now());
}
pub fn record_ais_message(&self, mut msg: AisMessage) {
if msg.ts_ms.is_none() {
msg.ts_ms = Some(current_timestamp_ms());
}
let mut h = lock_or_recover(&self.ais, "ais_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_ais(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_ais_history(&self) -> Vec<AisMessage> {
let mut h = lock_or_recover(&self.ais, "ais_history");
let before = h.len();
Self::prune_ais(&mut h);
self.adjust_total_count(before, h.len());
h.iter().map(|(_, msg)| msg.clone()).collect()
}
// --- VDES ---
fn prune_vdes(history: &mut VecDeque<(Instant, VdesMessage)>) {
prune_by_age(history, VDES_HISTORY_RETENTION, Instant::now());
}
pub fn record_vdes_message(&self, mut msg: VdesMessage) {
if msg.ts_ms.is_none() {
msg.ts_ms = Some(current_timestamp_ms());
}
let mut h = lock_or_recover(&self.vdes, "vdes_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_vdes(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_vdes_history(&self) -> Vec<VdesMessage> {
let mut h = lock_or_recover(&self.vdes, "vdes_history");
let before = h.len();
Self::prune_vdes(&mut h);
self.adjust_total_count(before, h.len());
h.iter().map(|(_, msg)| msg.clone()).collect()
}
// --- APRS ---
fn prune_aprs(history: &mut VecDeque<(Instant, AprsPacket)>) {
prune_by_age(history, APRS_HISTORY_RETENTION, Instant::now());
}
pub fn record_aprs_packet(&self, mut pkt: AprsPacket) {
if !pkt.crc_ok {
return;
}
if pkt.ts_ms.is_none() {
pkt.ts_ms = Some(current_timestamp_ms());
}
let mut h = lock_or_recover(&self.aprs, "aprs_history");
let before = h.len();
h.push_back((Instant::now(), pkt));
Self::prune_aprs(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_aprs_history(&self) -> Vec<AprsPacket> {
let mut h = lock_or_recover(&self.aprs, "aprs_history");
let before = h.len();
Self::prune_aprs(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, pkt): &(Instant, AprsPacket)| pkt.clone())
.collect()
}
pub fn clear_aprs_history(&self) {
let mut h = lock_or_recover(&self.aprs, "aprs_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- HF APRS ---
fn prune_hf_aprs(history: &mut VecDeque<(Instant, AprsPacket)>) {
prune_by_age(history, HF_APRS_HISTORY_RETENTION, Instant::now());
}
pub fn record_hf_aprs_packet(&self, mut pkt: AprsPacket) {
if !pkt.crc_ok {
return;
}
if pkt.ts_ms.is_none() {
pkt.ts_ms = Some(current_timestamp_ms());
}
let mut h = lock_or_recover(&self.hf_aprs, "hf_aprs_history");
let before = h.len();
h.push_back((Instant::now(), pkt));
Self::prune_hf_aprs(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_hf_aprs_history(&self) -> Vec<AprsPacket> {
let mut h = lock_or_recover(&self.hf_aprs, "hf_aprs_history");
let before = h.len();
Self::prune_hf_aprs(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, pkt): &(Instant, AprsPacket)| pkt.clone())
.collect()
}
pub fn clear_hf_aprs_history(&self) {
let mut h = lock_or_recover(&self.hf_aprs, "hf_aprs_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- CW ---
fn prune_cw(history: &mut VecDeque<(Instant, CwEvent)>) {
prune_by_age(history, CW_HISTORY_RETENTION, Instant::now());
}
pub fn record_cw_event(&self, evt: CwEvent) {
let mut h = lock_or_recover(&self.cw, "cw_history");
let before = h.len();
h.push_back((Instant::now(), evt));
Self::prune_cw(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_cw_history(&self) -> Vec<CwEvent> {
let mut h = lock_or_recover(&self.cw, "cw_history");
let before = h.len();
Self::prune_cw(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, evt): &(Instant, CwEvent)| evt.clone())
.collect()
}
pub fn clear_cw_history(&self) {
let mut h = lock_or_recover(&self.cw, "cw_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- FT8 ---
fn prune_ft8(history: &mut VecDeque<(Instant, Ft8Message)>) {
prune_by_age(history, FT8_HISTORY_RETENTION, Instant::now());
}
pub fn record_ft8_message(&self, msg: Ft8Message) {
let mut h = lock_or_recover(&self.ft8, "ft8_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_ft8(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_ft8_history(&self) -> Vec<Ft8Message> {
let mut h = lock_or_recover(&self.ft8, "ft8_history");
let before = h.len();
Self::prune_ft8(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, msg): &(Instant, Ft8Message)| msg.clone())
.collect()
}
pub fn clear_ft8_history(&self) {
let mut h = lock_or_recover(&self.ft8, "ft8_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- FT4 ---
fn prune_ft4(history: &mut VecDeque<(Instant, Ft8Message)>) {
prune_by_age(history, FT8_HISTORY_RETENTION, Instant::now());
}
pub fn record_ft4_message(&self, msg: Ft8Message) {
let mut h = lock_or_recover(&self.ft4, "ft4_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_ft4(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_ft4_history(&self) -> Vec<Ft8Message> {
let mut h = lock_or_recover(&self.ft4, "ft4_history");
let before = h.len();
Self::prune_ft4(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, msg): &(Instant, Ft8Message)| msg.clone())
.collect()
}
pub fn clear_ft4_history(&self) {
let mut h = lock_or_recover(&self.ft4, "ft4_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- FT2 ---
#[cfg_attr(not(feature = "ft2"), allow(dead_code))]
fn prune_ft2(history: &mut VecDeque<(Instant, Ft8Message)>) {
prune_by_age(history, FT8_HISTORY_RETENTION, Instant::now());
}
#[cfg_attr(not(feature = "ft2"), allow(dead_code))]
pub fn record_ft2_message(&self, msg: Ft8Message) {
let mut h = lock_or_recover(&self.ft2, "ft2_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_ft2(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
#[cfg_attr(not(feature = "ft2"), allow(dead_code))]
pub fn snapshot_ft2_history(&self) -> Vec<Ft8Message> {
let mut h = lock_or_recover(&self.ft2, "ft2_history");
let before = h.len();
Self::prune_ft2(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, msg): &(Instant, Ft8Message)| msg.clone())
.collect()
}
pub fn clear_ft2_history(&self) {
let mut h = lock_or_recover(&self.ft2, "ft2_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- WSPR ---
fn prune_wspr(history: &mut VecDeque<(Instant, WsprMessage)>) {
prune_by_age(history, WSPR_HISTORY_RETENTION, Instant::now());
}
pub fn record_wspr_message(&self, msg: WsprMessage) {
let mut h = lock_or_recover(&self.wspr, "wspr_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_wspr(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_wspr_history(&self) -> Vec<WsprMessage> {
let mut h = lock_or_recover(&self.wspr, "wspr_history");
let before = h.len();
Self::prune_wspr(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, msg): &(Instant, WsprMessage)| msg.clone())
.collect()
}
pub fn clear_wspr_history(&self) {
let mut h = lock_or_recover(&self.wspr, "wspr_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- LRPT ---
fn prune_lrpt(history: &mut VecDeque<(Instant, LrptImage)>) {
prune_by_age(history, LRPT_HISTORY_RETENTION, Instant::now());
}
pub fn record_lrpt_image(&self, mut img: LrptImage) {
if img.ts_ms.is_none() {
img.ts_ms = Some(current_timestamp_ms());
}
let mut h = lock_or_recover(&self.lrpt, "lrpt_history");
let before = h.len();
h.push_back((Instant::now(), img));
Self::prune_lrpt(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_lrpt_history(&self) -> Vec<LrptImage> {
let mut h = lock_or_recover(&self.lrpt, "lrpt_history");
let before = h.len();
Self::prune_lrpt(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, img): &(Instant, LrptImage)| img.clone())
.collect()
}
pub fn clear_lrpt_history(&self) {
let mut h = lock_or_recover(&self.lrpt, "lrpt_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- WEFAX ---
fn prune_wefax(history: &mut VecDeque<(Instant, WefaxMessage)>) {
prune_by_age(history, WEFAX_HISTORY_RETENTION, Instant::now());
}
pub fn record_wefax_message(&self, mut msg: WefaxMessage) {
if msg.ts_ms.is_none() {
msg.ts_ms = Some(current_timestamp_ms());
}
// Strip bulk PNG data before storing in memory/persistence.
msg.png_data = None;
let mut h = lock_or_recover(&self.wefax, "wefax_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_wefax(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_wefax_history(&self) -> Vec<WefaxMessage> {
let mut h = lock_or_recover(&self.wefax, "wefax_history");
let before = h.len();
Self::prune_wefax(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, msg)| {
let mut m = msg.clone();
// Re-read PNG from disk so remote clients can save a local copy.
if m.png_data.is_none() {
if let Some(ref path) = m.path {
if let Ok(bytes) = std::fs::read(path) {
m.png_data =
Some(base64::engine::general_purpose::STANDARD.encode(&bytes));
}
}
}
m
})
.collect()
}
pub fn clear_wefax_history(&self) {
let mut h = lock_or_recover(&self.wefax, "wefax_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
/// Returns a quick (non-pruning) estimate of the total number of history
/// entries across all decoders, used for pre-allocating the replay blob.
///
/// Uses an `AtomicUsize` counter maintained by record/prune/clear methods,
/// avoiding 9 separate mutex acquisitions.
pub fn estimated_total_count(&self) -> usize {
self.total_count.load(Ordering::Relaxed)
}
}
pub use crate::decoder_history::DecoderHistories;
/// Spawn the audio capture thread.
///
+489
View File
@@ -0,0 +1,489 @@
// SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
//
// SPDX-License-Identifier: GPL-2.0-or-later
//! Per-rig storage and lifecycle operations for decoded-message histories.
use std::collections::VecDeque;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::time::Instant;
use base64::Engine as _;
use trx_core::decode::{
AisMessage, AprsPacket, CwEvent, Ft8Message, LrptImage, VdesMessage, WefaxMessage, WsprMessage,
};
use crate::history_policy::{
current_timestamp_ms, enforce_capacity, lock_or_recover, prune_by_age, HISTORY_RETENTION,
MAX_HISTORY_ENTRIES,
};
/// Per-rig decoder history store.
///
/// Replaces the previous process-wide `OnceLock` statics so that each rig
/// instance can maintain its own independent history. Pass an
/// `Arc<DecoderHistories>` into every decoder task and into the audio listener.
pub struct DecoderHistories {
pub ais: Mutex<VecDeque<(Instant, AisMessage)>>,
pub vdes: Mutex<VecDeque<(Instant, VdesMessage)>>,
pub aprs: Mutex<VecDeque<(Instant, AprsPacket)>>,
pub hf_aprs: Mutex<VecDeque<(Instant, AprsPacket)>>,
pub cw: Mutex<VecDeque<(Instant, CwEvent)>>,
pub ft8: Mutex<VecDeque<(Instant, Ft8Message)>>,
pub ft4: Mutex<VecDeque<(Instant, Ft8Message)>>,
pub ft2: Mutex<VecDeque<(Instant, Ft8Message)>>,
pub wspr: Mutex<VecDeque<(Instant, WsprMessage)>>,
pub lrpt: Mutex<VecDeque<(Instant, LrptImage)>>,
pub wefax: Mutex<VecDeque<(Instant, WefaxMessage)>>,
/// Approximate total entry count across all decoders, maintained
/// atomically so `estimated_total_count()` avoids 11 lock acquisitions.
total_count: AtomicUsize,
}
impl DecoderHistories {
pub fn new() -> Arc<Self> {
Arc::new(Self {
ais: Mutex::new(VecDeque::new()),
vdes: Mutex::new(VecDeque::new()),
aprs: Mutex::new(VecDeque::new()),
hf_aprs: Mutex::new(VecDeque::new()),
cw: Mutex::new(VecDeque::new()),
ft8: Mutex::new(VecDeque::new()),
ft4: Mutex::new(VecDeque::new()),
ft2: Mutex::new(VecDeque::new()),
wspr: Mutex::new(VecDeque::new()),
lrpt: Mutex::new(VecDeque::new()),
wefax: Mutex::new(VecDeque::new()),
total_count: AtomicUsize::new(0),
})
}
/// Adjust the atomic total count after a record/prune/clear operation.
///
/// Uses a CAS loop for decrements to prevent underflow wrapping the
/// counter to `usize::MAX` (which would cause a capacity-overflow panic
/// when pre-allocating the history replay blob).
pub(crate) fn adjust_total_count(&self, old_len: usize, new_len: usize) {
if new_len > old_len {
self.total_count
.fetch_add(new_len - old_len, Ordering::Relaxed);
} else if old_len > new_len {
let delta = old_len - new_len;
let mut current = self.total_count.load(Ordering::Relaxed);
loop {
let next = current.saturating_sub(delta);
match self.total_count.compare_exchange_weak(
current,
next,
Ordering::Relaxed,
Ordering::Relaxed,
) {
Ok(_) => break,
Err(actual) => current = actual,
}
}
}
}
// --- AIS ---
fn prune_ais(history: &mut VecDeque<(Instant, AisMessage)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
pub fn record_ais_message(&self, mut msg: AisMessage) {
if msg.ts_ms.is_none() {
msg.ts_ms = Some(current_timestamp_ms());
}
let mut h = lock_or_recover(&self.ais, "ais_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_ais(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_ais_history(&self) -> Vec<AisMessage> {
let mut h = lock_or_recover(&self.ais, "ais_history");
let before = h.len();
Self::prune_ais(&mut h);
self.adjust_total_count(before, h.len());
h.iter().map(|(_, msg)| msg.clone()).collect()
}
// --- VDES ---
fn prune_vdes(history: &mut VecDeque<(Instant, VdesMessage)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
pub fn record_vdes_message(&self, mut msg: VdesMessage) {
if msg.ts_ms.is_none() {
msg.ts_ms = Some(current_timestamp_ms());
}
let mut h = lock_or_recover(&self.vdes, "vdes_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_vdes(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_vdes_history(&self) -> Vec<VdesMessage> {
let mut h = lock_or_recover(&self.vdes, "vdes_history");
let before = h.len();
Self::prune_vdes(&mut h);
self.adjust_total_count(before, h.len());
h.iter().map(|(_, msg)| msg.clone()).collect()
}
// --- APRS ---
fn prune_aprs(history: &mut VecDeque<(Instant, AprsPacket)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
pub fn record_aprs_packet(&self, mut pkt: AprsPacket) {
if !pkt.crc_ok {
return;
}
if pkt.ts_ms.is_none() {
pkt.ts_ms = Some(current_timestamp_ms());
}
let mut h = lock_or_recover(&self.aprs, "aprs_history");
let before = h.len();
h.push_back((Instant::now(), pkt));
Self::prune_aprs(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_aprs_history(&self) -> Vec<AprsPacket> {
let mut h = lock_or_recover(&self.aprs, "aprs_history");
let before = h.len();
Self::prune_aprs(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, pkt): &(Instant, AprsPacket)| pkt.clone())
.collect()
}
pub fn clear_aprs_history(&self) {
let mut h = lock_or_recover(&self.aprs, "aprs_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- HF APRS ---
fn prune_hf_aprs(history: &mut VecDeque<(Instant, AprsPacket)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
pub fn record_hf_aprs_packet(&self, mut pkt: AprsPacket) {
if !pkt.crc_ok {
return;
}
if pkt.ts_ms.is_none() {
pkt.ts_ms = Some(current_timestamp_ms());
}
let mut h = lock_or_recover(&self.hf_aprs, "hf_aprs_history");
let before = h.len();
h.push_back((Instant::now(), pkt));
Self::prune_hf_aprs(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_hf_aprs_history(&self) -> Vec<AprsPacket> {
let mut h = lock_or_recover(&self.hf_aprs, "hf_aprs_history");
let before = h.len();
Self::prune_hf_aprs(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, pkt): &(Instant, AprsPacket)| pkt.clone())
.collect()
}
pub fn clear_hf_aprs_history(&self) {
let mut h = lock_or_recover(&self.hf_aprs, "hf_aprs_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- CW ---
fn prune_cw(history: &mut VecDeque<(Instant, CwEvent)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
pub fn record_cw_event(&self, evt: CwEvent) {
let mut h = lock_or_recover(&self.cw, "cw_history");
let before = h.len();
h.push_back((Instant::now(), evt));
Self::prune_cw(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_cw_history(&self) -> Vec<CwEvent> {
let mut h = lock_or_recover(&self.cw, "cw_history");
let before = h.len();
Self::prune_cw(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, evt): &(Instant, CwEvent)| evt.clone())
.collect()
}
pub fn clear_cw_history(&self) {
let mut h = lock_or_recover(&self.cw, "cw_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- FT8 ---
fn prune_ft8(history: &mut VecDeque<(Instant, Ft8Message)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
pub fn record_ft8_message(&self, msg: Ft8Message) {
let mut h = lock_or_recover(&self.ft8, "ft8_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_ft8(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_ft8_history(&self) -> Vec<Ft8Message> {
let mut h = lock_or_recover(&self.ft8, "ft8_history");
let before = h.len();
Self::prune_ft8(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, msg): &(Instant, Ft8Message)| msg.clone())
.collect()
}
pub fn clear_ft8_history(&self) {
let mut h = lock_or_recover(&self.ft8, "ft8_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- FT4 ---
fn prune_ft4(history: &mut VecDeque<(Instant, Ft8Message)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
pub fn record_ft4_message(&self, msg: Ft8Message) {
let mut h = lock_or_recover(&self.ft4, "ft4_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_ft4(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_ft4_history(&self) -> Vec<Ft8Message> {
let mut h = lock_or_recover(&self.ft4, "ft4_history");
let before = h.len();
Self::prune_ft4(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, msg): &(Instant, Ft8Message)| msg.clone())
.collect()
}
pub fn clear_ft4_history(&self) {
let mut h = lock_or_recover(&self.ft4, "ft4_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- FT2 ---
#[cfg_attr(not(feature = "ft2"), allow(dead_code))]
fn prune_ft2(history: &mut VecDeque<(Instant, Ft8Message)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
#[cfg_attr(not(feature = "ft2"), allow(dead_code))]
pub fn record_ft2_message(&self, msg: Ft8Message) {
let mut h = lock_or_recover(&self.ft2, "ft2_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_ft2(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
#[cfg_attr(not(feature = "ft2"), allow(dead_code))]
pub fn snapshot_ft2_history(&self) -> Vec<Ft8Message> {
let mut h = lock_or_recover(&self.ft2, "ft2_history");
let before = h.len();
Self::prune_ft2(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, msg): &(Instant, Ft8Message)| msg.clone())
.collect()
}
pub fn clear_ft2_history(&self) {
let mut h = lock_or_recover(&self.ft2, "ft2_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- WSPR ---
fn prune_wspr(history: &mut VecDeque<(Instant, WsprMessage)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
pub fn record_wspr_message(&self, msg: WsprMessage) {
let mut h = lock_or_recover(&self.wspr, "wspr_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_wspr(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_wspr_history(&self) -> Vec<WsprMessage> {
let mut h = lock_or_recover(&self.wspr, "wspr_history");
let before = h.len();
Self::prune_wspr(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, msg): &(Instant, WsprMessage)| msg.clone())
.collect()
}
pub fn clear_wspr_history(&self) {
let mut h = lock_or_recover(&self.wspr, "wspr_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- LRPT ---
fn prune_lrpt(history: &mut VecDeque<(Instant, LrptImage)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
pub fn record_lrpt_image(&self, mut img: LrptImage) {
if img.ts_ms.is_none() {
img.ts_ms = Some(current_timestamp_ms());
}
let mut h = lock_or_recover(&self.lrpt, "lrpt_history");
let before = h.len();
h.push_back((Instant::now(), img));
Self::prune_lrpt(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_lrpt_history(&self) -> Vec<LrptImage> {
let mut h = lock_or_recover(&self.lrpt, "lrpt_history");
let before = h.len();
Self::prune_lrpt(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, img): &(Instant, LrptImage)| img.clone())
.collect()
}
pub fn clear_lrpt_history(&self) {
let mut h = lock_or_recover(&self.lrpt, "lrpt_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
// --- WEFAX ---
fn prune_wefax(history: &mut VecDeque<(Instant, WefaxMessage)>) {
prune_by_age(history, HISTORY_RETENTION, Instant::now());
}
pub fn record_wefax_message(&self, mut msg: WefaxMessage) {
if msg.ts_ms.is_none() {
msg.ts_ms = Some(current_timestamp_ms());
}
// Strip bulk PNG data before storing in memory/persistence.
msg.png_data = None;
let mut h = lock_or_recover(&self.wefax, "wefax_history");
let before = h.len();
h.push_back((Instant::now(), msg));
Self::prune_wefax(&mut h);
enforce_capacity(&mut h, MAX_HISTORY_ENTRIES);
self.adjust_total_count(before, h.len());
}
pub fn snapshot_wefax_history(&self) -> Vec<WefaxMessage> {
let mut h = lock_or_recover(&self.wefax, "wefax_history");
let before = h.len();
Self::prune_wefax(&mut h);
self.adjust_total_count(before, h.len());
h.iter()
.map(|(_, msg)| {
let mut m = msg.clone();
// Re-read PNG from disk so remote clients can save a local copy.
if m.png_data.is_none() {
if let Some(ref path) = m.path {
if let Ok(bytes) = std::fs::read(path) {
m.png_data =
Some(base64::engine::general_purpose::STANDARD.encode(&bytes));
}
}
}
m
})
.collect()
}
pub fn clear_wefax_history(&self) {
let mut h = lock_or_recover(&self.wefax, "wefax_history");
let before = h.len();
h.clear();
self.adjust_total_count(before, 0);
}
/// Returns a quick (non-pruning) estimate of the total number of history
/// entries across all decoders, used for pre-allocating the replay blob.
///
/// Uses an `AtomicUsize` counter maintained by record/prune/clear methods,
/// avoiding 11 separate mutex acquisitions.
pub fn estimated_total_count(&self) -> usize {
self.total_count.load(Ordering::Relaxed)
}
/// Rebuild the aggregate count after bulk restoration bypasses the normal
/// record methods.
pub(crate) fn recalculate_total_count(&self) {
let total = lock_or_recover(&self.ais, "ais_history").len()
+ lock_or_recover(&self.vdes, "vdes_history").len()
+ lock_or_recover(&self.aprs, "aprs_history").len()
+ lock_or_recover(&self.hf_aprs, "hf_aprs_history").len()
+ lock_or_recover(&self.cw, "cw_history").len()
+ lock_or_recover(&self.ft8, "ft8_history").len()
+ lock_or_recover(&self.ft4, "ft4_history").len()
+ lock_or_recover(&self.ft2, "ft2_history").len()
+ lock_or_recover(&self.wspr, "wspr_history").len()
+ lock_or_recover(&self.lrpt, "lrpt_history").len()
+ lock_or_recover(&self.wefax, "wefax_history").len();
self.total_count.store(total, Ordering::Relaxed);
}
}
+56
View File
@@ -0,0 +1,56 @@
// SPDX-FileCopyrightText: 2026 Stan Grams <sjg@haxx.space>
//
// SPDX-License-Identifier: GPL-2.0-or-later
//! Shared retention and synchronization policy for decoder histories.
use std::collections::VecDeque;
use std::sync::{Mutex, MutexGuard};
use std::time::{Duration, Instant};
pub(crate) const HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60);
/// Maximum entries per decoder history queue. Oldest entries are evicted on
/// busy channels independently of time-based pruning.
pub(crate) const MAX_HISTORY_ENTRIES: usize = 10_000;
pub(crate) fn current_timestamp_ms() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|duration| duration.as_millis() as i64)
.unwrap_or(0)
}
pub(crate) fn lock_or_recover<'a, T>(mutex: &'a Mutex<T>, label: &str) -> MutexGuard<'a, T> {
mutex.lock().unwrap_or_else(|error| {
tracing::warn!(
"Mutex for {} was poisoned (prior panic); recovering with potentially inconsistent data",
label
);
error.into_inner()
})
}
pub(crate) fn enforce_capacity<T>(deque: &mut VecDeque<T>, max: usize) {
while deque.len() > max {
deque.pop_front();
}
}
/// Drop entries older than `retention` from an ordered history queue.
pub(crate) fn prune_by_age<T>(
deque: &mut VecDeque<(Instant, T)>,
retention: Duration,
now: Instant,
) {
let Some(cutoff) = now.checked_sub(retention) else {
return;
};
while let Some((timestamp, _)) = deque.front() {
if *timestamp < cutoff {
deque.pop_front();
} else {
break;
}
}
}
+139 -5
View File
@@ -19,7 +19,7 @@ use pickledb::{PickleDb, PickleDbDumpPolicy, SerializationMethod};
use serde::{de::DeserializeOwned, Deserialize, Serialize};
use trx_core::decode::{
AisMessage, AprsPacket, CwEvent, Ft8Message, VdesMessage, WefaxMessage, WsprMessage,
AisMessage, AprsPacket, CwEvent, Ft8Message, LrptImage, VdesMessage, WefaxMessage, WsprMessage,
};
use crate::audio::DecoderHistories;
@@ -118,6 +118,11 @@ pub fn load_all(db: &PickleDb, rig_id: &str, histories: &Arc<DecoderHistories>)
h.push_back(e);
}
}
if let Ok(mut h) = histories.hf_aprs.lock() {
for e in load_key::<AprsPacket>(db, &k("hf_aprs")) {
h.push_back(e);
}
}
if let Ok(mut h) = histories.cw.lock() {
for e in load_key::<CwEvent>(db, &k("cw")) {
h.push_back(e);
@@ -128,16 +133,32 @@ pub fn load_all(db: &PickleDb, rig_id: &str, histories: &Arc<DecoderHistories>)
h.push_back(e);
}
}
if let Ok(mut h) = histories.ft4.lock() {
for e in load_key::<Ft8Message>(db, &k("ft4")) {
h.push_back(e);
}
}
if let Ok(mut h) = histories.ft2.lock() {
for e in load_key::<Ft8Message>(db, &k("ft2")) {
h.push_back(e);
}
}
if let Ok(mut h) = histories.wspr.lock() {
for e in load_key::<WsprMessage>(db, &k("wspr")) {
h.push_back(e);
}
}
if let Ok(mut h) = histories.lrpt.lock() {
for e in load_key::<LrptImage>(db, &k("lrpt")) {
h.push_back(e);
}
}
if let Ok(mut h) = histories.wefax.lock() {
for e in load_key::<WefaxMessage>(db, &k("wefax")) {
h.push_back(e);
}
}
histories.recalculate_total_count();
}
/// Flush `histories` to the database under `rig_id`-prefixed keys and sync.
@@ -162,6 +183,11 @@ pub fn flush_all(db: &mut PickleDb, rig_id: &str, histories: &Arc<DecoderHistori
drop(h);
save_key(db, &k("aprs"), &snapshot);
}
if let Ok(h) = histories.hf_aprs.lock() {
let snapshot = h.clone();
drop(h);
save_key(db, &k("hf_aprs"), &snapshot);
}
if let Ok(h) = histories.cw.lock() {
let snapshot = h.clone();
drop(h);
@@ -172,11 +198,26 @@ pub fn flush_all(db: &mut PickleDb, rig_id: &str, histories: &Arc<DecoderHistori
drop(h);
save_key(db, &k("ft8"), &snapshot);
}
if let Ok(h) = histories.ft4.lock() {
let snapshot = h.clone();
drop(h);
save_key(db, &k("ft4"), &snapshot);
}
if let Ok(h) = histories.ft2.lock() {
let snapshot = h.clone();
drop(h);
save_key(db, &k("ft2"), &snapshot);
}
if let Ok(h) = histories.wspr.lock() {
let snapshot = h.clone();
drop(h);
save_key(db, &k("wspr"), &snapshot);
}
if let Ok(h) = histories.lrpt.lock() {
let snapshot = h.clone();
drop(h);
save_key(db, &k("lrpt"), &snapshot);
}
if let Ok(h) = histories.wefax.lock() {
let snapshot = h.clone();
drop(h);
@@ -185,20 +226,38 @@ pub fn flush_all(db: &mut PickleDb, rig_id: &str, histories: &Arc<DecoderHistori
let _ = db.dump();
}
fn flush_all_rigs(db: &Mutex<PickleDb>, rig_histories: &[(String, Arc<DecoderHistories>)]) {
let Ok(mut guard) = db.lock() else {
tracing::warn!("history database mutex poisoned; skipping periodic flush");
return;
};
for (rig_id, histories) in rig_histories {
flush_all(&mut guard, rig_id, histories);
}
}
/// Spawn a Tokio task that flushes all rigs' histories to disk every 60 seconds.
///
/// Snapshot cloning, JSON serialization, and disk I/O run on Tokio's blocking
/// pool so a large history database cannot stall an async runtime worker.
pub fn spawn_flush_task(
db: Arc<Mutex<PickleDb>>,
rig_histories: Vec<(String, Arc<DecoderHistories>)>,
) {
tokio::spawn(async move {
let rig_histories = Arc::new(rig_histories);
let mut interval = tokio::time::interval(Duration::from_secs(60));
interval.tick().await; // consume the immediate first tick
loop {
interval.tick().await;
if let Ok(mut guard) = db.lock() {
for (rig_id, histories) in &rig_histories {
flush_all(&mut guard, rig_id, histories);
}
let db = Arc::clone(&db);
let rig_histories = Arc::clone(&rig_histories);
if let Err(err) = tokio::task::spawn_blocking(move || {
flush_all_rigs(&db, rig_histories.as_slice());
})
.await
{
tracing::warn!(error = %err, "history flush worker failed");
}
}
});
@@ -302,4 +361,79 @@ mod tests {
let _ = std::fs::remove_file(&db_file);
let _ = std::fs::remove_dir(&dir);
}
#[test]
fn flush_and_load_all_restores_previously_omitted_histories() {
let db_file = std::env::temp_dir().join(format!(
"trx_history_all_{}_{}.db",
std::process::id(),
now_unix_ms()
));
let mut db = PickleDb::new(
&db_file,
PickleDbDumpPolicy::DumpUponRequest,
SerializationMethod::Json,
);
let source = DecoderHistories::new();
let now = Instant::now();
source.hf_aprs.lock().unwrap().push_back((
now,
AprsPacket {
rig_id: Some("rig-a".into()),
ts_ms: Some(now_unix_ms()),
src_call: "TEST".into(),
dest_call: "APRS".into(),
path: String::new(),
info: "history".into(),
info_bytes: Vec::new(),
packet_type: "position".into(),
crc_ok: true,
lat: None,
lon: None,
symbol_table: None,
symbol_code: None,
},
));
for (queue, mode) in [(&source.ft4, "FT4"), (&source.ft2, "FT2")] {
queue.lock().unwrap().push_back((
now,
Ft8Message {
rig_id: Some("rig-a".into()),
ts_ms: now_unix_ms(),
snr_db: -10.0,
dt_s: 0.1,
freq_hz: 1_000.0,
message: mode.into(),
},
));
}
source.lrpt.lock().unwrap().push_back((
now,
LrptImage {
rig_id: Some("rig-a".into()),
pass_start_ms: now_unix_ms(),
pass_end_ms: now_unix_ms(),
mcu_count: 1,
path: "/tmp/lrpt.png".into(),
ts_ms: Some(now_unix_ms()),
satellite: None,
channels: None,
geo_bounds: None,
ground_track: None,
},
));
flush_all(&mut db, "rig-a", &source);
let restored = DecoderHistories::new();
load_all(&db, "rig-a", &restored);
assert_eq!(restored.hf_aprs.lock().unwrap().len(), 1);
assert_eq!(restored.ft4.lock().unwrap().len(), 1);
assert_eq!(restored.ft2.lock().unwrap().len(), 1);
assert_eq!(restored.lrpt.lock().unwrap().len(), 1);
assert_eq!(restored.estimated_total_count(), 4);
let _ = std::fs::remove_file(db_file);
}
}
+2
View File
@@ -4,7 +4,9 @@
mod audio;
mod config;
mod decoder_history;
mod error;
mod history_policy;
mod history_store;
mod listener;
mod rig_handle;
@@ -251,9 +251,9 @@ fn mul_freq_domain(buf: &mut [FftComplex<f32>], h_freq: &[FftComplex<f32>], scal
unsafe {
mul_freq_domain_neon(buf, h_freq, scale);
}
return;
}
#[cfg(not(target_arch = "aarch64"))]
mul_freq_domain_scalar(buf, h_freq, scale);
}