1use futures::Stream;
2use 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#[derive(Debug, Clone)]
21pub struct StreamingProcessorConfig {
22 pub trade_channel_capacity: usize,
24 pub bar_channel_capacity: usize,
26 pub memory_threshold_bytes: usize,
28 pub backpressure_timeout: Duration,
30 pub circuit_breaker_threshold: f64,
32 pub circuit_breaker_timeout: Duration,
34}
35
36impl StreamingProcessorConfig {
37 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, bar_channel_capacity: StreamingProcessorConfig::get_bar_channel_capacity(), memory_threshold_bytes: 100_000_000, backpressure_timeout: Duration::from_millis(100),
54 circuit_breaker_threshold: 0.5, circuit_breaker_timeout: Duration::from_secs(30),
56 }
57 }
58}
59
60pub struct StreamingProcessor {
62 processor: ExportOpenDeviationBarProcessor,
64
65 _threshold_decimal_bps: u32,
67
68 trade_sender: Option<mpsc::Sender<Tick>>,
70 trade_receiver: mpsc::Receiver<Tick>,
71
72 bar_sender: mpsc::Sender<OpenDeviationBar>,
74 bar_receiver: Option<mpsc::Receiver<OpenDeviationBar>>,
75
76 config: StreamingProcessorConfig,
78
79 metrics: Arc<StreamingMetrics>,
81
82 circuit_breaker: CircuitBreaker,
84}
85
86#[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#[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, pub total_block_time_ms: AtomicU64, }
117
118impl StreamingProcessor {
119 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 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 pub fn trade_sender(&mut self) -> Option<mpsc::Sender<Tick>> {
155 self.trade_sender.take()
156 }
157
158 pub fn bar_receiver(&mut self) -> Option<mpsc::Receiver<OpenDeviationBar>> {
160 self.bar_receiver.take()
161 }
162
163 pub async fn start_processing(&mut self) -> Result<(), StreamingError> {
165 loop {
166 if !self.circuit_breaker.can_process() {
168 tokio::time::sleep(Duration::from_millis(100)).await;
169 continue;
170 }
171
172 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 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, };
191
192 match self.process_single_trade(&trade).await {
194 Ok(bar_opt) => {
195 self.circuit_breaker.record_success();
196
197 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 async fn process_single_trade(
219 &mut self,
220 trade: &Tick,
221 ) -> Result<Option<OpenDeviationBar>, StreamingError> {
222 self.metrics
224 .trades_processed
225 .fetch_add(1, Ordering::Relaxed);
226
227 self.processor
229 .process_trades_continuously(std::slice::from_ref(trade));
230
231 let mut completed_bars = self.processor.get_all_completed_bars();
233
234 if !completed_bars.is_empty() {
235 let completed_bar = completed_bars.remove(0);
238
239 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 async fn send_bar_with_backpressure(
259 &self,
260 bar: OpenDeviationBar,
261 ) -> Result<(), StreamingError> {
262 match self.bar_sender.try_send(bar.clone()) {
264 Ok(()) => Ok(()),
265 Err(mpsc::error::TrySendError::Full(_)) => {
266 println!("Bar channel full, applying backpressure");
268 self.metrics
269 .backpressure_events
270 .fetch_add(1, Ordering::Relaxed);
271
272 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 pub fn metrics(&self) -> &StreamingMetrics {
284 &self.metrics
285 }
286
287 pub fn get_final_incomplete_bar(&mut self) -> Option<OpenDeviationBar> {
289 self.processor.get_incomplete_bar()
290 }
291
292 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 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 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
356pub 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#[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 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#[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 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 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 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(); let initial_metrics = processor.metrics().summary();
475
476 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 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 assert!(circuit_breaker.can_process());
496
497 for _ in 0..20 {
499 circuit_breaker.record_failure();
500 }
501
502 assert_eq!(circuit_breaker.state, CircuitBreakerState::Open);
504 assert!(!circuit_breaker.can_process());
505
506 tokio::time::sleep(Duration::from_millis(150)).await;
508
509 assert!(circuit_breaker.can_process());
511
512 circuit_breaker.record_success();
514
515 assert_eq!(circuit_breaker.state, CircuitBreakerState::Closed);
517 }
518
519 #[test]
522 fn test_circuit_breaker_stays_closed_below_threshold() {
523 let mut cb = CircuitBreaker::new(0.5, Duration::from_secs(10));
524
525 for _ in 0..8 {
527 cb.record_success();
528 }
529 for _ in 0..2 {
530 cb.record_failure();
531 }
532
533 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 for _ in 0..9 {
544 cb.record_failure();
545 }
546
547 assert_eq!(cb.state, CircuitBreakerState::Closed);
549 assert!(cb.can_process());
550
551 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 for _ in 0..10 {
562 cb.record_failure();
563 }
564 assert_eq!(cb.state, CircuitBreakerState::Open);
565
566 assert!(cb.can_process());
568 assert_eq!(cb.state, CircuitBreakerState::HalfOpen);
569
570 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 for _ in 0..10 {
582 cb.record_failure();
583 }
584 assert_eq!(cb.state, CircuitBreakerState::Open);
585
586 assert!(cb.can_process());
588 assert_eq!(cb.state, CircuitBreakerState::HalfOpen);
589
590 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)); for _ in 0..10 {
602 cb.record_failure();
603 }
604
605 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 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 #[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 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 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 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 let sender = processor.trade_sender();
709 assert!(
710 sender.is_some(),
711 "First trade_sender() call must return Some"
712 );
713
714 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 let receiver = processor.bar_receiver();
728 assert!(
729 receiver.is_some(),
730 "First bar_receiver() call must return Some"
731 );
732
733 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 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 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 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}