opendeviationbar-streaming 13.78.0

Real-time streaming engine for open deviation bar processing
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
//! Replay buffer for storing and replaying recent trade data
//!
//! This module provides a circular buffer that stores recent trades and allows
//! replaying them at different speeds for testing and analysis.

// Tick has Copy when built locally (workspace path dep) but NOT when
// packaged for crates.io (published dep). Use clone() for compatibility.
#![allow(clippy::clone_on_copy)]

use opendeviationbar_core::Tick;
use std::collections::VecDeque;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
use tokio::time::sleep;
use tokio_stream::Stream;

/// A circular buffer that stores recent trades and provides replay functionality
#[derive(Debug, Clone)]
pub struct ReplayBuffer {
    inner: Arc<Mutex<ReplayBufferInner>>,
}

#[derive(Debug)]
struct ReplayBufferInner {
    capacity: Duration,
    trades: VecDeque<Tick>,
    start_time: Option<Instant>,
}

impl ReplayBuffer {
    /// Create a new replay buffer with the specified capacity
    pub fn new(capacity: Duration) -> Self {
        Self {
            inner: Arc::new(Mutex::new(ReplayBufferInner {
                capacity,
                trades: VecDeque::new(),
                start_time: None,
            })),
        }
    }

    /// Add a new trade to the buffer
    pub fn push(&self, trade: Tick) {
        let mut inner = self
            .inner
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());

        // Set start time on first trade
        if inner.start_time.is_none() {
            inner.start_time = Some(Instant::now());
        }

        // Remove old trades beyond capacity (using microseconds)
        let cutoff_timestamp = trade.timestamp - (inner.capacity.as_micros() as i64);

        while let Some(front_trade) = inner.trades.front() {
            if front_trade.timestamp < cutoff_timestamp {
                inner.trades.pop_front();
            } else {
                break;
            }
        }

        inner.trades.push_back(trade);
    }

    /// Get the number of trades currently in the buffer
    pub fn len(&self) -> usize {
        self.inner
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner())
            .trades
            .len()
    }

    /// Check if the buffer is empty
    pub fn is_empty(&self) -> bool {
        self.inner
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner())
            .trades
            .is_empty()
    }

    /// Get the time span of trades in the buffer
    pub fn time_span(&self) -> Option<Duration> {
        let inner = self
            .inner
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());
        if let (Some(first), Some(last)) = (inner.trades.front(), inner.trades.back()) {
            let span_microseconds = last.timestamp - first.timestamp;
            if span_microseconds > 0 {
                Some(Duration::from_micros(span_microseconds as u64))
            } else {
                None
            }
        } else {
            None
        }
    }

    /// Get trades from the buffer starting from N minutes ago
    /// Issue #96 Task #73: Binary search for cutoff position + collect-only approach (3-8% speedup)
    pub fn get_trades_from(&self, minutes_ago: u32) -> Vec<Tick> {
        let inner = self
            .inner
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());

        // Fast-path: empty buffer (1-2% speedup on empty/single-trade case)
        if inner.trades.is_empty() {
            return Vec::new();
        }

        // Calculate cutoff timestamp
        let latest_timestamp = inner.trades.back().unwrap().timestamp;
        let cutoff_timestamp = latest_timestamp - (minutes_ago as i64 * 60 * 1000);

        // Issue #96 Task #73: Find first trade >= cutoff using early-exit linear scan
        // (VecDeque doesn't have binary_search, so linear scan with early break is optimal)
        let mut start_idx = 0;
        for (idx, trade) in inner.trades.iter().enumerate() {
            if trade.timestamp >= cutoff_timestamp {
                start_idx = idx;
                break;
            }
        }

        // Collect trades from cutoff position forward, avoiding filter overhead
        let mut result = Vec::with_capacity(inner.trades.len().saturating_sub(start_idx));
        result.extend(inner.trades.iter().skip(start_idx).cloned());
        result
    }

    /// Create a replay stream that emits trades at the specified speed
    pub fn replay_from(&self, minutes_ago: u32, speed_multiplier: f32) -> ReplayStream {
        let trades = self.get_trades_from(minutes_ago);
        ReplayStream::new(trades, speed_multiplier)
    }

    /// Get statistics about the buffer
    pub fn stats(&self) -> ReplayBufferStats {
        let inner = self
            .inner
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner());

        let (first_timestamp, last_timestamp) =
            if let (Some(first), Some(last)) = (inner.trades.front(), inner.trades.back()) {
                (Some(first.timestamp), Some(last.timestamp))
            } else {
                (None, None)
            };

        ReplayBufferStats {
            capacity: inner.capacity,
            trade_count: inner.trades.len(),
            first_timestamp,
            last_timestamp,
            memory_usage_bytes: inner.trades.len() * std::mem::size_of::<Tick>(),
        }
    }
}

/// Statistics about the replay buffer
#[derive(Debug, Clone)]
pub struct ReplayBufferStats {
    pub capacity: Duration,
    pub trade_count: usize,
    pub first_timestamp: Option<i64>,
    pub last_timestamp: Option<i64>,
    pub memory_usage_bytes: usize,
}

