diff --git a/Cargo.toml b/Cargo.toml index 43d2537..be90a80 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,3 +1,3 @@ [workspace] -members = ["packages/core"] +members = ["packages/core", "packages/devkit"] resolver = "2" \ No newline at end of file diff --git a/packages/core/src/alerts/mod.rs b/packages/core/src/alerts/mod.rs index 8cb8aed..724bf78 100644 --- a/packages/core/src/alerts/mod.rs +++ b/packages/core/src/alerts/mod.rs @@ -107,8 +107,8 @@ mod tests { use super::*; use chrono::{DateTime, Duration, Utc}; use wiremock::{ - Mock, MockServer, ResponseTemplate, matchers::{method, path}, + Mock, MockServer, ResponseTemplate, }; use crate::insights::{ @@ -246,4 +246,4 @@ mod tests { manager.check_and_dispatch(&update).await; tokio::time::sleep(std::time::Duration::from_millis(100)).await; } -} \ No newline at end of file +} diff --git a/packages/core/src/alerts/webhook.rs b/packages/core/src/alerts/webhook.rs index fdd5749..bbd1b51 100644 --- a/packages/core/src/alerts/webhook.rs +++ b/packages/core/src/alerts/webhook.rs @@ -89,8 +89,8 @@ impl WebhookDelivery { mod tests { use super::*; use wiremock::{ - Mock, MockServer, ResponseTemplate, matchers::{body_json, method, path}, + Mock, MockServer, ResponseTemplate, }; fn build_payload() -> AlertPayload { diff --git a/packages/core/src/api/alerts.rs b/packages/core/src/api/alerts.rs index 3d65616..4aa3e17 100644 --- a/packages/core/src/api/alerts.rs +++ b/packages/core/src/api/alerts.rs @@ -176,7 +176,6 @@ pub async fn delete_alert( } } - // ---- Alert history ---- #[derive(Debug, Deserialize)] @@ -241,8 +240,8 @@ mod tests { use axum::{ body::Body, http::{Method, Request}, - Router, routing::{delete, get, patch, post}, + Router, }; use http_body_util::BodyExt; use tower::ServiceExt; @@ -273,7 +272,9 @@ mod tests { .method(Method::POST) .uri("/alerts/config") .header("content-type", "application/json") - .body(Body::from(r#"{"webhook_url":"https://example.com/hook","threshold":"Major"}"#)) + .body(Body::from( + r#"{"webhook_url":"https://example.com/hook","threshold":"Major"}"#, + )) .unwrap(); let resp = app.oneshot(req).await.unwrap(); @@ -289,7 +290,9 @@ mod tests { .method(Method::POST) .uri("/alerts/config") .header("content-type", "application/json") - .body(Body::from(r#"{"webhook_url":"https://example.com/hook","threshold":"Catastrophic"}"#)) + .body(Body::from( + r#"{"webhook_url":"https://example.com/hook","threshold":"Catastrophic"}"#, + )) .unwrap(); let resp = app.oneshot(req).await.unwrap(); @@ -325,7 +328,10 @@ mod tests { async fn patch_updates_alert_config() { let pool = create_pool("sqlite::memory:").await.unwrap(); let repo = Arc::new(FeeRepository::new(pool)); - let id = repo.insert_alert_config("https://example.com/hook", "Minor").await.unwrap(); + let id = repo + .insert_alert_config("https://example.com/hook", "Minor") + .await + .unwrap(); let app = Router::new() .route("/alerts/config/:id", patch(update_alert)) @@ -346,7 +352,10 @@ mod tests { async fn patch_invalid_threshold_returns_400() { let pool = create_pool("sqlite::memory:").await.unwrap(); let repo = Arc::new(FeeRepository::new(pool)); - let id = repo.insert_alert_config("https://example.com/hook", "Minor").await.unwrap(); + let id = repo + .insert_alert_config("https://example.com/hook", "Minor") + .await + .unwrap(); let app = Router::new() .route("/alerts/config/:id", patch(update_alert)) @@ -367,7 +376,10 @@ mod tests { async fn delete_soft_deletes_alert_config() { let pool = create_pool("sqlite::memory:").await.unwrap(); let repo = Arc::new(FeeRepository::new(pool.clone())); - let id = repo.insert_alert_config("https://example.com/hook", "Major").await.unwrap(); + let id = repo + .insert_alert_config("https://example.com/hook", "Major") + .await + .unwrap(); let app = Router::new() .route("/alerts/config/:id", delete(delete_alert)) @@ -408,8 +420,8 @@ mod history_tests { use axum::{ body::Body, http::{Method, Request}, - Router, routing::get, + Router, }; use http_body_util::BodyExt; use tower::ServiceExt; diff --git a/packages/core/src/api/fees.rs b/packages/core/src/api/fees.rs index 16c4b50..9d2dee4 100644 --- a/packages/core/src/api/fees.rs +++ b/packages/core/src/api/fees.rs @@ -4,7 +4,7 @@ use async_trait::async_trait; use axum::{ body::Body, extract::{Query, State}, - http::{HeaderMap, StatusCode, header}, + http::{header, HeaderMap, StatusCode}, response::Response, Json, }; @@ -13,12 +13,12 @@ use serde::{Deserialize, Serialize}; use serde_json::{json, Value}; use tokio::sync::{Mutex, RwLock}; +use super::headers::{cache_control, compute_etag, if_none_match_matches, last_modified}; use crate::cache::ResponseCache; use crate::error::AppError; use crate::insights::{FeeDataPoint, FeeInsightsEngine, TrendIndicator, TrendStrength}; use crate::services::horizon::HorizonClient; use crate::store::FeeHistoryStore; -use super::headers::{cache_control, compute_etag, if_none_match_matches, last_modified}; /// Shared state type for the fees route. pub type FeesState = Arc; @@ -38,18 +38,18 @@ impl FeeStatsProvider for HorizonClient { max_fee: stats.fee_charged.max, avg_fee: stats.fee_charged.avg, percentiles: PercentileFees { - p10: stats.fee_charged.p10, - p20: stats.fee_charged.p20, - p30: stats.fee_charged.p30, - p40: stats.fee_charged.p40, - p50: stats.fee_charged.p50, - p60: stats.fee_charged.p60, - p70: stats.fee_charged.p70, - p80: stats.fee_charged.p80, - p90: stats.fee_charged.p90, - p95: stats.fee_charged.p95, - p99: stats.fee_charged.p99, -}, + p10: stats.fee_charged.p10, + p20: stats.fee_charged.p20, + p30: stats.fee_charged.p30, + p40: stats.fee_charged.p40, + p50: stats.fee_charged.p50, + p60: stats.fee_charged.p60, + p70: stats.fee_charged.p70, + p80: stats.fee_charged.p80, + p90: stats.fee_charged.p90, + p95: stats.fee_charged.p95, + p99: stats.fee_charged.p99, + }, }) } } @@ -324,9 +324,7 @@ pub struct FeeTrendResponse { pub last_updated: DateTime, } -pub async fn fee_trend( - State(state): State, -) -> Result, AppError> { +pub async fn fee_trend(State(state): State) -> Result, AppError> { let engine = state .insights_engine .as_ref() @@ -389,13 +387,13 @@ mod tests { use std::sync::Mutex as StdMutex; use std::time::Duration as StdDuration; + use crate::insights::InsightsConfig; use axum::{ body::{to_bytes, Body}, http::{Request, StatusCode}, routing::get, Router, }; - use crate::insights::InsightsConfig; use chrono::Duration as ChronoDuration; use tower::ServiceExt; @@ -477,19 +475,19 @@ mod tests { min_fee: "100".to_string(), max_fee: "5000".to_string(), avg_fee: "213".to_string(), - percentiles: PercentileFees { - p10: "100".to_string(), - p20: "100".to_string(), - p30: "100".to_string(), - p40: "100".to_string(), - p50: "150".to_string(), - p60: "200".to_string(), - p70: "250".to_string(), - p80: "300".to_string(), - p90: "500".to_string(), - p95: "800".to_string(), - p99: "1000".to_string(), -}, + percentiles: PercentileFees { + p10: "100".to_string(), + p20: "100".to_string(), + p30: "100".to_string(), + p40: "100".to_string(), + p50: "150".to_string(), + p60: "200".to_string(), + p70: "250".to_string(), + p80: "300".to_string(), + p90: "500".to_string(), + p95: "800".to_string(), + p99: "1000".to_string(), + }, } } @@ -551,7 +549,8 @@ mod tests { make_current_fee_response("100"), make_current_fee_response("200"), ]); - let state = make_fee_state_with_provider(Arc::new(mock.clone()), StdDuration::from_secs(60)); + let state = + make_fee_state_with_provider(Arc::new(mock.clone()), StdDuration::from_secs(60)); let app = Router::new() .route("/fees/current", get(current_fees)) diff --git a/packages/core/src/api/headers.rs b/packages/core/src/api/headers.rs index 48ad056..b5c0166 100644 --- a/packages/core/src/api/headers.rs +++ b/packages/core/src/api/headers.rs @@ -1,7 +1,7 @@ use std::collections::hash_map::DefaultHasher; use std::hash::{Hash, Hasher}; -use axum::http::{HeaderMap, HeaderValue, header}; +use axum::http::{header, HeaderMap, HeaderValue}; use chrono::{DateTime, Utc}; /// Compute a weakly-stable quoted ETag from response bytes. diff --git a/packages/core/src/api/health.rs b/packages/core/src/api/health.rs index 0849617..0d715c6 100644 --- a/packages/core/src/api/health.rs +++ b/packages/core/src/api/health.rs @@ -1,6 +1,6 @@ use axum::{ body::Body, - http::{HeaderValue, StatusCode, header}, + http::{header, HeaderValue, StatusCode}, response::{IntoResponse, Response}, }; diff --git a/packages/core/src/api/insights.rs b/packages/core/src/api/insights.rs index 3c66123..ddc8fe6 100644 --- a/packages/core/src/api/insights.rs +++ b/packages/core/src/api/insights.rs @@ -3,7 +3,7 @@ use axum::{ body::Body, extract::State, - http::{HeaderMap, StatusCode, header}, + http::{header, HeaderMap, StatusCode}, response::{Json, Response}, routing::get, Router, @@ -12,8 +12,8 @@ use serde_json::Value; use std::sync::Arc; use tokio::sync::RwLock; -use crate::insights::{CongestionTrends, FeeExtremes, FeeInsightsEngine, RollingAverages}; use super::headers::{cache_control, compute_etag, if_none_match_matches, last_modified}; +use crate::insights::{CongestionTrends, FeeExtremes, FeeInsightsEngine, RollingAverages}; /// Shared state for the insights API pub type InsightsState = Arc>; @@ -97,7 +97,7 @@ async fn get_insights_health( State(engine): State, ) -> Result, (StatusCode, Json)> { let engine = engine.read().await; - + let health_info = serde_json::json!({ "status": "healthy", "last_update": engine.get_last_update(), @@ -107,6 +107,6 @@ async fn get_insights_health( "spike_threshold": engine.get_config().spike_detection.threshold_multiplier } }); - + Ok(Json(health_info)) } diff --git a/packages/core/src/api/mod.rs b/packages/core/src/api/mod.rs index a635e43..d3c2ce1 100644 --- a/packages/core/src/api/mod.rs +++ b/packages/core/src/api/mod.rs @@ -1,6 +1,5 @@ -pub mod health; -pub mod fees; -pub mod insights; pub mod alerts; +pub mod fees; pub mod headers; - +pub mod health; +pub mod insights; diff --git a/packages/core/src/config.rs b/packages/core/src/config.rs index b6b7b86..b5017f5 100644 --- a/packages/core/src/config.rs +++ b/packages/core/src/config.rs @@ -145,8 +145,8 @@ impl Config { .unwrap_or(1000); // -------- Database URL -------- - let database_url = get("DATABASE_URL") - .unwrap_or_else(|| "sqlite://stellar_fees.db".to_string()); + let database_url = + get("DATABASE_URL").unwrap_or_else(|| "sqlite://stellar_fees.db".to_string()); // -------- Storage retention -------- let storage_retention_days = get("STORAGE_RETENTION_DAYS") @@ -369,9 +369,10 @@ mod tests { #[test] fn allowed_origins_parses_comma_separated_list() { let cli = make_cli("testnet", None); - let env = HashMap::from([ - ("ALLOWED_ORIGINS", "http://localhost:3000,https://app.example.com"), - ]); + let env = HashMap::from([( + "ALLOWED_ORIGINS", + "http://localhost:3000,https://app.example.com", + )]); let config = Config::from_sources_with_overrides(&cli, &env).unwrap(); assert_eq!( config.allowed_origins, @@ -382,9 +383,10 @@ mod tests { #[test] fn allowed_origins_trims_whitespace() { let cli = make_cli("testnet", None); - let env = HashMap::from([ - ("ALLOWED_ORIGINS", "http://localhost:3000 , https://app.example.com "), - ]); + let env = HashMap::from([( + "ALLOWED_ORIGINS", + "http://localhost:3000 , https://app.example.com ", + )]); let config = Config::from_sources_with_overrides(&cli, &env).unwrap(); assert_eq!( config.allowed_origins, diff --git a/packages/core/src/db.rs b/packages/core/src/db.rs index 6a33512..8c0f2da 100644 --- a/packages/core/src/db.rs +++ b/packages/core/src/db.rs @@ -37,7 +37,11 @@ mod tests { // Run migrations a second time explicitly let result = sqlx::migrate!("./migrations").run(&pool).await; - assert!(result.is_ok(), "Second migration run failed: {:?}", result.err()); + assert!( + result.is_ok(), + "Second migration run failed: {:?}", + result.err() + ); } #[tokio::test] @@ -69,4 +73,4 @@ mod tests { assert!(result.is_ok(), "Insert failed: {:?}", result.err()); } -} \ No newline at end of file +} diff --git a/packages/core/src/error.rs b/packages/core/src/error.rs index 9baebe1..3e0bb12 100644 --- a/packages/core/src/error.rs +++ b/packages/core/src/error.rs @@ -1,5 +1,5 @@ -use std::fmt; use std::error::Error; +use std::fmt; use axum::{ http::StatusCode, @@ -109,4 +109,4 @@ mod tests { "Unknown error: ???" ); } -} \ No newline at end of file +} diff --git a/packages/core/src/insights/calculator.rs b/packages/core/src/insights/calculator.rs index f267995..549e547 100644 --- a/packages/core/src/insights/calculator.rs +++ b/packages/core/src/insights/calculator.rs @@ -3,11 +3,7 @@ use chrono::{DateTime, Utc}; use std::collections::{HashMap, VecDeque}; -use crate::insights::{ - types::*, - error::InsightsError, - config::AverageConfig, -}; +use crate::insights::{config::AverageConfig, error::InsightsError, types::*}; /// Circular buffer for efficient storage of fee data points #[derive(Debug, Clone)] @@ -23,22 +19,22 @@ impl CircularBuffer { max_size, } } - + fn push(&mut self, item: T) { if self.data.len() >= self.max_size { self.data.pop_front(); } self.data.push_back(item); } - + fn iter(&self) -> impl Iterator { self.data.iter() } - + fn len(&self) -> usize { self.data.len() } - + fn is_empty(&self) -> bool { self.data.is_empty() } @@ -55,23 +51,23 @@ impl RollingAverageCalculator { /// Create a new rolling average calculator pub fn new(config: AverageConfig, time_windows: Vec) -> Self { let mut windows = HashMap::new(); - + // Initialize circular buffers for each time window for window in &time_windows { windows.insert(window.clone(), CircularBuffer::new(config.max_buffer_size)); } - + Self { config, windows, time_windows, } } - + /// Add a new data point to all relevant time windows pub fn add_data_point(&mut self, point: FeeDataPoint) { let now = Utc::now(); - + // Add to each time window if the point is within the window duration for window in &self.time_windows { if let Some(buffer) = self.windows.get_mut(window) { @@ -82,16 +78,16 @@ impl RollingAverageCalculator { } } } - + // Clean old data points from all buffers self.clean_old_data(now); } - + /// Clean old data points that are outside their respective time windows fn clean_old_data(&mut self, current_time: DateTime) { for (window, buffer) in &mut self.windows { let window_start = current_time - window.duration; - + // Remove data points that are too old while let Some(front) = buffer.data.front() { if front.timestamp < window_start { @@ -102,33 +98,42 @@ impl RollingAverageCalculator { } } } - + /// Calculate averages for all time windows pub fn calculate_averages(&self) -> Result { let now = Utc::now(); - + // Calculate average for each predefined window type let short_term = self.calculate_average_for_window("short_term", now)?; let medium_term = self.calculate_average_for_window("medium_term", now)?; let long_term = self.calculate_average_for_window("long_term", now)?; - + Ok(RollingAverages { short_term, medium_term, long_term, }) } - + /// Calculate average for a specific time window by name - fn calculate_average_for_window(&self, window_name: &str, calculated_at: DateTime) -> Result { + fn calculate_average_for_window( + &self, + window_name: &str, + calculated_at: DateTime, + ) -> Result { // Find the time window by name - let time_window = self.time_windows.iter() + let time_window = self + .time_windows + .iter() .find(|w| w.name == window_name) - .ok_or_else(|| InsightsError::config_error(format!("Time window '{}' not found", window_name)))?; - - let buffer = self.windows.get(time_window) - .ok_or_else(|| InsightsError::config_error(format!("Buffer for window '{}' not found", window_name)))?; - + .ok_or_else(|| { + InsightsError::config_error(format!("Time window '{}' not found", window_name)) + })?; + + let buffer = self.windows.get(time_window).ok_or_else(|| { + InsightsError::config_error(format!("Buffer for window '{}' not found", window_name)) + })?; + if buffer.is_empty() { return Ok(AverageResult { value: 0.0, @@ -138,15 +143,15 @@ impl RollingAverageCalculator { time_window: time_window.clone(), }); } - + // Calculate the average fee let total_fee: u64 = buffer.iter().map(|point| point.fee_amount).sum(); let sample_count = buffer.len(); let average = total_fee as f64 / sample_count as f64; - + // Determine if this is a partial result (insufficient samples) let is_partial = sample_count < time_window.min_samples; - + Ok(AverageResult { value: average, sample_count, @@ -155,20 +160,20 @@ impl RollingAverageCalculator { time_window: time_window.clone(), }) } - + /// Get average for a specific time window pub fn get_average_for_window(&self, window: &TimeWindow) -> Option { let buffer = self.windows.get(window)?; - + if buffer.is_empty() { return None; } - + let total_fee: u64 = buffer.iter().map(|point| point.fee_amount).sum(); let sample_count = buffer.len(); let average = total_fee as f64 / sample_count as f64; let is_partial = sample_count < window.min_samples; - + Some(AverageResult { value: average, sample_count, @@ -177,15 +182,18 @@ impl RollingAverageCalculator { time_window: window.clone(), }) } - + /// Get the number of data points in a specific time window pub fn get_sample_count(&self, window: &TimeWindow) -> usize { - self.windows.get(window).map(|buffer| buffer.len()).unwrap_or(0) + self.windows + .get(window) + .map(|buffer| buffer.len()) + .unwrap_or(0) } - + /// Check if a time window has sufficient data for reliable calculations pub fn has_sufficient_data(&self, window: &TimeWindow) -> bool { let sample_count = self.get_sample_count(window); sample_count >= window.min_samples } -} \ No newline at end of file +} diff --git a/packages/core/src/insights/config.rs b/packages/core/src/insights/config.rs index 1c04513..7a38dbe 100644 --- a/packages/core/src/insights/config.rs +++ b/packages/core/src/insights/config.rs @@ -1,8 +1,8 @@ //! Configuration for fee insights system +use crate::insights::types::TimeWindow; use chrono::Duration; use serde::{Deserialize, Serialize}; -use crate::insights::types::TimeWindow; /// Configuration for the fee insights engine #[derive(Debug, Clone, Serialize, Deserialize)] @@ -88,4 +88,4 @@ impl Default for ExtremesConfig { historical_periods_to_keep: 30, } } -} \ No newline at end of file +} diff --git a/packages/core/src/insights/detector.rs b/packages/core/src/insights/detector.rs index 346a07b..e2b486f 100644 --- a/packages/core/src/insights/detector.rs +++ b/packages/core/src/insights/detector.rs @@ -1,13 +1,9 @@ //! Congestion Detection System -use chrono::{Utc, Duration}; +use chrono::{Duration, Utc}; use std::collections::VecDeque; -use crate::insights::{ - types::*, - error::InsightsError, - config::SpikeConfig, -}; +use crate::insights::{config::SpikeConfig, error::InsightsError, types::*}; /// Analyzer for trend patterns #[derive(Debug, Clone)] @@ -23,15 +19,15 @@ impl TrendAnalyzer { congestion_window, } } - + fn add_spike(&mut self, spike: FeeSpike) { self.recent_spikes.push_back(spike); self.clean_old_spikes(); } - + fn clean_old_spikes(&mut self) { let cutoff_time = Utc::now() - self.congestion_window; - + while let Some(front_spike) = self.recent_spikes.front() { if front_spike.start_time < cutoff_time { self.recent_spikes.pop_front(); @@ -40,10 +36,12 @@ impl TrendAnalyzer { } } } - + fn calculate_trend_strength(&self) -> TrendStrength { let spike_count = self.recent_spikes.len(); - let total_severity_score: f64 = self.recent_spikes.iter() + let total_severity_score: f64 = self + .recent_spikes + .iter() .map(|spike| match spike.severity { SpikeSeverity::Minor => 1.0, SpikeSeverity::Moderate => 2.0, @@ -51,10 +49,10 @@ impl TrendAnalyzer { SpikeSeverity::Critical => 8.0, }) .sum(); - + // Calculate strength based on both count and severity let strength_score = total_severity_score + (spike_count as f64 * 0.5); - + match strength_score { s if s >= 10.0 => TrendStrength::Strong, s if s >= 4.0 => TrendStrength::Moderate, @@ -62,29 +60,33 @@ impl TrendAnalyzer { _ => TrendStrength::Weak, } } - + fn determine_trend_indicator(&self) -> TrendIndicator { if self.recent_spikes.is_empty() { return TrendIndicator::Normal; } - + let recent_spikes: Vec<_> = self.recent_spikes.iter().collect(); let spike_count = recent_spikes.len(); - + // Analyze trend based on recent spike patterns if spike_count >= 3 { // Check if spikes are increasing in severity - let avg_recent_ratio = recent_spikes.iter() + let avg_recent_ratio = recent_spikes + .iter() .rev() .take(2) .map(|s| s.spike_ratio) - .sum::() / 2.0; - - let avg_older_ratio = recent_spikes.iter() + .sum::() + / 2.0; + + let avg_older_ratio = recent_spikes + .iter() .take(recent_spikes.len() - 2) .map(|s| s.spike_ratio) - .sum::() / (recent_spikes.len() - 2) as f64; - + .sum::() + / (recent_spikes.len() - 2) as f64; + if avg_recent_ratio > avg_older_ratio * 1.2 { TrendIndicator::Rising } else if avg_recent_ratio < avg_older_ratio * 0.8 { @@ -103,24 +105,27 @@ impl TrendAnalyzer { TrendIndicator::Normal } } - + fn predict_duration(&self) -> Option { if self.recent_spikes.is_empty() { return None; } - + // Simple prediction based on recent spike patterns - let avg_duration: i64 = self.recent_spikes.iter() + let avg_duration: i64 = self + .recent_spikes + .iter() .map(|spike| spike.duration.num_minutes()) - .sum::() / self.recent_spikes.len() as i64; - + .sum::() + / self.recent_spikes.len() as i64; + // Predict based on trend strength let multiplier = match self.calculate_trend_strength() { TrendStrength::Strong => 2.0, TrendStrength::Moderate => 1.5, TrendStrength::Weak => 1.0, }; - + Some(Duration::minutes((avg_duration as f64 * multiplier) as i64)) } } @@ -141,34 +146,39 @@ impl CongestionDetector { historical_spikes: VecDeque::new(), } } - + /// Analyze congestion patterns - pub fn analyze_congestion(&mut self, current_fees: &[FeeDataPoint], baseline: f64) -> Result { + pub fn analyze_congestion( + &mut self, + current_fees: &[FeeDataPoint], + baseline: f64, + ) -> Result { // Detect new spikes in the current fee data let new_spikes = self.detect_spikes(current_fees, baseline)?; - + // Add new spikes to the trend analyzer for spike in &new_spikes { self.trend_analyzer.add_spike(spike.clone()); self.historical_spikes.push_back(spike.clone()); } - + // Maintain historical spike buffer (keep last 1000 spikes) while self.historical_spikes.len() > 1000 { self.historical_spikes.pop_front(); } - + // Clean old spikes from trend analyzer self.trend_analyzer.clean_old_spikes(); - + // Calculate current trend indicators let current_trend = self.trend_analyzer.determine_trend_indicator(); let trend_strength = self.trend_analyzer.calculate_trend_strength(); let predicted_duration = self.trend_analyzer.predict_duration(); - + // Get recent spikes for the response - let recent_spikes: Vec = self.trend_analyzer.recent_spikes.iter().cloned().collect(); - + let recent_spikes: Vec = + self.trend_analyzer.recent_spikes.iter().cloned().collect(); + Ok(CongestionTrends { current_trend, recent_spikes, @@ -176,29 +186,33 @@ impl CongestionDetector { predicted_duration, }) } - + /// Detect fee spikes in the given fee data - pub fn detect_spikes(&self, fees: &[FeeDataPoint], baseline: f64) -> Result, InsightsError> { + pub fn detect_spikes( + &self, + fees: &[FeeDataPoint], + baseline: f64, + ) -> Result, InsightsError> { if fees.is_empty() { return Ok(Vec::new()); } - + if baseline <= 0.0 { return Err(InsightsError::invalid_data("Baseline must be positive")); } - + let mut spikes = Vec::new(); let threshold = baseline * self.config.threshold_multiplier; - + // Sort fees by timestamp to process in chronological order let mut sorted_fees = fees.to_vec(); sorted_fees.sort_by(|a, b| a.timestamp.cmp(&b.timestamp)); - + let mut current_spike: Option = None; - + for fee_point in &sorted_fees { let fee_amount = fee_point.fee_amount as f64; - + if fee_amount >= threshold { // This is a spike match &mut current_spike { @@ -227,7 +241,7 @@ impl CongestionDetector { // Not a spike, end current spike if it exists if let Some(mut spike) = current_spike.take() { spike.duration = fee_point.timestamp - spike.start_time; - + // Only include spikes that meet minimum duration if spike.duration >= self.config.minimum_spike_duration { spikes.push(spike); @@ -235,21 +249,21 @@ impl CongestionDetector { } } } - + // Handle case where spike continues to the end of data if let Some(mut spike) = current_spike { if let Some(last_fee) = sorted_fees.last() { spike.duration = last_fee.timestamp - spike.start_time; - + if spike.duration >= self.config.minimum_spike_duration { spikes.push(spike); } } } - + Ok(spikes) } - + /// Classify the severity of a spike based on its ratio to baseline pub fn classify_spike_severity(&self, spike_ratio: f64) -> SpikeSeverity { match spike_ratio { @@ -259,25 +273,25 @@ impl CongestionDetector { _ => SpikeSeverity::Minor, } } - + /// Calculate trend strength based on recent spike activity pub fn calculate_trend_strength(&self) -> TrendStrength { self.trend_analyzer.calculate_trend_strength() } - + /// Get recent spikes within the congestion window pub fn get_recent_spikes(&self) -> Vec { self.trend_analyzer.recent_spikes.iter().cloned().collect() } - + /// Get all historical spikes pub fn get_historical_spikes(&self) -> Vec { self.historical_spikes.iter().cloned().collect() } - + /// Clear all spike history (useful for testing or reset scenarios) pub fn clear_history(&mut self) { self.trend_analyzer.recent_spikes.clear(); self.historical_spikes.clear(); } -} \ No newline at end of file +} diff --git a/packages/core/src/insights/engine.rs b/packages/core/src/insights/engine.rs index 4a62c21..6635cc1 100644 --- a/packages/core/src/insights/engine.rs +++ b/packages/core/src/insights/engine.rs @@ -4,12 +4,12 @@ use chrono::{DateTime, Utc}; use std::time::Instant; use crate::insights::{ - types::*, - error::InsightsError, - config::{InsightsConfig, AverageConfig, ExtremesConfig}, calculator::RollingAverageCalculator, - tracker::ExtremesTracker, + config::{AverageConfig, ExtremesConfig, InsightsConfig}, detector::CongestionDetector, + error::InsightsError, + tracker::ExtremesTracker, + types::*, }; /// Central fee insights engine that orchestrates all analysis operations @@ -28,15 +28,12 @@ impl FeeInsightsEngine { // Create component configurations let average_config = AverageConfig::default(); let extremes_config = ExtremesConfig::default(); - + // Initialize components - let calculator = RollingAverageCalculator::new( - average_config, - config.time_windows.clone() - ); + let calculator = RollingAverageCalculator::new(average_config, config.time_windows.clone()); let tracker = ExtremesTracker::new(extremes_config); let detector = CongestionDetector::new(config.spike_detection.clone()); - + Self { config, calculator, @@ -46,41 +43,46 @@ impl FeeInsightsEngine { last_insights: None, } } - + /// Process new fee data and update insights - pub async fn process_fee_data(&mut self, data: &[FeeDataPoint]) -> Result { + pub async fn process_fee_data( + &mut self, + data: &[FeeDataPoint], + ) -> Result { let start_time = Instant::now(); let processing_start = Utc::now(); - + if data.is_empty() { return Err(InsightsError::invalid_data("No fee data provided")); } - + // Validate fee data self.validate_fee_data(data)?; - + // Update rolling averages for fee_point in data { self.calculator.add_data_point(fee_point.clone()); } - + // Update extremes tracking self.tracker.update_with_fees(data)?; - + // Calculate rolling averages to get baseline for congestion detection let rolling_averages = self.calculator.calculate_averages()?; let baseline = rolling_averages.medium_term.value; // Use medium-term as baseline - + // Update congestion detection let congestion_trends = self.detector.analyze_congestion(data, baseline)?; - + // Get current extremes - let extremes = self.tracker.get_current_extremes() + let extremes = self + .tracker + .get_current_extremes() .unwrap_or_else(|_| self.create_default_extremes()); - + // Calculate data quality let data_quality = self.calculate_data_quality(data, processing_start); - + // Create current insights let insights = CurrentInsights { rolling_averages, @@ -89,86 +91,97 @@ impl FeeInsightsEngine { last_updated: processing_start, data_quality, }; - + // Update last update time self.last_update = Some(processing_start); self.last_insights = Some(insights.clone()); - + // Calculate processing time let processing_time = chrono::Duration::from_std(start_time.elapsed()) .unwrap_or_else(|_| chrono::Duration::zero()); - + Ok(InsightsUpdate { insights, processing_time, data_points_processed: data.len(), }) } - + /// Validate fee data for basic correctness pub fn validate_fee_data(&self, data: &[FeeDataPoint]) -> Result<(), InsightsError> { for (i, fee_point) in data.iter().enumerate() { // Check for reasonable fee amounts (not zero, not excessively large) if fee_point.fee_amount == 0 { - return Err(InsightsError::invalid_data( - format!("Zero fee amount at index {}", i) - )); + return Err(InsightsError::invalid_data(format!( + "Zero fee amount at index {}", + i + ))); } - + // Check for reasonable fee amounts (Stellar fees are typically in stroops) - if fee_point.fee_amount > 1_000_000_000 { // 1000 XLM in stroops - return Err(InsightsError::invalid_data( - format!("Unreasonably large fee amount {} at index {}", fee_point.fee_amount, i) - )); + if fee_point.fee_amount > 1_000_000_000 { + // 1000 XLM in stroops + return Err(InsightsError::invalid_data(format!( + "Unreasonably large fee amount {} at index {}", + fee_point.fee_amount, i + ))); } - + // Check for valid transaction hash if fee_point.transaction_hash.is_empty() { - return Err(InsightsError::invalid_data( - format!("Empty transaction hash at index {}", i) - )); + return Err(InsightsError::invalid_data(format!( + "Empty transaction hash at index {}", + i + ))); } - + // Check for reasonable timestamp (not too far in the future) let now = Utc::now(); if fee_point.timestamp > now + chrono::Duration::hours(1) { - return Err(InsightsError::invalid_data( - format!("Future timestamp at index {}: {}", i, fee_point.timestamp) - )); + return Err(InsightsError::invalid_data(format!( + "Future timestamp at index {}: {}", + i, fee_point.timestamp + ))); } } - + Ok(()) } - + /// Calculate data quality metrics - fn calculate_data_quality(&self, data: &[FeeDataPoint], processing_time: DateTime) -> DataQuality { + fn calculate_data_quality( + &self, + data: &[FeeDataPoint], + processing_time: DateTime, + ) -> DataQuality { let expected_points = self.estimate_expected_data_points(); let actual_points = data.len(); - + // Calculate completeness (0.0 to 1.0) let completeness = if expected_points > 0 { (actual_points as f64 / expected_points as f64).min(1.0) } else { 1.0 }; - + // Calculate freshness (time since last update) let freshness = match self.last_update { Some(last) => processing_time - last, None => chrono::Duration::zero(), }; - + // Check for gaps in data (simplified - just check if we have recent data) - let has_gaps = data.is_empty() || - data.iter().any(|point| processing_time - point.timestamp > chrono::Duration::hours(1)); - + let has_gaps = data.is_empty() + || data + .iter() + .any(|point| processing_time - point.timestamp > chrono::Duration::hours(1)); + let last_gap = if has_gaps { Some(processing_time) } else { None }; - + DataQuality { completeness, freshness, @@ -176,7 +189,7 @@ impl FeeInsightsEngine { last_gap, } } - + /// Estimate expected number of data points based on polling interval fn estimate_expected_data_points(&self) -> usize { // Simple estimation: assume 1 data point per polling interval @@ -184,13 +197,14 @@ impl FeeInsightsEngine { match self.last_update { Some(last) => { let time_diff = Utc::now() - last; - let intervals = time_diff.num_seconds() / self.config.polling_interval.num_seconds(); + let intervals = + time_diff.num_seconds() / self.config.polling_interval.num_seconds(); intervals.max(1) as usize } None => 1, } } - + /// Create default extremes when no data is available fn create_default_extremes(&self) -> FeeExtremes { let now = Utc::now(); @@ -199,7 +213,7 @@ impl FeeInsightsEngine { timestamp: now, transaction_hash: "unknown".to_string(), }; - + FeeExtremes { current_min: default_extreme.clone(), current_max: default_extreme, @@ -207,33 +221,37 @@ impl FeeInsightsEngine { period_end: now, } } - + /// Get current insights pub fn get_current_insights(&self) -> CurrentInsights { if let Some(insights) = &self.last_insights { return insights.clone(); } - let rolling_averages = self.calculator.calculate_averages() + let rolling_averages = self + .calculator + .calculate_averages() .unwrap_or_else(|_| self.create_default_rolling_averages()); - - let extremes = self.tracker.get_current_extremes() + + let extremes = self + .tracker + .get_current_extremes() .unwrap_or_else(|_| self.create_default_extremes()); - + let congestion_trends = CongestionTrends { current_trend: TrendIndicator::Normal, recent_spikes: self.detector.get_recent_spikes(), trend_strength: self.detector.calculate_trend_strength(), predicted_duration: None, }; - + let data_quality = DataQuality { completeness: 1.0, freshness: chrono::Duration::zero(), has_gaps: false, last_gap: None, }; - + CurrentInsights { rolling_averages, extremes, @@ -242,7 +260,7 @@ impl FeeInsightsEngine { data_quality, } } - + /// Create default rolling averages when no data is available fn create_default_rolling_averages(&self) -> RollingAverages { let now = Utc::now(); @@ -257,26 +275,28 @@ impl FeeInsightsEngine { min_samples: 1, }, }; - + RollingAverages { short_term: default_result.clone(), medium_term: default_result.clone(), long_term: default_result, } } - + /// Get rolling averages pub fn get_rolling_averages(&self) -> RollingAverages { - self.calculator.calculate_averages() + self.calculator + .calculate_averages() .unwrap_or_else(|_| self.create_default_rolling_averages()) } - + /// Get fee extremes pub fn get_extremes(&self) -> FeeExtremes { - self.tracker.get_current_extremes() + self.tracker + .get_current_extremes() .unwrap_or_else(|_| self.create_default_extremes()) } - + /// Get congestion trends pub fn get_congestion_trends(&self) -> CongestionTrends { CongestionTrends { @@ -286,37 +306,35 @@ impl FeeInsightsEngine { predicted_duration: None, } } - + /// Get engine configuration pub fn get_config(&self) -> &InsightsConfig { &self.config } - + /// Get last update time pub fn get_last_update(&self) -> Option> { self.last_update } - + /// Reset all components (useful for testing or maintenance) pub fn reset(&mut self) -> Result<(), InsightsError> { // Reset calculator by creating a new one let average_config = AverageConfig::default(); - self.calculator = RollingAverageCalculator::new( - average_config, - self.config.time_windows.clone() - ); - + self.calculator = + RollingAverageCalculator::new(average_config, self.config.time_windows.clone()); + // Reset tracker let extremes_config = ExtremesConfig::default(); self.tracker = ExtremesTracker::new(extremes_config); - + // Reset detector self.detector.clear_history(); - + // Reset update time self.last_update = None; self.last_insights = None; - + Ok(()) } } diff --git a/packages/core/src/insights/error.rs b/packages/core/src/insights/error.rs index 52820b8..a44af41 100644 --- a/packages/core/src/insights/error.rs +++ b/packages/core/src/insights/error.rs @@ -7,22 +7,24 @@ use thiserror::Error; pub enum InsightsError { #[error("Invalid fee data: {message}")] InvalidData { message: String }, - + #[error("Calculation error: {message}")] CalculationError { message: String }, - + #[error("Configuration error: {message}")] ConfigError { message: String }, - + #[error("Storage error: {message}")] StorageError { message: String }, - + #[error("Data provider error: {source}")] - ProviderError { source: Box }, - + ProviderError { + source: Box, + }, + #[error("Insufficient data for calculation: {operation}")] InsufficientData { operation: String }, - + #[error("Numerical overflow in calculation: {operation}")] NumericalOverflow { operation: String }, } @@ -32,42 +34,54 @@ pub enum InsightsError { pub enum ProviderError { #[error("Network error: {message}")] NetworkError { message: String }, - + #[error("Data format error: {message}")] FormatError { message: String }, - + #[error("Authentication error: {message}")] AuthError { message: String }, - + #[error("Rate limit exceeded")] RateLimitExceeded, - + #[error("Service unavailable")] ServiceUnavailable, } impl InsightsError { pub fn invalid_data(message: impl Into) -> Self { - Self::InvalidData { message: message.into() } + Self::InvalidData { + message: message.into(), + } } - + pub fn calculation_error(message: impl Into) -> Self { - Self::CalculationError { message: message.into() } + Self::CalculationError { + message: message.into(), + } } - + pub fn config_error(message: impl Into) -> Self { - Self::ConfigError { message: message.into() } + Self::ConfigError { + message: message.into(), + } } - + pub fn storage_error(message: impl Into) -> Self { - Self::StorageError { message: message.into() } + Self::StorageError { + message: message.into(), + } } - + pub fn insufficient_data(operation: impl Into) -> Self { - Self::InsufficientData { operation: operation.into() } + Self::InsufficientData { + operation: operation.into(), + } } - + pub fn numerical_overflow(operation: impl Into) -> Self { - Self::NumericalOverflow { operation: operation.into() } + Self::NumericalOverflow { + operation: operation.into(), + } } -} \ No newline at end of file +} diff --git a/packages/core/src/insights/horizon_adapter.rs b/packages/core/src/insights/horizon_adapter.rs index 9a55060..9024c23 100644 --- a/packages/core/src/insights/horizon_adapter.rs +++ b/packages/core/src/insights/horizon_adapter.rs @@ -1,5 +1,5 @@ //! Horizon Fee Data Provider Adapter -//! +//! //! Adapts the HorizonClient to implement the FeeDataProvider trait use async_trait::async_trait; @@ -8,9 +8,9 @@ use serde::Deserialize; use std::str::FromStr; use crate::insights::{ + error::ProviderError, provider::{FeeDataProvider, ProviderMetadata, ProviderResult}, types::FeeDataPoint, - error::ProviderError, }; use crate::services::horizon::HorizonClient; @@ -46,65 +46,73 @@ impl HorizonFeeDataProvider { pub fn new(client: HorizonClient) -> Self { let metadata = ProviderMetadata { supports_historical: true, - max_batch_size: 200, // Horizon's default limit + max_batch_size: 200, // Horizon's default limit rate_limit_per_minute: Some(3600), // Horizon's rate limit - data_freshness_seconds: 5, // Stellar ledger close time + data_freshness_seconds: 5, // Stellar ledger close time }; - - Self { - client, - metadata, - } + + Self { client, metadata } } - + /// Fetch recent transactions from Horizon - async fn fetch_recent_transactions(&self, limit: u32) -> ProviderResult> { - let url = format!("{}/transactions?order=desc&limit={}", self.client.base_url(), limit); - + async fn fetch_recent_transactions( + &self, + limit: u32, + ) -> ProviderResult> { + let url = format!( + "{}/transactions?order=desc&limit={}", + self.client.base_url(), + limit + ); + let response = reqwest::get(&url) .await - .map_err(|e| ProviderError::NetworkError { - message: format!("Failed to fetch transactions: {}", e) + .map_err(|e| ProviderError::NetworkError { + message: format!("Failed to fetch transactions: {}", e), })?; - + if !response.status().is_success() { return Err(ProviderError::NetworkError { message: format!("Horizon returned HTTP {}", response.status()), }); } - - let transaction_response: HorizonTransactionResponse = response - .json() - .await - .map_err(|e| ProviderError::FormatError { - message: format!("Failed to parse transaction response: {}", e), - })?; - + + let transaction_response: HorizonTransactionResponse = + response + .json() + .await + .map_err(|e| ProviderError::FormatError { + message: format!("Failed to parse transaction response: {}", e), + })?; + Ok(transaction_response.embedded.records) } - + /// Convert Horizon transaction record to FeeDataPoint - fn convert_to_fee_data_point(&self, record: HorizonTransactionRecord) -> ProviderResult { + fn convert_to_fee_data_point( + &self, + record: HorizonTransactionRecord, + ) -> ProviderResult { // Only include successful transactions if !record.successful { return Err(ProviderError::FormatError { message: "Transaction was not successful".to_string(), }); } - + // Parse fee amount - let fee_amount = u64::from_str(&record.fee_charged) - .map_err(|e| ProviderError::FormatError { + let fee_amount = + u64::from_str(&record.fee_charged).map_err(|e| ProviderError::FormatError { message: format!("Invalid fee amount '{}': {}", record.fee_charged, e), })?; - + // Parse timestamp let timestamp = DateTime::parse_from_rfc3339(&record.created_at) .map_err(|e| ProviderError::FormatError { message: format!("Invalid timestamp '{}': {}", record.created_at, e), })? .with_timezone(&Utc); - + Ok(FeeDataPoint { fee_amount, timestamp, @@ -119,7 +127,7 @@ impl FeeDataProvider for HorizonFeeDataProvider { async fn fetch_latest_fees(&self) -> ProviderResult> { // Fetch recent transactions (last 100 by default) let transactions = self.fetch_recent_transactions(100).await?; - + // Convert to fee data points, filtering out failed conversions let mut fee_data_points = Vec::new(); for transaction in transactions { @@ -131,32 +139,33 @@ impl FeeDataProvider for HorizonFeeDataProvider { } } } - + if fee_data_points.is_empty() { return Err(ProviderError::FormatError { message: "No valid fee data points found in recent transactions".to_string(), }); } - + Ok(fee_data_points) } - + fn provider_name(&self) -> &str { "Horizon" } - + async fn health_check(&self) -> ProviderResult<()> { // Use the existing fee_stats endpoint for health check - self.client.fetch_fee_stats() + self.client + .fetch_fee_stats() .await .map_err(|e| ProviderError::NetworkError { message: format!("Horizon health check failed: {}", e), })?; - + Ok(()) } - + fn get_metadata(&self) -> ProviderMetadata { self.metadata.clone() } -} \ No newline at end of file +} diff --git a/packages/core/src/insights/mod.rs b/packages/core/src/insights/mod.rs index 5fa6c18..d77a0ff 100644 --- a/packages/core/src/insights/mod.rs +++ b/packages/core/src/insights/mod.rs @@ -1,5 +1,5 @@ //! Fee Insights Module -//! +//! //! This module provides analytical insights from raw blockchain fee data, //! including rolling averages, extremes tracking, and congestion detection. @@ -7,22 +7,22 @@ // and API server are wired up. Suppress dead-code warnings until then. #![allow(unused_imports)] -pub mod engine; pub mod calculator; -pub mod tracker; +pub mod config; pub mod detector; -pub mod types; +pub mod engine; pub mod error; -pub mod config; -pub mod provider; pub mod horizon_adapter; +pub mod provider; +pub mod tracker; +pub mod types; #[cfg(test)] mod tests; +pub use config::InsightsConfig; pub use engine::FeeInsightsEngine; -pub use types::*; pub use error::InsightsError; -pub use config::InsightsConfig; +pub use horizon_adapter::HorizonFeeDataProvider; pub use provider::{FeeDataProvider, ProviderMetadata}; -pub use horizon_adapter::HorizonFeeDataProvider; \ No newline at end of file +pub use types::*; diff --git a/packages/core/src/insights/provider.rs b/packages/core/src/insights/provider.rs index e5e1a08..696ae5b 100644 --- a/packages/core/src/insights/provider.rs +++ b/packages/core/src/insights/provider.rs @@ -1,28 +1,25 @@ //! Fee Data Provider Interface -//! +//! //! Provides abstraction layer for different fee data sources +use crate::insights::{error::ProviderError, types::FeeDataPoint}; use async_trait::async_trait; -use crate::insights::{ - types::FeeDataPoint, - error::ProviderError, -}; /// Trait for fee data providers to ensure data source independence #[async_trait] pub trait FeeDataProvider { /// Fetch the latest fee data from the provider async fn fetch_latest_fees(&self) -> Result, ProviderError>; - + /// Get the name of this provider for logging/debugging fn provider_name(&self) -> &str; - + /// Check if the provider is currently available async fn health_check(&self) -> Result<(), ProviderError> { // Default implementation - just try to fetch data self.fetch_latest_fees().await.map(|_| ()) } - + /// Get provider-specific configuration or metadata fn get_metadata(&self) -> ProviderMetadata { ProviderMetadata::default() @@ -50,4 +47,4 @@ impl Default for ProviderMetadata { } /// Result type for provider operations -pub type ProviderResult = Result; \ No newline at end of file +pub type ProviderResult = Result; diff --git a/packages/core/src/insights/tests.rs b/packages/core/src/insights/tests.rs index d65a462..908576c 100644 --- a/packages/core/src/insights/tests.rs +++ b/packages/core/src/insights/tests.rs @@ -1,5 +1,5 @@ //! Comprehensive tests for fee metrics core calculations -//! +//! //! This module contains all unit tests and property-based tests for the fee insights engine //! core calculations, including rolling averages, extremes tracking, and congestion detection. @@ -8,15 +8,15 @@ mod tests { use super::super::*; use crate::insights::{ - engine::FeeInsightsEngine, calculator::RollingAverageCalculator, - tracker::ExtremesTracker, + config::{AverageConfig, ExtremesConfig, InsightsConfig, SpikeConfig}, detector::CongestionDetector, - config::{AverageConfig, ExtremesConfig, SpikeConfig, InsightsConfig}, - types::*, + engine::FeeInsightsEngine, error::InsightsError, + tracker::ExtremesTracker, + types::*, }; - use chrono::{Utc, Duration}; + use chrono::{Duration, Utc}; use proptest::prelude::*; // Test data generators @@ -27,29 +27,29 @@ mod tests { Utc::now() - Duration::seconds(secs.abs() % (86400 * 30)) // Within last 30 days }), "[a-f0-9]{64}".prop_map(|s| s), // Transaction hash - 1u64..1_000_000u64, // Ledger sequence - ).prop_map(|(fee_amount, timestamp, transaction_hash, ledger_sequence)| { - FeeDataPoint { - fee_amount, - timestamp, - transaction_hash, - ledger_sequence, - } - }) + 1u64..1_000_000u64, // Ledger sequence + ) + .prop_map( + |(fee_amount, timestamp, transaction_hash, ledger_sequence)| FeeDataPoint { + fee_amount, + timestamp, + transaction_hash, + ledger_sequence, + }, + ) } fn time_window_strategy() -> impl Strategy { ( prop::collection::vec("[a-z]+", 1..10).prop_map(|words| words.join("_")), - 1i64..86400i64, // Duration in seconds (1 second to 1 day) + 1i64..86400i64, // Duration in seconds (1 second to 1 day) 1usize..100usize, // Min samples - ).prop_map(|(name, duration_secs, min_samples)| { - TimeWindow { + ) + .prop_map(|(name, duration_secs, min_samples)| TimeWindow { name, duration: Duration::seconds(duration_secs), min_samples, - } - }) + }) } // ============================================================================= @@ -59,16 +59,14 @@ mod tests { #[test] fn test_rolling_average_calculator_creation() { let config = AverageConfig::default(); - let time_windows = vec![ - TimeWindow { - name: "test".to_string(), - duration: Duration::hours(1), - min_samples: 5, - } - ]; - + let time_windows = vec![TimeWindow { + name: "test".to_string(), + duration: Duration::hours(1), + min_samples: 5, + }]; + let _calculator = RollingAverageCalculator::new(config, time_windows); - + // Should be created successfully (just verify no panic) } @@ -92,9 +90,9 @@ mod tests { min_samples: 2, }, ]; - + let mut calculator = RollingAverageCalculator::new(config, time_windows); - + // Add some test data let now = Utc::now(); let fee_points = vec![ @@ -111,13 +109,13 @@ mod tests { ledger_sequence: 2, }, ]; - + for point in fee_points { calculator.add_data_point(point); } - + let averages = calculator.calculate_averages().unwrap(); - + // Should calculate correct average: (100 + 200) / 2 = 150 assert_eq!(averages.short_term.value, 150.0); assert_eq!(averages.short_term.sample_count, 2); @@ -144,9 +142,9 @@ mod tests { min_samples: 5, }, ]; - + let mut calculator = RollingAverageCalculator::new(config, time_windows); - + // Add only 2 data points (less than required 5) let now = Utc::now(); calculator.add_data_point(FeeDataPoint { @@ -161,9 +159,9 @@ mod tests { transaction_hash: "hash2".to_string(), ledger_sequence: 2, }); - + let averages = calculator.calculate_averages().unwrap(); - + // Should be marked as partial assert!(averages.short_term.is_partial); assert_eq!(averages.short_term.sample_count, 2); @@ -189,11 +187,11 @@ mod tests { min_samples: 1, }, ]; - + let mut calculator = RollingAverageCalculator::new(config, time_windows); - + let now = Utc::now(); - + // Add old data point (outside window) calculator.add_data_point(FeeDataPoint { fee_amount: 100, @@ -201,7 +199,7 @@ mod tests { transaction_hash: "hash1".to_string(), ledger_sequence: 1, }); - + // Add recent data point (inside window) calculator.add_data_point(FeeDataPoint { fee_amount: 200, @@ -209,9 +207,9 @@ mod tests { transaction_hash: "hash2".to_string(), ledger_sequence: 2, }); - + let averages = calculator.calculate_averages().unwrap(); - + // Should only include the recent data point assert_eq!(averages.short_term.value, 200.0); assert_eq!(averages.short_term.sample_count, 1); @@ -240,11 +238,11 @@ mod tests { min_samples: 1, }, ]; - + let mut calculator = RollingAverageCalculator::new(config, time_windows); - + let now = Utc::now(); - + // Add 5 data points (more than buffer capacity of 3) for i in 0..5 { calculator.add_data_point(FeeDataPoint { @@ -254,12 +252,12 @@ mod tests { ledger_sequence: i + 1, }); } - + let averages = calculator.calculate_averages().unwrap(); - + // Should only have 3 data points (buffer capacity) assert_eq!(averages.short_term.sample_count, 3); - + // Should contain the most recent 3 points: 500, 400, 300 // Average should be (500 + 400 + 300) / 3 = 400 assert_eq!(averages.short_term.value, 400.0); @@ -285,10 +283,10 @@ mod tests { min_samples: 1, }, ]; - + let calculator = RollingAverageCalculator::new(config, time_windows); let averages = calculator.calculate_averages().unwrap(); - + // Should return zero values with appropriate metadata assert_eq!(averages.short_term.value, 0.0); assert_eq!(averages.short_term.sample_count, 0); @@ -303,7 +301,7 @@ mod tests { fn test_extremes_tracker_creation() { let config = ExtremesConfig::default(); let tracker = ExtremesTracker::new(config); - + // Should be created successfully assert!(!tracker.has_current_data()); } @@ -312,7 +310,7 @@ mod tests { fn test_extremes_identification() { let config = ExtremesConfig::default(); let mut tracker = ExtremesTracker::new(config); - + let now = Utc::now(); let fee_data = vec![ FeeDataPoint { @@ -329,15 +327,15 @@ mod tests { }, FeeDataPoint { fee_amount: 300, // Maximum - timestamp: now, // Use current time + timestamp: now, // Use current time transaction_hash: "hash3".to_string(), ledger_sequence: 3, }, ]; - + tracker.update_with_fees(&fee_data).unwrap(); let extremes = tracker.get_current_extremes().unwrap(); - + // Should correctly identify min and max assert_eq!(extremes.current_min.value, 50); assert_eq!(extremes.current_min.transaction_hash, "hash2"); @@ -349,26 +347,26 @@ mod tests { fn test_extremes_tie_breaking() { let config = ExtremesConfig::default(); let mut tracker = ExtremesTracker::new(config); - + let now = Utc::now(); let fee_data = vec![ FeeDataPoint { - fee_amount: 100, // First occurrence of min + fee_amount: 100, // First occurrence of min timestamp: now - Duration::seconds(1), // Slightly earlier transaction_hash: "hash1".to_string(), ledger_sequence: 1, }, FeeDataPoint { fee_amount: 100, // Second occurrence of min (more recent) - timestamp: now, // More recent + timestamp: now, // More recent transaction_hash: "hash2".to_string(), ledger_sequence: 2, }, ]; - + tracker.update_with_fees(&fee_data).unwrap(); let extremes = tracker.get_current_extremes().unwrap(); - + // Should store the most recent occurrence assert_eq!(extremes.current_min.value, 100); assert_eq!(extremes.current_min.transaction_hash, "hash2"); @@ -378,20 +376,18 @@ mod tests { fn test_extremes_metadata_preservation() { let config = ExtremesConfig::default(); let mut tracker = ExtremesTracker::new(config); - + let now = Utc::now(); - let fee_data = vec![ - FeeDataPoint { - fee_amount: 200, - timestamp: now, - transaction_hash: "test_hash_123".to_string(), - ledger_sequence: 12345, - }, - ]; - + let fee_data = vec![FeeDataPoint { + fee_amount: 200, + timestamp: now, + transaction_hash: "test_hash_123".to_string(), + ledger_sequence: 12345, + }]; + tracker.update_with_fees(&fee_data).unwrap(); let extremes = tracker.get_current_extremes().unwrap(); - + // Should preserve all metadata assert_eq!(extremes.current_min.transaction_hash, "test_hash_123"); assert_eq!(extremes.current_max.transaction_hash, "test_hash_123"); @@ -407,7 +403,7 @@ mod tests { fn test_congestion_detector_creation() { let config = SpikeConfig::default(); let detector = CongestionDetector::new(config); - + // Should be created successfully assert_eq!(detector.get_recent_spikes().len(), 0); } @@ -420,10 +416,10 @@ mod tests { congestion_window: Duration::hours(1), }; let detector = CongestionDetector::new(config); - + let now = Utc::now(); let baseline = 100.0; - + let fee_data = vec![ FeeDataPoint { fee_amount: 100, // Normal fee @@ -450,9 +446,9 @@ mod tests { ledger_sequence: 4, }, ]; - + let spikes = detector.detect_spikes(&fee_data, baseline).unwrap(); - + // Should detect one spike assert_eq!(spikes.len(), 1); assert_eq!(spikes[0].peak_fee, 300); @@ -464,12 +460,18 @@ mod tests { fn test_spike_severity_classification() { let config = SpikeConfig::default(); let detector = CongestionDetector::new(config); - + // Test different severity levels assert_eq!(detector.classify_spike_severity(1.5), SpikeSeverity::Minor); - assert_eq!(detector.classify_spike_severity(3.5), SpikeSeverity::Moderate); + assert_eq!( + detector.classify_spike_severity(3.5), + SpikeSeverity::Moderate + ); assert_eq!(detector.classify_spike_severity(7.0), SpikeSeverity::Major); - assert_eq!(detector.classify_spike_severity(15.0), SpikeSeverity::Critical); + assert_eq!( + detector.classify_spike_severity(15.0), + SpikeSeverity::Critical + ); } #[test] @@ -480,10 +482,10 @@ mod tests { congestion_window: Duration::hours(1), }; let detector = CongestionDetector::new(config); - + let now = Utc::now(); let baseline = 200.0; - + let fee_data = vec![ FeeDataPoint { fee_amount: 600, // 3x baseline, should exceed threshold of 2.0 @@ -498,9 +500,9 @@ mod tests { ledger_sequence: 2, }, ]; - + let spikes = detector.detect_spikes(&fee_data, baseline).unwrap(); - + assert_eq!(spikes.len(), 1); assert_eq!(spikes[0].spike_ratio, 3.0); assert_eq!(spikes[0].baseline_fee, baseline); @@ -514,45 +516,42 @@ mod tests { fn test_fee_amount_validation() { let config = InsightsConfig::default(); let engine = FeeInsightsEngine::new(config); - + // Test zero fee amount (should fail) - let invalid_data = vec![ - FeeDataPoint { - fee_amount: 0, - timestamp: Utc::now(), - transaction_hash: "hash1".to_string(), - ledger_sequence: 1, - } - ]; - + let invalid_data = vec![FeeDataPoint { + fee_amount: 0, + timestamp: Utc::now(), + transaction_hash: "hash1".to_string(), + ledger_sequence: 1, + }]; + let result = engine.validate_fee_data(&invalid_data); assert!(result.is_err()); assert!(result.unwrap_err().to_string().contains("Zero fee amount")); - + // Test excessively large fee amount (should fail) - let invalid_data = vec![ - FeeDataPoint { - fee_amount: 2_000_000_000, // > 1 billion stroops - timestamp: Utc::now(), - transaction_hash: "hash1".to_string(), - ledger_sequence: 1, - } - ]; - + let invalid_data = vec![FeeDataPoint { + fee_amount: 2_000_000_000, // > 1 billion stroops + timestamp: Utc::now(), + transaction_hash: "hash1".to_string(), + ledger_sequence: 1, + }]; + let result = engine.validate_fee_data(&invalid_data); assert!(result.is_err()); - assert!(result.unwrap_err().to_string().contains("Unreasonably large fee amount")); - + assert!(result + .unwrap_err() + .to_string() + .contains("Unreasonably large fee amount")); + // Test valid fee amount (should pass) - let valid_data = vec![ - FeeDataPoint { - fee_amount: 100, - timestamp: Utc::now(), - transaction_hash: "hash1".to_string(), - ledger_sequence: 1, - } - ]; - + let valid_data = vec![FeeDataPoint { + fee_amount: 100, + timestamp: Utc::now(), + transaction_hash: "hash1".to_string(), + ledger_sequence: 1, + }]; + let result = engine.validate_fee_data(&valid_data); assert!(result.is_ok()); } @@ -561,31 +560,27 @@ mod tests { fn test_timestamp_validation() { let config = InsightsConfig::default(); let engine = FeeInsightsEngine::new(config); - + // Test future timestamp (should fail) - let invalid_data = vec![ - FeeDataPoint { - fee_amount: 100, - timestamp: Utc::now() + Duration::hours(2), // 2 hours in future - transaction_hash: "hash1".to_string(), - ledger_sequence: 1, - } - ]; - + let invalid_data = vec![FeeDataPoint { + fee_amount: 100, + timestamp: Utc::now() + Duration::hours(2), // 2 hours in future + transaction_hash: "hash1".to_string(), + ledger_sequence: 1, + }]; + let result = engine.validate_fee_data(&invalid_data); assert!(result.is_err()); assert!(result.unwrap_err().to_string().contains("Future timestamp")); - + // Test valid timestamp (should pass) - let valid_data = vec![ - FeeDataPoint { - fee_amount: 100, - timestamp: Utc::now() - Duration::minutes(30), - transaction_hash: "hash1".to_string(), - ledger_sequence: 1, - } - ]; - + let valid_data = vec![FeeDataPoint { + fee_amount: 100, + timestamp: Utc::now() - Duration::minutes(30), + transaction_hash: "hash1".to_string(), + ledger_sequence: 1, + }]; + let result = engine.validate_fee_data(&valid_data); assert!(result.is_ok()); } @@ -594,31 +589,30 @@ mod tests { fn test_transaction_hash_validation() { let config = InsightsConfig::default(); let engine = FeeInsightsEngine::new(config); - + // Test empty transaction hash (should fail) - let invalid_data = vec![ - FeeDataPoint { - fee_amount: 100, - timestamp: Utc::now(), - transaction_hash: "".to_string(), - ledger_sequence: 1, - } - ]; - + let invalid_data = vec![FeeDataPoint { + fee_amount: 100, + timestamp: Utc::now(), + transaction_hash: "".to_string(), + ledger_sequence: 1, + }]; + let result = engine.validate_fee_data(&invalid_data); assert!(result.is_err()); - assert!(result.unwrap_err().to_string().contains("Empty transaction hash")); - + assert!(result + .unwrap_err() + .to_string() + .contains("Empty transaction hash")); + // Test valid transaction hash (should pass) - let valid_data = vec![ - FeeDataPoint { - fee_amount: 100, - timestamp: Utc::now(), - transaction_hash: "valid_hash_123".to_string(), - ledger_sequence: 1, - } - ]; - + let valid_data = vec![FeeDataPoint { + fee_amount: 100, + timestamp: Utc::now(), + transaction_hash: "valid_hash_123".to_string(), + ledger_sequence: 1, + }]; + let result = engine.validate_fee_data(&valid_data); assert!(result.is_ok()); } @@ -647,11 +641,11 @@ mod tests { min_samples: 1, }, ]; - + let mut calculator = RollingAverageCalculator::new(config, time_windows); - + let now = Utc::now(); - + // Test with large fee amounts (but within valid range) let large_fees = vec![ FeeDataPoint { @@ -667,13 +661,13 @@ mod tests { ledger_sequence: 2, }, ]; - + for fee in large_fees { calculator.add_data_point(fee); } - + let averages = calculator.calculate_averages().unwrap(); - + // Should maintain accuracy with large numbers let expected_average = (999_999_999.0 + 999_999_998.0) / 2.0; assert_eq!(averages.short_term.value, expected_average); @@ -685,15 +679,15 @@ mod tests { let fee_amount: u64 = 123_456_789; let converted: f64 = fee_amount as f64; let back_converted: u64 = converted as u64; - + // Should preserve value for reasonable fee amounts assert_eq!(fee_amount, back_converted); - + // Test with maximum safe integer in f64 let max_safe: u64 = (1u64 << 53) - 1; // 2^53 - 1 let converted_max: f64 = max_safe as f64; let back_converted_max: u64 = converted_max as u64; - + assert_eq!(max_safe, back_converted_max); } @@ -701,47 +695,49 @@ mod tests { fn test_division_by_zero_handling() { let config = SpikeConfig::default(); let detector = CongestionDetector::new(config); - + let now = Utc::now(); - let fee_data = vec![ - FeeDataPoint { - fee_amount: 100, - timestamp: now, - transaction_hash: "hash1".to_string(), - ledger_sequence: 1, - } - ]; - + let fee_data = vec![FeeDataPoint { + fee_amount: 100, + timestamp: now, + transaction_hash: "hash1".to_string(), + ledger_sequence: 1, + }]; + // Test with zero baseline (should return error) let result = detector.detect_spikes(&fee_data, 0.0); assert!(result.is_err()); - assert!(result.unwrap_err().to_string().contains("Baseline must be positive")); - + assert!(result + .unwrap_err() + .to_string() + .contains("Baseline must be positive")); + // Test with negative baseline (should return error) let result = detector.detect_spikes(&fee_data, -100.0); assert!(result.is_err()); - assert!(result.unwrap_err().to_string().contains("Baseline must be positive")); + assert!(result + .unwrap_err() + .to_string() + .contains("Baseline must be positive")); } #[test] fn test_decimal_precision_maintenance() { let config = SpikeConfig::default(); let detector = CongestionDetector::new(config); - + let now = Utc::now(); let baseline = 100.0; - - let fee_data = vec![ - FeeDataPoint { - fee_amount: 333, // Should give ratio of 3.33 - timestamp: now, - transaction_hash: "hash1".to_string(), - ledger_sequence: 1, - } - ]; - + + let fee_data = vec![FeeDataPoint { + fee_amount: 333, // Should give ratio of 3.33 + timestamp: now, + transaction_hash: "hash1".to_string(), + ledger_sequence: 1, + }]; + let spikes = detector.detect_spikes(&fee_data, baseline).unwrap(); - + if !spikes.is_empty() { // Should maintain decimal precision assert_eq!(spikes[0].spike_ratio, 3.33); @@ -762,13 +758,13 @@ mod tests { ) { let config = AverageConfig::default(); let now = Utc::now(); - + // Adjust all timestamps to be within the time window let adjusted_fee_points: Vec = fee_points.into_iter().map(|mut point| { point.timestamp = now - Duration::minutes((rand::random::() % 60) as i64); // Within last hour point }).collect(); - + let time_windows = vec![ TimeWindow { name: "short_term".to_string(), @@ -786,20 +782,20 @@ mod tests { min_samples: 1, }, ]; - + let mut calculator = RollingAverageCalculator::new(config, time_windows); - + // Add all fee points for point in &adjusted_fee_points { calculator.add_data_point(point.clone()); } - + let averages = calculator.calculate_averages().unwrap(); - + // Calculate expected average manually let total: u64 = adjusted_fee_points.iter().map(|p| p.fee_amount).sum(); let expected_average = total as f64 / adjusted_fee_points.len() as f64; - + // Should match calculated average prop_assert_eq!(averages.short_term.value, expected_average); prop_assert_eq!(averages.short_term.sample_count, adjusted_fee_points.len()); @@ -812,14 +808,14 @@ mod tests { ) { let config = ExtremesConfig::default(); let mut tracker = ExtremesTracker::new(config); - + tracker.update_with_fees(&fee_points).unwrap(); - + if let Ok(extremes) = tracker.get_current_extremes() { // Find actual min and max let actual_min = fee_points.iter().map(|p| p.fee_amount).min().unwrap(); let actual_max = fee_points.iter().map(|p| p.fee_amount).max().unwrap(); - + // Should match tracked extremes prop_assert_eq!(extremes.current_min.value, actual_min); prop_assert_eq!(extremes.current_max.value, actual_max); @@ -838,9 +834,9 @@ mod tests { congestion_window: Duration::hours(1), }; let detector = CongestionDetector::new(config); - + let spikes = detector.detect_spikes(&fee_points, baseline).unwrap(); - + // Verify all spike ratios are calculated correctly for spike in spikes { let expected_ratio = spike.peak_fee as f64 / baseline; @@ -856,7 +852,7 @@ mod tests { ) { let config = InsightsConfig::default(); let engine = FeeInsightsEngine::new(config); - + let fee_data = vec![ FeeDataPoint { fee_amount, @@ -865,7 +861,7 @@ mod tests { ledger_sequence: 1, } ]; - + // Should pass validation for valid fee amounts let result = engine.validate_fee_data(&fee_data); prop_assert!(result.is_ok()); @@ -879,7 +875,7 @@ mod tests { // Convert u64 to f64 and back let as_float: f64 = fee_amount as f64; let back_to_int: u64 = as_float as u64; - + // Should preserve the original value prop_assert_eq!(fee_amount, back_to_int); } @@ -893,7 +889,7 @@ mod tests { fn test_full_engine_integration() { let config = InsightsConfig::default(); let mut engine = FeeInsightsEngine::new(config); - + let now = Utc::now(); let fee_data = vec![ FeeDataPoint { @@ -921,13 +917,13 @@ mod tests { ledger_sequence: 4, }, ]; - + // Process the fee data let result = tokio_test::block_on(engine.process_fee_data(&fee_data)); assert!(result.is_ok()); - + let update = result.unwrap(); - + // Verify insights were calculated assert!(update.insights.rolling_averages.short_term.value > 0.0); assert!(update.insights.extremes.current_min.value > 0); @@ -939,24 +935,22 @@ mod tests { fn test_engine_reset_functionality() { let config = InsightsConfig::default(); let mut engine = FeeInsightsEngine::new(config); - + let now = Utc::now(); - let fee_data = vec![ - FeeDataPoint { - fee_amount: 200, - timestamp: now - Duration::minutes(30), - transaction_hash: "hash1".to_string(), - ledger_sequence: 1, - } - ]; - + let fee_data = vec![FeeDataPoint { + fee_amount: 200, + timestamp: now - Duration::minutes(30), + transaction_hash: "hash1".to_string(), + ledger_sequence: 1, + }]; + // Process some data let _result = tokio_test::block_on(engine.process_fee_data(&fee_data)); - + // Reset the engine let reset_result = engine.reset(); assert!(reset_result.is_ok()); - + // Verify reset worked assert!(engine.get_last_update().is_none()); let insights = engine.get_current_insights(); diff --git a/packages/core/src/insights/tracker.rs b/packages/core/src/insights/tracker.rs index f5068e8..76e78f7 100644 --- a/packages/core/src/insights/tracker.rs +++ b/packages/core/src/insights/tracker.rs @@ -3,11 +3,7 @@ use chrono::{DateTime, Utc}; use std::collections::VecDeque; -use crate::insights::{ - types::*, - error::InsightsError, - config::ExtremesConfig, -}; +use crate::insights::{config::ExtremesConfig, error::InsightsError, types::*}; /// Represents a tracking period for extremes #[derive(Debug, Clone)] @@ -27,14 +23,14 @@ impl ExtremePeriod { period_end: end, } } - + fn update_with_fee(&mut self, fee_point: &FeeDataPoint) { let extreme_value = ExtremeValue { value: fee_point.fee_amount, timestamp: fee_point.timestamp, transaction_hash: fee_point.transaction_hash.clone(), }; - + // Update minimum match &self.min_value { None => self.min_value = Some(extreme_value.clone()), @@ -44,7 +40,7 @@ impl ExtremePeriod { } } } - + // Update maximum match &self.max_value { None => self.max_value = Some(extreme_value), @@ -55,7 +51,7 @@ impl ExtremePeriod { } } } - + fn to_fee_extremes(&self) -> Option { match (&self.min_value, &self.max_value) { (Some(min), Some(max)) => Some(FeeExtremes { @@ -82,59 +78,61 @@ impl ExtremesTracker { let now = Utc::now(); let period_start = now; let period_end = now + config.tracking_period; - + Self { config, current_period: ExtremePeriod::new(period_start, period_end), historical_periods: VecDeque::new(), } } - + /// Update with new fee data pub fn update_with_fees(&mut self, fees: &[FeeDataPoint]) -> Result<(), InsightsError> { let now = Utc::now(); - + // Check if we need to rotate to a new period if now >= self.current_period.period_end { self.rotate_period(now)?; } - + // Update current period with new fees for fee_point in fees { // Only process fees that are within the current tracking period - if fee_point.timestamp >= self.current_period.period_start - && fee_point.timestamp <= self.current_period.period_end { + if fee_point.timestamp >= self.current_period.period_start + && fee_point.timestamp <= self.current_period.period_end + { self.current_period.update_with_fee(fee_point); } } - + Ok(()) } - + /// Rotate to a new tracking period, preserving the current period as historical fn rotate_period(&mut self, current_time: DateTime) -> Result<(), InsightsError> { // Move current period to historical periods let completed_period = std::mem::replace( &mut self.current_period, - ExtremePeriod::new(current_time, current_time + self.config.tracking_period) + ExtremePeriod::new(current_time, current_time + self.config.tracking_period), ); - + self.historical_periods.push_back(completed_period); - + // Maintain the configured number of historical periods while self.historical_periods.len() > self.config.historical_periods_to_keep { self.historical_periods.pop_front(); } - + Ok(()) } - + /// Get current extremes pub fn get_current_extremes(&self) -> Result { - self.current_period.to_fee_extremes() - .ok_or_else(|| InsightsError::insufficient_data("No fee data available for current period")) + self.current_period.to_fee_extremes().ok_or_else(|| { + InsightsError::insufficient_data("No fee data available for current period") + }) } - + /// Get historical extremes for a specific number of periods back pub fn get_historical_extremes(&self, periods_back: usize) -> Vec { self.historical_periods @@ -144,7 +142,7 @@ impl ExtremesTracker { .filter_map(|period| period.to_fee_extremes()) .collect() } - + /// Get all historical extremes pub fn get_all_historical_extremes(&self) -> Vec { self.historical_periods @@ -152,20 +150,20 @@ impl ExtremesTracker { .filter_map(|period| period.to_fee_extremes()) .collect() } - + /// Reset the current tracking period while preserving historical data pub fn reset_current_period(&mut self) -> Result<(), InsightsError> { let now = Utc::now(); - + // If the current period has data, preserve it as historical if self.current_period.min_value.is_some() || self.current_period.max_value.is_some() { let completed_period = std::mem::replace( &mut self.current_period, - ExtremePeriod::new(now, now + self.config.tracking_period) + ExtremePeriod::new(now, now + self.config.tracking_period), ); - + self.historical_periods.push_back(completed_period); - + // Maintain the configured number of historical periods while self.historical_periods.len() > self.config.historical_periods_to_keep { self.historical_periods.pop_front(); @@ -174,22 +172,25 @@ impl ExtremesTracker { // Just reset the current period if it has no data self.current_period = ExtremePeriod::new(now, now + self.config.tracking_period); } - + Ok(()) } - + /// Check if the current period has any data pub fn has_current_data(&self) -> bool { self.current_period.min_value.is_some() || self.current_period.max_value.is_some() } - + /// Get the number of historical periods stored pub fn historical_period_count(&self) -> usize { self.historical_periods.len() } - + /// Get the current tracking period information pub fn get_current_period_info(&self) -> (DateTime, DateTime) { - (self.current_period.period_start, self.current_period.period_end) + ( + self.current_period.period_start, + self.current_period.period_end, + ) } -} \ No newline at end of file +} diff --git a/packages/core/src/insights/types.rs b/packages/core/src/insights/types.rs index 0598053..ed357a3 100644 --- a/packages/core/src/insights/types.rs +++ b/packages/core/src/insights/types.rs @@ -1,6 +1,6 @@ //! Core data types for fee insights -use chrono::{DateTime, Utc, Duration}; +use chrono::{DateTime, Duration, Utc}; use serde::{Deserialize, Serialize}; /// A single fee data point from the blockchain @@ -25,9 +25,9 @@ pub struct CurrentInsights { /// Rolling averages across different time windows #[derive(Debug, Clone, Serialize, Deserialize)] pub struct RollingAverages { - pub short_term: AverageResult, // 1 hour - pub medium_term: AverageResult, // 6 hours - pub long_term: AverageResult, // 24 hours + pub short_term: AverageResult, // 1 hour + pub medium_term: AverageResult, // 6 hours + pub long_term: AverageResult, // 24 hours } /// Result of a rolling average calculation @@ -114,7 +114,7 @@ pub enum SpikeSeverity { /// Data quality indicators #[derive(Debug, Clone, Serialize, Deserialize)] pub struct DataQuality { - pub completeness: f64, // 0.0 to 1.0 + pub completeness: f64, // 0.0 to 1.0 pub freshness: Duration, pub has_gaps: bool, pub last_gap: Option>, @@ -126,4 +126,4 @@ pub struct InsightsUpdate { pub insights: CurrentInsights, pub processing_time: Duration, pub data_points_processed: usize, -} \ No newline at end of file +} diff --git a/packages/core/src/logging.rs b/packages/core/src/logging.rs index f901f83..977044b 100644 --- a/packages/core/src/logging.rs +++ b/packages/core/src/logging.rs @@ -5,8 +5,7 @@ use tracing_subscriber::{fmt, EnvFilter}; /// /// This must be called once at startup (in main.rs). pub fn init_logging() { - let filter = EnvFilter::try_from_default_env() - .unwrap_or_else(|_| EnvFilter::new("info")); + let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info")); fmt() .with_env_filter(filter) @@ -15,4 +14,4 @@ pub fn init_logging() { .init(); info!("Logging initialized"); -} \ No newline at end of file +} diff --git a/packages/core/src/main.rs b/packages/core/src/main.rs index d027bc9..45d3cd8 100644 --- a/packages/core/src/main.rs +++ b/packages/core/src/main.rs @@ -5,39 +5,39 @@ mod alerts; mod api; mod cache; -mod metrics; mod cli; mod config; mod db; mod error; mod insights; mod logging; +mod metrics; mod middleware; mod repository; -mod services; mod scheduler; +mod services; mod store; use std::sync::Arc; -use axum::{routing::get, Router}; use axum::http::{HeaderName, Method}; +use axum::{routing::get, Router}; use clap::Parser; use dotenvy::dotenv; use std::time::Duration; use tokio::sync::{Mutex, RwLock}; use tower_http::cors::{AllowOrigin, CorsLayer}; +use crate::alerts::AlertManager; use crate::cache::ResponseCache; use crate::cli::Cli; use crate::config::Config; use crate::error::AppError; -use crate::insights::{FeeInsightsEngine, InsightsConfig, HorizonFeeDataProvider}; +use crate::insights::{FeeInsightsEngine, HorizonFeeDataProvider, InsightsConfig}; use crate::logging::init_logging; use crate::metrics::AppMetrics; -use crate::alerts::AlertManager; use crate::middleware::auth::require_api_key; -use crate::middleware::rate_limit::{RateLimitState, enforce_rate_limit}; +use crate::middleware::rate_limit::{enforce_rate_limit, RateLimitState}; use crate::repository::FeeRepository; use crate::scheduler::run_fee_polling_with_retry; use crate::services::horizon::HorizonClient; @@ -90,12 +90,10 @@ async fn main() { tracing::info!("Database initialised: {}", config.database_url); // ---- Metrics ---- - let app_metrics = Arc::new( - AppMetrics::new().unwrap_or_else(|err| { - tracing::error!("Failed to initialise Prometheus metrics: {}", err); - std::process::exit(1); - }), - ); + let app_metrics = Arc::new(AppMetrics::new().unwrap_or_else(|err| { + tracing::error!("Failed to initialise Prometheus metrics: {}", err); + std::process::exit(1); + })); let repository = Arc::new(FeeRepository::new(db_pool)); @@ -105,9 +103,9 @@ async fn main() { let fee_store = Arc::new(RwLock::new(FeeHistoryStore::new(DEFAULT_CAPACITY))); - let insights_engine = Arc::new(RwLock::new( - FeeInsightsEngine::new(InsightsConfig::default()), - )); + let insights_engine = Arc::new(RwLock::new(FeeInsightsEngine::new( + InsightsConfig::default(), + ))); let current_fees_cache = Arc::new(Mutex::new(ResponseCache::new(Duration::from_secs( config.cache_ttl_seconds, )))); @@ -134,9 +132,7 @@ async fn main() { Ok(_) => tracing::info!("No historical fee data found — starting cold"), Err(err) => tracing::warn!("Failed to rehydrate store from database: {}", err), } - let horizon_provider = Arc::new(HorizonFeeDataProvider::new( - (*horizon_client).clone(), - )); + let horizon_provider = Arc::new(HorizonFeeDataProvider::new((*horizon_client).clone())); let fee_stats_provider: Arc = horizon_client.clone(); let alert_manager = Arc::new(AlertManager::new( @@ -195,8 +191,7 @@ async fn main() { // Clone for metrics endpoint closure let metrics_for_handler = app_metrics.clone(); - let app = Router::new() - .route("/health", get(api::health::health)); + let app = Router::new().route("/health", get(api::health::health)); let protected_routes = Router::new() .route( @@ -225,14 +220,31 @@ async fn main() { }), ) .merge(fees_router) - .merge(api::insights::create_insights_router(insights_engine.clone())) + .merge(api::insights::create_insights_router( + insights_engine.clone(), + )) .merge( Router::new() - .route("/alerts/config", axum::routing::post(api::alerts::create_alert)) - .route("/alerts/config", axum::routing::get(api::alerts::list_alerts)) - .route("/alerts/config/:id", axum::routing::patch(api::alerts::update_alert)) - .route("/alerts/config/:id", axum::routing::delete(api::alerts::delete_alert)) - .route("/alerts/history", axum::routing::get(api::alerts::get_alert_history)) + .route( + "/alerts/config", + axum::routing::post(api::alerts::create_alert), + ) + .route( + "/alerts/config", + axum::routing::get(api::alerts::list_alerts), + ) + .route( + "/alerts/config/:id", + axum::routing::patch(api::alerts::update_alert), + ) + .route( + "/alerts/config/:id", + axum::routing::delete(api::alerts::delete_alert), + ) + .route( + "/alerts/history", + axum::routing::get(api::alerts::get_alert_history), + ) .with_state(repository.clone()), ); @@ -273,8 +285,8 @@ async fn main() { listener, app.into_make_service_with_connect_info::(), ) - .await - .unwrap_or_else(|err| tracing::error!("Server error: {}", err)); + .await + .unwrap_or_else(|err| tracing::error!("Server error: {}", err)); }, run_fee_polling_with_retry( horizon_provider, diff --git a/packages/core/src/metrics.rs b/packages/core/src/metrics.rs index 9eaab06..c28e3fd 100644 --- a/packages/core/src/metrics.rs +++ b/packages/core/src/metrics.rs @@ -8,9 +8,7 @@ //! (`text/plain; version=0.0.4`). The endpoint is intentionally excluded //! from API-key auth so it can be scraped by Prometheus / Grafana agents. -use prometheus::{ - Counter, CounterVec, Gauge, Histogram, HistogramOpts, Opts, Registry, -}; +use prometheus::{Counter, CounterVec, Gauge, Histogram, HistogramOpts, Opts, Registry}; /// All application-level Prometheus metrics. pub struct AppMetrics { @@ -76,7 +74,9 @@ impl AppMetrics { "stellar_fee_tracker_http_request_duration_seconds", "HTTP request latency in seconds", ) - .buckets(vec![0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0]), + .buckets(vec![ + 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, + ]), )?; registry.register(Box::new(polls_total.clone()))?; @@ -117,7 +117,11 @@ mod tests { #[test] fn all_metrics_register_without_error() { let metrics = AppMetrics::new(); - assert!(metrics.is_ok(), "AppMetrics::new() failed: {:?}", metrics.err()); + assert!( + metrics.is_ok(), + "AppMetrics::new() failed: {:?}", + metrics.err() + ); } #[test] @@ -220,7 +224,12 @@ mod integration_tests { .body(Body::empty()) .unwrap(); let resp = app.oneshot(req).await.unwrap(); - let ct = resp.headers().get("content-type").unwrap().to_str().unwrap(); + let ct = resp + .headers() + .get("content-type") + .unwrap() + .to_str() + .unwrap(); assert_eq!(ct, "text/plain; version=0.0.4"); } @@ -275,4 +284,4 @@ mod integration_tests { // Prometheus text format: metric_name value\n assert!(body.contains("stellar_fee_tracker_polls_total 5")); } -} \ No newline at end of file +} diff --git a/packages/core/src/middleware/auth.rs b/packages/core/src/middleware/auth.rs index f7017ce..c627106 100644 --- a/packages/core/src/middleware/auth.rs +++ b/packages/core/src/middleware/auth.rs @@ -75,8 +75,8 @@ fn constant_time_eq(a: &[u8], b: &[u8]) -> bool { mod tests { use super::*; use axum::{ - body::{Body, to_bytes}, - http::{HeaderValue, header}, + body::{to_bytes, Body}, + http::{header, HeaderValue}, middleware::from_fn_with_state, routing::get, Router, @@ -149,10 +149,7 @@ mod tests { assert_eq!(response.status(), StatusCode::UNAUTHORIZED); let body = to_bytes(response.into_body(), usize::MAX).await.unwrap(); let payload: serde_json::Value = serde_json::from_slice(&body).unwrap(); - assert_eq!( - payload["error"], - "Unauthorized: missing or invalid API key" - ); + assert_eq!(payload["error"], "Unauthorized: missing or invalid API key"); } #[tokio::test] @@ -172,10 +169,7 @@ mod tests { assert_eq!(response.status(), StatusCode::UNAUTHORIZED); let body = to_bytes(response.into_body(), usize::MAX).await.unwrap(); let payload: serde_json::Value = serde_json::from_slice(&body).unwrap(); - assert_eq!( - payload["error"], - "Unauthorized: missing or invalid API key" - ); + assert_eq!(payload["error"], "Unauthorized: missing or invalid API key"); } #[tokio::test] diff --git a/packages/core/src/middleware/rate_limit.rs b/packages/core/src/middleware/rate_limit.rs index 50478da..33e87e1 100644 --- a/packages/core/src/middleware/rate_limit.rs +++ b/packages/core/src/middleware/rate_limit.rs @@ -3,11 +3,11 @@ use std::sync::Arc; use std::time::Instant; use axum::{ - Json, - extract::{Request, State, connect_info::ConnectInfo}, - http::{HeaderValue, StatusCode, header}, + extract::{connect_info::ConnectInfo, Request, State}, + http::{header, HeaderValue, StatusCode}, middleware::Next, response::{IntoResponse, Response}, + Json, }; use dashmap::DashMap; use serde_json::json; @@ -158,7 +158,11 @@ fn parse_x_forwarded_for(value: &str) -> Option { fn attach_rate_limit_headers(response: &mut Response, limit: u32, remaining: u32, reset_secs: u64) { insert_number_header(response, X_RATE_LIMIT_LIMIT_HEADER, u64::from(limit)); - insert_number_header(response, X_RATE_LIMIT_REMAINING_HEADER, u64::from(remaining)); + insert_number_header( + response, + X_RATE_LIMIT_REMAINING_HEADER, + u64::from(remaining), + ); insert_number_header(response, X_RATE_LIMIT_RESET_HEADER, reset_secs); } @@ -172,7 +176,7 @@ fn insert_number_header(response: &mut Response, name: &'static str, value: u64) mod tests { use super::*; use axum::{ - body::{Body, to_bytes}, + body::{to_bytes, Body}, extract::connect_info::ConnectInfo, http::Request, middleware::from_fn_with_state, @@ -207,10 +211,7 @@ mod tests { } async fn request_with_connect_info(app: &Router, addr: SocketAddr) -> Response { - let mut request = Request::builder() - .uri("/test") - .body(Body::empty()) - .unwrap(); + let mut request = Request::builder().uri("/test").body(Body::empty()).unwrap(); request.extensions_mut().insert(ConnectInfo(addr)); app.clone().oneshot(request).await.unwrap() @@ -262,10 +263,7 @@ mod tests { let payload: serde_json::Value = serde_json::from_slice(&body).unwrap(); assert_eq!( payload["error"], - format!( - "Rate limit exceeded. Try again in {} seconds.", - retry_after - ) + format!("Rate limit exceeded. Try again in {} seconds.", retry_after) ); } diff --git a/packages/core/src/repository.rs b/packages/core/src/repository.rs index 4d5c886..8f7f128 100644 --- a/packages/core/src/repository.rs +++ b/packages/core/src/repository.rs @@ -27,7 +27,6 @@ pub struct AlertConfig { pub created_at: String, } - /// A single fired-alert log entry. #[derive(Debug, Clone, Serialize, Deserialize)] pub struct AlertEvent { @@ -149,12 +148,10 @@ impl FeeRepository { pub async fn prune_older_than(&self, cutoff: DateTime) -> Result { let cutoff_str = cutoff.to_rfc3339(); - let result = sqlx::query( - "DELETE FROM fee_data_points WHERE timestamp < ?", - ) - .bind(&cutoff_str) - .execute(&self.pool) - .await?; + let result = sqlx::query("DELETE FROM fee_data_points WHERE timestamp < ?") + .bind(&cutoff_str) + .execute(&self.pool) + .await?; Ok(result.rows_affected()) } @@ -167,13 +164,12 @@ impl FeeRepository { webhook_url: &str, threshold: &str, ) -> Result { - let result = sqlx::query( - "INSERT INTO alert_configs (webhook_url, threshold) VALUES (?, ?)", - ) - .bind(webhook_url) - .bind(threshold) - .execute(&self.pool) - .await?; + let result = + sqlx::query("INSERT INTO alert_configs (webhook_url, threshold) VALUES (?, ?)") + .bind(webhook_url) + .bind(threshold) + .execute(&self.pool) + .await?; Ok(result.last_insert_rowid()) } @@ -468,7 +464,10 @@ mod tests { assert_eq!(deleted, 1); - let remaining = repo.fetch_since(Utc::now() - Duration::days(1)).await.unwrap(); + let remaining = repo + .fetch_since(Utc::now() - Duration::days(1)) + .await + .unwrap(); assert_eq!(remaining.len(), 2); } @@ -510,7 +509,10 @@ mod tests { #[tokio::test] async fn fetch_since_returns_empty_when_no_data() { let repo = make_repo().await; - let fetched = repo.fetch_since(Utc::now() - Duration::hours(24)).await.unwrap(); + let fetched = repo + .fetch_since(Utc::now() - Duration::hours(24)) + .await + .unwrap(); assert!(fetched.is_empty()); } } @@ -546,7 +548,10 @@ mod alert_tests { .insert_alert_config("https://hooks.example.com/a", "Minor") .await .unwrap(); - let updated = repo.update_alert_config(id, "Critical", false).await.unwrap(); + let updated = repo + .update_alert_config(id, "Critical", false) + .await + .unwrap(); assert!(updated); let configs = repo.list_alert_configs().await.unwrap(); assert_eq!(configs[0].threshold, "Critical"); @@ -627,7 +632,9 @@ mod alert_event_tests { async fn log_and_query_five_events() { let repo = make_repo().await; for _ in 0..5 { - repo.log_alert_event(&make_event("Major", true)).await.unwrap(); + repo.log_alert_event(&make_event("Major", true)) + .await + .unwrap(); } let events = repo.query_alert_history(20, None, None).await.unwrap(); assert_eq!(events.len(), 5); @@ -636,15 +643,27 @@ mod alert_event_tests { #[tokio::test] async fn filter_by_severity() { let repo = make_repo().await; - repo.log_alert_event(&make_event("Minor", true)).await.unwrap(); - repo.log_alert_event(&make_event("Major", true)).await.unwrap(); - repo.log_alert_event(&make_event("Critical", false)).await.unwrap(); + repo.log_alert_event(&make_event("Minor", true)) + .await + .unwrap(); + repo.log_alert_event(&make_event("Major", true)) + .await + .unwrap(); + repo.log_alert_event(&make_event("Critical", false)) + .await + .unwrap(); - let major = repo.query_alert_history(20, Some("Major"), None).await.unwrap(); + let major = repo + .query_alert_history(20, Some("Major"), None) + .await + .unwrap(); assert_eq!(major.len(), 1); assert_eq!(major[0].severity, "Major"); - let critical = repo.query_alert_history(20, Some("Critical"), None).await.unwrap(); + let critical = repo + .query_alert_history(20, Some("Critical"), None) + .await + .unwrap(); assert_eq!(critical.len(), 1); assert_eq!(critical[0].severity, "Critical"); } @@ -652,14 +671,26 @@ mod alert_event_tests { #[tokio::test] async fn filter_by_delivered() { let repo = make_repo().await; - repo.log_alert_event(&make_event("Major", true)).await.unwrap(); - repo.log_alert_event(&make_event("Major", false)).await.unwrap(); - repo.log_alert_event(&make_event("Major", true)).await.unwrap(); + repo.log_alert_event(&make_event("Major", true)) + .await + .unwrap(); + repo.log_alert_event(&make_event("Major", false)) + .await + .unwrap(); + repo.log_alert_event(&make_event("Major", true)) + .await + .unwrap(); - let delivered = repo.query_alert_history(20, None, Some(true)).await.unwrap(); + let delivered = repo + .query_alert_history(20, None, Some(true)) + .await + .unwrap(); assert_eq!(delivered.len(), 2); - let failed = repo.query_alert_history(20, None, Some(false)).await.unwrap(); + let failed = repo + .query_alert_history(20, None, Some(false)) + .await + .unwrap(); assert_eq!(failed.len(), 1); } @@ -667,7 +698,9 @@ mod alert_event_tests { async fn limit_clamped_to_100() { let repo = make_repo().await; for _ in 0..5 { - repo.log_alert_event(&make_event("Major", true)).await.unwrap(); + repo.log_alert_event(&make_event("Major", true)) + .await + .unwrap(); } // Requesting 999 should be clamped to 100; still only 5 rows in DB let events = repo.query_alert_history(999, None, None).await.unwrap(); @@ -678,7 +711,9 @@ mod alert_event_tests { async fn count_alert_events_total() { let repo = make_repo().await; for _ in 0..5 { - repo.log_alert_event(&make_event("Major", true)).await.unwrap(); + repo.log_alert_event(&make_event("Major", true)) + .await + .unwrap(); } let total = repo.count_alert_events(None, None).await.unwrap(); assert_eq!(total, 5); @@ -687,9 +722,15 @@ mod alert_event_tests { #[tokio::test] async fn count_alert_events_filtered() { let repo = make_repo().await; - repo.log_alert_event(&make_event("Minor", true)).await.unwrap(); - repo.log_alert_event(&make_event("Major", true)).await.unwrap(); - repo.log_alert_event(&make_event("Critical", false)).await.unwrap(); + repo.log_alert_event(&make_event("Minor", true)) + .await + .unwrap(); + repo.log_alert_event(&make_event("Major", true)) + .await + .unwrap(); + repo.log_alert_event(&make_event("Critical", false)) + .await + .unwrap(); let major_count = repo.count_alert_events(Some("Major"), None).await.unwrap(); assert_eq!(major_count, 1); @@ -697,16 +738,21 @@ mod alert_event_tests { let delivered_count = repo.count_alert_events(None, Some(true)).await.unwrap(); assert_eq!(delivered_count, 2); - let critical_failed = repo.count_alert_events(Some("Critical"), Some(false)).await.unwrap(); + let critical_failed = repo + .count_alert_events(Some("Critical"), Some(false)) + .await + .unwrap(); assert_eq!(critical_failed, 1); } #[tokio::test] async fn logged_event_has_assigned_id() { let repo = make_repo().await; - repo.log_alert_event(&make_event("Major", true)).await.unwrap(); + repo.log_alert_event(&make_event("Major", true)) + .await + .unwrap(); let events = repo.query_alert_history(1, None, None).await.unwrap(); assert!(events[0].id.is_some()); assert!(events[0].id.unwrap() > 0); } -} \ No newline at end of file +} diff --git a/packages/core/src/scheduler.rs b/packages/core/src/scheduler.rs index 40321c3..2066899 100644 --- a/packages/core/src/scheduler.rs +++ b/packages/core/src/scheduler.rs @@ -17,14 +17,12 @@ use tokio::sync::RwLock; use tokio::time; use crate::alerts::AlertManager; -use crate::insights::{ - FeeDataProvider, FeeInsightsEngine, -}; use crate::insights::error::ProviderError; use crate::insights::types::FeeDataPoint; +use crate::insights::{FeeDataProvider, FeeInsightsEngine}; +use crate::metrics::AppMetrics; use crate::repository::FeeRepository; use crate::store::FeeHistoryStore; -use crate::metrics::AppMetrics; /// Run the fee polling loop until Ctrl+C is received. /// Uses defaults for retry and retention — prefer `run_fee_polling_with_retry` in production. @@ -164,10 +162,10 @@ async fn poll_once( update.insights.rolling_averages.short_term.value, ); if let Some(m) = metrics { - m.current_avg_fee.set(update.insights.rolling_averages.short_term.value); - m.spikes_detected_total.inc_by( - update.insights.congestion_trends.recent_spikes.len() as f64 - ); + m.current_avg_fee + .set(update.insights.rolling_averages.short_term.value); + m.spikes_detected_total + .inc_by(update.insights.congestion_trends.recent_spikes.len() as f64); } if let Some(manager) = alert_manager { manager.check_and_dispatch(&update).await; @@ -253,9 +251,9 @@ mod tests { use super::*; use chrono::Utc; - use crate::insights::{FeeInsightsEngine, InsightsConfig}; use crate::insights::error::ProviderError; use crate::insights::types::FeeDataPoint; + use crate::insights::{FeeInsightsEngine, InsightsConfig}; use crate::services::mock_horizon::MockHorizonClient; use crate::store::{FeeHistoryStore, DEFAULT_CAPACITY}; @@ -273,7 +271,9 @@ mod tests { } fn make_shared_engine() -> Arc> { - Arc::new(RwLock::new(FeeInsightsEngine::new(InsightsConfig::default()))) + Arc::new(RwLock::new(FeeInsightsEngine::new( + InsightsConfig::default(), + ))) } // ---- poll_once tests ---- @@ -306,9 +306,8 @@ mod tests { #[tokio::test] async fn poll_once_on_provider_error_does_not_push_to_store() { - let provider: Arc = Arc::new( - MockHorizonClient::new().with_error(ProviderError::ServiceUnavailable), - ); + let provider: Arc = + Arc::new(MockHorizonClient::new().with_error(ProviderError::ServiceUnavailable)); let store = make_shared_store(); let engine = make_shared_engine(); @@ -333,8 +332,7 @@ mod tests { #[tokio::test] async fn poll_once_with_empty_provider_response_leaves_store_unchanged() { - let provider: Arc = - Arc::new(MockHorizonClient::new()); + let provider: Arc = Arc::new(MockHorizonClient::new()); let store = make_shared_store(); let engine = make_shared_engine(); diff --git a/packages/core/src/services/horizon.rs b/packages/core/src/services/horizon.rs index 8ad2de1..18100ee 100644 --- a/packages/core/src/services/horizon.rs +++ b/packages/core/src/services/horizon.rs @@ -1,9 +1,8 @@ -use serde::Deserialize; use reqwest::Client; +use serde::Deserialize; use crate::error::AppError; - #[derive(Clone)] pub struct HorizonClient { base_url: String, @@ -16,10 +15,7 @@ impl HorizonClient { .no_proxy() .build() .unwrap_or_else(|_| Client::new()); - Self { - base_url, - http, - } + Self { base_url, http } } pub fn base_url(&self) -> &str { @@ -27,7 +23,6 @@ impl HorizonClient { } } - #[derive(Debug, Deserialize)] pub struct HorizonTransaction { pub hash: String, @@ -102,7 +97,6 @@ impl HorizonClient { } } - /// Wrapper structs for deserialising Horizon's `_embedded.records` envelope. #[derive(Debug, Deserialize)] struct HorizonTransactionResponse { @@ -307,4 +301,4 @@ mod tests { assert_eq!(stats.fee_charged.p50, "150"); assert_eq!(stats.fee_charged.p95, "800"); } -} \ No newline at end of file +} diff --git a/packages/core/src/services/mock_horizon.rs b/packages/core/src/services/mock_horizon.rs index b01a223..b8d0305 100644 --- a/packages/core/src/services/mock_horizon.rs +++ b/packages/core/src/services/mock_horizon.rs @@ -163,7 +163,10 @@ mod tests { let result = mock.fetch_latest_fees().await; assert!(result.is_err()); - assert!(matches!(result.unwrap_err(), ProviderError::NetworkError { .. })); + assert!(matches!( + result.unwrap_err(), + ProviderError::NetworkError { .. } + )); } #[tokio::test] @@ -196,7 +199,10 @@ mod tests { async fn health_check_fails_when_unhealthy() { let mock = MockHorizonClient::new().with_healthy(false); let result = mock.health_check().await; - assert!(matches!(result.unwrap_err(), ProviderError::ServiceUnavailable)); + assert!(matches!( + result.unwrap_err(), + ProviderError::ServiceUnavailable + )); } #[tokio::test] @@ -217,4 +223,4 @@ mod tests { let mock = MockHorizonClient::new(); assert_eq!(mock.provider_name(), "MockHorizon"); } -} \ No newline at end of file +} diff --git a/packages/core/src/services/mod.rs b/packages/core/src/services/mod.rs index 6300a2a..86e0859 100644 --- a/packages/core/src/services/mod.rs +++ b/packages/core/src/services/mod.rs @@ -1,4 +1,4 @@ pub mod horizon; #[cfg(test)] -pub mod mock_horizon; \ No newline at end of file +pub mod mock_horizon; diff --git a/packages/core/src/store.rs b/packages/core/src/store.rs index c5e51f8..f8891d4 100644 --- a/packages/core/src/store.rs +++ b/packages/core/src/store.rs @@ -147,7 +147,7 @@ mod tests { let mut store = FeeHistoryStore::new(10); store.push(make_point(100, 60)); // 60 min ago store.push(make_point(200, 30)); // 30 min ago - store.push(make_point(300, 5)); // 5 min ago + store.push(make_point(300, 5)); // 5 min ago let cutoff = Utc::now() - Duration::minutes(31); let result = store.get_since(cutoff); @@ -214,4 +214,4 @@ mod tests { let store = FeeHistoryStore::new(10); assert!(store.get_last_n(5).is_empty()); } -} \ No newline at end of file +} diff --git a/packages/core/tests/api_integration.rs b/packages/core/tests/api_integration.rs index c87a5fe..2ee0e3f 100644 --- a/packages/core/tests/api_integration.rs +++ b/packages/core/tests/api_integration.rs @@ -35,8 +35,8 @@ use stellar_fee_tracker::{ api, cache::ResponseCache, db, - insights::{FeeInsightsEngine, InsightsConfig}, insights::types::FeeDataPoint, + insights::{FeeInsightsEngine, InsightsConfig}, metrics::AppMetrics, repository::FeeRepository, services::horizon::HorizonClient, @@ -117,10 +117,7 @@ async fn build_test_app() -> (Router, MockServer) { let mock_server = MockServer::start().await; Mock::given(method("GET")) .and(path("/fee_stats")) - .respond_with( - ResponseTemplate::new(200) - .set_body_raw(FAKE_FEE_STATS, "application/json"), - ) + .respond_with(ResponseTemplate::new(200).set_body_raw(FAKE_FEE_STATS, "application/json")) .mount(&mock_server) .await; @@ -145,9 +142,9 @@ async fn build_test_app() -> (Router, MockServer) { } } - let insights_engine = Arc::new(RwLock::new( - FeeInsightsEngine::new(InsightsConfig::default()), - )); + let insights_engine = Arc::new(RwLock::new(FeeInsightsEngine::new( + InsightsConfig::default(), + ))); { let mut engine = insights_engine.write().await; engine.process_fee_data(&points).await.unwrap(); @@ -195,14 +192,31 @@ async fn build_test_app() -> (Router, MockServer) { }), ) .merge(fees_router) - .merge(api::insights::create_insights_router(insights_engine.clone())) + .merge(api::insights::create_insights_router( + insights_engine.clone(), + )) .merge( Router::new() - .route("/alerts/config", axum::routing::post(api::alerts::create_alert)) - .route("/alerts/config", axum::routing::get(api::alerts::list_alerts)) - .route("/alerts/config/:id", axum::routing::patch(api::alerts::update_alert)) - .route("/alerts/config/:id", axum::routing::delete(api::alerts::delete_alert)) - .route("/alerts/history", axum::routing::get(api::alerts::get_alert_history)) + .route( + "/alerts/config", + axum::routing::post(api::alerts::create_alert), + ) + .route( + "/alerts/config", + axum::routing::get(api::alerts::list_alerts), + ) + .route( + "/alerts/config/:id", + axum::routing::patch(api::alerts::update_alert), + ) + .route( + "/alerts/config/:id", + axum::routing::delete(api::alerts::delete_alert), + ) + .route( + "/alerts/history", + axum::routing::get(api::alerts::get_alert_history), + ) .with_state(repository), ); @@ -221,7 +235,12 @@ async fn json_body(body: Body) -> Value { async fn health_returns_200_with_ok_body() { let (app, _mock) = build_test_app().await; let resp = app - .oneshot(Request::builder().uri("/health").body(Body::empty()).unwrap()) + .oneshot( + Request::builder() + .uri("/health") + .body(Body::empty()) + .unwrap(), + ) .await .unwrap(); @@ -254,9 +273,18 @@ async fn fees_current_returns_200_with_required_fields() { assert_eq!(resp.status(), StatusCode::OK); let json = json_body(resp.into_body()).await; assert!(json["base_fee"].is_string(), "missing base_fee"); - assert!(json["percentiles"]["p50"].is_string(), "missing percentiles.p50"); - assert!(json["percentiles"]["p10"].is_string(), "missing percentiles.p10"); - assert!(json["percentiles"]["p95"].is_string(), "missing percentiles.p95"); + assert!( + json["percentiles"]["p50"].is_string(), + "missing percentiles.p50" + ); + assert!( + json["percentiles"]["p10"].is_string(), + "missing percentiles.p10" + ); + assert!( + json["percentiles"]["p95"].is_string(), + "missing percentiles.p95" + ); assert!(json["min_fee"].is_string(), "missing min_fee"); assert!(json["max_fee"].is_string(), "missing max_fee"); assert!(json["avg_fee"].is_string(), "missing avg_fee"); @@ -473,7 +501,10 @@ async fn fees_trend_returns_200_with_status_and_changes() { assert!(json["status"].is_string(), "missing status"); assert!(json["changes"].is_object(), "missing changes"); assert!(json["trend_strength"].is_string(), "missing trend_strength"); - assert!(json["recent_spike_count"].is_number(), "missing recent_spike_count"); + assert!( + json["recent_spike_count"].is_number(), + "missing recent_spike_count" + ); assert!(json["last_updated"].is_string(), "missing last_updated"); } @@ -516,7 +547,10 @@ async fn insights_returns_200_with_rolling_averages() { assert_eq!(resp.status(), StatusCode::OK); let json = json_body(resp.into_body()).await; - assert!(json["rolling_averages"].is_object(), "missing rolling_averages"); + assert!( + json["rolling_averages"].is_object(), + "missing rolling_averages" + ); } #[tokio::test] @@ -641,7 +675,10 @@ async fn insights_extremes_returns_200() { assert_eq!(resp.status(), StatusCode::OK); // Response is a JSON object let json = json_body(resp.into_body()).await; - assert!(json.is_object(), "expected JSON object from /insights/extremes"); + assert!( + json.is_object(), + "expected JSON object from /insights/extremes" + ); } // ---- GET /insights/congestion ----------------------------------------------- @@ -661,7 +698,10 @@ async fn insights_congestion_returns_200() { assert_eq!(resp.status(), StatusCode::OK); let json = json_body(resp.into_body()).await; - assert!(json.is_object(), "expected JSON object from /insights/congestion"); + assert!( + json.is_object(), + "expected JSON object from /insights/congestion" + ); } // ---- GET /insights/health --------------------------------------------------- @@ -782,4 +822,4 @@ async fn alerts_history_returns_200_empty() { let json = json_body(resp.into_body()).await; assert_eq!(json["total"], 0); assert!(json["items"].as_array().unwrap().is_empty()); -} \ No newline at end of file +} diff --git a/packages/devkit/Cargo.toml b/packages/devkit/Cargo.toml new file mode 100644 index 0000000..4731dfc --- /dev/null +++ b/packages/devkit/Cargo.toml @@ -0,0 +1,7 @@ +[package] +name = "stellar-devkit" +version = "0.1.0" +edition = "2021" +description = "Developer toolkit for testing and simulating the Stellar fee tracker" + +[dependencies] diff --git a/packages/devkit/README.md b/packages/devkit/README.md new file mode 100644 index 0000000..85a418c --- /dev/null +++ b/packages/devkit/README.md @@ -0,0 +1,17 @@ +# stellar-devkit + +Developer toolkit for the Stellar Fee Tracker. Provides utilities for testing, mocking, and simulating Stellar network behaviour without hitting live infrastructure. + +## Modules + +- **harness** — Test harness with a Horizon mock server and pre-built scenario runners. +- **simulation** — Fee models, network-load generators, and congestion predictors for local simulation. + +## Usage + +Add `stellar-devkit` to your `[dev-dependencies]` and import the modules you need. + +```toml +[dev-dependencies] +stellar-devkit = { path = "../devkit" } +``` diff --git a/packages/devkit/src/harness/horizon_mock.rs b/packages/devkit/src/harness/horizon_mock.rs new file mode 100644 index 0000000..92d02ac --- /dev/null +++ b/packages/devkit/src/harness/horizon_mock.rs @@ -0,0 +1,2 @@ +/// Mock implementation of the Horizon API server for use in tests. +pub struct HorizonMock; diff --git a/packages/devkit/src/harness/mod.rs b/packages/devkit/src/harness/mod.rs new file mode 100644 index 0000000..cac6d90 --- /dev/null +++ b/packages/devkit/src/harness/mod.rs @@ -0,0 +1,2 @@ +pub mod horizon_mock; +pub mod scenarios; diff --git a/packages/devkit/src/harness/scenarios/mod.rs b/packages/devkit/src/harness/scenarios/mod.rs new file mode 100644 index 0000000..a77f2b4 --- /dev/null +++ b/packages/devkit/src/harness/scenarios/mod.rs @@ -0,0 +1 @@ +//! Pre-built test scenarios for the Stellar fee tracker harness. diff --git a/packages/devkit/src/lib.rs b/packages/devkit/src/lib.rs new file mode 100644 index 0000000..c3a6ac2 --- /dev/null +++ b/packages/devkit/src/lib.rs @@ -0,0 +1,2 @@ +pub mod harness; +pub mod simulation; diff --git a/packages/devkit/src/simulation/congestion_predictor.rs b/packages/devkit/src/simulation/congestion_predictor.rs new file mode 100644 index 0000000..7ec0d0a --- /dev/null +++ b/packages/devkit/src/simulation/congestion_predictor.rs @@ -0,0 +1,2 @@ +/// Predicts network congestion based on simulated load and fee models. +pub struct CongestionPredictor; diff --git a/packages/devkit/src/simulation/fee_model.rs b/packages/devkit/src/simulation/fee_model.rs new file mode 100644 index 0000000..cb64205 --- /dev/null +++ b/packages/devkit/src/simulation/fee_model.rs @@ -0,0 +1,2 @@ +/// Models for simulating Stellar transaction fee behaviour. +pub struct FeeModel; diff --git a/packages/devkit/src/simulation/mod.rs b/packages/devkit/src/simulation/mod.rs new file mode 100644 index 0000000..342fca7 --- /dev/null +++ b/packages/devkit/src/simulation/mod.rs @@ -0,0 +1,3 @@ +pub mod congestion_predictor; +pub mod fee_model; +pub mod network_load; diff --git a/packages/devkit/src/simulation/network_load.rs b/packages/devkit/src/simulation/network_load.rs new file mode 100644 index 0000000..13c02d6 --- /dev/null +++ b/packages/devkit/src/simulation/network_load.rs @@ -0,0 +1,2 @@ +/// Generates synthetic network load profiles for simulation. +pub struct NetworkLoad;