From 8106f4987bbd7d463b2caa9ad1a6d4c5d5b930e5 Mon Sep 17 00:00:00 2001 From: jgrusewski Date: Mon, 23 Feb 2026 13:17:52 +0100 Subject: [PATCH] feat(common): add QuestDB client with ring buffer and health monitoring Non-critical path: if QuestDB is unavailable, metrics buffer locally (up to 10,000 entries) and flush when connection is restored. Feature-gated under `questdb` feature. 6 tests. Co-Authored-By: Claude Opus 4.6 --- Cargo.toml | 1 + common/Cargo.toml | 2 + common/src/lib.rs | 2 + common/src/questdb.rs | 302 ++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 307 insertions(+) create mode 100644 common/src/questdb.rs diff --git a/Cargo.toml b/Cargo.toml index 58bcef197..d60e01aa0 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -222,6 +222,7 @@ dashmap = { version = "6.0", features = ["serde"] } bytes = "1.5" smallvec = { version = "1.11", features = ["serde", "const_generics"] } prometheus = "0.14" +questdb-rs = { version = "4", default-features = false } # MINIMAL numerical libraries only - ALL ML/GPU frameworks REMOVED from workspace nalgebra = { version = "0.33", features = ["serde", "rand"] } diff --git a/common/Cargo.toml b/common/Cargo.toml index 3ff6e8ad0..1c3f920c3 100644 --- a/common/Cargo.toml +++ b/common/Cargo.toml @@ -60,6 +60,7 @@ config = { path = "../config" } uuid = { workspace = true, features = ["v4", "serde"] } once_cell.workspace = true rand.workspace = true +questdb-rs = { workspace = true, optional = true } # Metrics (Wave 5) prometheus.workspace = true @@ -76,6 +77,7 @@ tempfile = "3.8" [features] default = ["database"] database = ["sqlx"] +questdb = ["questdb-rs"] [[bench]] name = "ml_strategy_bench" diff --git a/common/src/lib.rs b/common/src/lib.rs index 357defbec..84eaf7236 100644 --- a/common/src/lib.rs +++ b/common/src/lib.rs @@ -39,6 +39,8 @@ pub mod thresholds; pub mod tls; pub mod traits; pub mod types; +#[cfg(feature = "questdb")] +pub mod questdb; // Re-export database types for external use pub use database::{DatabaseConfig, DatabaseError, DatabasePool}; diff --git a/common/src/questdb.rs b/common/src/questdb.rs new file mode 100644 index 000000000..360eac592 --- /dev/null +++ b/common/src/questdb.rs @@ -0,0 +1,302 @@ +//! QuestDB client with ring buffer for graceful degradation +//! +//! Non-critical path: if QuestDB is unavailable, metrics buffer locally +//! and flush when connection is restored. Trading never depends on QuestDB. + +use std::collections::VecDeque; +use std::sync::atomic::{AtomicBool, Ordering}; +use std::sync::Arc; +use std::time::{Duration, Instant}; +use tokio::sync::Mutex; +use tracing::{debug, info, warn}; + +/// Maximum entries in the ring buffer before oldest are dropped +const DEFAULT_BUFFER_CAPACITY: usize = 10_000; + +/// Health check interval +const HEALTH_CHECK_INTERVAL: Duration = Duration::from_secs(30); + +/// A metric entry waiting to be flushed to QuestDB +#[derive(Debug, Clone)] +pub struct MetricEntry { + pub table: String, + pub columns: Vec<(String, MetricValue)>, + pub timestamp_ns: i64, + pub created_at: Instant, +} + +/// Supported metric value types for ILP +#[derive(Debug, Clone)] +pub enum MetricValue { + Symbol(String), + Float(f64), + Int(i64), + Bool(bool), +} + +/// QuestDB connection configuration +#[derive(Debug, Clone)] +pub struct QuestDBConfig { + /// ILP host for writes (default: localhost) + pub ilp_host: String, + /// ILP port for writes (default: 9009) + pub ilp_port: u16, + /// PostgreSQL host for queries (default: localhost) + pub pg_host: String, + /// PostgreSQL port for queries (default: 8812) + pub pg_port: u16, + /// Ring buffer capacity (default: 10,000) + pub buffer_capacity: usize, +} + +impl Default for QuestDBConfig { + fn default() -> Self { + Self { + ilp_host: "localhost".into(), + ilp_port: 9009, + pg_host: "localhost".into(), + pg_port: 8812, + buffer_capacity: DEFAULT_BUFFER_CAPACITY, + } + } +} + +/// QuestDB client with graceful degradation via ring buffer +pub struct QuestDBClient { + config: QuestDBConfig, + buffer: Arc>>, + connected: Arc, + last_health_check: Arc>, +} + +impl QuestDBClient { + /// Create a new QuestDB client (does not connect immediately) + pub fn new(config: QuestDBConfig) -> Self { + Self { + buffer: Arc::new(Mutex::new(VecDeque::with_capacity(config.buffer_capacity))), + connected: Arc::new(AtomicBool::new(false)), + last_health_check: Arc::new(Mutex::new(Instant::now())), + config, + } + } + + /// Check if QuestDB is currently connected + pub fn is_connected(&self) -> bool { + self.connected.load(Ordering::Relaxed) + } + + /// Get current buffer size (metrics waiting to be flushed) + pub async fn buffer_size(&self) -> usize { + self.buffer.lock().await.len() + } + + /// Get age of oldest buffered entry + pub async fn buffer_age(&self) -> Duration { + let buf = self.buffer.lock().await; + buf.front() + .map(|e| e.created_at.elapsed()) + .unwrap_or(Duration::ZERO) + } + + /// Write a metric entry (buffers if QuestDB is unavailable) + pub async fn write(&self, entry: MetricEntry) { + let mut buf = self.buffer.lock().await; + + // If buffer is full, drop oldest entry + if buf.len() >= self.config.buffer_capacity { + buf.pop_front(); + debug!("QuestDB buffer full, dropped oldest entry"); + } + + buf.push_back(entry); + + // Try to flush if connected + if self.connected.load(Ordering::Relaxed) { + drop(buf); // Release lock before flush + if let Err(e) = self.try_flush().await { + warn!("QuestDB flush failed: {}", e); + self.connected.store(false, Ordering::Relaxed); + } + } + } + + /// Attempt to flush buffered entries to QuestDB via ILP + async fn try_flush(&self) -> Result<(), Box> { + let mut buf = self.buffer.lock().await; + if buf.is_empty() { + return Ok(()); + } + + let mut sender = questdb::ingress::SenderBuilder::new( + questdb::ingress::Protocol::Tcp, + &self.config.ilp_host, + self.config.ilp_port, + ) + .build()?; + + let mut ilp_buf = questdb::ingress::Buffer::new(); + let entries_to_flush: Vec = buf.drain(..).collect(); + let count = entries_to_flush.len(); + + for entry in &entries_to_flush { + ilp_buf.table(entry.table.as_str())?; + for (name, value) in &entry.columns { + match value { + MetricValue::Symbol(s) => { + ilp_buf.symbol(name.as_str(), s.as_str())?; + } + MetricValue::Float(f) => { + ilp_buf.column_f64(name.as_str(), *f)?; + } + MetricValue::Int(i) => { + ilp_buf.column_i64(name.as_str(), *i)?; + } + MetricValue::Bool(b) => { + ilp_buf.column_bool(name.as_str(), *b)?; + } + } + } + ilp_buf.at(questdb::ingress::TimestampNanos::new(entry.timestamp_ns))?; + } + + sender.flush(&mut ilp_buf)?; + info!("Flushed {} entries to QuestDB", count); + + Ok(()) + } + + /// Periodic health check — call from a background task + pub async fn health_check(&self) -> bool { + let mut last = self.last_health_check.lock().await; + if last.elapsed() < HEALTH_CHECK_INTERVAL { + return self.connected.load(Ordering::Relaxed); + } + *last = Instant::now(); + drop(last); + + // Try connecting via ILP + let connected = match questdb::ingress::SenderBuilder::new( + questdb::ingress::Protocol::Tcp, + &self.config.ilp_host, + self.config.ilp_port, + ) + .build() + { + Ok(_sender) => { + if !self.connected.load(Ordering::Relaxed) { + info!("QuestDB connection restored"); + } + true + } + Err(e) => { + if self.connected.load(Ordering::Relaxed) { + warn!("QuestDB connection lost: {}", e); + } + false + } + }; + + self.connected.store(connected, Ordering::Relaxed); + + // If reconnected, try flushing buffer + if connected { + if let Err(e) = self.try_flush().await { + warn!("QuestDB reconnect flush failed: {}", e); + } + } + + connected + } + + /// Get ILP connection string for the configured host + pub fn ilp_address(&self) -> String { + format!("{}:{}", self.config.ilp_host, self.config.ilp_port) + } + + /// Get PostgreSQL connection string for queries + pub fn pg_connection_string(&self) -> String { + format!( + "postgresql://admin:quest@{}:{}/qdb", + self.config.pg_host, self.config.pg_port + ) + } +} + +impl std::fmt::Debug for QuestDBClient { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("QuestDBClient") + .field("connected", &self.connected.load(Ordering::Relaxed)) + .field("config", &self.config) + .finish() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[tokio::test] + async fn test_client_creation() { + let client = QuestDBClient::new(QuestDBConfig::default()); + assert!(!client.is_connected()); + assert_eq!(client.buffer_size().await, 0); + } + + #[tokio::test] + async fn test_buffer_stores_entries() { + let client = QuestDBClient::new(QuestDBConfig::default()); + let entry = MetricEntry { + table: "test".into(), + columns: vec![("value".into(), MetricValue::Float(1.0))], + timestamp_ns: 1_234_567_890, + created_at: Instant::now(), + }; + client.write(entry).await; + assert_eq!(client.buffer_size().await, 1); + } + + #[tokio::test] + async fn test_buffer_capacity_limit() { + let config = QuestDBConfig { + buffer_capacity: 3, + ..Default::default() + }; + let client = QuestDBClient::new(config); + for i in 0..5 { + let entry = MetricEntry { + table: "test".into(), + columns: vec![("i".into(), MetricValue::Int(i))], + timestamp_ns: i, + created_at: Instant::now(), + }; + client.write(entry).await; + } + // Buffer capacity is 3, so oldest 2 should be dropped + assert_eq!(client.buffer_size().await, 3); + } + + #[tokio::test] + async fn test_buffer_age_empty() { + let client = QuestDBClient::new(QuestDBConfig::default()); + assert_eq!(client.buffer_age().await, Duration::ZERO); + } + + #[tokio::test] + async fn test_default_config() { + let config = QuestDBConfig::default(); + assert_eq!(config.ilp_host, "localhost"); + assert_eq!(config.ilp_port, 9009); + assert_eq!(config.pg_host, "localhost"); + assert_eq!(config.pg_port, 8812); + assert_eq!(config.buffer_capacity, 10_000); + } + + #[test] + fn test_pg_connection_string() { + let client = QuestDBClient::new(QuestDBConfig::default()); + assert_eq!( + client.pg_connection_string(), + "postgresql://admin:quest@localhost:8812/qdb" + ); + } +}