From 730fd7a1a3a852d1789cb4e419b1bed9cd4ab5b0 Mon Sep 17 00:00:00 2001 From: Stan Grams Date: Sat, 1 Aug 2026 01:42:42 +0200 Subject: [PATCH] [refactor](trx-server): extract decoder history store --- src/trx-server/src/audio.rs | 466 +------------------------ src/trx-server/src/decoder_history.rs | 472 ++++++++++++++++++++++++++ src/trx-server/src/history_policy.rs | 7 + src/trx-server/src/main.rs | 1 + 4 files changed, 484 insertions(+), 462 deletions(-) create mode 100644 src/trx-server/src/decoder_history.rs diff --git a/src/trx-server/src/audio.rs b/src/trx-server/src/audio.rs index bf9745c8..6f191fe7 100644 --- a/src/trx-server/src/audio.rs +++ b/src/trx-server/src/audio.rs @@ -8,11 +8,10 @@ use std::collections::HashMap; use std::collections::{HashSet, VecDeque}; use std::net::SocketAddr; -use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; -use base64::Engine as _; use bytes::Bytes; use flate2::write::GzEncoder; use flate2::Compression; @@ -36,7 +35,7 @@ use trx_core::audio::{ }; use trx_core::decode::{ AisMessage, AprsPacket, CwEvent, DecodedMessage, Ft8Message, LrptImage, LrptProgress, - VdesMessage, WefaxMessage, WsprMessage, + WsprMessage, }; use trx_core::rig::state::{RigMode, RigState}; use trx_core::vchan::SharedVChanManager; @@ -48,7 +47,7 @@ use uuid::Uuid; use crate::config::AudioConfig; use crate::history_policy::{ - enforce_capacity, lock_or_recover, prune_by_age, HISTORY_RETENTION, MAX_HISTORY_ENTRIES, + current_timestamp_ms, enforce_capacity, lock_or_recover, prune_by_age, MAX_HISTORY_ENTRIES, }; use trx_decode_log::DecoderLoggers; @@ -65,13 +64,6 @@ const DECODE_AUDIO_GATE_RMS: f32 = 2.5e-4; const AUDIO_STREAM_ERROR_LOG_INTERVAL: Duration = Duration::from_secs(60); const AUDIO_STREAM_RECOVERY_DELAY: Duration = Duration::from_secs(1); -fn current_timestamp_ms() -> i64 { - match std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH) { - Ok(dur) => dur.as_millis() as i64, - Err(_) => 0, - } -} - #[cfg(feature = "ft2")] fn retain_ft2_window(buf: &mut Vec) { if buf.len() > FT2_ASYNC_BUFFER_SAMPLES { @@ -350,457 +342,7 @@ fn classify_stream_error(err: &str) -> &'static str { } } -/// Per-rig decoder history store. -/// -/// Replaces the previous process-wide `OnceLock` statics so that each rig -/// instance can maintain its own independent history. Pass an -/// `Arc` into every decoder task and into the audio listener. -pub struct DecoderHistories { - pub ais: Mutex>, - pub vdes: Mutex>, - pub aprs: Mutex>, - pub hf_aprs: Mutex>, - pub cw: Mutex>, - pub ft8: Mutex>, - pub ft4: Mutex>, - pub ft2: Mutex>, - pub wspr: Mutex>, - pub lrpt: Mutex>, - pub wefax: Mutex>, - /// 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 { - 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 { - 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 { - 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 { - 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 { - 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 { - 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 { - 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 { - 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 { - 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 { - 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 { - 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 { - let mut h = lock_or_recover(&self.wefax, "wefax_history"); - let before = h.len(); - Self::prune_wefax(&mut h); - self.adjust_total_count(before, h.len()); - h.iter() - .map(|(_, msg)| { - let mut m = msg.clone(); - // Re-read PNG from disk so remote clients can save a local copy. - if m.png_data.is_none() { - if let Some(ref path) = m.path { - if let Ok(bytes) = std::fs::read(path) { - m.png_data = - Some(base64::engine::general_purpose::STANDARD.encode(&bytes)); - } - } - } - m - }) - .collect() - } - - pub fn clear_wefax_history(&self) { - let mut h = lock_or_recover(&self.wefax, "wefax_history"); - let before = h.len(); - h.clear(); - self.adjust_total_count(before, 0); - } - - /// Returns a quick (non-pruning) estimate of the total number of history - /// entries across all decoders, used for pre-allocating the replay blob. - /// - /// Uses an `AtomicUsize` counter maintained by record/prune/clear methods, - /// avoiding 9 separate mutex acquisitions. - pub fn estimated_total_count(&self) -> usize { - self.total_count.load(Ordering::Relaxed) - } -} +pub use crate::decoder_history::DecoderHistories; /// Spawn the audio capture thread. /// diff --git a/src/trx-server/src/decoder_history.rs b/src/trx-server/src/decoder_history.rs new file mode 100644 index 00000000..0f1e9f3b --- /dev/null +++ b/src/trx-server/src/decoder_history.rs @@ -0,0 +1,472 @@ +// SPDX-FileCopyrightText: 2026 Stan Grams +// +// 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` into every decoder task and into the audio listener. +pub struct DecoderHistories { + pub ais: Mutex>, + pub vdes: Mutex>, + pub aprs: Mutex>, + pub hf_aprs: Mutex>, + pub cw: Mutex>, + pub ft8: Mutex>, + pub ft4: Mutex>, + pub ft2: Mutex>, + pub wspr: Mutex>, + pub lrpt: Mutex>, + pub wefax: Mutex>, + /// 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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 { + 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) + } +} diff --git a/src/trx-server/src/history_policy.rs b/src/trx-server/src/history_policy.rs index 2ee54e3a..5186fdd7 100644 --- a/src/trx-server/src/history_policy.rs +++ b/src/trx-server/src/history_policy.rs @@ -14,6 +14,13 @@ pub(crate) const HISTORY_RETENTION: Duration = Duration::from_secs(24 * 60 * 60) /// busy channels independently of time-based pruning. pub(crate) const MAX_HISTORY_ENTRIES: usize = 10_000; +pub(crate) fn current_timestamp_ms() -> i64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|duration| duration.as_millis() as i64) + .unwrap_or(0) +} + pub(crate) fn lock_or_recover<'a, T>(mutex: &'a Mutex, label: &str) -> MutexGuard<'a, T> { mutex.lock().unwrap_or_else(|error| { tracing::warn!( diff --git a/src/trx-server/src/main.rs b/src/trx-server/src/main.rs index 3110095f..3501e3e1 100644 --- a/src/trx-server/src/main.rs +++ b/src/trx-server/src/main.rs @@ -4,6 +4,7 @@ mod audio; mod config; +mod decoder_history; mod error; mod history_policy; mod history_store;