/// A stream that replays trades at a specified speed
pub struct ReplayStream {
    trades: Vec<Tick>,
    current_index: usize,
    speed_multiplier: f32,
    base_timestamp: Option<i64>,
    start_time: Option<Instant>,
}

impl ReplayStream {
    /// Create a new replay stream
    pub fn new(trades: Vec<Tick>, speed_multiplier: f32) -> Self {
        let base_timestamp = trades.first().map(|t| t.timestamp);

        Self {
            trades,
            current_index: 0,
            speed_multiplier: speed_multiplier.max(0.1), // Minimum 0.1x speed
            base_timestamp,
            start_time: None,
        }
    }

    /// Set the replay speed
    pub fn set_speed(&mut self, speed_multiplier: f32) {
        self.speed_multiplier = speed_multiplier.max(0.1);
    }

    /// Get the current replay speed
    pub fn speed(&self) -> f32 {
        self.speed_multiplier
    }

    /// Get the number of remaining trades
    pub fn remaining(&self) -> usize {
        self.trades.len().saturating_sub(self.current_index)
    }

    /// Get the total number of trades
    pub fn total(&self) -> usize {
        self.trades.len()
    }

    /// Get the progress as a percentage (0.0 to 1.0)
    pub fn progress(&self) -> f32 {
        if self.trades.is_empty() {
            1.0
        } else {
            self.current_index as f32 / self.trades.len() as f32
        }
    }
}

impl Stream for ReplayStream {
    type Item = Tick;

