use super::types::*;
use crate::ring_buffer::ConcurrentRingBuffer;
use opendeviationbar_core::processor::OpenDeviationBarProcessor;
use opendeviationbar_core::{OpenDeviationBar, Tick};
use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::Ordering;
use tokio::sync::watch;
#[cfg(feature = "clickhouse-sink")]
pub(crate) type ExtraSinks =
Option<std::sync::Arc<std::sync::Mutex<Vec<Box<dyn crate::engine::traits::BarSink>>>>>;
#[cfg(not(feature = "clickhouse-sink"))]
pub(crate) type ExtraSinks = Option<()>;
#[allow(clippy::too_many_arguments)]
pub(crate) fn process_trade_through_processors(
trade: &Tick,
processors: &mut [(u32, OpenDeviationBarProcessor)],
last_days: &mut [i64],
last_trade_timestamps_us: &mut [i64], bar_buffer: &ConcurrentRingBuffer<CompletedBar>,
metrics: &LiveEngineMetrics,
symbol_arc: &Arc<str>,
symbol: &str,
forming_tx_map: &HashMap<u32, watch::Sender<Option<FormingBar>>>,
committed_floors: &mut HashMap<u32, i64>,
last_forming_update: &mut HashMap<u32, i64>, ouroboros_mode: OuroborosMode, extra_sinks: &ExtraSinks, ) {
let _ = &extra_sinks; for (idx, (threshold, processor)) in processors.iter_mut().enumerate() {
if let OuroborosMode::Week { max_gap_us } = ouroboros_mode
&& let Some(orphan) = maybe_reset_at_week_gap(
processor,
trade.timestamp,
&mut last_trade_timestamps_us[idx],
max_gap_us,
)
{
if let Some(&floor) = committed_floors.get(threshold)
&& orphan.last_agg_trade_id > 0
&& orphan.last_agg_trade_id <= floor
{
tracing::debug!(%symbol, threshold, last_tid = orphan.last_agg_trade_id, floor, "rust-dedup: skipping week-gap orphan (within committed range)");
continue;
}
metrics.bars_emitted.fetch_add(1, Ordering::Relaxed);
let completed = CompletedBar {
symbol: symbol_arc.clone(),
threshold_decimal_bps: *threshold,
bar: orphan,
};
committed_floors.insert(*threshold, completed.bar.last_agg_trade_id);
let was_added = bar_buffer.push(completed.clone());
if !was_added {
metrics.dropped_bars.fetch_add(1, Ordering::Relaxed);
metrics.backpressure_events.fetch_add(1, Ordering::Relaxed);
tracing::warn!(%symbol, threshold, "ring buffer full, old bar dropped (week-gap orphan)");
}
#[cfg(feature = "clickhouse-sink")]
if let Some(sinks) = extra_sinks {
crate::clickhouse_writer::guards::dispatch_to_sinks(&completed, sinks);
}
}
if let Some(orphan) = maybe_reset_at_midnight(
processor,
trade.timestamp,
&mut last_days[idx],
ouroboros_mode,
) {
if let Some(&floor) = committed_floors.get(threshold)
&& orphan.last_agg_trade_id > 0
&& orphan.last_agg_trade_id <= floor
{
tracing::debug!(%symbol, threshold, last_tid = orphan.last_agg_trade_id, floor, "rust-dedup: skipping orphan (within committed range)");
continue;
}
metrics.bars_emitted.fetch_add(1, Ordering::Relaxed);
let completed = CompletedBar {
symbol: symbol_arc.clone(),
threshold_decimal_bps: *threshold,
bar: orphan,
};
committed_floors.insert(*threshold, completed.bar.last_agg_trade_id);
let was_added = bar_buffer.push(completed.clone());
if !was_added {
metrics.dropped_bars.fetch_add(1, Ordering::Relaxed);
metrics.backpressure_events.fetch_add(1, Ordering::Relaxed);
tracing::warn!(%symbol, threshold, "ring buffer full, old bar dropped (midnight orphan)");
}
#[cfg(feature = "clickhouse-sink")]
if let Some(sinks) = extra_sinks {
crate::clickhouse_writer::guards::dispatch_to_sinks(&completed, sinks);
}
}
match processor.process_single_trade(trade) {
Ok(Some(bar)) => {
if let Some(&floor) = committed_floors.get(threshold)
&& bar.last_agg_trade_id > 0
&& bar.last_agg_trade_id <= floor
{
tracing::warn!(
%symbol, threshold,
first_tid = bar.first_agg_trade_id,
last_tid = bar.last_agg_trade_id,
floor,
"rust-dedup: SUPPRESSING bar (last_tid <= committed floor) — #345 telemetry"
);
metrics.bars_suppressed.fetch_add(1, Ordering::Relaxed);
if let Some(tx) = forming_tx_map.get(threshold) {
let _ = tx.send(None);
}
continue;
}
if let Some(tx) = forming_tx_map.get(threshold) {
let _ = tx.send(None);
}
metrics.bars_emitted.fetch_add(1, Ordering::Relaxed);
let completed = CompletedBar {
symbol: symbol_arc.clone(),
threshold_decimal_bps: *threshold,
bar,
};
if let Some(&floor) = committed_floors.get(threshold) {
let expected_first = floor + 1;
if completed.bar.first_agg_trade_id > expected_first
&& completed.bar.first_agg_trade_id > 0
&& floor > 0
{
let gap = completed.bar.first_agg_trade_id - expected_first;
tracing::warn!(
%symbol, threshold,
expected_first,
actual_first = completed.bar.first_agg_trade_id,
prev_last = floor,
gap,
bar_last = completed.bar.last_agg_trade_id,
bar_trades = completed.bar.last_agg_trade_id - completed.bar.first_agg_trade_id + 1,
"STATHERA-GAP-AT-EMISSION: bar first_tid != prev_last+1 — trades missing from processor"
);
}
}
if *threshold == 100 {
tracing::info!(
%symbol,
first = completed.bar.first_agg_trade_id,
last = completed.bar.last_agg_trade_id,
trades = completed.bar.last_agg_trade_id - completed.bar.first_agg_trade_id + 1,
"BAR100-EMIT"
);
}
committed_floors.insert(*threshold, completed.bar.last_agg_trade_id);
let was_added = bar_buffer.push(completed.clone());
if !was_added {
metrics.dropped_bars.fetch_add(1, Ordering::Relaxed);
metrics.backpressure_events.fetch_add(1, Ordering::Relaxed);
tracing::warn!(%symbol, threshold, "ring buffer full, old bar dropped");
}
#[cfg(feature = "clickhouse-sink")]
if let Some(sinks) = extra_sinks {
crate::clickhouse_writer::guards::dispatch_to_sinks(&completed, sinks);
}
let current_depth = bar_buffer.len() as u64;
let _ = metrics
.max_queue_depth
.fetch_max(current_depth, Ordering::Relaxed);
last_forming_update.insert(*threshold, 0);
}
Ok(None) => {
let last_us = last_forming_update.entry(*threshold).or_insert(0);
if trade.timestamp - *last_us >= 1_000_000 {
if let Some(tx) = forming_tx_map.get(threshold)
&& let Some(incomplete) = processor.get_incomplete_bar()
{
let _ = tx.send(Some(FormingBar {
symbol: symbol_arc.clone(),
threshold_decimal_bps: *threshold,
bar: incomplete,
last_trade_timestamp_us: trade.timestamp,
}));
*last_us = trade.timestamp;
}
}
}
Err(e) => {
tracing::warn!(%symbol, threshold = *threshold, ?e, "trade processing error");
}
}
}
}
pub(crate) fn maybe_reset_at_midnight(
processor: &mut OpenDeviationBarProcessor,
trade_timestamp_us: i64,
last_day: &mut i64,
ouroboros_mode: OuroborosMode,
) -> Option<OpenDeviationBar> {
if matches!(
ouroboros_mode,
OuroborosMode::Aion | OuroborosMode::Week { .. }
) {
let trade_day = trade_timestamp_us / DAY_US;
*last_day = trade_day;
return None;
}
let trade_day = trade_timestamp_us / DAY_US;
if *last_day >= 0 && trade_day != *last_day {
*last_day = trade_day;
processor.reset_at_ouroboros()
} else {
*last_day = trade_day;
None
}
}
pub(crate) fn maybe_reset_at_week_gap(
processor: &mut OpenDeviationBarProcessor,
trade_timestamp_us: i64,
last_trade_ts: &mut i64,
max_gap_us: i64,
) -> Option<OpenDeviationBar> {
if *last_trade_ts > 0 {
let gap = trade_timestamp_us - *last_trade_ts;
if gap > max_gap_us {
*last_trade_ts = trade_timestamp_us;
return processor.reset_at_ouroboros();
}
}
*last_trade_ts = trade_timestamp_us;
None
}
#[cfg(test)]
mod tests {
use super::*;
use opendeviationbar_core::processor::OpenDeviationBarProcessor;
use opendeviationbar_core::{FixedPoint, Tick};
const FOUR_HOURS_US: i64 = 14_400_000_000;
fn make_tick(id: i64, price_f64: f64, timestamp_us: i64) -> Tick {
let price_str = format!("{price_f64:.8}");
Tick {
ref_id: id,
price: FixedPoint::from_str(&price_str).unwrap(),
volume: FixedPoint::from_str("1.00000000").unwrap(),
first_sub_id: id,
last_sub_id: id,
timestamp: timestamp_us,
is_buyer_maker: false,
is_best_match: None,
best_bid: None,
best_ask: None,
}
}
#[test]
fn test_week_variant_construction_and_traits() {
let mode = OuroborosMode::Week {
max_gap_us: FOUR_HOURS_US,
};
let mode2 = mode; let mode3 = mode.clone(); assert_eq!(mode, mode2); assert_eq!(mode, mode3);
assert_ne!(mode, OuroborosMode::Aion);
assert_ne!(mode, OuroborosMode::Day);
let mode_different = OuroborosMode::Week {
max_gap_us: 1_000_000,
};
assert_ne!(mode, mode_different);
}
#[test]
fn test_week_gap_first_trade_no_reset() {
let mut processor = OpenDeviationBarProcessor::new(250).unwrap();
let mut last_ts: i64 = 0;
let tick = make_tick(1, 100.0, 1_000_000_000_000);
let _ = processor.process_single_trade(&tick);
assert!(processor.get_incomplete_bar().is_some());
let result = maybe_reset_at_week_gap(
&mut processor,
1_000_000_000_000,
&mut last_ts,
FOUR_HOURS_US,
);
assert!(result.is_none(), "first trade should not trigger reset");
assert_eq!(last_ts, 1_000_000_000_000);
}
#[test]
fn test_week_gap_small_gap_no_reset() {
let mut processor = OpenDeviationBarProcessor::new(250).unwrap();
let base_ts = 1_000_000_000_000i64;
let tick1 = make_tick(1, 100.0, base_ts);
let _ = processor.process_single_trade(&tick1);
let mut last_ts = base_ts;
let next_ts = base_ts + 3_600_000_000;
let result = maybe_reset_at_week_gap(&mut processor, next_ts, &mut last_ts, FOUR_HOURS_US);
assert!(result.is_none(), "gap < threshold should not trigger reset");
assert_eq!(last_ts, next_ts, "last_ts should be updated");
assert!(processor.get_incomplete_bar().is_some());
}
#[test]
fn test_week_gap_large_gap_triggers_reset() {
let mut processor = OpenDeviationBarProcessor::new(250).unwrap();
let base_ts = 1_000_000_000_000i64;
let tick1 = make_tick(1, 100.0, base_ts);
let _ = processor.process_single_trade(&tick1);
let tick2 = make_tick(2, 100.001, base_ts + 1_000_000);
let _ = processor.process_single_trade(&tick2);
assert!(processor.get_incomplete_bar().is_some());
let mut last_ts = base_ts + 1_000_000;
let weekend_ts = base_ts + 18_000_000_000; let result =
maybe_reset_at_week_gap(&mut processor, weekend_ts, &mut last_ts, FOUR_HOURS_US);
assert!(result.is_some(), "gap > threshold should trigger reset");
let orphan = result.unwrap();
assert_eq!(orphan.first_agg_trade_id, 1);
assert_eq!(orphan.last_agg_trade_id, 2);
assert_eq!(last_ts, weekend_ts);
assert!(processor.get_incomplete_bar().is_none());
}
#[test]
fn test_week_gap_updates_last_ts_on_both_paths() {
let mut processor = OpenDeviationBarProcessor::new(250).unwrap();
let mut last_ts = 1_000_000_000_000i64;
let next_ts = last_ts + 1_000_000; let tick = make_tick(1, 100.0, next_ts);
let _ = processor.process_single_trade(&tick);
let result = maybe_reset_at_week_gap(&mut processor, next_ts, &mut last_ts, FOUR_HOURS_US);
assert!(result.is_none());
assert_eq!(last_ts, next_ts, "non-reset path should update last_ts");
let weekend_ts = next_ts + 20_000_000_000; let result = maybe_reset_at_week_gap(
&mut processor,
weekend_ts,
&mut last_ts,
1_000_000, );
assert!(
result.is_some(),
"should trigger reset with small threshold"
);
assert_eq!(last_ts, weekend_ts, "reset path should update last_ts");
}
#[test]
fn test_midnight_returns_none_for_week_mode() {
let mut processor = OpenDeviationBarProcessor::new(250).unwrap();
let before_midnight = 86_399_000_000i64; let tick1 = make_tick(1, 100.0, before_midnight);
let _ = processor.process_single_trade(&tick1);
let mut last_day = -1i64;
let mode = OuroborosMode::Week {
max_gap_us: FOUR_HOURS_US,
};
let result = maybe_reset_at_midnight(&mut processor, before_midnight, &mut last_day, mode);
assert!(result.is_none());
let after_midnight = 86_401_000_000i64; let result = maybe_reset_at_midnight(&mut processor, after_midnight, &mut last_day, mode);
assert!(
result.is_none(),
"Week mode should not trigger midnight reset"
);
assert!(processor.get_incomplete_bar().is_some());
}
#[test]
fn test_week_checkpoint_roundtrip() {
let mut processor = OpenDeviationBarProcessor::new(250).unwrap();
let base_ts = 1_000_000_000_000i64;
for i in 0..5i64 {
let tick = make_tick(i + 1, 100.0 + (i as f64) * 0.001, base_ts + i * 1_000_000);
let _ = processor.process_single_trade(&tick);
}
assert!(processor.get_incomplete_bar().is_some());
let bar_before = processor.get_incomplete_bar().unwrap();
assert_eq!(bar_before.first_agg_trade_id, 1);
assert_eq!(bar_before.last_agg_trade_id, 5);
let mut last_ts = base_ts + 4_000_000;
let weekend_ts = base_ts + 20_000_000_000; let orphan = maybe_reset_at_week_gap(
&mut processor,
weekend_ts,
&mut last_ts,
1_000_000, );
assert!(orphan.is_some(), "should emit orphan at week gap");
let orphan_bar = orphan.unwrap();
assert_eq!(orphan_bar.first_agg_trade_id, 1);
assert_eq!(orphan_bar.last_agg_trade_id, 5);
let checkpoint = processor.create_checkpoint("TESTUSDT");
assert!(
!checkpoint.has_incomplete_bar(),
"processor should be fresh after reset"
);
let restored = OpenDeviationBarProcessor::from_checkpoint(checkpoint).unwrap();
let mut restored = restored;
let new_tick = make_tick(100, 200.0, weekend_ts + 1_000_000);
let result = restored.process_single_trade(&new_tick).unwrap();
assert!(result.is_none(), "first trade opens bar, no breach");
let new_bar = restored.get_incomplete_bar().unwrap();
assert_eq!(
new_bar.first_agg_trade_id, 100,
"new bar starts from new trade"
);
assert_eq!(new_bar.last_agg_trade_id, 100);
}
}