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
18 changed files with 1039 additions and 578 deletions
+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 ## License
GPL-2.0-or-later. See [`LICENSES`](LICENSES) for the full license text and 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 bundled third-party license files. Bundled third-party components retain their
the Leaflet AIS tracksymbol plugin under `assets/web/vendor/`) retain their original licenses: Leaflet is BSD-2-Clause, DSEG is OFL-1.1, and opus-decoder
original BSD-2-Clause license. is MIT.
+8
View File
@@ -42,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"] 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-FileCopyrightText = "2020 The DSEG Authors (https://github.com/keshikan/DSEG)"
SPDX-License-Identifier = "OFL-1.1" 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"
+2 -1
View File
@@ -274,7 +274,8 @@ mod tests {
// Pseudo-random noise vs gradient — correlation should be low. // Pseudo-random noise vs gradient — correlation should be low.
let noise: Vec<u8> = (0..256) 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(); .collect();
let r = asm.correlation_with_last(&noise).expect("r"); let r = asm.correlation_with_last(&noise).expect("r");
assert!( assert!(
@@ -645,7 +645,7 @@ function flushDeferredDecodeMapSync() {
if (!decodeMapSyncPending || decodeHistoryReplayActive || !window.trx?.map?.aprsMap) return; if (!decodeMapSyncPending || decodeHistoryReplayActive || !window.trx?.map?.aprsMap) return;
decodeMapSyncPending = false; decodeMapSyncPending = false;
scheduleUiFrameJob("decode-map-maintenance", () => { 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 bmPopulateScopePicker === "function") bmPopulateScopePicker();
if (typeof bmFetch === "function") bmFetch(document.getElementById("bm-category-filter")?.value || ""); 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; serverRigs = rigs;
serverActiveRigId = data.active_remote || null; serverActiveRigId = data.active_remote || null;
applyRigList(data.active_remote, rigIds, displayNames); applyRigList(data.active_remote, rigIds, displayNames);
window.trx.map?.syncAprsReceiverMarker(); window.trx.modules.map?.syncAprsReceiverMarker();
} catch (e) { } catch (e) {
// Non-fatal: SSE/status path still drives main UI. // Non-fatal: SSE/status path still drives main UI.
} }
@@ -3143,9 +3143,9 @@ function render(update) {
const grid = latLonToMaidenhead(serverLat, serverLon); const grid = latLonToMaidenhead(serverLat, serverLon);
locationSubtitle.textContent = `Location: ${grid}`; locationSubtitle.textContent = `Location: ${grid}`;
locationSubtitle.style.display = ""; 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)) { if (typeof update.initial_map_zoom === "number" && Number.isFinite(update.initial_map_zoom)) {
initialMapZoom = Math.max(1, Math.round(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 setSchedulerRig === "function") setSchedulerRig(lastActiveRigId);
if (typeof setBackgroundDecodeRig === "function") setBackgroundDecodeRig(lastActiveRigId); if (typeof setBackgroundDecodeRig === "function") setBackgroundDecodeRig(lastActiveRigId);
if (typeof bmFetch === "function") bmFetch(document.getElementById("bm-category-filter")?.value || ""); 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 // Switch this session's rig and reconnect SSE to the new rig's
// state channel. // state channel.
try { try {
@@ -4480,23 +4480,23 @@ function updateTabHistory(name, replaceHistory = false) {
} }
// Initialise the Leaflet map, waiting for both Leaflet (L) and map-core.js // 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; let _mapInitTimer = null;
function _initMapWhenReady() { function _initMapWhenReady() {
const loadingEl = document.getElementById("map-loading"); 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 (_mapInitTimer) { clearInterval(_mapInitTimer); _mapInitTimer = null; }
if (loadingEl) loadingEl.classList.add("is-hidden"); if (loadingEl) loadingEl.classList.add("is-hidden");
window.trx.map.initAprsMap(); window.trx.modules.map.initAprsMap();
window.trx.map.sizeAprsMapToViewport(); window.trx.modules.map.sizeAprsMapToViewport();
// The map panel was just made visible (display:none → ""); the browser // The map panel was just made visible (display:none → ""); the browser
// may not have laid it out yet, so getBoundingClientRect() can return // may not have laid it out yet, so getBoundingClientRect() can return
// stale/zero dimensions. Double-rAF ensures a full layout pass has // stale/zero dimensions. Double-rAF ensures a full layout pass has
// completed before we re-measure and tell Leaflet about its real size. // completed before we re-measure and tell Leaflet about its real size.
requestAnimationFrame(() => { requestAnimationFrame(() => {
requestAnimationFrame(() => { requestAnimationFrame(() => {
window.trx.map.sizeAprsMapToViewport(); window.trx.modules.map.sizeAprsMapToViewport();
if (window.trx.map.aprsMap) window.trx.map.aprsMap.invalidateSize(); if (window.trx.modules.map.aprsMap) window.trx.modules.map.aprsMap.invalidateSize();
}); });
}); });
return; return;
@@ -4544,7 +4544,7 @@ function navigateToTab(name, options = {}) {
_initMapWhenReady(); _initMapWhenReady();
} }
if (name === "statistics") { if (name === "statistics") {
window.trx.map?.scheduleStatsRender(); window.trx.modules.map?.scheduleStatsRender();
} }
if (name === "recorder") { if (name === "recorder") {
refreshRecorderStatus(); refreshRecorderStatus();
@@ -4721,10 +4721,11 @@ if (headerAuthBtn) {
// ── Shared namespace for lazy-loaded modules ──────────────────────────────── // ── Shared namespace for lazy-loaded modules ────────────────────────────────
// Modules (map-core.js, screenshot.js) access core state and utilities via // Modules (map-core.js, screenshot.js) access core state and utilities via
// window.trx. Modules register their own APIs as sub-namespaces // window.trx. Modules register their own APIs as sub-namespaces
// (e.g. window.trx.map, window.trx.screenshot). // (e.g. window.trx.modules.map, window.trx.modules.screenshot).
window.trx = Object.create(null); const trxState = Object.create(null);
const trxModules = Object.create(null);
// -- State getters (backed by core-scoped variables) -- // -- State getters (backed by core-scoped variables) --
Object.defineProperties(window.trx, { Object.defineProperties(trxState, {
serverLat: { get() { return serverLat; }, set(v) { serverLat = v; } }, serverLat: { get() { return serverLat; }, set(v) { serverLat = v; } },
serverLon: { get() { return serverLon; }, set(v) { serverLon = v; } }, serverLon: { get() { return serverLon; }, set(v) { serverLon = v; } },
lastFreqHz: { get() { return lastFreqHz; } }, lastFreqHz: { get() { return lastFreqHz; } },
@@ -4762,7 +4763,7 @@ Object.defineProperties(window.trx, {
signalOverlayGl: { get() { return signalOverlayGl; } }, signalOverlayGl: { get() { return signalOverlayGl; } },
}); });
// -- Shared utility functions -- // -- Shared utility functions --
Object.assign(window.trx, { const trxCore = Object.freeze({
saveSetting, loadSetting, showHint, escapeMapHtml, formatFreq, formatFreqForHumans, saveSetting, loadSetting, showHint, escapeMapHtml, formatFreq, formatFreqForHumans,
formatWavelength, formatBwLabel, formatUptime, formatSigStrength, formatSignal, formatWavelength, formatBwLabel, formatUptime, formatSigStrength, formatSignal,
postPath, scheduleUiFrameJob, navigateToTab, rigBadgeColor, postPath, scheduleUiFrameJob, navigateToTab, rigBadgeColor,
@@ -4772,18 +4773,19 @@ Object.assign(window.trx, {
currentTheme, canvasPalette, currentStyle, currentTheme, canvasPalette, currentStyle,
cssColorToRgba, rgbaWithAlpha, isBinsArray, estimateNoiseFloorDb, cssColorToRgba, rgbaWithAlpha, isBinsArray, estimateNoiseFloorDb,
spectrumVisibleRange, drawSpectrum, spectrumVisibleRange, drawSpectrum,
bandForHz: function(hz) { return window.trx.map?.bandForHz?.(hz); }, bandForHz: function(hz) { return trxModules.map?.bandForHz?.(hz); },
markDecodeMapSyncPending, markDecodeMapSyncPending,
decodeHistoryMapRenderingDeferred, decodeHistoryMapRenderingDeferred,
updateDocumentTitle, updateDocumentTitle,
activeChannelRds, activeChannelRds,
}); });
Object.defineProperties(window.trx, { Object.defineProperties(trxState, {
decodeHistoryReplayActive: { get() { return decodeHistoryReplayActive; } }, decodeHistoryReplayActive: { get() { return decodeHistoryReplayActive; } },
decodeMapSyncPending: { get() { return decodeMapSyncPending; } }, decodeMapSyncPending: { get() { return decodeMapSyncPending; } },
_activeTab: { get() { return _activeTab; } }, _activeTab: { get() { return _activeTab; } },
locationSubtitle: { get() { return locationSubtitle; } }, 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 // 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. // 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) ────────────────────── // ── Map module (extracted to map-core.js, lazy-loaded) ──────────────────────
// The map, statistics, and geolocation code (~3,450 lines) has been moved to // 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. // 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) ───────────────── // ── Geo utilities (shared with map-core.js via window.trx) ─────────────────
function haversineKm(lat1, lon1, lat2, lon2) { function haversineKm(lat1, lon1, lat2, lon2) {
@@ -4932,7 +4934,7 @@ document.querySelectorAll(".sub-tab-bar").forEach(_wireSubTabBar);
window.addEventListener("resize", () => { window.addEventListener("resize", () => {
const mapTab = document.getElementById("tab-map"); const mapTab = document.getElementById("tab-map");
if (!mapTab || mapTab.style.display === "none") return; if (!mapTab || mapTab.style.display === "none") return;
window.trx.map?.sizeAprsMapToViewport(); window.trx.modules.map?.sizeAprsMapToViewport();
}); });
// --- Signal measurement --- // --- Signal measurement ---
@@ -6102,8 +6104,8 @@ function dispatchDecodeMessage(msg, skipStats) {
if (msg.type === "wefax" && window.onServerWefax) window.onServerWefax(msg); if (msg.type === "wefax" && window.onServerWefax) window.onServerWefax(msg);
if (msg.type === "wefax_progress" && window.onServerWefaxProgress) window.onServerWefaxProgress(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") { 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.modules.map?.statsRecordDecode(msg.type, msg.rig_id || msg.remote || null);
window.trx.map?.scheduleStatsRender(); window.trx.modules.map?.scheduleStatsRender();
} }
} }
@@ -6112,10 +6114,10 @@ function dispatchDecodeBatch(batch) {
// Record statistics for every message in the batch regardless of dispatch path. // Record statistics for every message in the batch regardless of dispatch path.
for (const msg of batch) { for (const msg of batch) {
if (msg.type && msg.type !== "lrpt_image" && msg.type !== "lrpt_progress" && msg.type !== "wefax" && msg.type !== "wefax_progress") { 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 type = String(batch[0]?.type || "");
const uniformType = batch.every((msg) => String(msg?.type || "") === type); const uniformType = batch.every((msg) => String(msg?.type || "") === type);
if (uniformType) { if (uniformType) {
@@ -6200,9 +6202,9 @@ function restoreDecodeHistoryGroup(kind, messages) {
// Record statistics for restored history messages. // Record statistics for restored history messages.
if (kind !== "lrpt_image" && kind !== "lrpt_progress" && kind !== "wefax" && kind !== "wefax_progress") { if (kind !== "lrpt_image" && kind !== "lrpt_progress" && kind !== "wefax" && kind !== "wefax_progress") {
for (const msg of messages) { 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 (kind === "ais") {
if (window.restoreAisHistory) { window.restoreAisHistory(messages); } if (window.restoreAisHistory) { window.restoreAisHistory(messages); }
@@ -8000,12 +8002,12 @@ window.addEventListener("keydown", (event) => {
// S — spectrum screenshot (lazy-loads screenshot.js on first use) // S — spectrum screenshot (lazy-loads screenshot.js on first use)
if (key === "s") { if (key === "s") {
event.preventDefault(); event.preventDefault();
if (window.trx.screenshot) { if (window.trx.modules.screenshot) {
void window.trx.screenshot.captureSpectrumScreenshot(); void window.trx.modules.screenshot.captureSpectrumScreenshot();
} else { } else {
const s = document.createElement("script"); const s = document.createElement("script");
s.src = "/screenshot.js"; s.src = "/screenshot.js";
s.onload = () => { void window.trx.screenshot?.captureSpectrumScreenshot(); }; s.onload = () => { void window.trx.modules.screenshot?.captureSpectrumScreenshot(); };
document.body.appendChild(s); document.body.appendChild(s);
} }
return; return;
@@ -1623,7 +1623,7 @@ 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 id="decode-history-overlay-sub" class="decode-history-overlay-sub">Preparing recent decodes for the UI</div>
</div> </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="/vendor/leaflet.js"></script>
<script defer src="/leaflet-ais-tracksymbol.js"></script> <script defer src="/leaflet-ais-tracksymbol.js"></script>
<script defer src="/webgl-renderer.js"></script> <script defer src="/webgl-renderer.js"></script>
@@ -3,10 +3,10 @@
// SPDX-License-Identifier: GPL-2.0-or-later // SPDX-License-Identifier: GPL-2.0-or-later
// Map, statistics, and geolocation module (lazy-loaded on map tab activation). // 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 () { (function () {
"use strict"; "use strict";
const T = window.trx; const { state: T, core: C, modules } = window.trx;
// Destructure shared utility functions for convenience // Destructure shared utility functions for convenience
const { saveSetting, loadSetting, showHint, escapeMapHtml, formatFreq, formatFreqForHumans, const { saveSetting, loadSetting, showHint, escapeMapHtml, formatFreq, formatFreqForHumans,
@@ -14,7 +14,7 @@
formatUptime, latLonToMaidenhead, locatorToLatLon, haversineKm, formatUptime, latLonToMaidenhead, locatorToLatLon, haversineKm,
formatDistanceKm, formatTimeAgo, currentDecodeHistoryRetentionMs, formatDistanceKm, formatTimeAgo, currentDecodeHistoryRetentionMs,
formatWavelength, bookmarkDistanceText, buildBookmarkTooltipText, formatWavelength, bookmarkDistanceText, buildBookmarkTooltipText,
nearestBookmarkForHz } = T; nearestBookmarkForHz } = C;
function updateMapRigFilter() { function updateMapRigFilter() {
const el = document.getElementById("map-rig-filter"); const el = document.getElementById("map-rig-filter");
@@ -261,7 +261,7 @@
if (canRenderMap) { if (canRenderMap) {
refreshAprsTrack(call, entry); refreshAprsTrack(call, entry);
} else { } else {
T.markDecodeMapSyncPending(); C.markDecodeMapSyncPending();
} }
if (!visible) { if (!visible) {
if (canRenderMap && selectedAprsTrackCall && String(selectedAprsTrackCall) === String(call)) { if (canRenderMap && selectedAprsTrackCall && String(selectedAprsTrackCall) === String(call)) {
@@ -294,7 +294,7 @@
if (canRenderMap) { if (canRenderMap) {
refreshAisTrack(key, entry); refreshAisTrack(key, entry);
} else { } else {
T.markDecodeMapSyncPending(); C.markDecodeMapSyncPending();
} }
if (!visible) { if (!visible) {
if (canRenderMap && selectedAisTrackMmsi && String(selectedAisTrackMmsi) === String(key)) { if (canRenderMap && selectedAisTrackMmsi && String(selectedAisTrackMmsi) === String(key)) {
@@ -337,7 +337,7 @@
entry.stations = new Set(); entry.stations = new Set();
entry.bandMeta = new Map(); entry.bandMeta = new Map();
if (canRenderMap) setRetainedMapMarkerVisible(entry.marker, false); if (canRenderMap) setRetainedMapMarkerVisible(entry.marker, false);
else T.markDecodeMapSyncPending(); else C.markDecodeMapSyncPending();
return false; return false;
} }
const nextStations = new Set(); const nextStations = new Set();
@@ -352,7 +352,7 @@
); );
const count = Math.max(nextDetails.size, nextStations.size || 0, 1); const count = Math.max(nextDetails.size, nextStations.size || 0, 1);
if (!canRenderMap) { if (!canRenderMap) {
T.markDecodeMapSyncPending(); C.markDecodeMapSyncPending();
return true; return true;
} }
ensureDecodeLocatorMarker(entry); ensureDecodeLocatorMarker(entry);
@@ -392,7 +392,7 @@
pruneLocatorEntry(key, entry, cutoffMs); pruneLocatorEntry(key, entry, cutoffMs);
} }
if (!aprsMap || T.decodeHistoryReplayActive) { if (!aprsMap || T.decodeHistoryReplayActive) {
T.markDecodeMapSyncPending(); C.markDecodeMapSyncPending();
return; return;
} }
rebuildDecodeContactPaths(); rebuildDecodeContactPaths();
@@ -415,7 +415,7 @@
function locatorFilterColor(type) { function locatorFilterColor(type) {
const hues = locatorThemeHues(); const hues = locatorThemeHues();
const lightTheme = T.currentTheme() === "light"; const lightTheme = C.currentTheme() === "light";
const sat = lightTheme ? 66 : 76; const sat = lightTheme ? 66 : 76;
const light = lightTheme ? 42 : 56; const light = lightTheme ? 42 : 56;
const hue = type === "bookmark" const hue = type === "bookmark"
@@ -539,7 +539,7 @@
} }
function locatorThemeHues() { function locatorThemeHues() {
const pal = T.canvasPalette(); const pal = C.canvasPalette();
const baseHue = paletteHue(pal?.spectrumLine, 145); const baseHue = paletteHue(pal?.spectrumLine, 145);
const waveHue = paletteHue(pal?.waveformLine, baseHue + 34); const waveHue = paletteHue(pal?.waveformLine, baseHue + 34);
const peakHue = paletteHue(pal?.waveformPeak, baseHue - 42); const peakHue = paletteHue(pal?.waveformPeak, baseHue - 42);
@@ -560,7 +560,7 @@
function locatorBandChipColor(label) { function locatorBandChipColor(label) {
const hues = locatorThemeHues(); const hues = locatorThemeHues();
const lightTheme = T.currentTheme() === "light"; const lightTheme = C.currentTheme() === "light";
const hue = wrapHue(hues.bandBase + locatorBandIndex(label) * 137.508); const hue = wrapHue(hues.bandBase + locatorBandIndex(label) * 137.508);
const sat = lightTheme ? 68 : 78; const sat = lightTheme ? 68 : 78;
const light = lightTheme ? 44 : 58; const light = lightTheme ? 44 : 58;
@@ -606,7 +606,7 @@
const safeCount = Math.max(1, Number.isFinite(count) ? count : 1); const safeCount = Math.max(1, Number.isFinite(count) ? count : 1);
const intensity = Math.min(1, Math.log2(safeCount + 1) / 5); const intensity = Math.min(1, Math.log2(safeCount + 1) / 5);
const hue = locatorHueForEntry(entry); const hue = locatorHueForEntry(entry);
const lightTheme = T.currentTheme() === "light"; const lightTheme = C.currentTheme() === "light";
const strokeSat = lightTheme ? 62 : 74; const strokeSat = lightTheme ? 62 : 74;
const fillSat = lightTheme ? 68 : 78; const fillSat = lightTheme ? 68 : 78;
const strokeLight = lightTheme ? 40 : 56; const strokeLight = lightTheme ? 40 : 56;
@@ -1573,7 +1573,7 @@
stageResizeObserver = new ResizeObserver(() => sizeAprsMapToViewport()); stageResizeObserver = new ResizeObserver(() => sizeAprsMapToViewport());
stageResizeObserver.observe(stage); stageResizeObserver.observe(stage);
} }
updateMapBaseLayerForTheme(T.currentTheme()); updateMapBaseLayerForTheme(C.currentTheme());
syncAprsReceiverMarker(); syncAprsReceiverMarker();
// Rebuild popup content on open (keeps age/distance/rig list fresh) // Rebuild popup content on open (keeps age/distance/rig list fresh)
@@ -2307,18 +2307,18 @@
existing.rigIds.add(msgRigId); existing.rigIds.add(msgRigId);
} }
if (!visible) { if (!visible) {
if (!T.decodeHistoryMapRenderingDeferred()) { if (!C.decodeHistoryMapRenderingDeferred()) {
setRetainedMapMarkerVisible(existing.marker, false); setRetainedMapMarkerVisible(existing.marker, false);
} else { } else {
T.markDecodeMapSyncPending(); C.markDecodeMapSyncPending();
} }
return; return;
} }
if (!T.decodeHistoryMapRenderingDeferred()) { if (!C.decodeHistoryMapRenderingDeferred()) {
ensureVdesMarker(key, existing); ensureVdesMarker(key, existing);
setRetainedMapMarkerVisible(existing.marker, true); setRetainedMapMarkerVisible(existing.marker, true);
} else { } else {
T.markDecodeMapSyncPending(); C.markDecodeMapSyncPending();
} }
if (aprsMap && existing.marker && !T.decodeHistoryReplayActive) { if (aprsMap && existing.marker && !T.decodeHistoryReplayActive) {
existing.marker.setLatLng([msg.lat, msg.lon]); existing.marker.setLatLng([msg.lat, msg.lon]);
@@ -2334,11 +2334,11 @@
}; };
vdesMarkers.set(key, entry); vdesMarkers.set(key, entry);
if (!visible) return; if (!visible) return;
if (!T.decodeHistoryMapRenderingDeferred()) { if (!C.decodeHistoryMapRenderingDeferred()) {
ensureVdesMarker(key, entry); ensureVdesMarker(key, entry);
setRetainedMapMarkerVisible(entry.marker, true); setRetainedMapMarkerVisible(entry.marker, true);
} else { } else {
T.markDecodeMapSyncPending(); C.markDecodeMapSyncPending();
} }
if (aprsMap && entry.marker && !T.decodeHistoryReplayActive) { if (aprsMap && entry.marker && !T.decodeHistoryReplayActive) {
entry.marker.setPopupContent(popupHtml); entry.marker.setPopupContent(popupHtml);
@@ -2365,7 +2365,7 @@
if (T.locationSubtitle) { if (T.locationSubtitle) {
T.locationSubtitle.textContent = `Location: ${grid} · ${label}`; T.locationSubtitle.textContent = `Location: ${grid} · ${label}`;
} }
T.updateDocumentTitle(T.activeChannelRds()); C.updateDocumentTitle(C.activeChannelRds());
}) })
.catch(() => {}); .catch(() => {});
} }
@@ -2443,8 +2443,8 @@
} }
function scheduleDecodeMapMaintenance() { function scheduleDecodeMapMaintenance() {
if (T.decodeHistoryMapRenderingDeferred()) { if (C.decodeHistoryMapRenderingDeferred()) {
T.markDecodeMapSyncPending(); C.markDecodeMapSyncPending();
return; return;
} }
scheduleUiFrameJob("decode-map-maintenance", () => { scheduleUiFrameJob("decode-map-maintenance", () => {
@@ -3461,7 +3461,7 @@
} }
// Register module API for core to call // Register module API for core to call
window.trx.map = { modules.map = {
initAprsMap, initAprsMap,
sizeAprsMapToViewport, sizeAprsMapToViewport,
syncAprsReceiverMarker, syncAprsReceiverMarker,
@@ -257,7 +257,7 @@
} }
// Register module API // Register module API
window.trx.screenshot = { window.trx.modules.screenshot = {
captureSpectrumScreenshot, captureSpectrumScreenshot,
buildSpectrumSnapshotCanvas, buildSpectrumSnapshotCanvas,
saveCanvasAsPng, 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 // Vendored DSEG14 Classic font
// (binary woff2 — served directly, not through gz_cache) // (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 // Vendored Leaflet 1.9.4
define_gz_cache!(gz_leaflet_js, status::LEAFLET_JS, "leaflet.js"); define_gz_cache!(gz_leaflet_js, status::LEAFLET_JS, "leaflet.js");
define_gz_cache!(gz_leaflet_css, status::LEAFLET_CSS, "leaflet.css"); 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) .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 // Vendored Leaflet 1.9.4
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
@@ -668,6 +668,8 @@ pub fn configure(cfg: &mut web::ServiceConfig) {
.service(assets::bandplan_json) .service(assets::bandplan_json)
// Vendored DSEG14 Classic font // Vendored DSEG14 Classic font
.service(assets::dseg14_classic_woff2) .service(assets::dseg14_classic_woff2)
// Vendored opus-decoder 0.7.11
.service(assets::opus_decoder_js)
// Vendored Leaflet 1.9.4 // Vendored Leaflet 1.9.4
.service(assets::leaflet_js) .service(assets::leaflet_js)
.service(assets::leaflet_css) .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] = pub const DSEG14_CLASSIC_WOFF2: &[u8] =
include_bytes!("../assets/web/vendor/dseg14-classic-latin-400-normal.woff2"); 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 // Vendored Leaflet 1.9.4
pub const LEAFLET_JS: &str = include_str!("../assets/web/vendor/leaflet.js"); pub const LEAFLET_JS: &str = include_str!("../assets/web/vendor/leaflet.js");
pub const LEAFLET_CSS: &str = include_str!("../assets/web/vendor/leaflet.css"); pub const LEAFLET_CSS: &str = include_str!("../assets/web/vendor/leaflet.css");
+11 -513
View File
@@ -6,13 +6,14 @@
#[cfg(feature = "ft2")] #[cfg(feature = "ft2")]
use std::collections::HashMap; 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::net::SocketAddr;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex}; use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant}; use std::time::{Duration, Instant};
use base64::Engine as _;
use bytes::Bytes; use bytes::Bytes;
use flate2::write::GzEncoder; use flate2::write::GzEncoder;
use flate2::Compression; 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_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, AUDIO_MSG_VDES_DECODE, AUDIO_MSG_WEFAX_DECODE, AUDIO_MSG_WEFAX_PROGRESS, AUDIO_MSG_WSPR_DECODE,
}; };
use trx_core::decode::{ #[cfg(test)]
AisMessage, AprsPacket, CwEvent, DecodedMessage, Ft8Message, LrptImage, LrptProgress, use trx_core::decode::{AisMessage, AprsPacket, CwEvent};
VdesMessage, WefaxMessage, WsprMessage, use trx_core::decode::{DecodedMessage, Ft8Message, LrptImage, LrptProgress, WsprMessage};
};
use trx_core::rig::state::{RigMode, RigState}; use trx_core::rig::state::{RigMode, RigState};
use trx_core::vchan::SharedVChanManager; use trx_core::vchan::SharedVChanManager;
use trx_cw::CwDecoder; use trx_cw::CwDecoder;
@@ -47,21 +47,11 @@ use trx_wspr::WsprDecoder;
use uuid::Uuid; use uuid::Uuid;
use crate::config::AudioConfig; 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; 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). /// Silence timeout before auto-finalising an LRPT pass (30 s without new MCUs).
const LRPT_PASS_SILENCE_TIMEOUT: Duration = Duration::from_secs(30); const LRPT_PASS_SILENCE_TIMEOUT: Duration = Duration::from_secs(30);
const FT8_SAMPLE_RATE: u32 = 12_000; 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_ERROR_LOG_INTERVAL: Duration = Duration::from_secs(60);
const AUDIO_STREAM_RECOVERY_DELAY: Duration = Duration::from_secs(1); 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")] #[cfg(feature = "ft2")]
fn retain_ft2_window(buf: &mut Vec<f32>) { fn retain_ft2_window(buf: &mut Vec<f32>) {
if buf.len() > FT2_ASYNC_BUFFER_SAMPLES { if buf.len() > FT2_ASYNC_BUFFER_SAMPLES {
@@ -360,492 +343,7 @@ fn classify_stream_error(err: &str) -> &'static str {
} }
} }
/// Per-rig decoder history store. pub use crate::decoder_history::DecoderHistories;
///
/// 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)
}
}
/// Spawn the audio capture thread. /// 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 serde::{de::DeserializeOwned, Deserialize, Serialize};
use trx_core::decode::{ use trx_core::decode::{
AisMessage, AprsPacket, CwEvent, Ft8Message, VdesMessage, WefaxMessage, WsprMessage, AisMessage, AprsPacket, CwEvent, Ft8Message, LrptImage, VdesMessage, WefaxMessage, WsprMessage,
}; };
use crate::audio::DecoderHistories; use crate::audio::DecoderHistories;
@@ -118,6 +118,11 @@ pub fn load_all(db: &PickleDb, rig_id: &str, histories: &Arc<DecoderHistories>)
h.push_back(e); 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() { if let Ok(mut h) = histories.cw.lock() {
for e in load_key::<CwEvent>(db, &k("cw")) { for e in load_key::<CwEvent>(db, &k("cw")) {
h.push_back(e); h.push_back(e);
@@ -128,16 +133,32 @@ pub fn load_all(db: &PickleDb, rig_id: &str, histories: &Arc<DecoderHistories>)
h.push_back(e); 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() { if let Ok(mut h) = histories.wspr.lock() {
for e in load_key::<WsprMessage>(db, &k("wspr")) { for e in load_key::<WsprMessage>(db, &k("wspr")) {
h.push_back(e); 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() { if let Ok(mut h) = histories.wefax.lock() {
for e in load_key::<WefaxMessage>(db, &k("wefax")) { for e in load_key::<WefaxMessage>(db, &k("wefax")) {
h.push_back(e); h.push_back(e);
} }
} }
histories.recalculate_total_count();
} }
/// Flush `histories` to the database under `rig_id`-prefixed keys and sync. /// 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); drop(h);
save_key(db, &k("aprs"), &snapshot); 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() { if let Ok(h) = histories.cw.lock() {
let snapshot = h.clone(); let snapshot = h.clone();
drop(h); drop(h);
@@ -172,11 +198,26 @@ pub fn flush_all(db: &mut PickleDb, rig_id: &str, histories: &Arc<DecoderHistori
drop(h); drop(h);
save_key(db, &k("ft8"), &snapshot); 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() { if let Ok(h) = histories.wspr.lock() {
let snapshot = h.clone(); let snapshot = h.clone();
drop(h); drop(h);
save_key(db, &k("wspr"), &snapshot); 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() { if let Ok(h) = histories.wefax.lock() {
let snapshot = h.clone(); let snapshot = h.clone();
drop(h); drop(h);
@@ -185,20 +226,38 @@ pub fn flush_all(db: &mut PickleDb, rig_id: &str, histories: &Arc<DecoderHistori
let _ = db.dump(); 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. /// 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( pub fn spawn_flush_task(
db: Arc<Mutex<PickleDb>>, db: Arc<Mutex<PickleDb>>,
rig_histories: Vec<(String, Arc<DecoderHistories>)>, rig_histories: Vec<(String, Arc<DecoderHistories>)>,
) { ) {
tokio::spawn(async move { tokio::spawn(async move {
let rig_histories = Arc::new(rig_histories);
let mut interval = tokio::time::interval(Duration::from_secs(60)); let mut interval = tokio::time::interval(Duration::from_secs(60));
interval.tick().await; // consume the immediate first tick interval.tick().await; // consume the immediate first tick
loop { loop {
interval.tick().await; interval.tick().await;
if let Ok(mut guard) = db.lock() { let db = Arc::clone(&db);
for (rig_id, histories) in &rig_histories { let rig_histories = Arc::clone(&rig_histories);
flush_all(&mut guard, rig_id, 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_file(&db_file);
let _ = std::fs::remove_dir(&dir); 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 audio;
mod config; mod config;
mod decoder_history;
mod error; mod error;
mod history_policy;
mod history_store; mod history_store;
mod listener; mod listener;
mod rig_handle; mod rig_handle;
@@ -251,9 +251,9 @@ fn mul_freq_domain(buf: &mut [FftComplex<f32>], h_freq: &[FftComplex<f32>], scal
unsafe { unsafe {
mul_freq_domain_neon(buf, h_freq, scale); mul_freq_domain_neon(buf, h_freq, scale);
} }
return;
} }
#[cfg(not(target_arch = "aarch64"))]
mul_freq_domain_scalar(buf, h_freq, scale); mul_freq_domain_scalar(buf, h_freq, scale);
} }