use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use crate::engine::traits::{BarSink, SinkError};
use crate::live_engine::CompletedBar;
use super::config::ClickHouseWriterConfig;
use super::flush_thread::{FlushCommand, FlushThreadMetrics, spawn_flush_thread};
use super::row::ClickHouseBarRow;
pub struct ClickHouseWriterSink {
tx: std::sync::mpsc::SyncSender<FlushCommand>,
flush_thread: Option<std::thread::JoinHandle<()>>,
metrics: Arc<FlushThreadMetrics>,
bars_sent: AtomicU64,
bars_dropped: AtomicU64,
bars_rejected_ghost: AtomicU64,
bars_rejected_megaspan: AtomicU64,
}
fn assert_not_prod_write_under_test(url: &str) {
let test_ctx = std::env::var_os("NEXTEST").is_some()
|| std::env::var("OPENDEVIATIONBAR_ENV").as_deref() == Ok("test");
if !test_ctx {
return;
}
let loopback = url.contains("localhost") || url.contains("127.0.0.1") || url.contains("[::1]");
assert!(
loopback,
"#660 prod-write fence: refusing ClickHouse sink to '{url}' under test context \
(NEXTEST / OPENDEVIATIONBAR_ENV=test). Point the test at a loopback mock \
(wiremock) or unset OPENDEVIATIONBAR_CH_HOSTS for the test."
);
}
impl ClickHouseWriterSink {
pub fn new(config: ClickHouseWriterConfig) -> Self {
assert_not_prod_write_under_test(&config.url);
let (tx, rx) = std::sync::mpsc::sync_channel(config.channel_capacity);
let metrics = Arc::new(FlushThreadMetrics::default());
let flush_thread = spawn_flush_thread(rx, config, Arc::clone(&metrics));
Self {
tx,
flush_thread: Some(flush_thread),
metrics,
bars_sent: AtomicU64::new(0),
bars_dropped: AtomicU64::new(0),
bars_rejected_ghost: AtomicU64::new(0),
bars_rejected_megaspan: AtomicU64::new(0),
}
}
pub fn metrics(&self) -> &FlushThreadMetrics {
&self.metrics
}
pub fn shared_metrics(&self) -> Arc<FlushThreadMetrics> {
Arc::clone(&self.metrics)
}
pub fn bars_sent(&self) -> u64 {
self.bars_sent.load(Ordering::Relaxed)
}
pub fn bars_dropped(&self) -> u64 {
self.bars_dropped.load(Ordering::Relaxed)
}
pub fn bars_rejected_ghost(&self) -> u64 {
self.bars_rejected_ghost.load(Ordering::Relaxed)
}
pub fn bars_rejected_megaspan(&self) -> u64 {
self.bars_rejected_megaspan.load(Ordering::Relaxed)
}
}
impl BarSink for ClickHouseWriterSink {
fn on_bar(&mut self, bar: &CompletedBar) -> Result<(), SinkError> {
if bar.bar.first_agg_trade_id == 0 || bar.bar.last_agg_trade_id == 0 {
self.bars_rejected_ghost.fetch_add(1, Ordering::Relaxed);
tracing::error!(
symbol = %bar.symbol,
threshold_decimal_bps = bar.threshold_decimal_bps,
first_agg_trade_id = bar.bar.first_agg_trade_id,
last_agg_trade_id = bar.bar.last_agg_trade_id,
open_time_us = bar.bar.open_time,
close_time_us = bar.bar.close_time,
open = bar.bar.open.to_f64(),
high = bar.bar.high.to_f64(),
low = bar.bar.low.to_f64(),
close = bar.bar.close.to_f64(),
volume_raw = bar.bar.volume,
individual_trade_count = bar.bar.individual_trade_count,
agg_record_count = bar.bar.agg_record_count,
"#363 ghost bar rejected at CH sink boundary (first_agg_trade_id or last_agg_trade_id == 0 — live engine must never emit TID=0)"
);
return Ok(());
}
const MAX_BAR_SPAN_TRADES: i64 = 1_000_000;
let span = bar.bar.last_agg_trade_id - bar.bar.first_agg_trade_id;
if span > MAX_BAR_SPAN_TRADES {
self.bars_rejected_megaspan.fetch_add(1, Ordering::Relaxed);
tracing::error!(
symbol = %bar.symbol,
threshold_decimal_bps = bar.threshold_decimal_bps,
first_agg_trade_id = bar.bar.first_agg_trade_id,
last_agg_trade_id = bar.bar.last_agg_trade_id,
span = span,
open_time_us = bar.bar.open_time,
close_time_us = bar.bar.close_time,
"#641 mega-span bar rejected at CH sink (span > 1M trades — stale floor suspected, checkpoint restore before seed order?)"
);
if let Err(dl_err) = super::dead_letter::write_dead_letter_single(
ClickHouseBarRow::from_completed_bar(bar),
) {
tracing::error!(error = %dl_err, "dead-letter write failed for mega-span bar");
}
return Ok(());
}
let row = ClickHouseBarRow::from_completed_bar(bar);
match self.tx.try_send(FlushCommand::Bar(Box::new(row))) {
Ok(()) => {
self.bars_sent.fetch_add(1, Ordering::Relaxed);
Ok(())
}
Err(std::sync::mpsc::TrySendError::Full(cmd)) => {
self.bars_dropped.fetch_add(1, Ordering::Relaxed);
if let FlushCommand::Bar(boxed_row) = cmd {
match super::dead_letter::write_dead_letter_single(*boxed_row) {
Ok(path) => {
tracing::warn!(
path = %path.display(),
"backpressure - bar dead-lettered to Parquet"
);
}
Err(dl_err) => {
tracing::error!(
error = %dl_err,
"CRITICAL: backpressure AND dead-letter write failed -- bar lost"
);
}
}
} else {
tracing::warn!(
"clickhouse flush thread backpressure - non-bar command dropped"
);
}
Err(SinkError::Recoverable(
"clickhouse flush thread backpressure".into(),
))
}
Err(std::sync::mpsc::TrySendError::Disconnected(_)) => Err(SinkError::Unrecoverable(
"clickhouse flush thread died".into(),
)),
}
}
fn flush(&mut self) -> Result<(), SinkError> {
self.tx
.send(FlushCommand::Flush)
.map_err(|_| SinkError::Unrecoverable("clickhouse flush thread died".into()))
}
fn name(&self) -> &str {
"clickhouse"
}
}
impl Drop for ClickHouseWriterSink {
fn drop(&mut self) {
let _ = self.tx.send(FlushCommand::Shutdown);
if let Some(handle) = self.flush_thread.take() {
let _ = handle.join();
}
}
}
#[cfg(test)]
mod tests {
#[test]
#[should_panic(expected = "#660 prod-write fence")]
fn fence_panics_on_prod_host_under_test_context() {
unsafe { std::env::set_var("OPENDEVIATIONBAR_ENV", "test") };
super::assert_not_prod_write_under_test("http://bigblack:8123");
}
#[test]
fn fence_allows_loopback_under_test_context() {
unsafe { std::env::set_var("OPENDEVIATIONBAR_ENV", "test") };
super::assert_not_prod_write_under_test("http://127.0.0.1:18123");
super::assert_not_prod_write_under_test("http://localhost:8123");
}
#[test]
#[should_panic(expected = "#660 prod-write fence")]
fn fence_fires_via_sink_constructor() {
unsafe { std::env::set_var("OPENDEVIATIONBAR_ENV", "test") };
let config = ClickHouseWriterConfig {
url: "http://bigblack:8123".to_string(),
..ClickHouseWriterConfig::default()
};
let _sink = ClickHouseWriterSink::new(config);
}
use super::*;
use opendeviationbar_core::{FixedPoint, OpenDeviationBar};
use std::sync::Arc;
fn make_bar() -> CompletedBar {
make_bar_with_ref_id(1)
}
fn make_bar_with_ref_id(ref_id: i64) -> CompletedBar {
let trade = opendeviationbar_core::Tick {
ref_id,
price: FixedPoint::from_str("50000.0").unwrap(),
volume: FixedPoint::from_str("1.0").unwrap(),
first_sub_id: 1,
last_sub_id: 1,
timestamp: opendeviationbar_core::normalize_timestamp(1_700_000_000_000),
is_buyer_maker: false,
is_best_match: None,
best_bid: None,
best_ask: None,
};
CompletedBar {
symbol: Arc::from("BTCUSDT"),
threshold_decimal_bps: 250,
bar: OpenDeviationBar::new(&trade),
}
}
#[test]
fn test_sink_name() {
let config = ClickHouseWriterConfig {
url: "http://127.0.0.1:1".to_string(),
channel_capacity: 10,
max_retries: 0,
flush_period_ms: 60_000,
..Default::default()
};
let mut sink = ClickHouseWriterSink::new(config);
assert_eq!(sink.name(), "clickhouse");
let _ = sink.flush();
}
#[test]
fn test_sink_on_bar_sends_to_channel() {
let config = ClickHouseWriterConfig {
url: "http://127.0.0.1:1".to_string(),
channel_capacity: 100,
max_retries: 0,
flush_period_ms: 60_000,
..Default::default()
};
let mut sink = ClickHouseWriterSink::new(config);
let bar = make_bar();
assert!(sink.on_bar(&bar).is_ok());
assert_eq!(sink.bars_sent(), 1);
assert_eq!(sink.bars_dropped(), 0);
}
#[test]
fn test_sink_backpressure() {
let config = ClickHouseWriterConfig {
url: "http://127.0.0.1:1".to_string(),
channel_capacity: 1, max_retries: 0,
flush_period_ms: 60_000,
max_rows: 10_000, ..Default::default()
};
let mut sink = ClickHouseWriterSink::new(config);
let bar = make_bar();
assert!(sink.on_bar(&bar).is_ok());
let result = sink.on_bar(&bar);
assert!(matches!(result, Err(SinkError::Recoverable(_))));
assert_eq!(sink.bars_dropped(), 1);
}
#[test]
fn test_sink_graceful_shutdown() {
let config = ClickHouseWriterConfig {
url: "http://127.0.0.1:1".to_string(),
channel_capacity: 100,
max_retries: 0,
flush_period_ms: 60_000,
..Default::default()
};
let mut sink = ClickHouseWriterSink::new(config);
let bar = make_bar();
sink.on_bar(&bar).unwrap();
sink.on_bar(&bar).unwrap();
assert_eq!(sink.bars_sent(), 2);
drop(sink);
}
#[test]
fn test_sink_rejects_ghost_bar_tid_zero() {
let config = ClickHouseWriterConfig {
url: "http://127.0.0.1:1".to_string(),
channel_capacity: 100,
max_retries: 0,
flush_period_ms: 60_000,
..Default::default()
};
let mut sink = ClickHouseWriterSink::new(config);
let ghost = make_bar_with_ref_id(0);
assert_eq!(ghost.bar.first_agg_trade_id, 0);
assert_eq!(ghost.bar.last_agg_trade_id, 0);
assert!(sink.on_bar(&ghost).is_ok());
assert_eq!(sink.bars_rejected_ghost(), 1);
assert_eq!(sink.bars_sent(), 0);
assert_eq!(sink.bars_dropped(), 0);
let real = make_bar_with_ref_id(42);
assert!(sink.on_bar(&real).is_ok());
assert_eq!(sink.bars_rejected_ghost(), 1);
assert_eq!(sink.bars_sent(), 1);
let ghost2 = make_bar_with_ref_id(0);
assert!(sink.on_bar(&ghost2).is_ok());
assert_eq!(sink.bars_rejected_ghost(), 2);
assert_eq!(sink.bars_sent(), 1);
let _ = sink.flush();
}
}