Compare commits
12
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6676b66993 | ||
|
|
412651d612 | ||
|
|
d347b493d9 | ||
|
|
39e59dca96 | ||
|
|
5654520901 | ||
|
|
ff75fcc692 | ||
|
|
4728b578ae | ||
|
|
7db7fad9b0 | ||
|
|
830f7299fe | ||
|
|
061738a63b | ||
|
|
e8bd97655f | ||
|
|
7d0b36450d |
@@ -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!(
|
||||||
|
|||||||
@@ -840,6 +840,8 @@ let jogMult = loadSetting("jogMult", 1); // divisor: 1, 10, 100
|
|||||||
let jogStep = Math.max(Math.round(jogUnit / jogMult), 1);
|
let jogStep = Math.max(Math.round(jogUnit / jogMult), 1);
|
||||||
let minFreqStepHz = 1;
|
let minFreqStepHz = 1;
|
||||||
let lastModeName = "";
|
let lastModeName = "";
|
||||||
|
let lastWfmCci = 0;
|
||||||
|
let lastWfmAci = 0;
|
||||||
const VFO_COLORS = ["var(--accent-green)", "var(--accent-yellow)"];
|
const VFO_COLORS = ["var(--accent-green)", "var(--accent-yellow)"];
|
||||||
function vfoColor(idx) {
|
function vfoColor(idx) {
|
||||||
if (idx < VFO_COLORS.length) return VFO_COLORS[idx];
|
if (idx < VFO_COLORS.length) return VFO_COLORS[idx];
|
||||||
@@ -3293,8 +3295,14 @@ function render(update) {
|
|||||||
wfmStFlagEl.classList.toggle("wfm-st-flag-stereo", detected);
|
wfmStFlagEl.classList.toggle("wfm-st-flag-stereo", detected);
|
||||||
wfmStFlagEl.classList.toggle("wfm-st-flag-mono", !detected);
|
wfmStFlagEl.classList.toggle("wfm-st-flag-mono", !detected);
|
||||||
}
|
}
|
||||||
if (typeof update.filter.wfm_cci === "number") updateIntfBar(wfmCciFillEl, wfmCciValEl, update.filter.wfm_cci);
|
if (typeof update.filter.wfm_cci === "number") {
|
||||||
if (typeof update.filter.wfm_aci === "number") updateIntfBar(wfmAciFillEl, wfmAciValEl, update.filter.wfm_aci);
|
lastWfmCci = Math.max(0, Math.min(100, update.filter.wfm_cci));
|
||||||
|
updateIntfBar(wfmCciFillEl, wfmCciValEl, lastWfmCci);
|
||||||
|
}
|
||||||
|
if (typeof update.filter.wfm_aci === "number") {
|
||||||
|
lastWfmAci = Math.max(0, Math.min(100, update.filter.wfm_aci));
|
||||||
|
updateIntfBar(wfmAciFillEl, wfmAciValEl, lastWfmAci);
|
||||||
|
}
|
||||||
if (samStereoWidthEl && typeof update.filter.sam_stereo_width === "number") {
|
if (samStereoWidthEl && typeof update.filter.sam_stereo_width === "number") {
|
||||||
samStereoWidthEl.value = String(Math.round(update.filter.sam_stereo_width * 100));
|
samStereoWidthEl.value = String(Math.round(update.filter.sam_stereo_width * 100));
|
||||||
}
|
}
|
||||||
@@ -4268,7 +4276,7 @@ const MODE_BW_DEFAULTS = {
|
|||||||
FM: [12_500, 2_500, 25_000, 500],
|
FM: [12_500, 2_500, 25_000, 500],
|
||||||
AIS: [25_000, 12_500, 50_000, 500],
|
AIS: [25_000, 12_500, 50_000, 500],
|
||||||
VDES: [100_000, 25_000, 200_000, 1_000],
|
VDES: [100_000, 25_000, 200_000, 1_000],
|
||||||
WFM: [180_000, 50_000,300_000,5_000],
|
WFM: [180_000, 60_000,300_000,5_000],
|
||||||
DIG: [3_000, 300, 6_000, 100],
|
DIG: [3_000, 300, 6_000, 100],
|
||||||
PKT: [25_000, 300, 50_000, 500],
|
PKT: [25_000, 300, 50_000, 500],
|
||||||
};
|
};
|
||||||
@@ -4348,64 +4356,110 @@ async function applyBandwidthFromInput() {
|
|||||||
} catch (_) {}
|
} catch (_) {}
|
||||||
}
|
}
|
||||||
|
|
||||||
function estimateBandwidthAroundPeak(data, centerHz) {
|
function estimateOccupiedBandwidth(data, centerHz, interference = {}) {
|
||||||
if (!data || !isBinsArray(data.bins) || data.bins.length < 3 || !Number.isFinite(centerHz)) {
|
if (!data || !isBinsArray(data.bins) || data.bins.length < 3 || !Number.isFinite(centerHz)) {
|
||||||
return null;
|
return null;
|
||||||
}
|
}
|
||||||
|
|
||||||
const bins = data.bins;
|
const bins = data.bins;
|
||||||
const maxIdx = bins.length - 1;
|
const maxIdx = bins.length - 1;
|
||||||
|
const hzPerBin = data.sample_rate / maxIdx;
|
||||||
const fullLoHz = data.center_hz - data.sample_rate / 2;
|
const fullLoHz = data.center_hz - data.sample_rate / 2;
|
||||||
const centerIdx = Math.max(
|
const centerIdx = Math.max(
|
||||||
1,
|
1,
|
||||||
Math.min(maxIdx - 1, Math.round(((centerHz - fullLoHz) / data.sample_rate) * maxIdx)),
|
Math.min(maxIdx - 1, Math.round(((centerHz - fullLoHz) / data.sample_rate) * maxIdx)),
|
||||||
);
|
);
|
||||||
const searchRadius = Math.max(6, Math.min(120, Math.round(maxIdx * 0.03)));
|
const mode = (modeEl ? modeEl.value : "USB").toUpperCase();
|
||||||
const searchLo = Math.max(1, centerIdx - searchRadius);
|
const [defaultBw, minBw, maxBw, stepBw] = mwDefaultsForMode(mode);
|
||||||
const searchHi = Math.min(maxIdx - 1, centerIdx + searchRadius);
|
const oneSided = mode === "USB" || mode === "DIG" || mode === "CW"
|
||||||
|
? 1
|
||||||
let peakIdx = centerIdx;
|
: mode === "LSB" || mode === "CWR" ? -1 : 0;
|
||||||
for (let i = searchLo; i <= searchHi; i++) {
|
const isWfm = mode === "WFM";
|
||||||
if (bins[i] > bins[peakIdx]) peakIdx = i;
|
|
||||||
}
|
|
||||||
|
|
||||||
|
// Reduce single-bin peaks and holes before finding occupied-channel edges.
|
||||||
|
// WFM needs a wider smoothing window because its energy is noise-like and
|
||||||
|
// spread across the entire channel rather than concentrated at a carrier.
|
||||||
|
const smoothRadius = isWfm ? 3 : 1;
|
||||||
|
const smoothed = bins.map((_, i) => {
|
||||||
|
let sum = 0;
|
||||||
|
let count = 0;
|
||||||
|
for (let j = Math.max(0, i - smoothRadius); j <= Math.min(maxIdx, i + smoothRadius); j++) {
|
||||||
|
sum += bins[j];
|
||||||
|
count += 1;
|
||||||
|
}
|
||||||
|
return sum / count;
|
||||||
|
});
|
||||||
const sorted = [...bins].sort((a, b) => a - b);
|
const sorted = [...bins].sort((a, b) => a - b);
|
||||||
const noise = sorted[Math.floor(sorted.length * 0.2)];
|
const noise = sorted[Math.floor(sorted.length * 0.2)];
|
||||||
const peak = bins[peakIdx];
|
const maxSpanBins = Math.max(2, Math.ceil(maxBw / hzPerBin));
|
||||||
const threshold = Math.max(noise + 4, peak - Math.max(8, (peak - noise) * 0.35));
|
const searchHalfBins = oneSided === 0 ? Math.ceil(maxSpanBins / 2) : maxSpanBins;
|
||||||
|
const searchLo = Math.max(1, centerIdx - (oneSided > 0 ? 2 : searchHalfBins));
|
||||||
|
const searchHi = Math.min(maxIdx - 1, centerIdx + (oneSided < 0 ? 2 : searchHalfBins));
|
||||||
|
let peak = -Infinity;
|
||||||
|
for (let i = searchLo; i <= searchHi; i++) peak = Math.max(peak, smoothed[i]);
|
||||||
|
const snr = peak - noise;
|
||||||
|
if (!Number.isFinite(snr) || snr < (isWfm ? 5 : 4)) return isWfm ? minBw : defaultBw;
|
||||||
|
|
||||||
let left = peakIdx;
|
// A threshold relative to the noise floor finds occupied bandwidth much
|
||||||
let right = peakIdx;
|
// more reliably than one relative to the peak. The latter fails for WFM,
|
||||||
let belowCount = 0;
|
// whose multiplex spectrum has peaks, notches, and no narrow centre carrier.
|
||||||
for (let i = peakIdx; i > 1; i--) {
|
const threshold = noise + Math.max(3, Math.min(isWfm ? 6 : 10, snr * (isWfm ? 0.18 : 0.28)));
|
||||||
if (bins[i] < threshold) belowCount += 1;
|
const allowedGap = Math.max(isWfm ? 4 : 2, Math.ceil((isWfm ? 12_000 : stepBw) / hzPerBin));
|
||||||
else belowCount = 0;
|
|
||||||
if (belowCount >= 2) break;
|
function occupiedExtent(direction, limitBins) {
|
||||||
left = i;
|
let lastOccupied = centerIdx;
|
||||||
|
let gap = 0;
|
||||||
|
for (let n = 0; n <= limitBins; n++) {
|
||||||
|
const i = centerIdx + direction * n;
|
||||||
|
if (i <= 0 || i >= maxIdx) break;
|
||||||
|
if (smoothed[i] >= threshold) {
|
||||||
|
lastOccupied = i;
|
||||||
|
gap = 0;
|
||||||
|
} else if (++gap > allowedGap) {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return Math.abs(lastOccupied - centerIdx) * hzPerBin;
|
||||||
}
|
}
|
||||||
|
|
||||||
belowCount = 0;
|
let rawBw;
|
||||||
for (let i = peakIdx; i < maxIdx - 1; i++) {
|
if (oneSided !== 0) {
|
||||||
if (bins[i] < threshold) belowCount += 1;
|
rawBw = occupiedExtent(oneSided, maxSpanBins);
|
||||||
else belowCount = 0;
|
} else {
|
||||||
if (belowCount >= 2) break;
|
const leftHz = occupiedExtent(-1, searchHalfBins);
|
||||||
right = i;
|
const rightHz = occupiedExtent(1, searchHalfBins);
|
||||||
|
// A symmetric RF filter must contain the larger of the two sidebands.
|
||||||
|
rawBw = 2 * Math.max(leftHz, rightHz);
|
||||||
}
|
}
|
||||||
|
|
||||||
const shoulderPad = Math.max(1, Math.round((right - left) * 0.08));
|
// Add a transition-band margin. Weak WFM deliberately falls back to the
|
||||||
left = Math.max(0, left - shoulderPad);
|
// 60 kHz mode floor above: a narrower filter trades stereo/RDS content for
|
||||||
right = Math.min(maxIdx, right + shoulderPad);
|
// a useful improvement in intelligibility when the signal is very poor.
|
||||||
|
rawBw *= isWfm ? 1.08 : 1.12;
|
||||||
const hzPerBin = data.sample_rate / maxIdx;
|
if (isWfm) {
|
||||||
const rawBw = Math.max(hzPerBin, (right - left) * hzPerBin);
|
const aci = Math.max(0, Math.min(100, Number(interference.aci) || 0)) / 100;
|
||||||
const [, minBw, maxBw, stepBw] = mwDefaultsForMode(modeEl ? modeEl.value : "USB");
|
const cci = Math.max(0, Math.min(100, Number(interference.cci) || 0)) / 100;
|
||||||
|
// Adjacent-channel energy is outside the wanted modulation, so ACI can
|
||||||
|
// safely drive the cap all the way from the 300 kHz ceiling to 60 kHz.
|
||||||
|
const aciCap = maxBw - (maxBw - minBw) * aci;
|
||||||
|
// CCI overlaps the wanted station and cannot be removed by an RF filter.
|
||||||
|
// Only distrust the widest edge estimates, retaining at least 65% of the
|
||||||
|
// useful range between the weak-signal floor and nominal WFM bandwidth.
|
||||||
|
const cciFloor = minBw + (defaultBw - minBw) * 0.65;
|
||||||
|
const cciCap = maxBw - (maxBw - cciFloor) * cci;
|
||||||
|
rawBw = Math.min(rawBw, aciCap, cciCap);
|
||||||
|
}
|
||||||
const clamped = Math.max(minBw, Math.min(maxBw, rawBw));
|
const clamped = Math.max(minBw, Math.min(maxBw, rawBw));
|
||||||
return Math.max(stepBw, Math.round(clamped / stepBw) * stepBw);
|
return Math.max(stepBw, Math.round(clamped / stepBw) * stepBw);
|
||||||
}
|
}
|
||||||
|
|
||||||
async function applyAutoBandwidth() {
|
async function applyAutoBandwidth() {
|
||||||
if (!lastSpectrumData || lastFreqHz == null) return;
|
if (!lastSpectrumData || lastFreqHz == null) return;
|
||||||
const estimated = estimateBandwidthAroundPeak(lastSpectrumData, lastFreqHz);
|
// WFM interference telemetry belongs to the primary DSP channel. Do not
|
||||||
|
// apply it to a virtual channel, where it would describe the wrong signal.
|
||||||
|
const onVirtual = typeof vchanIsOnVirtual === "function" && vchanIsOnVirtual();
|
||||||
|
const interference = onVirtual ? {} : { cci: lastWfmCci, aci: lastWfmAci };
|
||||||
|
const estimated = estimateOccupiedBandwidth(lastSpectrumData, lastFreqHz, interference);
|
||||||
if (!Number.isFinite(estimated) || estimated <= 0) {
|
if (!Number.isFinite(estimated) || estimated <= 0) {
|
||||||
syncBandwidthInput(currentBandwidthHz);
|
syncBandwidthInput(currentBandwidthHz);
|
||||||
return;
|
return;
|
||||||
|
|||||||
@@ -1641,29 +1641,56 @@ SPDX-License-Identifier: GPL-2.0-or-later
|
|||||||
'settings': ['/vchan.js', '/scheduler.js']
|
'settings': ['/vchan.js', '/scheduler.js']
|
||||||
};
|
};
|
||||||
var loaded = new Set();
|
var loaded = new Set();
|
||||||
function loadPlugins(tab) {
|
var loading = new Map();
|
||||||
var scripts = pluginScripts[tab];
|
|
||||||
if (!scripts) return;
|
function loadScript(src) {
|
||||||
scripts.forEach(function(src) {
|
if (loaded.has(src)) return Promise.resolve();
|
||||||
if (loaded.has(src)) return;
|
if (loading.has(src)) return loading.get(src);
|
||||||
loaded.add(src);
|
|
||||||
|
var request = new Promise(function(resolve, reject) {
|
||||||
var s = document.createElement('script');
|
var s = document.createElement('script');
|
||||||
s.src = src;
|
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);
|
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)
|
// Eager plugin loading is triggered by app.js (after window.trx is set up)
|
||||||
// via window.loadEagerPlugins(). Dynamic scripts are effectively async, so
|
// via window.loadEagerPlugins(). Dynamic scripts are effectively async, so
|
||||||
// loading them before app.js would cause map-core.js to crash when
|
// loading them before app.js would cause map-core.js to crash when
|
||||||
// window.trx is not yet defined.
|
// window.trx is not yet defined.
|
||||||
window.loadEagerPlugins = function() {
|
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
|
// Load others on tab switch
|
||||||
document.addEventListener('click', function(e) {
|
document.addEventListener('click', function(e) {
|
||||||
var tab = e.target.closest('[data-tab]');
|
var tab = e.target.closest('[data-tab]');
|
||||||
if (tab) loadPlugins(tab.dataset.tab);
|
if (tab) requestPlugins(tab.dataset.tab);
|
||||||
});
|
});
|
||||||
window.loadPluginsForTab = loadPlugins;
|
window.loadPluginsForTab = loadPlugins;
|
||||||
})();
|
})();
|
||||||
|
|||||||
+11
-468
@@ -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,9 +47,9 @@ use trx_wspr::WsprDecoder;
|
|||||||
use uuid::Uuid;
|
use uuid::Uuid;
|
||||||
|
|
||||||
use crate::config::AudioConfig;
|
use crate::config::AudioConfig;
|
||||||
use crate::history_policy::{
|
use crate::history_policy::{current_timestamp_ms, lock_or_recover};
|
||||||
enforce_capacity, lock_or_recover, prune_by_age, HISTORY_RETENTION, MAX_HISTORY_ENTRIES,
|
#[cfg(test)]
|
||||||
};
|
use crate::history_policy::{enforce_capacity, prune_by_age, MAX_HISTORY_ENTRIES};
|
||||||
use trx_decode_log::DecoderLoggers;
|
use trx_decode_log::DecoderLoggers;
|
||||||
|
|
||||||
/// 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).
|
||||||
@@ -65,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 {
|
||||||
@@ -350,457 +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,
|
|
||||||
}
|
|
||||||
|
|
||||||
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, 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 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.
|
||||||
///
|
///
|
||||||
|
|||||||
@@ -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);
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -14,6 +14,13 @@ pub(crate) const HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60)
|
|||||||
/// busy channels independently of time-based pruning.
|
/// busy channels independently of time-based pruning.
|
||||||
pub(crate) const MAX_HISTORY_ENTRIES: usize = 10_000;
|
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> {
|
pub(crate) fn lock_or_recover<'a, T>(mutex: &'a Mutex<T>, label: &str) -> MutexGuard<'a, T> {
|
||||||
mutex.lock().unwrap_or_else(|error| {
|
mutex.lock().unwrap_or_else(|error| {
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
|
|||||||
@@ -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);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,6 +4,7 @@
|
|||||||
|
|
||||||
mod audio;
|
mod audio;
|
||||||
mod config;
|
mod config;
|
||||||
|
mod decoder_history;
|
||||||
mod error;
|
mod error;
|
||||||
mod history_policy;
|
mod history_policy;
|
||||||
mod history_store;
|
mod history_store;
|
||||||
|
|||||||
@@ -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);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user