Compare commits

..
Author SHA1 Message Date
sjg 3080a58c4a [fix](trx-frontend): serialize plugin loading
CI / lint (pull_request) Failing after 1s
CI / test (pull_request) Failing after 2s
CI / reuse (pull_request) Failing after 1s
2026-08-01 01:28:55 +02:00
4 changed files with 8 additions and 160 deletions
+1 -2
View File
@@ -274,8 +274,7 @@ 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 as u32).wrapping_mul(1_103_515_245).wrapping_add(12_345)) .map(|i| ((i * 1103515245 + 12345) as u32 >> 8 & 0xff) as u8)
.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!(
+1 -18
View File
@@ -378,7 +378,7 @@ pub struct DecoderHistories {
pub lrpt: Mutex<VecDeque<(Instant, LrptImage)>>, pub lrpt: Mutex<VecDeque<(Instant, LrptImage)>>,
pub wefax: Mutex<VecDeque<(Instant, WefaxMessage)>>, pub wefax: Mutex<VecDeque<(Instant, WefaxMessage)>>,
/// Approximate total entry count across all decoders, maintained /// Approximate total entry count across all decoders, maintained
/// atomically so `estimated_total_count()` avoids 11 lock acquisitions. /// atomically so `estimated_total_count()` avoids 9 lock acquisitions.
total_count: AtomicUsize, total_count: AtomicUsize,
} }
@@ -845,23 +845,6 @@ impl DecoderHistories {
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);
}
} }
/// Spawn the audio capture thread. /// Spawn the audio capture thread.
+5 -139
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, LrptImage, VdesMessage, WefaxMessage, WsprMessage, AisMessage, AprsPacket, CwEvent, Ft8Message, VdesMessage, WefaxMessage, WsprMessage,
}; };
use crate::audio::DecoderHistories; use crate::audio::DecoderHistories;
@@ -118,11 +118,6 @@ 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);
@@ -133,32 +128,16 @@ 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.
@@ -183,11 +162,6 @@ 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);
@@ -198,26 +172,11 @@ 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);
@@ -226,38 +185,20 @@ 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;
let db = Arc::clone(&db); if let Ok(mut guard) = db.lock() {
let rig_histories = Arc::clone(&rig_histories); for (rig_id, histories) in &rig_histories {
if let Err(err) = tokio::task::spawn_blocking(move || { flush_all(&mut guard, rig_id, histories);
flush_all_rigs(&db, rig_histories.as_slice()); }
})
.await
{
tracing::warn!(error = %err, "history flush worker failed");
} }
} }
}); });
@@ -361,79 +302,4 @@ 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);
} }