Compare commits

..
Author SHA1 Message Date
sjg bf819aa177 [fix](trx-frontend): narrow weak WFM to 60 kHz
CI / lint (pull_request) Failing after 1s
CI / test (pull_request) Failing after 2s
CI / reuse (pull_request) Failing after 2s
2026-08-01 02:03:53 +02:00
sjg 0f4afda843 [fix](trx-frontend): make auto bandwidth modulation-aware
CI / lint (pull_request) Failing after 0s
CI / test (pull_request) Failing after 0s
CI / reuse (pull_request) Failing after 0s
2026-08-01 02:02:13 +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
6 changed files with 259 additions and 51 deletions
+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!(
@@ -4268,7 +4268,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 +4348,93 @@ async function applyBandwidthFromInput() {
} catch (_) {} } catch (_) {}
} }
function estimateBandwidthAroundPeak(data, centerHz) { function estimateOccupiedBandwidth(data, centerHz) {
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;
const rawBw = Math.max(hzPerBin, (right - left) * hzPerBin);
const [, minBw, maxBw, stepBw] = mwDefaultsForMode(modeEl ? modeEl.value : "USB");
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); const estimated = estimateOccupiedBandwidth(lastSpectrumData, lastFreqHz);
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;
})(); })();
+18 -1
View File
@@ -465,8 +465,25 @@ impl DecoderHistories {
/// entries across all decoders, used for pre-allocating the replay blob. /// entries across all decoders, used for pre-allocating the replay blob.
/// ///
/// Uses an `AtomicUsize` counter maintained by record/prune/clear methods, /// Uses an `AtomicUsize` counter maintained by record/prune/clear methods,
/// avoiding 9 separate mutex acquisitions. /// avoiding 11 separate mutex acquisitions.
pub fn estimated_total_count(&self) -> usize { pub fn estimated_total_count(&self) -> usize {
self.total_count.load(Ordering::Relaxed) 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);
}
} }
+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);
}
} }
@@ -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);
} }