use std::path::{Path, PathBuf};
use arrow_array::RecordBatch;
use arrow_array::builder::{
BooleanBuilder, Float64Builder, Int64Builder, StringBuilder, UInt8Builder, UInt32Builder,
};
use arrow_schema::{DataType, Field, Schema};
use super::row::ClickHouseBarRow;
pub const DEAD_LETTER_DIR: &str = "/tmp/opendeviationbar-dead-letter";
#[derive(Debug)]
pub enum DeadLetterError {
Io(std::io::Error),
Arrow(String),
}
impl std::fmt::Display for DeadLetterError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
DeadLetterError::Io(e) => write!(f, "dead-letter I/O error: {e}"),
DeadLetterError::Arrow(e) => write!(f, "dead-letter Arrow error: {e}"),
}
}
}
impl From<std::io::Error> for DeadLetterError {
fn from(e: std::io::Error) -> Self {
DeadLetterError::Io(e)
}
}
impl From<arrow_schema::ArrowError> for DeadLetterError {
fn from(e: arrow_schema::ArrowError) -> Self {
DeadLetterError::Arrow(e.to_string())
}
}
impl From<parquet::errors::ParquetError> for DeadLetterError {
fn from(e: parquet::errors::ParquetError) -> Self {
DeadLetterError::Arrow(format!("Parquet error: {e}"))
}
}
pub fn dead_letter_schema() -> Schema {
Schema::new(vec![
Field::new("symbol", DataType::Utf8, false),
Field::new("threshold_decimal_bps", DataType::UInt32, false),
Field::new("close_time_us", DataType::Int64, false),
Field::new("open_time_us", DataType::Int64, false),
Field::new("open", DataType::Float64, false),
Field::new("high", DataType::Float64, false),
Field::new("low", DataType::Float64, false),
Field::new("close", DataType::Float64, false),
Field::new("volume", DataType::Float64, false),
Field::new("vwap", DataType::Float64, false),
Field::new("buy_volume", DataType::Float64, false),
Field::new("sell_volume", DataType::Float64, false),
Field::new("individual_trade_count", DataType::UInt32, false),
Field::new("agg_record_count", DataType::UInt32, false),
Field::new("duration_us", DataType::Int64, false),
Field::new("ofi", DataType::Float64, false),
Field::new("vwap_close_deviation", DataType::Float64, false),
Field::new("price_impact", DataType::Float64, false),
Field::new("kyle_lambda_proxy", DataType::Float64, false),
Field::new("trade_intensity", DataType::Float64, false),
Field::new("volume_per_trade", DataType::Float64, false),
Field::new("aggression_ratio", DataType::Float64, false),
Field::new("aggregation_density", DataType::Float64, false),
Field::new("turnover_imbalance", DataType::Float64, false),
Field::new("ouroboros_mode", DataType::Utf8, false),
Field::new("exchange_session_sydney", DataType::UInt8, false),
Field::new("exchange_session_tokyo", DataType::UInt8, false),
Field::new("exchange_session_london", DataType::UInt8, false),
Field::new("exchange_session_newyork", DataType::UInt8, false),
Field::new("lookback_trade_count", DataType::UInt32, true),
Field::new("lookback_ofi", DataType::Float64, true),
Field::new("lookback_duration_us", DataType::Int64, true),
Field::new("lookback_intensity", DataType::Float64, true),
Field::new("lookback_vwap_raw", DataType::Float64, true),
Field::new("lookback_vwap_position", DataType::Float64, true),
Field::new("lookback_count_imbalance", DataType::Float64, true),
Field::new("lookback_kyle_lambda", DataType::Float64, true),
Field::new("lookback_burstiness", DataType::Float64, true),
Field::new("lookback_volume_skew", DataType::Float64, true),
Field::new("lookback_volume_kurt", DataType::Float64, true),
Field::new("lookback_price_range", DataType::Float64, true),
Field::new("lookback_kaufman_er", DataType::Float64, true),
Field::new("lookback_garman_klass_vol", DataType::Float64, true),
Field::new("lookback_hurst", DataType::Float64, true),
Field::new("lookback_permutation_entropy", DataType::Float64, true),
Field::new("bar_petrosian_fd", DataType::Float64, true),
Field::new("bar_katz_fd", DataType::Float64, true),
Field::new("bar_dispersion_entropy", DataType::Float64, true),
Field::new("bar_cecp_velocity", DataType::Float64, true),
Field::new("bar_categorical_recurrence_rate", DataType::Float64, true),
Field::new("bar_sign_markov_flux", DataType::Float64, true),
Field::new("bar_ramsey_rothman_bicov_lag1", DataType::Float64, true),
Field::new("bar_ehlers_increment_asymmetry", DataType::Float64, true),
Field::new("bar_cox_stuart_trend_z", DataType::Float64, true),
Field::new("bar_groeneveld_meeden_b3_skewness", DataType::Float64, true),
Field::new("bar_l_kurtosis_tau4", DataType::Float64, true),
Field::new("bar_bartels_rank_vn_ratio", DataType::Float64, true),
Field::new(
"bar_hoeffding_phi_squared_midreturn_duration",
DataType::Float64,
true,
),
Field::new(
"bar_hvg_forward_visibility_horizon_mean",
DataType::Float64,
true,
),
Field::new(
"bar_vg_time_directed_clustering_meangap",
DataType::Float64,
true,
),
Field::new("intra_bull_epoch_density", DataType::Float64, true),
Field::new("intra_bear_epoch_density", DataType::Float64, true),
Field::new("intra_bull_excess_gain", DataType::Float64, true),
Field::new("intra_bear_excess_gain", DataType::Float64, true),
Field::new("intra_bull_cv", DataType::Float64, true),
Field::new("intra_bear_cv", DataType::Float64, true),
Field::new("intra_max_drawdown", DataType::Float64, true),
Field::new("intra_max_runup", DataType::Float64, true),
Field::new("intra_trade_count", DataType::UInt32, true),
Field::new("intra_ofi", DataType::Float64, true),
Field::new("intra_duration_us", DataType::Int64, true),
Field::new("intra_intensity", DataType::Float64, true),
Field::new("intra_vwap_position", DataType::Float64, true),
Field::new("intra_count_imbalance", DataType::Float64, true),
Field::new("intra_kyle_lambda", DataType::Float64, true),
Field::new("intra_burstiness", DataType::Float64, true),
Field::new("intra_volume_skew", DataType::Float64, true),
Field::new("intra_volume_kurt", DataType::Float64, true),
Field::new("intra_kaufman_er", DataType::Float64, true),
Field::new("intra_garman_klass_vol", DataType::Float64, true),
Field::new("intra_hurst", DataType::Float64, true),
Field::new("intra_permutation_entropy", DataType::Float64, true),
Field::new("has_gap", DataType::Boolean, false),
Field::new("gap_trade_count", DataType::Int64, false),
Field::new("max_gap_duration_us", DataType::Int64, false),
Field::new("is_exchange_gap", DataType::Boolean, false),
Field::new("first_agg_trade_id", DataType::Int64, false),
Field::new("last_agg_trade_id", DataType::Int64, false),
Field::new("is_orphan", DataType::UInt8, false),
Field::new("cache_key", DataType::Utf8, false),
Field::new("opendeviationbar_version", DataType::Utf8, false),
Field::new("source_start_ts", DataType::Int64, false),
Field::new("source_end_ts", DataType::Int64, false),
])
}
pub fn rows_to_record_batch_public(
rows: &[ClickHouseBarRow],
schema: &Schema,
) -> Result<RecordBatch, DeadLetterError> {
rows_to_record_batch(rows, schema)
}
fn rows_to_record_batch(
rows: &[ClickHouseBarRow],
schema: &Schema,
) -> Result<RecordBatch, DeadLetterError> {
let n = rows.len();
let mut symbol_b = StringBuilder::with_capacity(n, n * 10);
let mut ouroboros_mode_b = StringBuilder::with_capacity(n, n * 4);
let mut cache_key_b = StringBuilder::with_capacity(n, n * 32);
let mut version_b = StringBuilder::with_capacity(n, n * 10);
let mut threshold_b = UInt32Builder::with_capacity(n);
let mut close_time_b = Int64Builder::with_capacity(n);
let mut open_time_b = Int64Builder::with_capacity(n);
let mut open_b = Float64Builder::with_capacity(n);
let mut high_b = Float64Builder::with_capacity(n);
let mut low_b = Float64Builder::with_capacity(n);
let mut close_b = Float64Builder::with_capacity(n);
let mut volume_b = Float64Builder::with_capacity(n);
let mut vwap_b = Float64Builder::with_capacity(n);
let mut buy_vol_b = Float64Builder::with_capacity(n);
let mut sell_vol_b = Float64Builder::with_capacity(n);
let mut ind_count_b = UInt32Builder::with_capacity(n);
let mut agg_count_b = UInt32Builder::with_capacity(n);
let mut duration_b = Int64Builder::with_capacity(n);
let mut ofi_b = Float64Builder::with_capacity(n);
let mut vwap_dev_b = Float64Builder::with_capacity(n);
let mut price_impact_b = Float64Builder::with_capacity(n);
let mut kyle_b = Float64Builder::with_capacity(n);
let mut intensity_b = Float64Builder::with_capacity(n);
let mut vol_per_trade_b = Float64Builder::with_capacity(n);
let mut aggression_b = Float64Builder::with_capacity(n);
let mut agg_density_b = Float64Builder::with_capacity(n);
let mut turnover_b = Float64Builder::with_capacity(n);
let mut sess_sydney_b = UInt8Builder::with_capacity(n);
let mut sess_tokyo_b = UInt8Builder::with_capacity(n);
let mut sess_london_b = UInt8Builder::with_capacity(n);
let mut sess_ny_b = UInt8Builder::with_capacity(n);
let mut lb_trade_count_b = UInt32Builder::with_capacity(n);
let mut lb_ofi_b = Float64Builder::with_capacity(n);
let mut lb_duration_b = Int64Builder::with_capacity(n);
let mut lb_intensity_b = Float64Builder::with_capacity(n);
let mut lb_vwap_raw_b = Float64Builder::with_capacity(n);
let mut lb_vwap_pos_b = Float64Builder::with_capacity(n);
let mut lb_count_imb_b = Float64Builder::with_capacity(n);
let mut lb_kyle_b = Float64Builder::with_capacity(n);
let mut lb_burst_b = Float64Builder::with_capacity(n);
let mut lb_vol_skew_b = Float64Builder::with_capacity(n);
let mut lb_vol_kurt_b = Float64Builder::with_capacity(n);
let mut lb_price_range_b = Float64Builder::with_capacity(n);
let mut lb_kaufman_b = Float64Builder::with_capacity(n);
let mut lb_gk_vol_b = Float64Builder::with_capacity(n);
let mut lb_hurst_b = Float64Builder::with_capacity(n);
let mut lb_perm_ent_b = Float64Builder::with_capacity(n);
let mut bar_petrosian_fd_b = Float64Builder::with_capacity(n);
let mut bar_katz_fd_b = Float64Builder::with_capacity(n);
let mut bar_dispersion_entropy_b = Float64Builder::with_capacity(n);
let mut bar_cecp_velocity_b = Float64Builder::with_capacity(n);
let mut bar_categorical_recurrence_rate_b = Float64Builder::with_capacity(n);
let mut bar_sign_markov_flux_b = Float64Builder::with_capacity(n);
let mut bar_ramsey_rothman_bicov_lag1_b = Float64Builder::with_capacity(n);
let mut bar_ehlers_increment_asymmetry_b = Float64Builder::with_capacity(n);
let mut bar_cox_stuart_trend_z_b = Float64Builder::with_capacity(n);
let mut bar_groeneveld_meeden_b3_skewness_b = Float64Builder::with_capacity(n);
let mut bar_l_kurtosis_tau4_b = Float64Builder::with_capacity(n);
let mut bar_bartels_rank_vn_ratio_b = Float64Builder::with_capacity(n);
let mut bar_hoeffding_phi_squared_midreturn_duration_b = Float64Builder::with_capacity(n);
let mut bar_hvg_forward_visibility_horizon_mean_b = Float64Builder::with_capacity(n);
let mut bar_vg_time_directed_clustering_meangap_b = Float64Builder::with_capacity(n);
let mut intra_bull_density_b = Float64Builder::with_capacity(n);
let mut intra_bear_density_b = Float64Builder::with_capacity(n);
let mut intra_bull_gain_b = Float64Builder::with_capacity(n);
let mut intra_bear_gain_b = Float64Builder::with_capacity(n);
let mut intra_bull_cv_b = Float64Builder::with_capacity(n);
let mut intra_bear_cv_b = Float64Builder::with_capacity(n);
let mut intra_max_dd_b = Float64Builder::with_capacity(n);
let mut intra_max_ru_b = Float64Builder::with_capacity(n);
let mut intra_trade_count_b = UInt32Builder::with_capacity(n);
let mut intra_ofi_b = Float64Builder::with_capacity(n);
let mut intra_duration_b = Int64Builder::with_capacity(n);
let mut intra_intensity_b = Float64Builder::with_capacity(n);
let mut intra_vwap_pos_b = Float64Builder::with_capacity(n);
let mut intra_count_imb_b = Float64Builder::with_capacity(n);
let mut intra_kyle_b = Float64Builder::with_capacity(n);
let mut intra_burst_b = Float64Builder::with_capacity(n);
let mut intra_vol_skew_b = Float64Builder::with_capacity(n);
let mut intra_vol_kurt_b = Float64Builder::with_capacity(n);
let mut intra_kaufman_b = Float64Builder::with_capacity(n);
let mut intra_gk_vol_b = Float64Builder::with_capacity(n);
let mut intra_hurst_b = Float64Builder::with_capacity(n);
let mut intra_perm_ent_b = Float64Builder::with_capacity(n);
let mut has_gap_b = BooleanBuilder::with_capacity(n);
let mut gap_count_b = Int64Builder::with_capacity(n);
let mut gap_dur_b = Int64Builder::with_capacity(n);
let mut is_exch_gap_b = BooleanBuilder::with_capacity(n);
let mut first_tid_b = Int64Builder::with_capacity(n);
let mut last_tid_b = Int64Builder::with_capacity(n);
let mut is_orphan_b = UInt8Builder::with_capacity(n);
let mut src_start_b = Int64Builder::with_capacity(n);
let mut src_end_b = Int64Builder::with_capacity(n);
for r in rows {
symbol_b.append_value(&r.symbol);
threshold_b.append_value(r.threshold_decimal_bps);
close_time_b.append_value(r.close_time_us);
open_time_b.append_value(r.open_time_us);
open_b.append_value(r.open);
high_b.append_value(r.high);
low_b.append_value(r.low);
close_b.append_value(r.close);
volume_b.append_value(r.volume);
vwap_b.append_value(r.vwap);
buy_vol_b.append_value(r.buy_volume);
sell_vol_b.append_value(r.sell_volume);
ind_count_b.append_value(r.individual_trade_count);
agg_count_b.append_value(r.agg_record_count);
duration_b.append_value(r.duration_us);
ofi_b.append_value(r.ofi);
vwap_dev_b.append_value(r.vwap_close_deviation);
price_impact_b.append_value(r.price_impact);
kyle_b.append_value(r.kyle_lambda_proxy);
intensity_b.append_value(r.trade_intensity);
vol_per_trade_b.append_value(r.volume_per_trade);
aggression_b.append_value(r.aggression_ratio);
agg_density_b.append_value(r.aggregation_density);
turnover_b.append_value(r.turnover_imbalance);
ouroboros_mode_b.append_value(&r.ouroboros_mode);
sess_sydney_b.append_value(r.exchange_session_sydney);
sess_tokyo_b.append_value(r.exchange_session_tokyo);
sess_london_b.append_value(r.exchange_session_london);
sess_ny_b.append_value(r.exchange_session_newyork);
lb_trade_count_b.append_option(r.lookback_trade_count);
lb_ofi_b.append_option(r.lookback_ofi);
lb_duration_b.append_option(r.lookback_duration_us);
lb_intensity_b.append_option(r.lookback_intensity);
lb_vwap_raw_b.append_option(r.lookback_vwap_raw);
lb_vwap_pos_b.append_option(r.lookback_vwap_position);
lb_count_imb_b.append_option(r.lookback_count_imbalance);
lb_kyle_b.append_option(r.lookback_kyle_lambda);
lb_burst_b.append_option(r.lookback_burstiness);
lb_vol_skew_b.append_option(r.lookback_volume_skew);
lb_vol_kurt_b.append_option(r.lookback_volume_kurt);
lb_price_range_b.append_option(r.lookback_price_range);
lb_kaufman_b.append_option(r.lookback_kaufman_er);
lb_gk_vol_b.append_option(r.lookback_garman_klass_vol);
lb_hurst_b.append_option(r.lookback_hurst);
lb_perm_ent_b.append_option(r.lookback_permutation_entropy);
bar_petrosian_fd_b.append_option(r.bar_petrosian_fd);
bar_katz_fd_b.append_option(r.bar_katz_fd);
bar_dispersion_entropy_b.append_option(r.bar_dispersion_entropy);
bar_cecp_velocity_b.append_option(r.bar_cecp_velocity);
bar_categorical_recurrence_rate_b.append_option(r.bar_categorical_recurrence_rate);
bar_sign_markov_flux_b.append_option(r.bar_sign_markov_flux);
bar_ramsey_rothman_bicov_lag1_b.append_option(r.bar_ramsey_rothman_bicov_lag1);
bar_ehlers_increment_asymmetry_b.append_option(r.bar_ehlers_increment_asymmetry);
bar_cox_stuart_trend_z_b.append_option(r.bar_cox_stuart_trend_z);
bar_groeneveld_meeden_b3_skewness_b.append_option(r.bar_groeneveld_meeden_b3_skewness);
bar_l_kurtosis_tau4_b.append_option(r.bar_l_kurtosis_tau4);
bar_bartels_rank_vn_ratio_b.append_option(r.bar_bartels_rank_vn_ratio);
bar_hoeffding_phi_squared_midreturn_duration_b
.append_option(r.bar_hoeffding_phi_squared_midreturn_duration);
bar_hvg_forward_visibility_horizon_mean_b
.append_option(r.bar_hvg_forward_visibility_horizon_mean);
bar_vg_time_directed_clustering_meangap_b
.append_option(r.bar_vg_time_directed_clustering_meangap);
intra_bull_density_b.append_option(r.intra_bull_epoch_density);
intra_bear_density_b.append_option(r.intra_bear_epoch_density);
intra_bull_gain_b.append_option(r.intra_bull_excess_gain);
intra_bear_gain_b.append_option(r.intra_bear_excess_gain);
intra_bull_cv_b.append_option(r.intra_bull_cv);
intra_bear_cv_b.append_option(r.intra_bear_cv);
intra_max_dd_b.append_option(r.intra_max_drawdown);
intra_max_ru_b.append_option(r.intra_max_runup);
intra_trade_count_b.append_option(r.intra_trade_count);
intra_ofi_b.append_option(r.intra_ofi);
intra_duration_b.append_option(r.intra_duration_us);
intra_intensity_b.append_option(r.intra_intensity);
intra_vwap_pos_b.append_option(r.intra_vwap_position);
intra_count_imb_b.append_option(r.intra_count_imbalance);
intra_kyle_b.append_option(r.intra_kyle_lambda);
intra_burst_b.append_option(r.intra_burstiness);
intra_vol_skew_b.append_option(r.intra_volume_skew);
intra_vol_kurt_b.append_option(r.intra_volume_kurt);
intra_kaufman_b.append_option(r.intra_kaufman_er);
intra_gk_vol_b.append_option(r.intra_garman_klass_vol);
intra_hurst_b.append_option(r.intra_hurst);
intra_perm_ent_b.append_option(r.intra_permutation_entropy);
has_gap_b.append_value(r.has_gap);
gap_count_b.append_value(r.gap_trade_count);
gap_dur_b.append_value(r.max_gap_duration_us);
is_exch_gap_b.append_value(r.is_exchange_gap);
first_tid_b.append_value(r.first_agg_trade_id);
last_tid_b.append_value(r.last_agg_trade_id);
is_orphan_b.append_value(r.is_orphan);
cache_key_b.append_value(&r.cache_key);
version_b.append_value(&r.opendeviationbar_version);
src_start_b.append_value(r.source_start_ts);
src_end_b.append_value(r.source_end_ts);
}
let columns: Vec<arrow_array::ArrayRef> = vec![
std::sync::Arc::new(symbol_b.finish()),
std::sync::Arc::new(threshold_b.finish()),
std::sync::Arc::new(close_time_b.finish()),
std::sync::Arc::new(open_time_b.finish()),
std::sync::Arc::new(open_b.finish()),
std::sync::Arc::new(high_b.finish()),
std::sync::Arc::new(low_b.finish()),
std::sync::Arc::new(close_b.finish()),
std::sync::Arc::new(volume_b.finish()),
std::sync::Arc::new(vwap_b.finish()),
std::sync::Arc::new(buy_vol_b.finish()),
std::sync::Arc::new(sell_vol_b.finish()),
std::sync::Arc::new(ind_count_b.finish()),
std::sync::Arc::new(agg_count_b.finish()),
std::sync::Arc::new(duration_b.finish()),
std::sync::Arc::new(ofi_b.finish()),
std::sync::Arc::new(vwap_dev_b.finish()),
std::sync::Arc::new(price_impact_b.finish()),
std::sync::Arc::new(kyle_b.finish()),
std::sync::Arc::new(intensity_b.finish()),
std::sync::Arc::new(vol_per_trade_b.finish()),
std::sync::Arc::new(aggression_b.finish()),
std::sync::Arc::new(agg_density_b.finish()),
std::sync::Arc::new(turnover_b.finish()),
std::sync::Arc::new(ouroboros_mode_b.finish()),
std::sync::Arc::new(sess_sydney_b.finish()),
std::sync::Arc::new(sess_tokyo_b.finish()),
std::sync::Arc::new(sess_london_b.finish()),
std::sync::Arc::new(sess_ny_b.finish()),
std::sync::Arc::new(lb_trade_count_b.finish()),
std::sync::Arc::new(lb_ofi_b.finish()),
std::sync::Arc::new(lb_duration_b.finish()),
std::sync::Arc::new(lb_intensity_b.finish()),
std::sync::Arc::new(lb_vwap_raw_b.finish()),
std::sync::Arc::new(lb_vwap_pos_b.finish()),
std::sync::Arc::new(lb_count_imb_b.finish()),
std::sync::Arc::new(lb_kyle_b.finish()),
std::sync::Arc::new(lb_burst_b.finish()),
std::sync::Arc::new(lb_vol_skew_b.finish()),
std::sync::Arc::new(lb_vol_kurt_b.finish()),
std::sync::Arc::new(lb_price_range_b.finish()),
std::sync::Arc::new(lb_kaufman_b.finish()),
std::sync::Arc::new(lb_gk_vol_b.finish()),
std::sync::Arc::new(lb_hurst_b.finish()),
std::sync::Arc::new(lb_perm_ent_b.finish()),
std::sync::Arc::new(bar_petrosian_fd_b.finish()),
std::sync::Arc::new(bar_katz_fd_b.finish()),
std::sync::Arc::new(bar_dispersion_entropy_b.finish()),
std::sync::Arc::new(bar_cecp_velocity_b.finish()),
std::sync::Arc::new(bar_categorical_recurrence_rate_b.finish()),
std::sync::Arc::new(bar_sign_markov_flux_b.finish()),
std::sync::Arc::new(bar_ramsey_rothman_bicov_lag1_b.finish()),
std::sync::Arc::new(bar_ehlers_increment_asymmetry_b.finish()),
std::sync::Arc::new(bar_cox_stuart_trend_z_b.finish()),
std::sync::Arc::new(bar_groeneveld_meeden_b3_skewness_b.finish()),
std::sync::Arc::new(bar_l_kurtosis_tau4_b.finish()),
std::sync::Arc::new(bar_bartels_rank_vn_ratio_b.finish()),
std::sync::Arc::new(bar_hoeffding_phi_squared_midreturn_duration_b.finish()),
std::sync::Arc::new(bar_hvg_forward_visibility_horizon_mean_b.finish()),
std::sync::Arc::new(bar_vg_time_directed_clustering_meangap_b.finish()),
std::sync::Arc::new(intra_bull_density_b.finish()),
std::sync::Arc::new(intra_bear_density_b.finish()),
std::sync::Arc::new(intra_bull_gain_b.finish()),
std::sync::Arc::new(intra_bear_gain_b.finish()),
std::sync::Arc::new(intra_bull_cv_b.finish()),
std::sync::Arc::new(intra_bear_cv_b.finish()),
std::sync::Arc::new(intra_max_dd_b.finish()),
std::sync::Arc::new(intra_max_ru_b.finish()),
std::sync::Arc::new(intra_trade_count_b.finish()),
std::sync::Arc::new(intra_ofi_b.finish()),
std::sync::Arc::new(intra_duration_b.finish()),
std::sync::Arc::new(intra_intensity_b.finish()),
std::sync::Arc::new(intra_vwap_pos_b.finish()),
std::sync::Arc::new(intra_count_imb_b.finish()),
std::sync::Arc::new(intra_kyle_b.finish()),
std::sync::Arc::new(intra_burst_b.finish()),
std::sync::Arc::new(intra_vol_skew_b.finish()),
std::sync::Arc::new(intra_vol_kurt_b.finish()),
std::sync::Arc::new(intra_kaufman_b.finish()),
std::sync::Arc::new(intra_gk_vol_b.finish()),
std::sync::Arc::new(intra_hurst_b.finish()),
std::sync::Arc::new(intra_perm_ent_b.finish()),
std::sync::Arc::new(has_gap_b.finish()),
std::sync::Arc::new(gap_count_b.finish()),
std::sync::Arc::new(gap_dur_b.finish()),
std::sync::Arc::new(is_exch_gap_b.finish()),
std::sync::Arc::new(first_tid_b.finish()),
std::sync::Arc::new(last_tid_b.finish()),
std::sync::Arc::new(is_orphan_b.finish()),
std::sync::Arc::new(cache_key_b.finish()),
std::sync::Arc::new(version_b.finish()),
std::sync::Arc::new(src_start_b.finish()),
std::sync::Arc::new(src_end_b.finish()),
];
let schema_ref = std::sync::Arc::new(schema.clone());
RecordBatch::try_new(schema_ref, columns).map_err(|e| DeadLetterError::Arrow(e.to_string()))
}
pub fn write_dead_letter(rows: &[ClickHouseBarRow]) -> Result<PathBuf, DeadLetterError> {
if rows.is_empty() {
return Ok(PathBuf::new());
}
let dir = Path::new(DEAD_LETTER_DIR);
std::fs::create_dir_all(dir)?;
let symbol = &rows[0].symbol;
let threshold = rows[0].threshold_decimal_bps;
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_secs();
let filename = format!("{symbol}_{threshold}_{timestamp}.parquet");
let path = dir.join(&filename);
let schema = dead_letter_schema();
let batch = rows_to_record_batch(rows, &schema)?;
let props = parquet::file::properties::WriterProperties::builder()
.set_compression(parquet::basic::Compression::ZSTD(
parquet::basic::ZstdLevel::try_new(3)?,
))
.build();
let file = std::fs::File::create(&path)?;
let mut writer =
parquet::arrow::ArrowWriter::try_new(file, std::sync::Arc::new(schema), Some(props))?;
writer.write(&batch)?;
writer.close()?;
tracing::warn!(
symbol,
threshold,
rows = rows.len(),
path = %path.display(),
"dead-lettered failed bars to Parquet"
);
Ok(path)
}
pub fn write_dead_letter_single(row: ClickHouseBarRow) -> Result<PathBuf, DeadLetterError> {
write_dead_letter(&[row])
}
#[cfg(test)]
mod tests {
use super::super::row::{CORE_COLUMNS, ClickHouseBarRow};
use super::*;
use opendeviationbar_core::OpenDeviationBar;
use opendeviationbar_core::fixed_point::FixedPoint;
use parquet::file::reader::FileReader;
use std::sync::Arc;
use crate::live_engine::CompletedBar;
fn test_row(first_tid: i64, last_tid: i64) -> ClickHouseBarRow {
let mut bar = OpenDeviationBar::default();
bar.open = FixedPoint::from_str("50000.0").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.first_agg_trade_id = first_tid;
bar.last_agg_trade_id = last_tid;
bar.individual_trade_count = 100;
bar.agg_record_count = 50;
bar.duration_us = 100_000_000;
bar.lookback_trade_count = Some(200);
bar.lookback_ofi = Some(0.1);
let completed = CompletedBar {
symbol: Arc::from("BTCUSDT"),
threshold_decimal_bps: 250,
bar,
};
ClickHouseBarRow::from_completed_bar(&completed)
}
fn test_dead_letter_dir() -> PathBuf {
let pid = std::process::id();
let nanos = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_nanos();
let dir =
std::env::temp_dir().join(format!("opendeviationbar-dead-letter-test-{pid}-{nanos}"));
let _ = std::fs::remove_dir_all(&dir);
std::fs::create_dir_all(&dir).unwrap();
dir
}
#[test]
fn test_dead_letter_write_3_rows() {
let dir = test_dead_letter_dir();
let rows = vec![test_row(1, 10), test_row(11, 20), test_row(21, 30)];
let schema = dead_letter_schema();
let batch = rows_to_record_batch(&rows, &schema).unwrap();
let path = dir.join("BTCUSDT_250_12345.parquet");
let props = parquet::file::properties::WriterProperties::builder()
.set_compression(parquet::basic::Compression::ZSTD(
parquet::basic::ZstdLevel::try_new(3).unwrap(),
))
.build();
let file = std::fs::File::create(&path).unwrap();
let mut writer =
parquet::arrow::ArrowWriter::try_new(file, std::sync::Arc::new(schema), Some(props))
.unwrap();
writer.write(&batch).unwrap();
writer.close().unwrap();
assert!(path.exists());
let reader = parquet::arrow::arrow_reader::ParquetRecordBatchReader::try_new(
std::fs::File::open(&path).unwrap(),
1024,
)
.unwrap();
let batches: Vec<RecordBatch> = reader.into_iter().collect::<Result<_, _>>().unwrap();
let total_rows: usize = batches.iter().map(RecordBatch::num_rows).sum();
assert_eq!(total_rows, 3);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn test_dead_letter_write_single() {
let dir = test_dead_letter_dir();
let row = test_row(1, 10);
let schema = dead_letter_schema();
let batch = rows_to_record_batch(&[row], &schema).unwrap();
let path = dir.join("BTCUSDT_250_single.parquet");
let props = parquet::file::properties::WriterProperties::builder()
.set_compression(parquet::basic::Compression::ZSTD(
parquet::basic::ZstdLevel::try_new(3).unwrap(),
))
.build();
let file = std::fs::File::create(&path).unwrap();
let mut writer =
parquet::arrow::ArrowWriter::try_new(file, std::sync::Arc::new(schema), Some(props))
.unwrap();
writer.write(&batch).unwrap();
writer.close().unwrap();
assert!(path.exists());
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn test_dead_letter_schema_nullable_flags() {
let schema = dead_letter_schema();
let fields = schema.fields();
for field in fields.iter() {
let name = field.name();
if name.starts_with("lookback_") || name.starts_with("intra_") {
assert!(field.is_nullable(), "Field {name} should be nullable");
}
}
let non_nullable_names = [
"symbol",
"threshold_decimal_bps",
"open",
"high",
"low",
"close",
"volume",
"has_gap",
"first_agg_trade_id",
"last_agg_trade_id",
];
for name in &non_nullable_names {
let field = schema.field_with_name(name).unwrap();
assert!(!field.is_nullable(), "Field {name} should NOT be nullable");
}
}
#[test]
fn test_dead_letter_column_count() {
let schema = dead_letter_schema();
assert_eq!(
schema.fields().len(),
CORE_COLUMNS.len(),
"Schema column count must match CORE_COLUMNS ({})",
CORE_COLUMNS.len()
);
}
#[test]
fn test_dead_letter_zstd_compression() {
let dir = test_dead_letter_dir();
let row = test_row(1, 10);
let schema = dead_letter_schema();
let batch = rows_to_record_batch(&[row], &schema).unwrap();
let path = dir.join("BTCUSDT_250_zstd.parquet");
let props = parquet::file::properties::WriterProperties::builder()
.set_compression(parquet::basic::Compression::ZSTD(
parquet::basic::ZstdLevel::try_new(3).unwrap(),
))
.build();
let file = std::fs::File::create(&path).unwrap();
let mut writer =
parquet::arrow::ArrowWriter::try_new(file, std::sync::Arc::new(schema), Some(props))
.unwrap();
writer.write(&batch).unwrap();
writer.close().unwrap();
let file = std::fs::File::open(&path).unwrap();
let reader = parquet::file::reader::SerializedFileReader::new(file).unwrap();
let meta = reader.metadata();
let row_group = meta.row_group(0);
let col_meta = row_group.column(0);
assert!(
matches!(col_meta.compression(), parquet::basic::Compression::ZSTD(_)),
"Expected ZSTD compression, got {:?}",
col_meta.compression()
);
let _ = std::fs::remove_dir_all(&dir);
}
#[test]
fn test_dead_letter_empty_slice() {
let result = write_dead_letter(&[]);
assert!(result.is_ok());
let path = result.unwrap();
assert_eq!(path, PathBuf::new());
}
#[test]
fn test_write_roundtrip_artifact() {
let artifact_dir = std::env::temp_dir().join("opendeviationbar-dead-letter-roundtrip");
let _ = std::fs::remove_dir_all(&artifact_dir);
std::fs::create_dir_all(&artifact_dir).unwrap();
let mut row = test_row(1000, 1099);
row.lookback_duration_us = None;
row.lookback_intensity = None;
row.intra_bull_epoch_density = None;
row.intra_hurst = None;
row.lookback_trade_count = Some(500);
row.lookback_ofi = Some(-0.42);
let schema = dead_letter_schema();
let batch = rows_to_record_batch(&[row], &schema).unwrap();
let path = artifact_dir.join("BTCUSDT_250_9999999999.parquet");
let props = parquet::file::properties::WriterProperties::builder()
.set_compression(parquet::basic::Compression::ZSTD(
parquet::basic::ZstdLevel::try_new(3).unwrap(),
))
.build();
let file = std::fs::File::create(&path).unwrap();
let mut writer =
parquet::arrow::ArrowWriter::try_new(file, std::sync::Arc::new(schema), Some(props))
.unwrap();
writer.write(&batch).unwrap();
writer.close().unwrap();
assert!(path.exists());
eprintln!("ROUNDTRIP_ARTIFACT={}", path.display());
}
#[test]
fn test_dead_letter_column_names_match_core_columns() {
let schema = dead_letter_schema();
for (i, field) in schema.fields().iter().enumerate() {
assert_eq!(
field.name(),
CORE_COLUMNS[i],
"Column {i} name mismatch: schema has '{}', CORE_COLUMNS has '{}'",
field.name(),
CORE_COLUMNS[i]
);
}
}
}