use opendeviationbar_core::fixed_point::SCALE;
use crate::live_engine::CompletedBar;
const VOLUME_SCALE: f64 = SCALE as f64;
pub const CORE_COLUMNS: &[&str] = &[
"symbol",
"threshold_decimal_bps",
"close_time_us",
"open_time_us",
"open",
"high",
"low",
"close",
"volume",
"vwap",
"buy_volume",
"sell_volume",
"individual_trade_count",
"agg_record_count",
"duration_us",
"ofi",
"vwap_close_deviation",
"price_impact",
"kyle_lambda_proxy",
"trade_intensity",
"volume_per_trade",
"aggression_ratio",
"aggregation_density",
"turnover_imbalance",
"ouroboros_mode",
"exchange_session_sydney",
"exchange_session_tokyo",
"exchange_session_london",
"exchange_session_newyork",
"lookback_trade_count",
"lookback_ofi",
"lookback_duration_us",
"lookback_intensity",
"lookback_vwap_raw",
"lookback_vwap_position",
"lookback_count_imbalance",
"lookback_kyle_lambda",
"lookback_burstiness",
"lookback_volume_skew",
"lookback_volume_kurt",
"lookback_price_range",
"lookback_kaufman_er",
"lookback_garman_klass_vol",
"lookback_hurst",
"lookback_permutation_entropy",
"intra_bull_epoch_density",
"intra_bear_epoch_density",
"intra_bull_excess_gain",
"intra_bear_excess_gain",
"intra_bull_cv",
"intra_bear_cv",
"intra_max_drawdown",
"intra_max_runup",
"intra_trade_count",
"intra_ofi",
"intra_duration_us",
"intra_intensity",
"intra_vwap_position",
"intra_count_imbalance",
"intra_kyle_lambda",
"intra_burstiness",
"intra_volume_skew",
"intra_volume_kurt",
"intra_kaufman_er",
"intra_garman_klass_vol",
"intra_hurst",
"intra_permutation_entropy",
"has_gap",
"gap_trade_count",
"max_gap_duration_us",
"is_exchange_gap",
"first_agg_trade_id",
"last_agg_trade_id",
"is_orphan",
"cache_key",
"opendeviationbar_version",
"source_start_ts",
"source_end_ts",
];
#[derive(Debug, Clone, clickhouse::Row, serde::Serialize)]
pub struct ClickHouseBarRow {
pub symbol: String,
pub threshold_decimal_bps: u32,
pub close_time_us: i64,
pub open_time_us: i64,
pub open: f64,
pub high: f64,
pub low: f64,
pub close: f64,
pub volume: f64,
pub vwap: f64,
pub buy_volume: f64,
pub sell_volume: f64,
pub individual_trade_count: u32,
pub agg_record_count: u32,
pub duration_us: i64,
pub ofi: f64,
pub vwap_close_deviation: f64,
pub price_impact: f64,
pub kyle_lambda_proxy: f64,
pub trade_intensity: f64,
pub volume_per_trade: f64,
pub aggression_ratio: f64,
pub aggregation_density: f64,
pub turnover_imbalance: f64,
pub ouroboros_mode: String,
pub exchange_session_sydney: u8,
pub exchange_session_tokyo: u8,
pub exchange_session_london: u8,
pub exchange_session_newyork: u8,
pub lookback_trade_count: Option<u32>,
pub lookback_ofi: Option<f64>,
pub lookback_duration_us: Option<i64>,
pub lookback_intensity: Option<f64>,
pub lookback_vwap_raw: Option<f64>,
pub lookback_vwap_position: Option<f64>,
pub lookback_count_imbalance: Option<f64>,
pub lookback_kyle_lambda: Option<f64>,
pub lookback_burstiness: Option<f64>,
pub lookback_volume_skew: Option<f64>,
pub lookback_volume_kurt: Option<f64>,
pub lookback_price_range: Option<f64>,
pub lookback_kaufman_er: Option<f64>,
pub lookback_garman_klass_vol: Option<f64>,
pub lookback_hurst: Option<f64>,
pub lookback_permutation_entropy: Option<f64>,
pub intra_bull_epoch_density: Option<f64>,
pub intra_bear_epoch_density: Option<f64>,
pub intra_bull_excess_gain: Option<f64>,
pub intra_bear_excess_gain: Option<f64>,
pub intra_bull_cv: Option<f64>,
pub intra_bear_cv: Option<f64>,
pub intra_max_drawdown: Option<f64>,
pub intra_max_runup: Option<f64>,
pub intra_trade_count: Option<u32>,
pub intra_ofi: Option<f64>,
pub intra_duration_us: Option<i64>,
pub intra_intensity: Option<f64>,
pub intra_vwap_position: Option<f64>,
pub intra_count_imbalance: Option<f64>,
pub intra_kyle_lambda: Option<f64>,
pub intra_burstiness: Option<f64>,
pub intra_volume_skew: Option<f64>,
pub intra_volume_kurt: Option<f64>,
pub intra_kaufman_er: Option<f64>,
pub intra_garman_klass_vol: Option<f64>,
pub intra_hurst: Option<f64>,
pub intra_permutation_entropy: Option<f64>,
pub has_gap: bool,
pub gap_trade_count: i64,
pub max_gap_duration_us: i64,
pub is_exchange_gap: bool,
pub first_agg_trade_id: i64,
pub last_agg_trade_id: i64,
pub is_orphan: u8,
pub cache_key: String,
pub opendeviationbar_version: String,
pub source_start_ts: i64,
pub source_end_ts: i64,
}
impl ClickHouseBarRow {
pub fn from_completed_bar(completed: &CompletedBar) -> Self {
let bar = &completed.bar;
Self {
symbol: completed.symbol.to_string(),
threshold_decimal_bps: completed.threshold_decimal_bps,
close_time_us: bar.close_time,
open_time_us: bar.open_time,
open: bar.open.to_f64(),
high: bar.high.to_f64(),
low: bar.low.to_f64(),
close: bar.close.to_f64(),
volume: bar.volume as f64 / VOLUME_SCALE,
vwap: bar.vwap.to_f64(),
buy_volume: bar.buy_volume as f64 / VOLUME_SCALE,
sell_volume: bar.sell_volume as f64 / VOLUME_SCALE,
individual_trade_count: bar.individual_trade_count,
agg_record_count: bar.agg_record_count,
duration_us: bar.duration_us,
ofi: bar.ofi,
vwap_close_deviation: bar.vwap_close_deviation,
price_impact: bar.price_impact,
kyle_lambda_proxy: bar.kyle_lambda_proxy,
trade_intensity: bar.trade_intensity,
volume_per_trade: bar.volume_per_trade,
aggression_ratio: bar.aggression_ratio,
aggregation_density: bar.aggregation_density_f64, turnover_imbalance: bar.turnover_imbalance,
ouroboros_mode: "aion".to_string(),
exchange_session_sydney: 0,
exchange_session_tokyo: 0,
exchange_session_london: 0,
exchange_session_newyork: 0,
lookback_trade_count: bar.lookback_trade_count,
lookback_ofi: bar.lookback_ofi,
lookback_duration_us: bar.lookback_duration_us,
lookback_intensity: bar.lookback_intensity,
lookback_vwap_raw: bar.lookback_vwap_raw.map(|v| v as f64 / VOLUME_SCALE),
lookback_vwap_position: bar.lookback_vwap_position,
lookback_count_imbalance: bar.lookback_count_imbalance,
lookback_kyle_lambda: bar.lookback_kyle_lambda,
lookback_burstiness: bar.lookback_burstiness,
lookback_volume_skew: bar.lookback_volume_skew,
lookback_volume_kurt: bar.lookback_volume_kurt,
lookback_price_range: bar.lookback_price_range,
lookback_kaufman_er: bar.lookback_kaufman_er,
lookback_garman_klass_vol: bar.lookback_garman_klass_vol,
lookback_hurst: bar.lookback_hurst,
lookback_permutation_entropy: bar.lookback_permutation_entropy,
intra_bull_epoch_density: bar.intra_bull_epoch_density,
intra_bear_epoch_density: bar.intra_bear_epoch_density,
intra_bull_excess_gain: bar.intra_bull_excess_gain,
intra_bear_excess_gain: bar.intra_bear_excess_gain,
intra_bull_cv: bar.intra_bull_cv,
intra_bear_cv: bar.intra_bear_cv,
intra_max_drawdown: bar.intra_max_drawdown,
intra_max_runup: bar.intra_max_runup,
intra_trade_count: bar.intra_trade_count,
intra_ofi: bar.intra_ofi,
intra_duration_us: bar.intra_duration_us,
intra_intensity: bar.intra_intensity,
intra_vwap_position: bar.intra_vwap_position,
intra_count_imbalance: bar.intra_count_imbalance,
intra_kyle_lambda: bar.intra_kyle_lambda,
intra_burstiness: bar.intra_burstiness,
intra_volume_skew: bar.intra_volume_skew,
intra_volume_kurt: bar.intra_volume_kurt,
intra_kaufman_er: bar.intra_kaufman_er,
intra_garman_klass_vol: bar.intra_garman_klass_vol,
intra_hurst: bar.intra_hurst,
intra_permutation_entropy: bar.intra_permutation_entropy,
has_gap: bar.has_gap,
gap_trade_count: bar.gap_trade_count,
max_gap_duration_us: bar.max_gap_duration_us,
is_exchange_gap: bar.is_exchange_gap,
first_agg_trade_id: bar.first_agg_trade_id,
last_agg_trade_id: bar.last_agg_trade_id,
is_orphan: 0,
cache_key: String::new(),
opendeviationbar_version: env!("CARGO_PKG_VERSION").to_string(),
source_start_ts: 0,
source_end_ts: 0,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use opendeviationbar_core::fixed_point::FixedPoint;
use opendeviationbar_core::OpenDeviationBar;
use std::sync::Arc;
fn test_completed_bar() -> CompletedBar {
let mut bar = OpenDeviationBar::default();
bar.open = FixedPoint::from_str("50000.12345678").unwrap();
bar.high = FixedPoint::from_str("50100.0").unwrap();
bar.low = FixedPoint::from_str("49900.0").unwrap();
bar.close = FixedPoint::from_str("50050.0").unwrap();
bar.vwap = FixedPoint::from_str("50025.0").unwrap();
bar.open_time = 1_700_000_000_000_000;
bar.close_time = 1_700_000_100_000_000;
bar.volume = 5_000_000_000; bar.buy_volume = 3_000_000_000; bar.sell_volume = 2_000_000_000;
bar.individual_trade_count = 100;
bar.agg_record_count = 50;
bar.first_trade_id = 1000;
bar.last_trade_id = 1099;
bar.first_agg_trade_id = 500;
bar.last_agg_trade_id = 549;
bar.duration_us = 100_000_000;
bar.ofi = 0.2;
bar.vwap_close_deviation = 0.1;
bar.price_impact = 0.001;
bar.kyle_lambda_proxy = 0.5;
bar.trade_intensity = 1000.0;
bar.volume_per_trade = 0.5;
bar.aggression_ratio = 1.5;
bar.aggregation_density_f64 = 2.0;
bar.turnover_imbalance = 0.2;
bar.lookback_trade_count = Some(200);
bar.lookback_ofi = Some(0.1);
bar.lookback_vwap_raw = Some(5_002_500_000_000);
bar.has_gap = false;
bar.gap_trade_count = 0;
bar.max_gap_duration_us = 0;
bar.is_exchange_gap = false;
CompletedBar {
symbol: Arc::from("BTCUSDT"),
threshold_decimal_bps: 250,
bar,
}
}
#[test]
fn test_from_completed_bar_prices() {
let completed = test_completed_bar();
let row = ClickHouseBarRow::from_completed_bar(&completed);
let epsilon = 1e-8;
assert!((row.open - 50000.12345678).abs() < epsilon);
assert!((row.high - 50100.0).abs() < epsilon);
assert!((row.low - 49900.0).abs() < epsilon);
assert!((row.close - 50050.0).abs() < epsilon);
assert!((row.vwap - 50025.0).abs() < epsilon);
}
#[test]
fn test_from_completed_bar_volumes() {
let completed = test_completed_bar();
let row = ClickHouseBarRow::from_completed_bar(&completed);
let epsilon = 1e-8;
assert!((row.volume - 50.0).abs() < epsilon);
assert!((row.buy_volume - 30.0).abs() < epsilon);
assert!((row.sell_volume - 20.0).abs() < epsilon);
}
#[test]
fn test_from_completed_bar_lookback_vwap() {
let completed = test_completed_bar();
let row = ClickHouseBarRow::from_completed_bar(&completed);
let vwap = row.lookback_vwap_raw.expect("lookback_vwap_raw should be Some");
let expected = 5_002_500_000_000_f64 / 100_000_000.0;
assert!((vwap - expected).abs() < 1e-6);
}
#[test]
fn test_from_completed_bar_aggregation_density() {
let completed = test_completed_bar();
let row = ClickHouseBarRow::from_completed_bar(&completed);
assert!((row.aggregation_density - 2.0).abs() < f64::EPSILON);
}
#[test]
fn test_from_completed_bar_hardcoded_fields() {
let completed = test_completed_bar();
let row = ClickHouseBarRow::from_completed_bar(&completed);
assert_eq!(row.ouroboros_mode, "aion");
assert_eq!(row.exchange_session_sydney, 0);
assert_eq!(row.exchange_session_tokyo, 0);
assert_eq!(row.exchange_session_london, 0);
assert_eq!(row.exchange_session_newyork, 0);
assert_eq!(row.is_orphan, 0);
}
#[test]
fn test_computed_columns_not_in_struct() {
let completed = test_completed_bar();
let row = ClickHouseBarRow::from_completed_bar(&completed);
let _symbol = &row.symbol;
let _version = &row.opendeviationbar_version;
}
#[test]
fn test_from_completed_bar_timestamps() {
let completed = test_completed_bar();
let row = ClickHouseBarRow::from_completed_bar(&completed);
assert_eq!(row.open_time_us, 1_700_000_000_000_000);
assert_eq!(row.close_time_us, 1_700_000_100_000_000);
}
#[test]
fn test_from_completed_bar_trade_ids() {
let completed = test_completed_bar();
let row = ClickHouseBarRow::from_completed_bar(&completed);
assert_eq!(row.first_agg_trade_id, 500);
assert_eq!(row.last_agg_trade_id, 549);
}
#[test]
fn test_from_completed_bar_symbol_and_threshold() {
let completed = test_completed_bar();
let row = ClickHouseBarRow::from_completed_bar(&completed);
assert_eq!(row.symbol, "BTCUSDT");
assert_eq!(row.threshold_decimal_bps, 250);
}
#[test]
fn test_from_completed_bar_cache_metadata_defaults() {
let completed = test_completed_bar();
let row = ClickHouseBarRow::from_completed_bar(&completed);
assert!(row.cache_key.is_empty());
assert_eq!(row.source_start_ts, 0);
assert_eq!(row.source_end_ts, 0);
assert!(!row.opendeviationbar_version.is_empty());
}
#[test]
fn test_core_columns_count() {
assert!(CORE_COLUMNS.len() >= 75, "Expected 75+ core columns, got {}", CORE_COLUMNS.len());
}
}