1use std::cmp::Reverse;
12use std::collections::BinaryHeap;
13use std::marker::PhantomData;
14
15use chrono::NaiveDateTime;
16use qs_core::ExecutionPricer;
17use qs_core::types::PriceQuote;
18use serde::{Deserialize, Serialize};
19use thiserror::Error;
20
21#[derive(Debug, Clone, Serialize, Deserialize)]
25pub enum MarketEvent {
26 Tick {
28 symbol: String,
29 ts: NaiveDateTime,
30 bid: f64,
31 ask: f64,
32 },
33 Bar {
35 symbol: String,
36 ts: NaiveDateTime,
37 open: f64,
38 high: f64,
39 low: f64,
40 close: f64,
41 volume: i64,
42 },
43}
44
45impl MarketEvent {
46 pub fn ts(&self) -> NaiveDateTime {
48 match self {
49 MarketEvent::Tick { ts, .. } => *ts,
50 MarketEvent::Bar { ts, .. } => *ts,
51 }
52 }
53
54 pub fn symbol(&self) -> &str {
56 match self {
57 MarketEvent::Tick { symbol, .. } => symbol,
58 MarketEvent::Bar { symbol, .. } => symbol,
59 }
60 }
61
62 pub fn to_quote(&self) -> PriceQuote {
69 match self {
70 MarketEvent::Tick {
71 symbol,
72 ts,
73 bid,
74 ask,
75 } => PriceQuote {
76 symbol: symbol.clone(),
77 ts: *ts,
78 bid: *bid,
79 ask: *ask,
80 },
81 MarketEvent::Bar {
82 symbol, ts, close, ..
83 } => PriceQuote {
84 symbol: symbol.clone(),
85 ts: *ts,
86 bid: *close,
87 ask: *close,
88 },
89 }
90 }
91
92 pub fn to_valid_quote(&self) -> Option<PriceQuote> {
95 let quote = self.to_quote();
96 ExecutionPricer::validate_quote("e).ok().map(|()| quote)
97 }
98}
99
100#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
102pub struct SeriesRoles {
103 pub primary: bool,
104 pub conversion: bool,
105}
106
107impl SeriesRoles {
108 pub const PRIMARY: Self = Self {
109 primary: true,
110 conversion: false,
111 };
112 pub const CONVERSION: Self = Self {
113 primary: false,
114 conversion: true,
115 };
116 pub const PRIMARY_AND_CONVERSION: Self = Self {
117 primary: true,
118 conversion: true,
119 };
120}
121
122#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
124pub struct EventMetadata {
125 pub roles: SeriesRoles,
126 pub series_rank: u32,
127 pub row_sequence: u64,
129}
130
131impl EventMetadata {
132 pub const fn new(roles: SeriesRoles, series_rank: u32, row_sequence: u64) -> Self {
133 Self {
134 roles,
135 series_rank,
136 row_sequence,
137 }
138 }
139}
140
141#[derive(Debug, Clone, Serialize, Deserialize)]
143pub struct FeedEvent {
144 pub event: MarketEvent,
145 pub metadata: EventMetadata,
146}
147
148impl FeedEvent {
149 pub fn new(event: MarketEvent, metadata: EventMetadata) -> Self {
150 Self { event, metadata }
151 }
152
153 pub fn ordering_key(&self) -> (NaiveDateTime, u32, u64) {
155 (
156 self.event.ts(),
157 self.metadata.series_rank,
158 self.metadata.row_sequence,
159 )
160 }
161}
162
163#[derive(Debug, Clone)]
165pub struct SequencedMarketEvent {
166 pub event: MarketEvent,
167 pub source_row_ordinal: u64,
168}
169
170impl SequencedMarketEvent {
171 pub fn new(event: MarketEvent, source_row_ordinal: u64) -> Self {
172 Self {
173 event,
174 source_row_ordinal,
175 }
176 }
177}
178
179pub trait EventSourceItem {
181 fn into_event_and_ordinal(self) -> (MarketEvent, Option<u64>);
182}
183
184impl EventSourceItem for MarketEvent {
185 fn into_event_and_ordinal(self) -> (MarketEvent, Option<u64>) {
186 (self, None)
187 }
188}
189
190impl EventSourceItem for SequencedMarketEvent {
191 fn into_event_and_ordinal(self) -> (MarketEvent, Option<u64>) {
192 (self.event, Some(self.source_row_ordinal))
193 }
194}
195
196#[derive(Debug, Clone, Serialize, Deserialize)]
198pub struct TimestampBatch {
199 pub ts: NaiveDateTime,
200 pub events: Vec<FeedEvent>,
201}
202
203impl TimestampBatch {
204 pub fn len(&self) -> usize {
205 self.events.len()
206 }
207
208 pub fn is_empty(&self) -> bool {
209 self.events.is_empty()
210 }
211}
212
213pub trait DataFeed {
215 fn next_event(&mut self) -> Option<MarketEvent>;
217
218 fn peek(&self) -> Option<&MarketEvent>;
220
221 fn next_batch(&mut self) -> Option<TimestampBatch> {
226 let ts = self.peek()?.ts();
227 let mut events = Vec::new();
228 let mut row_sequence = 0;
229
230 while self.peek().is_some_and(|event| event.ts() == ts) {
231 let event = self.next_event()?;
232 events.push(FeedEvent::new(
233 event,
234 EventMetadata::new(SeriesRoles::PRIMARY, 0, row_sequence),
235 ));
236 row_sequence += 1;
237 }
238
239 Some(TimestampBatch { ts, events })
240 }
241
242 fn total_events(&self) -> Option<usize> {
244 None
245 }
246}
247
248pub trait FallibleBatchFeed {
250 type Error;
251
252 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error>;
254}
255
256impl<F> FallibleBatchFeed for Box<F>
257where
258 F: FallibleBatchFeed + ?Sized,
259{
260 type Error = F::Error;
261
262 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
263 (**self).next_batch()
264 }
265}
266
267#[derive(Debug, Error)]
269pub enum EventBatchFeedError<E> {
270 #[error("event source failed: {0}")]
271 Source(E),
272 #[error("event source moved backwards from {previous} to {current}")]
273 NonMonotonic {
274 previous: NaiveDateTime,
275 current: NaiveDateTime,
276 },
277 #[error("event row sequence is exhausted")]
278 RowSequenceExhausted,
279}
280
281pub struct EventBatchFeed<F, I = MarketEvent> {
283 next_event: F,
284 pending: Option<FeedEvent>,
285 roles: SeriesRoles,
286 series_rank: u32,
287 next_row_sequence: u64,
288 last_source_ts: Option<NaiveDateTime>,
289 exhausted: bool,
290 item: PhantomData<fn() -> I>,
291}
292
293impl<F, I> EventBatchFeed<F, I> {
294 pub fn new<E>(next_event: F, roles: SeriesRoles, series_rank: u32) -> Self
295 where
296 F: FnMut() -> Result<Option<I>, E>,
297 I: EventSourceItem,
298 {
299 Self {
300 next_event,
301 pending: None,
302 roles,
303 series_rank,
304 next_row_sequence: 0,
305 last_source_ts: None,
306 exhausted: false,
307 item: PhantomData,
308 }
309 }
310}
311
312impl<F, I, E> EventBatchFeed<F, I>
313where
314 F: FnMut() -> Result<Option<I>, E>,
315 I: EventSourceItem,
316{
317 pub fn next_timestamp_batch(
319 &mut self,
320 ) -> Result<Option<TimestampBatch>, EventBatchFeedError<E>> {
321 let first = match self.pending.take() {
322 Some(event) => event,
323 None => match self.pull_event()? {
324 Some(event) => event,
325 None => return Ok(None),
326 },
327 };
328 let ts = first.event.ts();
329 let mut events = vec![first];
330
331 loop {
332 match self.pull_event()? {
333 Some(event) if event.event.ts() == ts => events.push(event),
334 Some(event) => {
335 self.pending = Some(event);
336 break;
337 }
338 None => break,
339 }
340 }
341
342 Ok(Some(TimestampBatch { ts, events }))
343 }
344
345 fn pull_event(&mut self) -> Result<Option<FeedEvent>, EventBatchFeedError<E>> {
346 if self.exhausted {
347 return Ok(None);
348 }
349 let item = match (self.next_event)() {
350 Ok(Some(item)) => item,
351 Ok(None) => {
352 self.exhausted = true;
353 return Ok(None);
354 }
355 Err(error) => {
356 self.exhausted = true;
357 return Err(EventBatchFeedError::Source(error));
358 }
359 };
360 let (event, source_row_ordinal) = item.into_event_and_ordinal();
361 let current = event.ts();
362 if let Some(previous) = self.last_source_ts
363 && current < previous
364 {
365 self.exhausted = true;
366 return Err(EventBatchFeedError::NonMonotonic { previous, current });
367 }
368 let row_sequence = match source_row_ordinal {
369 Some(source_row_ordinal) => source_row_ordinal,
370 None => {
371 let row_sequence = self.next_row_sequence;
372 self.next_row_sequence = self
373 .next_row_sequence
374 .checked_add(1)
375 .ok_or(EventBatchFeedError::RowSequenceExhausted)?;
376 row_sequence
377 }
378 };
379 self.last_source_ts = Some(current);
380 Ok(Some(FeedEvent::new(
381 event,
382 EventMetadata::new(self.roles, self.series_rank, row_sequence),
383 )))
384 }
385}
386
387impl<F, I, E> FallibleBatchFeed for EventBatchFeed<F, I>
388where
389 F: FnMut() -> Result<Option<I>, E>,
390 I: EventSourceItem,
391{
392 type Error = EventBatchFeedError<E>;
393
394 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
395 self.next_timestamp_batch()
396 }
397}
398
399#[derive(Debug, Error)]
401pub enum KWayMergeError<E> {
402 #[error("series {series_rank} failed: {error}")]
403 Source { series_rank: u32, error: E },
404 #[error("series {series_rank} returned an empty timestamp batch at {ts}")]
405 EmptyBatch { series_rank: u32, ts: NaiveDateTime },
406 #[error("series {series_rank} batch at {batch_ts} contains an event at {event_ts}")]
407 TimestampMismatch {
408 series_rank: u32,
409 batch_ts: NaiveDateTime,
410 event_ts: NaiveDateTime,
411 },
412 #[error("series {series_rank} moved backwards from {previous} to {current}")]
413 NonMonotonic {
414 series_rank: u32,
415 previous: NaiveDateTime,
416 current: NaiveDateTime,
417 },
418 #[error("duplicate merge ordering key ({ts}, {series_rank}, {row_sequence})")]
419 DuplicateOrderingKey {
420 ts: NaiveDateTime,
421 series_rank: u32,
422 row_sequence: u64,
423 },
424}
425
426pub struct KWayMergeFeed<F>
428where
429 F: FallibleBatchFeed,
430{
431 feeds: Vec<F>,
432 heads: Vec<Option<TimestampBatch>>,
433 head_heap: BinaryHeap<Reverse<(NaiveDateTime, usize)>>,
434 exhausted: Vec<bool>,
435 last_batch_ts: Vec<Option<NaiveDateTime>>,
436 initialized: bool,
437}
438
439impl<F> KWayMergeFeed<F>
440where
441 F: FallibleBatchFeed,
442{
443 pub fn new(feeds: Vec<F>) -> Self {
445 assert!(
446 feeds.len() <= u32::MAX as usize,
447 "KWayMergeFeed supports at most u32::MAX series"
448 );
449 let series_count = feeds.len();
450 Self {
451 feeds,
452 heads: (0..series_count).map(|_| None).collect(),
453 head_heap: BinaryHeap::with_capacity(series_count),
454 exhausted: vec![false; series_count],
455 last_batch_ts: vec![None; series_count],
456 initialized: false,
457 }
458 }
459
460 pub fn series_count(&self) -> usize {
461 self.feeds.len()
462 }
463
464 pub fn is_empty(&self) -> bool {
465 self.feeds.is_empty()
466 }
467
468 pub fn next_timestamp_batch(
470 &mut self,
471 ) -> Result<Option<TimestampBatch>, KWayMergeError<F::Error>> {
472 if !self.initialized {
473 for index in 0..self.feeds.len() {
474 self.fill_head(index)?;
475 }
476 self.initialized = true;
477 }
478 let Some(Reverse((ts, _))) = self.head_heap.peek().copied() else {
479 return Ok(None);
480 };
481
482 let mut events = Vec::new();
483 while self
484 .head_heap
485 .peek()
486 .is_some_and(|Reverse((head_ts, _))| *head_ts == ts)
487 {
488 let Reverse((_, index)) = self.head_heap.pop().expect("matching heap head exists");
489 let batch = self.heads[index]
490 .take()
491 .expect("heap head retains its batch");
492 let series_rank = index as u32;
493 events.extend(batch.events.into_iter().map(|mut event| {
494 event.metadata.series_rank = series_rank;
495 event
496 }));
497 self.fill_head(index)?;
498 }
499
500 events.sort_by_key(FeedEvent::ordering_key);
501 for duplicate in events.windows(2) {
502 if duplicate[0].ordering_key() == duplicate[1].ordering_key() {
503 let (_, series_rank, row_sequence) = duplicate[0].ordering_key();
504 return Err(KWayMergeError::DuplicateOrderingKey {
505 ts,
506 series_rank,
507 row_sequence,
508 });
509 }
510 }
511 Ok(Some(TimestampBatch { ts, events }))
512 }
513
514 fn fill_head(&mut self, index: usize) -> Result<(), KWayMergeError<F::Error>> {
515 if self.heads[index].is_some() || self.exhausted[index] {
516 return Ok(());
517 }
518 let series_rank = index as u32;
519 let Some(batch) = self.feeds[index]
520 .next_batch()
521 .map_err(|error| KWayMergeError::Source { series_rank, error })?
522 else {
523 self.exhausted[index] = true;
524 return Ok(());
525 };
526 if batch.is_empty() {
527 return Err(KWayMergeError::EmptyBatch {
528 series_rank,
529 ts: batch.ts,
530 });
531 }
532 for event in &batch.events {
533 if event.event.ts() != batch.ts {
534 return Err(KWayMergeError::TimestampMismatch {
535 series_rank,
536 batch_ts: batch.ts,
537 event_ts: event.event.ts(),
538 });
539 }
540 }
541 if let Some(previous) = self.last_batch_ts[index]
542 && batch.ts < previous
543 {
544 return Err(KWayMergeError::NonMonotonic {
545 series_rank,
546 previous,
547 current: batch.ts,
548 });
549 }
550 self.last_batch_ts[index] = Some(batch.ts);
551 self.head_heap.push(Reverse((batch.ts, index)));
552 self.heads[index] = Some(batch);
553 Ok(())
554 }
555}
556
557impl<F> FallibleBatchFeed for KWayMergeFeed<F>
558where
559 F: FallibleBatchFeed,
560{
561 type Error = KWayMergeError<F::Error>;
562
563 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
564 self.next_timestamp_batch()
565 }
566}
567
568#[derive(Debug, Clone)]
573pub struct VecFeed {
574 events: Vec<FeedEvent>,
575 index: usize,
576}
577
578impl VecFeed {
579 pub fn new(events: Vec<MarketEvent>) -> Self {
581 let events = events
582 .into_iter()
583 .enumerate()
584 .map(|(row_sequence, event)| {
585 FeedEvent::new(
586 event,
587 EventMetadata::new(SeriesRoles::PRIMARY, 0, row_sequence as u64),
588 )
589 })
590 .collect();
591 Self { events, index: 0 }
592 }
593
594 pub fn from_feed_events(mut events: Vec<FeedEvent>) -> Self {
596 events.sort_by_key(FeedEvent::ordering_key);
597 Self { events, index: 0 }
598 }
599
600 pub fn next_feed_event(&mut self) -> Option<FeedEvent> {
602 let event = self.events.get(self.index)?.clone();
603 self.index += 1;
604 Some(event)
605 }
606
607 pub fn peek_metadata(&self) -> Option<&EventMetadata> {
609 self.events.get(self.index).map(|event| &event.metadata)
610 }
611
612 pub fn next_timestamp_batch(&mut self) -> Option<TimestampBatch> {
614 let ts = self.events.get(self.index)?.event.ts();
615 let start = self.index;
616 while self
617 .events
618 .get(self.index)
619 .is_some_and(|event| event.event.ts() == ts)
620 {
621 self.index += 1;
622 }
623
624 Some(TimestampBatch {
625 ts,
626 events: self.events[start..self.index].to_vec(),
627 })
628 }
629
630 pub fn remaining(&self) -> usize {
632 self.events.len().saturating_sub(self.index)
633 }
634
635 pub fn total(&self) -> usize {
637 self.events.len()
638 }
639
640 pub fn reset(&mut self) {
642 self.index = 0;
643 }
644}
645
646impl DataFeed for VecFeed {
647 fn next_event(&mut self) -> Option<MarketEvent> {
648 self.next_feed_event().map(|event| event.event)
649 }
650
651 fn peek(&self) -> Option<&MarketEvent> {
652 self.events.get(self.index).map(|event| &event.event)
653 }
654
655 fn next_batch(&mut self) -> Option<TimestampBatch> {
656 self.next_timestamp_batch()
657 }
658
659 fn total_events(&self) -> Option<usize> {
660 Some(self.events.len())
661 }
662}
663
664pub fn ticks_to_feed(ticks: Vec<data_preprocess::Tick>) -> VecFeed {
668 ticks_to_feed_with_metadata(ticks, SeriesRoles::PRIMARY, 0)
669}
670
671pub fn ticks_to_feed_with_metadata(
673 ticks: Vec<data_preprocess::Tick>,
674 roles: SeriesRoles,
675 series_rank: u32,
676) -> VecFeed {
677 let events = ticks
678 .into_iter()
679 .enumerate()
680 .filter_map(|(row_sequence, tick)| {
681 let bid = tick.bid?;
682 let ask = tick.ask?;
683 let event = MarketEvent::Tick {
684 symbol: tick.symbol,
685 ts: tick.ts,
686 bid,
687 ask,
688 };
689 event.to_valid_quote().map(|_| {
690 FeedEvent::new(
691 event,
692 EventMetadata::new(roles, series_rank, row_sequence as u64),
693 )
694 })
695 })
696 .collect();
697 VecFeed::from_feed_events(events)
698}
699
700pub fn bars_to_feed(bars: Vec<data_preprocess::Bar>) -> VecFeed {
702 bars_to_feed_with_metadata(bars, SeriesRoles::PRIMARY, 0)
703}
704
705pub fn bars_to_feed_with_metadata(
707 bars: Vec<data_preprocess::Bar>,
708 roles: SeriesRoles,
709 series_rank: u32,
710) -> VecFeed {
711 let events = bars
712 .into_iter()
713 .enumerate()
714 .map(|(row_sequence, bar)| {
715 let event = MarketEvent::Bar {
716 symbol: bar.symbol,
717 ts: bar.ts,
718 open: bar.open,
719 high: bar.high,
720 low: bar.low,
721 close: bar.close,
722 volume: bar.volume,
723 };
724 FeedEvent::new(
725 event,
726 EventMetadata::new(roles, series_rank, row_sequence as u64),
727 )
728 })
729 .collect();
730 VecFeed::from_feed_events(events)
731}
732
733pub fn merge_feeds(feeds: Vec<VecFeed>) -> VecFeed {
738 let mut all_events = Vec::new();
739 for (series_rank, feed) in feeds.into_iter().enumerate() {
740 let series_rank = series_rank.min(u32::MAX as usize) as u32;
741 all_events.extend(feed.events.into_iter().map(|mut event| {
742 event.metadata.series_rank = series_rank;
743 event
744 }));
745 }
746 VecFeed::from_feed_events(all_events)
747}
748
749#[cfg(test)]
752mod tests {
753 use std::collections::VecDeque;
754
755 use super::*;
756 use chrono::NaiveDate;
757
758 fn ts(h: u32, m: u32, s: u32) -> NaiveDateTime {
759 NaiveDate::from_ymd_opt(2026, 1, 1)
760 .unwrap()
761 .and_hms_opt(h, m, s)
762 .unwrap()
763 }
764
765 fn sample_events() -> Vec<MarketEvent> {
766 vec![
767 MarketEvent::Tick {
768 symbol: "EURUSD".into(),
769 ts: ts(10, 0, 0),
770 bid: 1.0848,
771 ask: 1.0850,
772 },
773 MarketEvent::Tick {
774 symbol: "EURUSD".into(),
775 ts: ts(10, 0, 1),
776 bid: 1.0849,
777 ask: 1.0851,
778 },
779 MarketEvent::Tick {
780 symbol: "EURUSD".into(),
781 ts: ts(10, 0, 2),
782 bid: 1.0847,
783 ask: 1.0849,
784 },
785 ]
786 }
787
788 #[test]
789 fn vec_feed_iterates() {
790 let mut feed = VecFeed::new(sample_events());
791 assert_eq!(feed.total(), 3);
792 assert_eq!(feed.remaining(), 3);
793
794 let e1 = feed.next_event().unwrap();
795 assert_eq!(e1.symbol(), "EURUSD");
796 assert_eq!(feed.remaining(), 2);
797
798 let _ = feed.next_event().unwrap();
799 let _ = feed.next_event().unwrap();
800 assert!(feed.next_event().is_none());
801 assert_eq!(feed.remaining(), 0);
802 }
803
804 #[test]
805 fn vec_feed_peek() {
806 let feed = VecFeed::new(sample_events());
807 let peeked = feed.peek().unwrap();
808 assert_eq!(peeked.ts(), ts(10, 0, 0));
809 }
810
811 #[test]
812 fn vec_feed_reset() {
813 let mut feed = VecFeed::new(sample_events());
814 let _ = feed.next_event();
815 let _ = feed.next_event();
816 feed.reset();
817 assert_eq!(feed.remaining(), 3);
818 }
819
820 #[test]
821 fn to_quote_tick() {
822 let event = MarketEvent::Tick {
823 symbol: "EURUSD".into(),
824 ts: ts(10, 0, 0),
825 bid: 1.0848,
826 ask: 1.0850,
827 };
828 let q = event.to_quote();
829 assert_eq!(q.symbol, "EURUSD");
830 assert!((q.bid - 1.0848).abs() < f64::EPSILON);
831 assert!((q.ask - 1.0850).abs() < f64::EPSILON);
832 }
833
834 #[test]
835 fn to_quote_bar() {
836 let event = MarketEvent::Bar {
837 symbol: "EURUSD".into(),
838 ts: ts(10, 0, 0),
839 open: 1.0840,
840 high: 1.0860,
841 low: 1.0830,
842 close: 1.0855,
843 volume: 1000,
844 };
845 let q = event.to_quote();
846 assert!((q.bid - 1.0855).abs() < f64::EPSILON);
848 assert!((q.ask - 1.0855).abs() < f64::EPSILON);
849 }
850
851 #[test]
852 fn merge_two_feeds() {
853 let feed_a = VecFeed::new(vec![
854 MarketEvent::Tick {
855 symbol: "EURUSD".into(),
856 ts: ts(10, 0, 0),
857 bid: 1.08,
858 ask: 1.09,
859 },
860 MarketEvent::Tick {
861 symbol: "EURUSD".into(),
862 ts: ts(10, 0, 2),
863 bid: 1.08,
864 ask: 1.09,
865 },
866 ]);
867 let feed_b = VecFeed::new(vec![MarketEvent::Tick {
868 symbol: "XAUUSD".into(),
869 ts: ts(10, 0, 1),
870 bid: 2000.0,
871 ask: 2001.0,
872 }]);
873
874 let mut merged = merge_feeds(vec![feed_a, feed_b]);
875 assert_eq!(merged.total(), 3);
876
877 let e1 = merged.next_event().unwrap();
878 assert_eq!(e1.ts(), ts(10, 0, 0));
879 assert_eq!(e1.symbol(), "EURUSD");
880
881 let e2 = merged.next_event().unwrap();
882 assert_eq!(e2.ts(), ts(10, 0, 1));
883 assert_eq!(e2.symbol(), "XAUUSD");
884
885 let e3 = merged.next_event().unwrap();
886 assert_eq!(e3.ts(), ts(10, 0, 2));
887 assert_eq!(e3.symbol(), "EURUSD");
888 }
889
890 #[test]
891 fn ticks_to_feed_skips_missing_prices() {
892 let ticks = vec![
893 data_preprocess::Tick {
894 exchange: "test".into(),
895 symbol: "EURUSD".into(),
896 ts: ts(10, 0, 0),
897 bid: Some(1.08),
898 ask: Some(1.09),
899 last: None,
900 volume: None,
901 flags: None,
902 },
903 data_preprocess::Tick {
904 exchange: "test".into(),
905 symbol: "EURUSD".into(),
906 ts: ts(10, 0, 1),
907 bid: Some(1.08),
908 ask: None, last: None,
910 volume: None,
911 flags: None,
912 },
913 data_preprocess::Tick {
914 exchange: "test".into(),
915 symbol: "EURUSD".into(),
916 ts: ts(10, 0, 2),
917 bid: None, ask: Some(1.09),
919 last: None,
920 volume: None,
921 flags: None,
922 },
923 ];
924
925 let feed = ticks_to_feed(ticks);
926 assert_eq!(feed.total(), 1); }
928
929 #[test]
930 fn validated_quote_contract_rejects_nonfinite_nonpositive_and_crossed_ticks() {
931 for (bid, ask) in [
932 (f64::NAN, 1.0),
933 (1.0, f64::INFINITY),
934 (0.0, 1.0),
935 (2.0, 1.0),
936 ] {
937 let event = MarketEvent::Tick {
938 symbol: "EURUSD".into(),
939 ts: ts(10, 0, 0),
940 bid,
941 ask,
942 };
943 assert!(event.to_valid_quote().is_none());
944 }
945
946 let ticks = vec![
947 data_preprocess::Tick {
948 exchange: "test".into(),
949 symbol: "EURUSD".into(),
950 ts: ts(10, 0, 0),
951 bid: Some(1.1),
952 ask: Some(1.0),
953 last: None,
954 volume: None,
955 flags: None,
956 },
957 data_preprocess::Tick {
958 exchange: "test".into(),
959 symbol: "EURUSD".into(),
960 ts: ts(10, 0, 1),
961 bid: Some(1.0),
962 ask: Some(1.1),
963 last: None,
964 volume: None,
965 flags: None,
966 },
967 ];
968 assert_eq!(ticks_to_feed(ticks).total(), 1);
969 }
970
971 #[test]
972 fn bars_to_feed_converts_all() {
973 let bars = vec![
974 data_preprocess::Bar {
975 exchange: "test".into(),
976 symbol: "EURUSD".into(),
977 timeframe: data_preprocess::Timeframe::M5,
978 ts: ts(10, 0, 0),
979 open: 1.0840,
980 high: 1.0860,
981 low: 1.0830,
982 close: 1.0855,
983 tick_vol: 100,
984 volume: 1000,
985 spread: 2,
986 },
987 data_preprocess::Bar {
988 exchange: "test".into(),
989 symbol: "EURUSD".into(),
990 timeframe: data_preprocess::Timeframe::M5,
991 ts: ts(10, 5, 0),
992 open: 1.0855,
993 high: 1.0870,
994 low: 1.0845,
995 close: 1.0865,
996 tick_vol: 120,
997 volume: 1200,
998 spread: 2,
999 },
1000 ];
1001
1002 let feed = bars_to_feed(bars);
1003 assert_eq!(feed.total(), 2);
1004 }
1005
1006 #[test]
1007 fn equal_timestamp_batch_uses_deterministic_metadata_order() {
1008 let batch_ts = ts(10, 0, 0);
1009 let later_ts = ts(10, 0, 1);
1010 let events = vec![
1011 FeedEvent::new(
1012 MarketEvent::Tick {
1013 symbol: "CONVERSION".into(),
1014 ts: batch_ts,
1015 bid: 2.0,
1016 ask: 2.1,
1017 },
1018 EventMetadata::new(SeriesRoles::CONVERSION, 1, 0),
1019 ),
1020 FeedEvent::new(
1021 MarketEvent::Bar {
1022 symbol: "PRIMARY".into(),
1023 ts: batch_ts,
1024 open: 1.0,
1025 high: 1.2,
1026 low: 0.9,
1027 close: 1.1,
1028 volume: 10,
1029 },
1030 EventMetadata::new(SeriesRoles::PRIMARY, 0, 1),
1031 ),
1032 FeedEvent::new(
1033 MarketEvent::Tick {
1034 symbol: "PRIMARY".into(),
1035 ts: batch_ts,
1036 bid: 1.0,
1037 ask: 1.1,
1038 },
1039 EventMetadata::new(SeriesRoles::PRIMARY, 0, 0),
1040 ),
1041 FeedEvent::new(
1042 MarketEvent::Tick {
1043 symbol: "PRIMARY".into(),
1044 ts: later_ts,
1045 bid: 1.1,
1046 ask: 1.2,
1047 },
1048 EventMetadata::new(SeriesRoles::PRIMARY, 0, 2),
1049 ),
1050 ];
1051 let mut feed = VecFeed::from_feed_events(events);
1052
1053 let batch = feed
1054 .next_timestamp_batch()
1055 .expect("equal timestamp batch should be available");
1056 assert_eq!(batch.ts, batch_ts);
1057 assert_eq!(batch.len(), 3);
1058 let keys: Vec<_> = batch.events.iter().map(FeedEvent::ordering_key).collect();
1059 assert_eq!(
1060 keys,
1061 vec![(batch_ts, 0, 0), (batch_ts, 0, 1), (batch_ts, 1, 0)]
1062 );
1063 assert!(matches!(batch.events[1].event, MarketEvent::Bar { .. }));
1064 assert!(matches!(batch.events[2].event, MarketEvent::Tick { .. }));
1065
1066 let later = DataFeed::next_batch(&mut feed).expect("later batch should remain");
1067 assert_eq!(later.ts, later_ts);
1068 assert_eq!(later.len(), 1);
1069 assert!(DataFeed::next_batch(&mut feed).is_none());
1070 }
1071
1072 #[test]
1073 fn one_tick_can_serve_primary_and_conversion_roles() {
1074 let ticks = vec![data_preprocess::Tick {
1075 exchange: "test".into(),
1076 symbol: "EURUSD".into(),
1077 ts: ts(10, 0, 0),
1078 bid: Some(1.08),
1079 ask: Some(1.09),
1080 last: None,
1081 volume: None,
1082 flags: None,
1083 }];
1084 let mut feed = ticks_to_feed_with_metadata(ticks, SeriesRoles::PRIMARY_AND_CONVERSION, 0);
1085
1086 let batch = feed.next_timestamp_batch().unwrap();
1087 assert_eq!(batch.len(), 1);
1088 assert_eq!(
1089 batch.events[0].metadata.roles,
1090 SeriesRoles::PRIMARY_AND_CONVERSION
1091 );
1092 assert_eq!(batch.events[0].metadata.row_sequence, 0);
1093 }
1094
1095 #[test]
1096 fn merge_feeds_ranks_series_and_batches_equal_timestamps_stably() {
1097 let event_ts = ts(10, 0, 0);
1098 let primary = VecFeed::new(vec![MarketEvent::Bar {
1099 symbol: "PRIMARY".into(),
1100 ts: event_ts,
1101 open: 1.0,
1102 high: 1.2,
1103 low: 0.9,
1104 close: 1.1,
1105 volume: 10,
1106 }]);
1107 let conversion = VecFeed::from_feed_events(vec![FeedEvent::new(
1108 MarketEvent::Tick {
1109 symbol: "CONVERSION".into(),
1110 ts: event_ts,
1111 bid: 2.0,
1112 ask: 2.1,
1113 },
1114 EventMetadata::new(SeriesRoles::CONVERSION, 0, 0),
1115 )]);
1116
1117 let mut merged = merge_feeds(vec![primary, conversion]);
1118 let batch = merged.next_timestamp_batch().unwrap();
1119 assert_eq!(batch.len(), 2);
1120 assert_eq!(batch.events[0].metadata.series_rank, 0);
1121 assert_eq!(batch.events[0].metadata.roles, SeriesRoles::PRIMARY);
1122 assert_eq!(batch.events[1].metadata.series_rank, 1);
1123 assert_eq!(batch.events[1].metadata.roles, SeriesRoles::CONVERSION);
1124 }
1125
1126 fn ranked_tick(
1127 symbol: &str,
1128 event_ts: NaiveDateTime,
1129 roles: SeriesRoles,
1130 row_sequence: u64,
1131 ) -> FeedEvent {
1132 FeedEvent::new(
1133 MarketEvent::Tick {
1134 symbol: symbol.into(),
1135 ts: event_ts,
1136 bid: 1.0,
1137 ask: 1.1,
1138 },
1139 EventMetadata::new(roles, 99, row_sequence),
1140 )
1141 }
1142
1143 struct ScriptedFeed {
1144 batches: VecDeque<Result<TimestampBatch, &'static str>>,
1145 }
1146
1147 impl ScriptedFeed {
1148 fn new(batches: Vec<Result<TimestampBatch, &'static str>>) -> Self {
1149 Self {
1150 batches: batches.into(),
1151 }
1152 }
1153 }
1154
1155 impl FallibleBatchFeed for ScriptedFeed {
1156 type Error = &'static str;
1157
1158 fn next_batch(&mut self) -> Result<Option<TimestampBatch>, Self::Error> {
1159 self.batches.pop_front().transpose()
1160 }
1161 }
1162
1163 #[test]
1164 fn event_batch_feed_groups_complete_timestamps_and_supports_shared_roles() {
1165 let first_ts = ts(10, 0, 0);
1166 let second_ts = ts(10, 0, 1);
1167 let mut source = vec![
1168 MarketEvent::Tick {
1169 symbol: "EURUSD".into(),
1170 ts: first_ts,
1171 bid: 1.0,
1172 ask: 1.1,
1173 },
1174 MarketEvent::Bar {
1175 symbol: "EURUSD".into(),
1176 ts: first_ts,
1177 open: 1.0,
1178 high: 1.2,
1179 low: 0.9,
1180 close: 1.1,
1181 volume: 10,
1182 },
1183 MarketEvent::Tick {
1184 symbol: "EURUSD".into(),
1185 ts: second_ts,
1186 bid: 1.1,
1187 ask: 1.2,
1188 },
1189 ]
1190 .into_iter();
1191 let mut feed = EventBatchFeed::new(
1192 move || Ok::<_, &'static str>(source.next()),
1193 SeriesRoles::PRIMARY_AND_CONVERSION,
1194 4,
1195 );
1196
1197 let first = feed.next_timestamp_batch().unwrap().unwrap();
1198 assert_eq!(first.ts, first_ts);
1199 assert_eq!(first.len(), 2);
1200 assert!(
1201 first
1202 .events
1203 .iter()
1204 .all(|event| event.metadata.roles == SeriesRoles::PRIMARY_AND_CONVERSION)
1205 );
1206 assert_eq!(first.events[0].metadata.row_sequence, 0);
1207 assert_eq!(first.events[1].metadata.row_sequence, 1);
1208
1209 let second = feed.next_timestamp_batch().unwrap().unwrap();
1210 assert_eq!(second.ts, second_ts);
1211 assert_eq!(second.events[0].metadata.row_sequence, 2);
1212 assert!(feed.next_timestamp_batch().unwrap().is_none());
1213 }
1214
1215 #[test]
1216 fn event_batch_feed_preserves_explicit_physical_source_ordinals() {
1217 let mut source = vec![
1218 SequencedMarketEvent::new(
1219 MarketEvent::Tick {
1220 symbol: "EURUSD".into(),
1221 ts: ts(10, 0, 0),
1222 bid: 1.0,
1223 ask: 1.1,
1224 },
1225 4,
1226 ),
1227 SequencedMarketEvent::new(
1228 MarketEvent::Tick {
1229 symbol: "EURUSD".into(),
1230 ts: ts(10, 0, 0),
1231 bid: 1.1,
1232 ask: 1.2,
1233 },
1234 9,
1235 ),
1236 ]
1237 .into_iter();
1238 let mut feed = EventBatchFeed::new(
1239 move || Ok::<_, &'static str>(source.next()),
1240 SeriesRoles::PRIMARY,
1241 2,
1242 );
1243
1244 let batch = feed.next_timestamp_batch().unwrap().unwrap();
1245 assert_eq!(
1246 batch
1247 .events
1248 .iter()
1249 .map(|event| event.metadata.row_sequence)
1250 .collect::<Vec<_>>(),
1251 vec![4, 9]
1252 );
1253 assert_eq!(batch.events[0].metadata.series_rank, 2);
1254 }
1255
1256 #[test]
1257 fn event_batch_feed_rejects_a_non_monotonic_event_source() {
1258 let mut source = vec![
1259 MarketEvent::Tick {
1260 symbol: "EURUSD".into(),
1261 ts: ts(10, 0, 1),
1262 bid: 1.0,
1263 ask: 1.1,
1264 },
1265 MarketEvent::Tick {
1266 symbol: "EURUSD".into(),
1267 ts: ts(10, 0, 0),
1268 bid: 1.0,
1269 ask: 1.1,
1270 },
1271 ]
1272 .into_iter();
1273 let mut feed = EventBatchFeed::new(
1274 move || Ok::<_, &'static str>(source.next()),
1275 SeriesRoles::PRIMARY,
1276 0,
1277 );
1278
1279 assert!(matches!(
1280 feed.next_timestamp_batch(),
1281 Err(EventBatchFeedError::NonMonotonic { previous, current })
1282 if previous == ts(10, 0, 1) && current == ts(10, 0, 0)
1283 ));
1284 }
1285
1286 #[test]
1287 fn k_way_merge_drains_all_equal_timestamp_batches_and_orders_by_full_key() {
1288 let first_ts = ts(10, 0, 0);
1289 let second_ts = ts(10, 0, 1);
1290 let first = ScriptedFeed::new(vec![
1291 Ok(TimestampBatch {
1292 ts: first_ts,
1293 events: vec![ranked_tick(
1294 "SHARED",
1295 first_ts,
1296 SeriesRoles::PRIMARY_AND_CONVERSION,
1297 0,
1298 )],
1299 }),
1300 Ok(TimestampBatch {
1301 ts: first_ts,
1302 events: vec![ranked_tick("PRIMARY", first_ts, SeriesRoles::PRIMARY, 1)],
1303 }),
1304 Ok(TimestampBatch {
1305 ts: second_ts,
1306 events: vec![ranked_tick("PRIMARY", second_ts, SeriesRoles::PRIMARY, 2)],
1307 }),
1308 ]);
1309 let second = ScriptedFeed::new(vec![Ok(TimestampBatch {
1310 ts: first_ts,
1311 events: vec![ranked_tick(
1312 "CONVERSION",
1313 first_ts,
1314 SeriesRoles::CONVERSION,
1315 0,
1316 )],
1317 })]);
1318 let mut merged = KWayMergeFeed::new(vec![first, second]);
1319
1320 let batch = merged.next_timestamp_batch().unwrap().unwrap();
1321 assert_eq!(batch.ts, first_ts);
1322 assert_eq!(batch.len(), 3);
1323 assert_eq!(
1324 batch
1325 .events
1326 .iter()
1327 .map(FeedEvent::ordering_key)
1328 .collect::<Vec<_>>(),
1329 vec![(first_ts, 0, 0), (first_ts, 0, 1), (first_ts, 1, 0)]
1330 );
1331 assert_eq!(
1332 batch.events[0].metadata.roles,
1333 SeriesRoles::PRIMARY_AND_CONVERSION
1334 );
1335
1336 let later = merged.next_timestamp_batch().unwrap().unwrap();
1337 assert_eq!(later.ts, second_ts);
1338 assert_eq!(later.len(), 1);
1339 assert!(merged.next_timestamp_batch().unwrap().is_none());
1340 }
1341
1342 #[test]
1343 fn k_way_merge_propagates_errors_while_draining_a_complete_timestamp() {
1344 let event_ts = ts(10, 0, 0);
1345 let first = ScriptedFeed::new(vec![
1346 Ok(TimestampBatch {
1347 ts: event_ts,
1348 events: vec![ranked_tick("PRIMARY", event_ts, SeriesRoles::PRIMARY, 0)],
1349 }),
1350 Err("same timestamp continuation failed"),
1351 ]);
1352 let second = ScriptedFeed::new(vec![Ok(TimestampBatch {
1353 ts: event_ts,
1354 events: vec![ranked_tick(
1355 "CONVERSION",
1356 event_ts,
1357 SeriesRoles::CONVERSION,
1358 0,
1359 )],
1360 })]);
1361 let mut merged = KWayMergeFeed::new(vec![first, second]);
1362
1363 assert!(matches!(
1364 merged.next_timestamp_batch(),
1365 Err(KWayMergeError::Source {
1366 series_rank: 0,
1367 error: "same timestamp continuation failed"
1368 })
1369 ));
1370 }
1371
1372 #[test]
1373 fn k_way_merge_propagates_ranked_source_failures() {
1374 let mut merged = KWayMergeFeed::new(vec![ScriptedFeed::new(vec![Err("read failed")])]);
1375 assert!(matches!(
1376 merged.next_timestamp_batch(),
1377 Err(KWayMergeError::Source {
1378 series_rank: 0,
1379 error: "read failed"
1380 })
1381 ));
1382 }
1383}