    fn poll_next(
        mut self: std::pin::Pin<&mut Self>,
        cx: &mut std::task::Context<'_>,
    ) -> std::task::Poll<Option<Self::Item>> {
        // Check if we have more trades
        if self.current_index >= self.trades.len() {
            return std::task::Poll::Ready(None);
        }

        // Initialize start time if this is the first trade
        if self.start_time.is_none() {
            let current_trade = self.trades[self.current_index].clone();
            self.start_time = Some(Instant::now());
            self.current_index += 1;
            return std::task::Poll::Ready(Some(current_trade));
        }

        let current_trade = &self.trades[self.current_index];

        // Calculate when this trade should be emitted based on timestamp differences
        if let (Some(base_timestamp), Some(start_time)) = (self.base_timestamp, self.start_time) {
            let time_diff_microseconds = current_trade.timestamp - base_timestamp;
            let real_time_diff = Duration::from_micros(time_diff_microseconds as u64);
            let scaled_time_diff = Duration::from_micros(
                (real_time_diff.as_micros() as f64 / self.speed_multiplier as f64) as u64,
            );

            let target_time = start_time + scaled_time_diff;
            let now = Instant::now();

            if now >= target_time {
                let trade = current_trade.clone();
                self.current_index += 1;
                std::task::Poll::Ready(Some(trade))
            } else {
                // Set up a timer to wake us when it's time
                let waker = cx.waker().clone();
                let sleep_duration = target_time - now;

                tokio::spawn(async move {
                    sleep(sleep_duration).await;
                    waker.wake();
                });

                std::task::Poll::Pending
            }
        } else {
            // Fallback: emit immediately
            let trade = current_trade.clone();
            self.current_index += 1;
            std::task::Poll::Ready(Some(trade))
        }
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use futures::StreamExt;
    use opendeviationbar_core::FixedPoint;

    fn create_test_trade(id: i64, timestamp: i64, price: f64) -> Tick {
        Tick {
            ref_id: id,
            price: FixedPoint::from_str(&price.to_string()).unwrap(),
            volume: FixedPoint::from_str("1.0").unwrap(),
            first_sub_id: id,
            last_sub_id: id,
            timestamp,
            is_buyer_maker: false,
            is_best_match: None,
            best_bid: None,
            best_ask: None,
        }
    }

    #[test]
    fn test_replay_buffer_capacity() {
        let buffer = ReplayBuffer::new(Duration::from_secs(60)); // 1 minute capacity

        let base_time = 1_704_067_200_000_000_i64; // 2024-01-01 00:00:00 in microseconds

        // Add trades spanning 2 minutes (1 trade per second in microseconds)
        for i in 0..120 {
            let trade = create_test_trade(i, base_time + (i * 1_000_000), 50000.0); // 1 second intervals in microseconds
            buffer.push(trade);
        }

        // Should only keep the last 60 seconds worth
        let stats = buffer.stats();
        assert!(
            stats.trade_count <= 120,
            "Expected <= 120 trades, got {}",
            stats.trade_count
        );

        let trades = buffer.get_trades_from(1); // Last 1 minute
        assert!(!trades.is_empty());
    }

    #[test]
    fn test_replay_buffer_time_span() {
        let buffer = ReplayBuffer::new(Duration::from_secs(300)); // 5 minutes

        // Use a more realistic timestamp (approximately 2024-01-01 in microseconds)
        let base_time = 1_704_067_200_000_000_i64; // 2024-01-01 00:00:00 in microseconds

        // Add a few trades over 1 minute (using microseconds)
        buffer.push(create_test_trade(1, base_time, 50000.0));
        buffer.push(create_test_trade(2, base_time + 30_000_000, 50100.0)); // 30s later in microseconds
        buffer.push(create_test_trade(3, base_time + 60_000_000, 50200.0)); // 60s later in microseconds

        let span = buffer.time_span().unwrap();
        assert_eq!(span.as_secs(), 60);
    }

    #[tokio::test]
    async fn test_replay_stream() {
        // Use a more realistic timestamp (approximately 2024-01-01 in microseconds)
        let base_time = 1_704_067_200_000_000_i64; // 2024-01-01 00:00:00 in microseconds
        let trades = vec![
            create_test_trade(1, base_time, 50000.0),
            create_test_trade(2, base_time + 1_000_000, 50100.0), // 1s later in microseconds
            create_test_trade(3, base_time + 2_000_000, 50200.0), // 2s later in microseconds
        ];

        let mut stream = ReplayStream::new(trades, 10.0); // 10x speed
        assert_eq!(stream.total(), 3);
        assert_eq!(stream.remaining(), 3);
        assert_eq!(stream.progress(), 0.0);

        // First trade should be immediate
        let first = StreamExt::next(&mut stream).await;
        assert!(first.is_some());
        assert_eq!(first.unwrap().ref_id, 1);
    }

    #[test]
    fn test_get_trades_from_empty_buffer() {
        let buffer = ReplayBuffer::new(Duration::from_secs(60));
        let trades = buffer.get_trades_from(1);
        assert_eq!(trades.len(), 0); // Should not panic
    }

    // === Issue #96: Expanded replay buffer + replay stream coverage ===

    #[test]
    fn test_push_len_is_empty() {
        let buffer = ReplayBuffer::new(Duration::from_secs(300));
        assert!(buffer.is_empty());
        assert_eq!(buffer.len(), 0);

        let base = 1_704_067_200_000_000_i64;
        buffer.push(create_test_trade(1, base, 50000.0));
        assert!(!buffer.is_empty());
        assert_eq!(buffer.len(), 1);

        buffer.push(create_test_trade(2, base + 1_000_000, 50100.0));
        assert_eq!(buffer.len(), 2);
    }

    #[test]
    fn test_push_eviction_by_capacity() {
        // 10 second capacity
        let buffer = ReplayBuffer::new(Duration::from_secs(10));
        let base = 1_704_067_200_000_000_i64;

        // Add 20 trades 1 second apart (20 sec span)
        for i in 0..20 {
            buffer.push(create_test_trade(i, base + (i * 1_000_000), 50000.0));
        }

        // Capacity is 10s, so ~last 10 trades should remain
        let len = buffer.len();
        assert!(len <= 12, "Expected <=12 trades after eviction, got {len}");
        assert!(len >= 10, "Expected >=10 trades retained, got {len}");
    }

    #[test]
    fn test_replay_stream_set_speed() {
        let trades = vec![
            create_test_trade(1, 1_704_067_200_000_000, 50000.0),
            create_test_trade(2, 1_704_067_201_000_000, 50100.0),
        ];
        let mut stream = ReplayStream::new(trades, 1.0);
        assert!((stream.speed() - 1.0).abs() < f32::EPSILON);

        stream.set_speed(5.0);
        assert!((stream.speed() - 5.0).abs() < f32::EPSILON);

        // Minimum speed clamp
        stream.set_speed(0.01);
        assert!((stream.speed() - 0.1).abs() < f32::EPSILON);
    }

    #[test]
    fn test_replay_stream_remaining_total_progress() {
        let base = 1_704_067_200_000_000_i64;
        let trades = vec![
            create_test_trade(1, base, 50000.0),
            create_test_trade(2, base + 1_000_000, 50100.0),
            create_test_trade(3, base + 2_000_000, 50200.0),
            create_test_trade(4, base + 3_000_000, 50300.0),
        ];
        let stream = ReplayStream::new(trades, 10.0);

        assert_eq!(stream.total(), 4);
        assert_eq!(stream.remaining(), 4);
        assert!((stream.progress() - 0.0).abs() < f32::EPSILON);
    }

    #[test]
    fn test_replay_stream_empty_progress() {
        let stream = ReplayStream::new(vec![], 1.0);
        assert_eq!(stream.total(), 0);
        assert_eq!(stream.remaining(), 0);
        assert!((stream.progress() - 1.0).abs() < f32::EPSILON);
    }

    #[test]
    fn test_buffer_stats() {
        let buffer = ReplayBuffer::new(Duration::from_secs(60));
        let base = 1_704_067_200_000_000_i64;

        let stats = buffer.stats();
        assert_eq!(stats.trade_count, 0);
        assert!(stats.first_timestamp.is_none());
        assert!(stats.last_timestamp.is_none());

        buffer.push(create_test_trade(1, base, 50000.0));
        buffer.push(create_test_trade(2, base + 5_000_000, 50100.0));

        let stats = buffer.stats();
        assert_eq!(stats.trade_count, 2);
        assert_eq!(stats.first_timestamp, Some(base));
        assert_eq!(stats.last_timestamp, Some(base + 5_000_000));
        assert!(stats.memory_usage_bytes > 0);
    }
}