Skip to main content

opendeviationbar_streaming/
processor.rs

1use futures::Stream;
2/// Production-ready streaming architecture with bounded memory and backpressure
3/// # FILE-SIZE-OK
4///
5/// This module implements true infinite streaming capabilities addressing critical failures:
6/// - Eliminates Vec<OpenDeviationBar> accumulation (unbounded memory growth)
7/// - Implements proper backpressure with bounded channels
8/// - Provides circuit breaker resilience patterns
9/// - Maintains temporal integrity for financial data
10use opendeviationbar_core::processor::ExportOpenDeviationBarProcessor;
11use opendeviationbar_core::{OpenDeviationBar, Tick};
12use std::pin::Pin;
13use std::sync::Arc;
14use std::sync::atomic::{AtomicU64, Ordering};
15use std::task::{Context, Poll};
16use tokio::sync::mpsc;
17use tokio::time::{Duration, Instant};
18
19/// Configuration for production streaming
20#[derive(Debug, Clone)]
21pub struct StreamingProcessorConfig {
22    /// Channel capacity for trade input
23    pub trade_channel_capacity: usize,
24    /// Channel capacity for completed bars
25    pub bar_channel_capacity: usize,
26    /// Memory usage threshold in bytes
27    pub memory_threshold_bytes: usize,
28    /// Backpressure timeout
29    pub backpressure_timeout: Duration,
30    /// Circuit breaker error rate threshold (0.0-1.0)
31    pub circuit_breaker_threshold: f64,
32    /// Circuit breaker timeout before retry
33    pub circuit_breaker_timeout: Duration,
34}
35
36impl StreamingProcessorConfig {
37    /// Get bar channel capacity from environment or use default (10K)
38    /// Issue #96 Task #6: OPENDEVIATIONBAR_MAX_PENDING_BARS env var support
39    fn get_bar_channel_capacity() -> usize {
40        std::env::var("OPENDEVIATIONBAR_MAX_PENDING_BARS")
41            .ok()
42            .and_then(|v| v.parse::<usize>().ok())
43            .unwrap_or(10_000)
44    }
45}
46
47impl Default for StreamingProcessorConfig {
48    fn default() -> Self {
49        Self {
50            trade_channel_capacity: 5_000, // Based on consensus analysis
51            bar_channel_capacity: StreamingProcessorConfig::get_bar_channel_capacity(), // Issue #96: 10K backpressure bound
52            memory_threshold_bytes: 100_000_000, // 100MB limit
53            backpressure_timeout: Duration::from_millis(100),
54            circuit_breaker_threshold: 0.5, // 50% error rate
55            circuit_breaker_timeout: Duration::from_secs(30),
56        }
57    }
58}
59
60/// Production streaming processor with bounded memory
61pub struct StreamingProcessor {
62    /// Open deviation bar processor (single instance, no accumulation)
63    processor: ExportOpenDeviationBarProcessor,
64
65    /// Threshold in decimal basis points for recreating processor
66    _threshold_decimal_bps: u32,
67
68    /// Bounded channel for incoming trades
69    trade_sender: Option<mpsc::Sender<Tick>>,
70    trade_receiver: mpsc::Receiver<Tick>,
71
72    /// Bounded channel for outgoing bars
73    bar_sender: mpsc::Sender<OpenDeviationBar>,
74    bar_receiver: Option<mpsc::Receiver<OpenDeviationBar>>,
75
76    /// Configuration
77    config: StreamingProcessorConfig,
78
79    /// Metrics
80    metrics: Arc<StreamingMetrics>,
81
82    /// Circuit breaker state
83    circuit_breaker: CircuitBreaker,
84}
85
86/// Circuit breaker implementation
87#[derive(Debug)]
88struct CircuitBreaker {
89    state: CircuitBreakerState,
90    failure_count: u64,
91    success_count: u64,
92    last_failure_time: Option<Instant>,
93    threshold: f64,
94    timeout: Duration,
95}
96
97#[derive(Debug, PartialEq)]
98enum CircuitBreakerState {
99    Closed,
100    Open,
101    HalfOpen,
102}
103
104/// Streaming metrics for observability
105/// Issue #96 Task #6: Extended with queue depth and block time tracking
106#[derive(Debug, Default)]
107pub struct StreamingMetrics {
108    pub trades_processed: AtomicU64,
109    pub bars_generated: AtomicU64,
110    pub errors_total: AtomicU64,
111    pub backpressure_events: AtomicU64,
112    pub circuit_breaker_trips: AtomicU64,
113    pub memory_usage_bytes: AtomicU64,
114    pub max_queue_depth: AtomicU64, // Issue #96 Task #6: Max observed queue depth
115    pub total_block_time_ms: AtomicU64, // Issue #96 Task #6: Accumulated block time
116}
117
118impl StreamingProcessor {
119    /// Create new production streaming processor
120    pub fn new(
121        threshold_decimal_bps: u32,
122    ) -> Result<Self, opendeviationbar_core::processor::ProcessingError> {
123        Self::with_config(threshold_decimal_bps, StreamingProcessorConfig::default())
124    }
125
126    /// Create with custom configuration
127    pub fn with_config(
128        threshold_decimal_bps: u32,
129        config: StreamingProcessorConfig,
130    ) -> Result<Self, opendeviationbar_core::processor::ProcessingError> {
131        let (trade_sender, trade_receiver) = mpsc::channel(config.trade_channel_capacity);
132        let (bar_sender, bar_receiver) = mpsc::channel(config.bar_channel_capacity);
133
134        let circuit_breaker_threshold = config.circuit_breaker_threshold;
135        let circuit_breaker_timeout = config.circuit_breaker_timeout;
136
137        Ok(Self {
138            processor: ExportOpenDeviationBarProcessor::new(threshold_decimal_bps)?,
139            _threshold_decimal_bps: threshold_decimal_bps,
140            trade_sender: Some(trade_sender),
141            trade_receiver,
142            bar_sender,
143            bar_receiver: Some(bar_receiver),
144            config,
145            metrics: Arc::new(StreamingMetrics::default()),
146            circuit_breaker: CircuitBreaker::new(
147                circuit_breaker_threshold,
148                circuit_breaker_timeout,
149            ),
150        })
151    }
152
153    /// Get trade sender for external components
154    pub fn trade_sender(&mut self) -> Option<mpsc::Sender<Tick>> {
155        self.trade_sender.take()
156    }
157
158    /// Get bar receiver for external components
159    pub fn bar_receiver(&mut self) -> Option<mpsc::Receiver<OpenDeviationBar>> {
160        self.bar_receiver.take()
161    }
162
163    /// Start processing loop (bounded memory, infinite capability)
164    pub async fn start_processing(&mut self) -> Result<(), StreamingError> {
165        loop {
166            // Check circuit breaker state
167            if !self.circuit_breaker.can_process() {
168                tokio::time::sleep(Duration::from_millis(100)).await;
169                continue;
170            }
171
172            // Receive trade with timeout (prevents blocking forever)
173            let trade = match tokio::time::timeout(
174                self.config.backpressure_timeout,
175                self.trade_receiver.recv(),
176            )
177            .await
178            {
179                Ok(Some(trade)) => trade,
180                Ok(None) => {
181                    // Channel closed - send final incomplete bar if exists
182                    if let Some(final_bar) = self.processor.get_incomplete_bar()
183                        && let Err(e) = self.send_bar_with_backpressure(final_bar).await
184                    {
185                        println!("Failed to send final incomplete bar: {:?}", e);
186                    }
187                    break;
188                }
189                Err(_) => continue, // Timeout, check circuit breaker again
190            };
191
192            // Process single trade (use borrowed reference per Issue #96 Task #78)
193            match self.process_single_trade(&trade).await {
194                Ok(bar_opt) => {
195                    self.circuit_breaker.record_success();
196
197                    // If bar completed, send with backpressure handling
198                    if let Some(bar) = bar_opt
199                        && let Err(e) = self.send_bar_with_backpressure(bar).await
200                    {
201                        println!("Failed to send bar: {:?}", e);
202                        self.circuit_breaker.record_failure();
203                    }
204                }
205                Err(e) => {
206                    println!("Trade processing error: {:?}", e);
207                    self.circuit_breaker.record_failure();
208                    self.metrics.errors_total.fetch_add(1, Ordering::Relaxed);
209                }
210            }
211        }
212
213        Ok(())
214    }
215
216    /// Process single trade - extracts completed bars without accumulation
217    // Issue #96 Task #78: Accept borrowed Tick reference
218    async fn process_single_trade(
219        &mut self,
220        trade: &Tick,
221    ) -> Result<Option<OpenDeviationBar>, StreamingError> {
222        // Update metrics
223        self.metrics
224            .trades_processed
225            .fetch_add(1, Ordering::Relaxed);
226
227        // Process trade using existing algorithm (single trade at a time)
228        self.processor
229            .process_trades_continuously(std::slice::from_ref(trade));
230
231        // Extract completed bars immediately (prevents accumulation)
232        let mut completed_bars = self.processor.get_all_completed_bars();
233
234        if !completed_bars.is_empty() {
235            // Bounded memory: only return first completed bar
236            // Additional bars would be rare edge cases but must be handled
237            let completed_bar = completed_bars.remove(0);
238
239            // Handle rare case of multiple completions
240            if !completed_bars.is_empty() {
241                println!(
242                    "Warning: {} additional bars completed, dropping for bounded memory",
243                    completed_bars.len()
244                );
245                self.metrics
246                    .backpressure_events
247                    .fetch_add(completed_bars.len() as u64, Ordering::Relaxed);
248            }
249
250            self.metrics.bars_generated.fetch_add(1, Ordering::Relaxed);
251            Ok(Some(completed_bar))
252        } else {
253            Ok(None)
254        }
255    }
256
257    /// Send bar with backpressure handling
258    async fn send_bar_with_backpressure(
259        &self,
260        bar: OpenDeviationBar,
261    ) -> Result<(), StreamingError> {
262        // Use try_send for immediate check, then send for blocking
263        match self.bar_sender.try_send(bar.clone()) {
264            Ok(()) => Ok(()),
265            Err(mpsc::error::TrySendError::Full(_)) => {
266                // Apply backpressure - channel is full
267                println!("Bar channel full, applying backpressure");
268                self.metrics
269                    .backpressure_events
270                    .fetch_add(1, Ordering::Relaxed);
271
272                // Wait for capacity with blocking send
273                self.bar_sender
274                    .send(bar)
275                    .await
276                    .map_err(|_| StreamingError::ChannelClosed)
277            }
278            Err(mpsc::error::TrySendError::Closed(_)) => Err(StreamingError::ChannelClosed),
279        }
280    }
281
282    /// Get current metrics
283    pub fn metrics(&self) -> &StreamingMetrics {
284        &self.metrics
285    }
286
287    /// Extract final incomplete bar when stream ends (for algorithmic consistency)
288    pub fn get_final_incomplete_bar(&mut self) -> Option<OpenDeviationBar> {
289        self.processor.get_incomplete_bar()
290    }
291
292    /// Check memory usage against threshold
293    pub fn check_memory_usage(&self) -> bool {
294        let current_usage = self.metrics.memory_usage_bytes.load(Ordering::Relaxed);
295        current_usage < self.config.memory_threshold_bytes as u64
296    }
297}
298
299impl CircuitBreaker {
300    fn new(threshold: f64, timeout: Duration) -> Self {
301        Self {
302            state: CircuitBreakerState::Closed,
303            failure_count: 0,
304            success_count: 0,
305            last_failure_time: None,
306            threshold,
307            timeout,
308        }
309    }
310
311    fn can_process(&mut self) -> bool {
312        match self.state {
313            CircuitBreakerState::Closed => true,
314            CircuitBreakerState::Open => {
315                if let Some(last_failure) = self.last_failure_time {
316                    if last_failure.elapsed() >= self.timeout {
317                        self.state = CircuitBreakerState::HalfOpen;
318                        true
319                    } else {
320                        false
321                    }
322                } else {
323                    true
324                }
325            }
326            CircuitBreakerState::HalfOpen => true,
327        }
328    }
329
330    fn record_success(&mut self) {
331        self.success_count += 1;
332
333        if self.state == CircuitBreakerState::HalfOpen {
334            // Successful request in half-open, close circuit
335            self.state = CircuitBreakerState::Closed;
336            self.failure_count = 0;
337        }
338    }
339
340    fn record_failure(&mut self) {
341        self.failure_count += 1;
342        self.last_failure_time = Some(Instant::now());
343
344        let total_requests = self.failure_count + self.success_count;
345        if total_requests >= 10 {
346            // Minimum sample size
347            let failure_rate = self.failure_count as f64 / total_requests as f64;
348
349            if failure_rate >= self.threshold {
350                self.state = CircuitBreakerState::Open;
351            }
352        }
353    }
354}
355
356/// Stream implementation for open deviation bars (true streaming)
357pub struct OpenDeviationBarStream {
358    receiver: mpsc::Receiver<OpenDeviationBar>,
359}
360
361impl OpenDeviationBarStream {
362    pub fn new(receiver: mpsc::Receiver<OpenDeviationBar>) -> Self {
363        Self { receiver }
364    }
365}
366
367impl Stream for OpenDeviationBarStream {
368    type Item = Result<OpenDeviationBar, StreamingError>;
369
370    fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
371        match self.receiver.poll_recv(cx) {
372            Poll::Ready(Some(bar)) => Poll::Ready(Some(Ok(bar))),
373            Poll::Ready(None) => Poll::Ready(None),
374            Poll::Pending => Poll::Pending,
375        }
376    }
377}
378
379/// Streaming errors
380#[derive(Debug, thiserror::Error)]
381pub enum StreamingError {
382    #[error("Channel closed")]
383    ChannelClosed,
384
385    #[error("Backpressure timeout")]
386    BackpressureTimeout,
387
388    #[error("Circuit breaker open")]
389    CircuitBreakerOpen,
390
391    #[error("Memory threshold exceeded")]
392    MemoryThresholdExceeded,
393
394    #[error("Processing error: {0}")]
395    ProcessingError(String),
396}
397
398impl StreamingMetrics {
399    /// Get metrics summary
400    pub fn summary(&self) -> MetricsSummary {
401        MetricsSummary {
402            trades_processed: self.trades_processed.load(Ordering::Relaxed),
403            bars_generated: self.bars_generated.load(Ordering::Relaxed),
404            errors_total: self.errors_total.load(Ordering::Relaxed),
405            backpressure_events: self.backpressure_events.load(Ordering::Relaxed),
406            circuit_breaker_trips: self.circuit_breaker_trips.load(Ordering::Relaxed),
407            memory_usage_bytes: self.memory_usage_bytes.load(Ordering::Relaxed),
408        }
409    }
410}
411
412/// Metrics snapshot
413#[derive(Debug, Clone)]
414pub struct MetricsSummary {
415    pub trades_processed: u64,
416    pub bars_generated: u64,
417    pub errors_total: u64,
418    pub backpressure_events: u64,
419    pub circuit_breaker_trips: u64,
420    pub memory_usage_bytes: u64,
421}
422
423impl MetricsSummary {
424    /// Calculate bars per aggTrade ratio
425    pub fn bars_per_aggtrade(&self) -> f64 {
426        if self.trades_processed > 0 {
427            self.bars_generated as f64 / self.trades_processed as f64
428        } else {
429            0.0
430        }
431    }
432
433    /// Calculate error rate
434    pub fn error_rate(&self) -> f64 {
435        if self.trades_processed > 0 {
436            self.errors_total as f64 / self.trades_processed as f64
437        } else {
438            0.0
439        }
440    }
441
442    /// Format memory usage
443    pub fn memory_usage_mb(&self) -> f64 {
444        self.memory_usage_bytes as f64 / 1_000_000.0
445    }
446}
447
448#[cfg(test)]
449mod tests {
450    use super::*;
451    use opendeviationbar_core::FixedPoint;
452
453    fn create_test_trade(id: u64, price: f64, timestamp: u64) -> Tick {
454        let price_str = format!("{:.8}", price);
455        Tick {
456            ref_id: id as i64,
457            price: FixedPoint::from_str(&price_str).unwrap(),
458            volume: FixedPoint::from_str("1.0").unwrap(),
459            first_sub_id: id as i64,
460            last_sub_id: id as i64,
461            timestamp: timestamp as i64,
462            is_buyer_maker: false,
463            is_best_match: None,
464            best_bid: None,
465            best_ask: None,
466        }
467    }
468
469    #[tokio::test]
470    async fn test_bounded_memory_streaming() {
471        let mut processor = StreamingProcessor::new(25).unwrap(); // 0.25% threshold
472
473        // Test that memory remains bounded
474        let initial_metrics = processor.metrics().summary();
475
476        // Send 1000 trades
477        for i in 0..1000 {
478            let trade = create_test_trade(i, 23000.0 + (i as f64), 1659312000000 + i);
479            if let Ok(bar_opt) = processor.process_single_trade(&trade).await {
480                // Verify no accumulation - at most one bar per aggTrade
481                assert!(bar_opt.is_none() || bar_opt.is_some());
482            }
483        }
484
485        let final_metrics = processor.metrics().summary();
486        assert!(final_metrics.trades_processed >= initial_metrics.trades_processed);
487        assert!(final_metrics.trades_processed <= 1000);
488    }
489
490    #[tokio::test]
491    async fn test_circuit_breaker() {
492        let mut circuit_breaker = CircuitBreaker::new(0.5, Duration::from_millis(100));
493
494        // Initially closed
495        assert!(circuit_breaker.can_process());
496
497        // Record failures
498        for _ in 0..20 {
499            circuit_breaker.record_failure();
500        }
501
502        // Should open after 50% failure rate
503        assert_eq!(circuit_breaker.state, CircuitBreakerState::Open);
504        assert!(!circuit_breaker.can_process());
505
506        // Wait for timeout
507        tokio::time::sleep(Duration::from_millis(150)).await;
508
509        // Should transition to half-open
510        assert!(circuit_breaker.can_process());
511
512        // Record success
513        circuit_breaker.record_success();
514
515        // Should close
516        assert_eq!(circuit_breaker.state, CircuitBreakerState::Closed);
517    }
518
519    // === Circuit Breaker State Machine Tests ===
520
521    #[test]
522    fn test_circuit_breaker_stays_closed_below_threshold() {
523        let mut cb = CircuitBreaker::new(0.5, Duration::from_secs(10));
524
525        // 8 successes, 2 failures = 20% failure rate, below 50% threshold
526        for _ in 0..8 {
527            cb.record_success();
528        }
529        for _ in 0..2 {
530            cb.record_failure();
531        }
532
533        // Should remain closed (20% < 50%)
534        assert_eq!(cb.state, CircuitBreakerState::Closed);
535        assert!(cb.can_process());
536    }
537
538    #[test]
539    fn test_circuit_breaker_minimum_sample_size() {
540        let mut cb = CircuitBreaker::new(0.5, Duration::from_secs(10));
541
542        // 9 failures, 0 successes = 100% failure rate, but only 9 requests (< 10 minimum)
543        for _ in 0..9 {
544            cb.record_failure();
545        }
546
547        // Should remain closed (minimum sample size not met)
548        assert_eq!(cb.state, CircuitBreakerState::Closed);
549        assert!(cb.can_process());
550
551        // 10th failure triggers open
552        cb.record_failure();
553        assert_eq!(cb.state, CircuitBreakerState::Open);
554    }
555
556    #[test]
557    fn test_circuit_breaker_halfopen_failure_reopens() {
558        let mut cb = CircuitBreaker::new(0.5, Duration::from_secs(0));
559
560        // Trip the breaker: 10 failures opens it
561        for _ in 0..10 {
562            cb.record_failure();
563        }
564        assert_eq!(cb.state, CircuitBreakerState::Open);
565
566        // Zero-second timeout → immediately transitions to HalfOpen on can_process
567        assert!(cb.can_process());
568        assert_eq!(cb.state, CircuitBreakerState::HalfOpen);
569
570        // Record failure in HalfOpen → should re-open
571        // (failure_count accumulates, total >= 10, rate >= threshold)
572        cb.record_failure();
573        assert_eq!(cb.state, CircuitBreakerState::Open);
574    }
575
576    #[test]
577    fn test_circuit_breaker_closed_resets_failure_count() {
578        let mut cb = CircuitBreaker::new(0.5, Duration::from_secs(0));
579
580        // Trip the breaker
581        for _ in 0..10 {
582            cb.record_failure();
583        }
584        assert_eq!(cb.state, CircuitBreakerState::Open);
585
586        // Transition to HalfOpen
587        assert!(cb.can_process());
588        assert_eq!(cb.state, CircuitBreakerState::HalfOpen);
589
590        // Record success → closes and resets failure_count
591        cb.record_success();
592        assert_eq!(cb.state, CircuitBreakerState::Closed);
593        assert_eq!(cb.failure_count, 0);
594    }
595
596    #[test]
597    fn test_circuit_breaker_open_blocks_until_timeout() {
598        let mut cb = CircuitBreaker::new(0.5, Duration::from_secs(3600)); // 1 hour timeout
599
600        // Trip the breaker
601        for _ in 0..10 {
602            cb.record_failure();
603        }
604
605        // Should be blocked — timeout hasn't elapsed
606        assert!(!cb.can_process());
607        assert_eq!(cb.state, CircuitBreakerState::Open);
608    }
609
610    #[test]
611    fn test_metrics_zero_trades() {
612        let metrics = MetricsSummary {
613            trades_processed: 0,
614            bars_generated: 0,
615            errors_total: 0,
616            backpressure_events: 0,
617            circuit_breaker_trips: 0,
618            memory_usage_bytes: 0,
619        };
620
621        // Division by zero guarded
622        assert_eq!(metrics.bars_per_aggtrade(), 0.0);
623        assert_eq!(metrics.error_rate(), 0.0);
624        assert_eq!(metrics.memory_usage_mb(), 0.0);
625    }
626
627    #[test]
628    fn test_metrics_calculations() {
629        let metrics = MetricsSummary {
630            trades_processed: 1000,
631            bars_generated: 50,
632            errors_total: 5,
633            backpressure_events: 2,
634            circuit_breaker_trips: 1,
635            memory_usage_bytes: 50_000_000,
636        };
637
638        assert_eq!(metrics.bars_per_aggtrade(), 0.05);
639        assert_eq!(metrics.error_rate(), 0.005);
640        assert_eq!(metrics.memory_usage_mb(), 50.0);
641    }
642
643    // === Metrics Snapshot & Take-Once Tests (Issue #96 Task #114) ===
644
645    #[test]
646    fn test_streaming_metrics_summary_snapshot() {
647        let metrics = StreamingMetrics::default();
648        metrics.trades_processed.store(500, Ordering::Relaxed);
649        metrics.bars_generated.store(25, Ordering::Relaxed);
650        metrics.errors_total.store(3, Ordering::Relaxed);
651        metrics.backpressure_events.store(1, Ordering::Relaxed);
652        metrics.circuit_breaker_trips.store(0, Ordering::Relaxed);
653        metrics
654            .memory_usage_bytes
655            .store(42_000_000, Ordering::Relaxed);
656
657        let summary = metrics.summary();
658
659        assert_eq!(summary.trades_processed, 500);
660        assert_eq!(summary.bars_generated, 25);
661        assert_eq!(summary.errors_total, 3);
662        assert_eq!(summary.backpressure_events, 1);
663        assert_eq!(summary.circuit_breaker_trips, 0);
664        assert_eq!(summary.memory_usage_bytes, 42_000_000);
665    }
666
667    #[test]
668    fn test_memory_usage_mb_conversion() {
669        // Exact MB boundary
670        let m1 = MetricsSummary {
671            trades_processed: 0,
672            bars_generated: 0,
673            errors_total: 0,
674            backpressure_events: 0,
675            circuit_breaker_trips: 0,
676            memory_usage_bytes: 1_000_000,
677        };
678        assert_eq!(m1.memory_usage_mb(), 1.0);
679
680        // Fractional MB
681        let m2 = MetricsSummary {
682            trades_processed: 0,
683            bars_generated: 0,
684            errors_total: 0,
685            backpressure_events: 0,
686            circuit_breaker_trips: 0,
687            memory_usage_bytes: 1_500_000,
688        };
689        assert_eq!(m2.memory_usage_mb(), 1.5);
690
691        // Large value (4 GB)
692        let m3 = MetricsSummary {
693            trades_processed: 0,
694            bars_generated: 0,
695            errors_total: 0,
696            backpressure_events: 0,
697            circuit_breaker_trips: 0,
698            memory_usage_bytes: 4_000_000_000,
699        };
700        assert_eq!(m3.memory_usage_mb(), 4000.0);
701    }
702
703    #[test]
704    fn test_trade_sender_take_once() {
705        let mut processor = StreamingProcessor::new(25).unwrap();
706
707        // First call returns Some
708        let sender = processor.trade_sender();
709        assert!(
710            sender.is_some(),
711            "First trade_sender() call must return Some"
712        );
713
714        // Second call returns None (already taken)
715        let sender2 = processor.trade_sender();
716        assert!(
717            sender2.is_none(),
718            "Second trade_sender() call must return None"
719        );
720    }
721
722    #[test]
723    fn test_bar_receiver_take_once() {
724        let mut processor = StreamingProcessor::new(25).unwrap();
725
726        // First call returns Some
727        let receiver = processor.bar_receiver();
728        assert!(
729            receiver.is_some(),
730            "First bar_receiver() call must return Some"
731        );
732
733        // Second call returns None (already taken)
734        let receiver2 = processor.bar_receiver();
735        assert!(
736            receiver2.is_none(),
737            "Second bar_receiver() call must return None"
738        );
739    }
740
741    #[test]
742    fn test_check_memory_usage_below_threshold() {
743        let processor = StreamingProcessor::new(25).unwrap();
744
745        // Default: memory_usage_bytes = 0, threshold = 100MB → within bounds
746        assert!(
747            processor.check_memory_usage(),
748            "Zero memory usage should be within threshold"
749        );
750    }
751
752    #[test]
753    fn test_check_memory_usage_above_threshold() {
754        let processor = StreamingProcessor::new(25).unwrap();
755
756        // Simulate exceeding threshold (100MB default)
757        processor
758            .metrics
759            .memory_usage_bytes
760            .store(200_000_000, Ordering::Relaxed);
761        assert!(
762            !processor.check_memory_usage(),
763            "200MB should exceed 100MB threshold"
764        );
765    }
766
767    #[test]
768    fn test_get_final_incomplete_bar_empty() {
769        let mut processor = StreamingProcessor::new(25).unwrap();
770
771        // No trades processed → no incomplete bar
772        let bar = processor.get_final_incomplete_bar();
773        assert!(bar.is_none(), "No incomplete bar before any trades");
774    }
775
776    #[test]
777    fn test_bars_per_aggtrade_ratio() {
778        let metrics = MetricsSummary {
779            trades_processed: 200,
780            bars_generated: 10,
781            errors_total: 0,
782            backpressure_events: 0,
783            circuit_breaker_trips: 0,
784            memory_usage_bytes: 0,
785        };
786
787        assert_eq!(metrics.bars_per_aggtrade(), 0.05);
788        assert_eq!(metrics.error_rate(), 0.0);
789    }
790}