1use std::{
19 any::{Any, TypeId},
20 fmt::Debug,
21};
22
23use ahash::AHashMap;
24use nautilus_core::UnixNanos;
25use nautilus_model::{
26 data::{Bar, BookOrder, DEPTH10_LEN, OrderBookDeltas, OrderBookDepth, QuoteTick, TradeTick},
27 types::{Price, Quantity, fixed::FIXED_PRECISION},
28};
29
30use crate::markers::DataClass;
31
32const QUOTE_FINGERPRINT_DOMAIN: &[u8] = b"nautilus-event-store/marker/fingerprint/quote/v1";
33const TRADE_FINGERPRINT_DOMAIN: &[u8] = b"nautilus-event-store/marker/fingerprint/trade/v1";
34const BAR_FINGERPRINT_DOMAIN: &[u8] = b"nautilus-event-store/marker/fingerprint/bar/v1";
35const DEPTH_FINGERPRINT_DOMAIN: &[u8] = b"nautilus-event-store/marker/fingerprint/depth/v1";
36const DEPTH10_FINGERPRINT_DOMAIN: &[u8] = b"nautilus-event-store/marker/fingerprint/depth10/v1";
37const DELTAS_FINGERPRINT_DOMAIN: &[u8] = b"nautilus-event-store/marker/fingerprint/deltas/v1";
38
39pub trait DataMarkerExtractor: Send + Sync {
44 fn data_class(&self) -> DataClass;
46
47 fn identifier(&self, msg: &dyn Any) -> Option<String>;
49
50 fn timestamps(&self, msg: &dyn Any) -> Option<(UnixNanos, UnixNanos)>;
52
53 fn fingerprint(&self, msg: &dyn Any) -> Option<[u8; 32]>;
55}
56
57pub struct DataMarkerExtractorRegistry {
62 by_type: AHashMap<TypeId, Box<dyn DataMarkerExtractor>>,
63}
64
65impl DataMarkerExtractorRegistry {
66 #[must_use]
68 pub fn new() -> Self {
69 Self {
70 by_type: AHashMap::new(),
71 }
72 }
73
74 pub fn register<T: 'static>(&mut self, ex: Box<dyn DataMarkerExtractor>) {
79 self.by_type.insert(TypeId::of::<T>(), ex);
80 }
81
82 #[must_use]
84 pub fn default_registry(classes: &[DataClass]) -> Self {
85 let mut registry = Self::new();
86
87 for class in classes {
88 match class {
89 DataClass::BookDeltas => {
90 registry.register::<OrderBookDeltas>(Box::new(OrderBookDeltasExtractor));
91 }
92 DataClass::BookDepth => {
93 registry.register::<OrderBookDepth>(Box::new(OrderBookDepthExtractor));
94 }
95 DataClass::Quote => {
96 registry.register::<QuoteTick>(Box::new(QuoteTickExtractor));
97 }
98 DataClass::Trade => {
99 registry.register::<TradeTick>(Box::new(TradeTickExtractor));
100 }
101 DataClass::Bar => {
102 registry.register::<Bar>(Box::new(BarExtractor));
103 }
104 }
105 }
106
107 registry
108 }
109
110 #[must_use]
112 pub fn lookup(&self, msg: &dyn Any) -> Option<&dyn DataMarkerExtractor> {
113 self.by_type.get(&msg.type_id()).map(Box::as_ref)
114 }
115}
116
117impl Default for DataMarkerExtractorRegistry {
118 fn default() -> Self {
119 Self::new()
120 }
121}
122
123impl Debug for DataMarkerExtractorRegistry {
124 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
125 f.debug_struct(stringify!(DataMarkerExtractorRegistry))
126 .field("len", &self.by_type.len())
127 .finish()
128 }
129}
130
131#[derive(Debug)]
132struct QuoteTickExtractor;
133
134impl DataMarkerExtractor for QuoteTickExtractor {
135 fn data_class(&self) -> DataClass {
136 DataClass::Quote
137 }
138
139 fn identifier(&self, msg: &dyn Any) -> Option<String> {
140 msg.downcast_ref::<QuoteTick>()
141 .map(|quote| quote.instrument_id.to_string())
142 }
143
144 fn timestamps(&self, msg: &dyn Any) -> Option<(UnixNanos, UnixNanos)> {
145 msg.downcast_ref::<QuoteTick>()
146 .map(|quote| (quote.ts_event, quote.ts_init))
147 }
148
149 fn fingerprint(&self, msg: &dyn Any) -> Option<[u8; 32]> {
150 msg.downcast_ref::<QuoteTick>().map(fingerprint_quote)
151 }
152}
153
154#[derive(Debug)]
155struct TradeTickExtractor;
156
157impl DataMarkerExtractor for TradeTickExtractor {
158 fn data_class(&self) -> DataClass {
159 DataClass::Trade
160 }
161
162 fn identifier(&self, msg: &dyn Any) -> Option<String> {
163 msg.downcast_ref::<TradeTick>()
164 .map(|trade| trade.instrument_id.to_string())
165 }
166
167 fn timestamps(&self, msg: &dyn Any) -> Option<(UnixNanos, UnixNanos)> {
168 msg.downcast_ref::<TradeTick>()
169 .map(|trade| (trade.ts_event, trade.ts_init))
170 }
171
172 fn fingerprint(&self, msg: &dyn Any) -> Option<[u8; 32]> {
173 msg.downcast_ref::<TradeTick>().map(fingerprint_trade)
174 }
175}
176
177#[derive(Debug)]
178struct BarExtractor;
179
180impl DataMarkerExtractor for BarExtractor {
181 fn data_class(&self) -> DataClass {
182 DataClass::Bar
183 }
184
185 fn identifier(&self, msg: &dyn Any) -> Option<String> {
186 msg.downcast_ref::<Bar>()
187 .map(|bar| bar.bar_type.to_string())
188 }
189
190 fn timestamps(&self, msg: &dyn Any) -> Option<(UnixNanos, UnixNanos)> {
191 msg.downcast_ref::<Bar>()
192 .map(|bar| (bar.ts_event, bar.ts_init))
193 }
194
195 fn fingerprint(&self, msg: &dyn Any) -> Option<[u8; 32]> {
196 msg.downcast_ref::<Bar>().map(fingerprint_bar)
197 }
198}
199
200#[derive(Debug)]
201struct OrderBookDepthExtractor;
202
203impl DataMarkerExtractor for OrderBookDepthExtractor {
204 fn data_class(&self) -> DataClass {
205 DataClass::BookDepth
206 }
207
208 fn identifier(&self, msg: &dyn Any) -> Option<String> {
209 msg.downcast_ref::<OrderBookDepth>()
210 .map(|depth| depth.instrument_id.to_string())
211 }
212
213 fn timestamps(&self, msg: &dyn Any) -> Option<(UnixNanos, UnixNanos)> {
214 msg.downcast_ref::<OrderBookDepth>()
215 .map(|depth| (depth.ts_event, depth.ts_init))
216 }
217
218 fn fingerprint(&self, msg: &dyn Any) -> Option<[u8; 32]> {
219 msg.downcast_ref::<OrderBookDepth>().map(fingerprint_depth)
220 }
221}
222
223#[derive(Debug)]
224struct OrderBookDeltasExtractor;
225
226impl DataMarkerExtractor for OrderBookDeltasExtractor {
227 fn data_class(&self) -> DataClass {
228 DataClass::BookDeltas
229 }
230
231 fn identifier(&self, msg: &dyn Any) -> Option<String> {
232 msg.downcast_ref::<OrderBookDeltas>()
233 .map(|deltas| deltas.instrument_id.to_string())
234 }
235
236 fn timestamps(&self, msg: &dyn Any) -> Option<(UnixNanos, UnixNanos)> {
237 msg.downcast_ref::<OrderBookDeltas>()
238 .map(|deltas| (deltas.ts_event, deltas.ts_init))
239 }
240
241 fn fingerprint(&self, msg: &dyn Any) -> Option<[u8; 32]> {
242 msg.downcast_ref::<OrderBookDeltas>()
243 .map(fingerprint_deltas)
244 }
245}
246
247fn fingerprint_quote(quote: &QuoteTick) -> [u8; 32] {
248 let mut hasher = blake3::Hasher::new();
249 hasher.update(QUOTE_FINGERPRINT_DOMAIN);
250 write_price_raw(&mut hasher, quote.bid_price);
251 write_price_raw(&mut hasher, quote.ask_price);
252 write_quantity_raw(&mut hasher, quote.bid_size);
253 write_quantity_raw(&mut hasher, quote.ask_size);
254 write_unix_nanos(&mut hasher, quote.ts_event);
255 *hasher.finalize().as_bytes()
256}
257
258fn fingerprint_trade(trade: &TradeTick) -> [u8; 32] {
259 let mut hasher = blake3::Hasher::new();
260 hasher.update(TRADE_FINGERPRINT_DOMAIN);
261 write_price_raw(&mut hasher, trade.price);
262 write_quantity_raw(&mut hasher, trade.size);
263 hasher.update(&[trade.aggressor_side as u8]);
264 write_str(&mut hasher, trade.trade_id.as_str());
265 write_unix_nanos(&mut hasher, trade.ts_event);
266 *hasher.finalize().as_bytes()
267}
268
269fn fingerprint_bar(bar: &Bar) -> [u8; 32] {
270 let mut hasher = blake3::Hasher::new();
271 hasher.update(BAR_FINGERPRINT_DOMAIN);
272 write_str(&mut hasher, &bar.bar_type.to_string());
273 write_price_raw(&mut hasher, bar.open);
274 write_price_raw(&mut hasher, bar.high);
275 write_price_raw(&mut hasher, bar.low);
276 write_price_raw(&mut hasher, bar.close);
277 write_quantity_raw(&mut hasher, bar.volume);
278 write_unix_nanos(&mut hasher, bar.ts_event);
279 *hasher.finalize().as_bytes()
280}
281
282fn fingerprint_depth(depth: &OrderBookDepth) -> [u8; 32] {
283 let mut hasher = blake3::Hasher::new();
284 if depth.bids.len() == DEPTH10_LEN && depth.asks.len() == DEPTH10_LEN {
285 hasher.update(DEPTH10_FINGERPRINT_DOMAIN);
286 } else {
287 hasher.update(DEPTH_FINGERPRINT_DOMAIN);
288 hasher.update(&(depth.bids.len() as u64).to_be_bytes());
289 hasher.update(&(depth.asks.len() as u64).to_be_bytes());
290 }
291
292 for (order, count) in depth.bids.iter().zip(depth.bid_counts.iter().copied()) {
293 write_depth_level(&mut hasher, order, count);
294 }
295
296 for (order, count) in depth.asks.iter().zip(depth.ask_counts.iter().copied()) {
297 write_depth_level(&mut hasher, order, count);
298 }
299 write_unix_nanos(&mut hasher, depth.ts_event);
300 *hasher.finalize().as_bytes()
301}
302
303fn fingerprint_deltas(deltas: &OrderBookDeltas) -> [u8; 32] {
304 let mut hasher = blake3::Hasher::new();
305 hasher.update(DELTAS_FINGERPRINT_DOMAIN);
306 hasher.update(&(deltas.deltas.len() as u64).to_be_bytes());
307 for delta in &deltas.deltas {
308 hasher.update(&[delta.action as u8]);
309 hasher.update(&[delta.order.side.map_or(0, |side| side as u8)]);
310 write_price_raw(&mut hasher, delta.order.price);
311 write_quantity_raw(&mut hasher, delta.order.size);
312 hasher.update(&delta.order.order_id.to_be_bytes());
313 hasher.update(&[delta.flags]);
314 }
315 write_unix_nanos(&mut hasher, deltas.ts_event);
316 *hasher.finalize().as_bytes()
317}
318
319fn write_depth_level(hasher: &mut blake3::Hasher, order: &BookOrder, count: u32) {
320 write_price_raw(hasher, order.price);
321 write_quantity_raw(hasher, order.size);
322 hasher.update(&count.to_be_bytes());
323}
324
325fn write_price_raw(hasher: &mut blake3::Hasher, price: Price) {
326 hasher.update(&[price.precision]);
327 hasher.update(&price_raw_at_precision(price).to_be_bytes());
328}
329
330fn write_quantity_raw(hasher: &mut blake3::Hasher, quantity: Quantity) {
331 hasher.update(&[quantity.precision]);
332 hasher.update(&quantity_raw_at_precision(quantity).to_be_bytes());
333}
334
335#[allow(
336 clippy::useless_conversion,
337 reason = "PriceRaw is i64 or i128 depending on feature unification; the conversion is only useless in high-precision builds"
338)]
339fn price_raw_at_precision(price: Price) -> i128 {
340 let scale_down = FIXED_PRECISION.saturating_sub(price.precision);
341 #[cfg(feature = "defi")]
342 let raw = price.raw();
343 #[cfg(not(feature = "defi"))]
344 let raw = i128::from(price.raw());
345
346 raw / 10_i128.pow(u32::from(scale_down))
347}
348
349#[allow(
350 clippy::useless_conversion,
351 reason = "QuantityRaw is u64 or u128 depending on feature unification; the conversion is only useless in high-precision builds"
352)]
353fn quantity_raw_at_precision(quantity: Quantity) -> u128 {
354 let scale_down = FIXED_PRECISION.saturating_sub(quantity.precision);
355 #[cfg(feature = "defi")]
356 let raw = quantity.raw();
357 #[cfg(not(feature = "defi"))]
358 let raw = u128::from(quantity.raw());
359
360 raw / 10_u128.pow(u32::from(scale_down))
361}
362
363fn write_unix_nanos(hasher: &mut blake3::Hasher, value: UnixNanos) {
364 hasher.update(&value.as_u64().to_be_bytes());
365}
366
367fn write_str(hasher: &mut blake3::Hasher, value: &str) {
368 let bytes = value.as_bytes();
369 hasher.update(&(bytes.len() as u64).to_be_bytes());
370 hasher.update(bytes);
371}
372
373#[cfg(test)]
374mod tests {
375 use std::{any::Any, fmt::Write};
376
377 use nautilus_core::UnixNanos;
378 use nautilus_model::{
379 data::{
380 Bar, BarType, BookOrder, OrderBookDelta, OrderBookDeltas, OrderBookDepth, QuoteTick,
381 TradeTick, depth::DEPTH10_LEN,
382 },
383 enums::{AggressorSide, BookAction, OrderSide},
384 identifiers::{InstrumentId, TradeId},
385 types::{Price, Quantity, price::PriceRaw, quantity::QuantityRaw},
386 };
387 use rstest::rstest;
388
389 use super::*;
390 use crate::markers::DataClass;
391
392 fn hex32(bytes: &[u8; 32]) -> String {
393 let mut out = String::with_capacity(64);
394 for byte in bytes {
395 write!(out, "{byte:02x}").expect("writing to a String is infallible");
396 }
397 out
398 }
399
400 fn quote_tick() -> QuoteTick {
401 QuoteTick::new(
402 InstrumentId::from("ETHUSDT.BINANCE"),
403 Price::from("3000.12"),
404 Price::from("3000.25"),
405 Quantity::from("1.25"),
406 Quantity::from("2.50"),
407 UnixNanos::from(1_700_000_000_000_000_100),
408 UnixNanos::from(1_700_000_000_000_000_200),
409 )
410 }
411
412 fn trade_tick() -> TradeTick {
413 TradeTick::new(
414 InstrumentId::from("ETHUSDT.BINANCE"),
415 Price::from("3000.18"),
416 Quantity::from("0.75"),
417 AggressorSide::Buy,
418 TradeId::new("T-ABC-123"),
419 UnixNanos::from(1_700_000_000_000_000_300),
420 UnixNanos::from(1_700_000_000_000_000_400),
421 )
422 }
423
424 fn bar() -> Bar {
425 Bar::new(
426 BarType::from("ETHUSDT.BINANCE-1-MINUTE-LAST-EXTERNAL"),
427 Price::from("3000.00"),
428 Price::from("3010.50"),
429 Price::from("2995.25"),
430 Price::from("3005.75"),
431 Quantity::from("42.25"),
432 UnixNanos::from(1_700_000_000_000_000_500),
433 UnixNanos::from(1_700_000_000_000_000_600),
434 )
435 }
436
437 fn price_from_cents(cents: i64) -> Price {
438 let scale_down = FIXED_PRECISION.saturating_sub(2);
439 let scale = PriceRaw::from(10_i64.pow(u32::from(scale_down)));
440 Price::from_raw(PriceRaw::from(cents) * scale, 2)
441 }
442
443 fn quantity_from_cents(cents: u64) -> Quantity {
444 let scale_down = FIXED_PRECISION.saturating_sub(2);
445 let scale = QuantityRaw::from(10_u64.pow(u32::from(scale_down)));
446 Quantity::from_raw(QuantityRaw::from(cents) * scale, 2)
447 }
448
449 #[rstest]
450 fn variable_depth_fingerprint_distinguishes_side_boundaries() {
451 let bid = BookOrder::new(OrderSide::Buy, Price::from("1.00"), Quantity::from("2"), 0);
452 let ask = BookOrder::new(OrderSide::Sell, bid.price, bid.size, 0);
453 let left = OrderBookDepth::new(
454 InstrumentId::from("ETHUSDT.BINANCE"),
455 vec![bid],
456 vec![ask; 2],
457 vec![1],
458 vec![1; 2],
459 0,
460 0,
461 UnixNanos::from(1),
462 UnixNanos::from(2),
463 );
464 let right = OrderBookDepth::new(
465 left.instrument_id,
466 vec![bid; 2],
467 vec![ask],
468 vec![1; 2],
469 vec![1],
470 left.flags,
471 left.sequence,
472 left.ts_event,
473 left.ts_init,
474 );
475 assert_ne!(fingerprint_depth(&left), fingerprint_depth(&right));
476 }
477
478 fn depth() -> OrderBookDepth {
479 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
480 let bids: [BookOrder; 10] = std::array::from_fn(|i| {
481 let level = i64::try_from(i).expect("depth index fits i64");
482 let order_offset = u64::try_from(i).expect("depth index fits u64");
483
484 BookOrder::new(
485 OrderSide::Buy,
486 price_from_cents(300_000 - level),
487 quantity_from_cents(10_000 + order_offset),
488 1_000 + order_offset,
489 )
490 });
491 let asks: [BookOrder; 10] = std::array::from_fn(|i| {
492 let level = i64::try_from(i).expect("depth index fits i64");
493 let order_offset = u64::try_from(i).expect("depth index fits u64");
494
495 BookOrder::new(
496 OrderSide::Sell,
497 price_from_cents(300_100 + level),
498 quantity_from_cents(20_000 + order_offset),
499 2_000 + order_offset,
500 )
501 });
502 let bid_counts: [u32; 10] =
503 std::array::from_fn(|i| 10 + u32::try_from(i).expect("depth index fits u32"));
504 let ask_counts: [u32; 10] =
505 std::array::from_fn(|i| 20 + u32::try_from(i).expect("depth index fits u32"));
506
507 OrderBookDepth::new(
508 instrument_id,
509 bids,
510 asks,
511 bid_counts,
512 ask_counts,
513 0x20,
514 42,
515 UnixNanos::from(1_700_000_000_000_000_700),
516 UnixNanos::from(1_700_000_000_000_000_800),
517 )
518 }
519
520 fn deltas() -> OrderBookDeltas {
521 let instrument_id = InstrumentId::from("ETHUSDT.BINANCE");
522 let first = OrderBookDelta::new(
523 instrument_id,
524 BookAction::Add,
525 BookOrder::new(
526 OrderSide::Buy,
527 Price::from("3000.00"),
528 Quantity::from("1.10"),
529 10,
530 ),
531 0x01,
532 41,
533 UnixNanos::from(1_700_000_000_000_000_900),
534 UnixNanos::from(1_700_000_000_000_001_000),
535 );
536 let second = OrderBookDelta::new(
537 instrument_id,
538 BookAction::Update,
539 BookOrder::new(
540 OrderSide::Sell,
541 Price::from("3001.00"),
542 Quantity::from("2.20"),
543 11,
544 ),
545 0x20,
546 42,
547 UnixNanos::from(1_700_000_000_000_000_900),
548 UnixNanos::from(1_700_000_000_000_001_000),
549 );
550
551 OrderBookDeltas::new(instrument_id, vec![first, second])
552 }
553
554 fn extractor_for<'a>(
555 registry: &'a DataMarkerExtractorRegistry,
556 msg: &dyn Any,
557 ) -> &'a dyn DataMarkerExtractor {
558 registry
559 .lookup(msg)
560 .expect("registered data marker extractor")
561 }
562
563 #[rstest]
564 fn quote_extractor_fields_and_fingerprint() {
565 let quote = quote_tick();
566 let registry = DataMarkerExtractorRegistry::default_registry(&[DataClass::Quote]);
567 let extractor = extractor_for(®istry, "e);
568
569 assert_eq!(extractor.data_class(), DataClass::Quote);
570 assert_eq!(
571 extractor.identifier("e),
572 Some("ETHUSDT.BINANCE".to_string())
573 );
574 assert_eq!(
575 extractor.timestamps("e),
576 Some((
577 UnixNanos::from(1_700_000_000_000_000_100),
578 UnixNanos::from(1_700_000_000_000_000_200),
579 ))
580 );
581 assert_eq!(
582 hex32(&extractor.fingerprint("e).expect("fingerprint")),
583 "7c6671e34f01b7b547ac8695c6d2cd19c1a37f6d2e3910d9195ed66fd4c02628"
584 );
585 }
586
587 #[rstest]
588 fn raw_writers_use_declared_precision_scale() {
589 assert_eq!(price_raw_at_precision(Price::from("3000.12")), 300_012);
590 assert_eq!(quantity_raw_at_precision(Quantity::from("1.25")), 125);
591 }
592
593 #[rstest]
594 fn trade_extractor_fields_and_fingerprint() {
595 let trade = trade_tick();
596 let registry = DataMarkerExtractorRegistry::default_registry(&[DataClass::Trade]);
597 let extractor = extractor_for(®istry, &trade);
598
599 assert_eq!(extractor.data_class(), DataClass::Trade);
600 assert_eq!(
601 extractor.identifier(&trade),
602 Some("ETHUSDT.BINANCE".to_string())
603 );
604 assert_eq!(
605 extractor.timestamps(&trade),
606 Some((
607 UnixNanos::from(1_700_000_000_000_000_300),
608 UnixNanos::from(1_700_000_000_000_000_400),
609 ))
610 );
611 assert_eq!(
612 hex32(&extractor.fingerprint(&trade).expect("fingerprint")),
613 "6b32a3187d353451a07d92a0d91051406ce4fe912010202b81616ed315f565cb"
614 );
615 }
616
617 #[rstest]
618 fn bar_extractor_fields_and_fingerprint() {
619 let bar = bar();
620 let registry = DataMarkerExtractorRegistry::default_registry(&[DataClass::Bar]);
621 let extractor = extractor_for(®istry, &bar);
622
623 assert_eq!(extractor.data_class(), DataClass::Bar);
624 assert_eq!(
625 extractor.identifier(&bar),
626 Some("ETHUSDT.BINANCE-1-MINUTE-LAST-EXTERNAL".to_string())
627 );
628 assert_eq!(
629 extractor.timestamps(&bar),
630 Some((
631 UnixNanos::from(1_700_000_000_000_000_500),
632 UnixNanos::from(1_700_000_000_000_000_600),
633 ))
634 );
635 assert_eq!(
636 hex32(&extractor.fingerprint(&bar).expect("fingerprint")),
637 "f2283ae7ed8d2e6a3874473b11935557b6bc2cbf20446419fbdff8fe51f91e84"
638 );
639 }
640
641 #[rstest]
642 fn depth_extractor_fields_and_fingerprint() {
643 let depth = depth();
644 let registry = DataMarkerExtractorRegistry::default_registry(&[DataClass::BookDepth]);
645 let extractor = extractor_for(®istry, &depth);
646
647 assert_eq!(extractor.data_class(), DataClass::BookDepth);
648 assert_eq!(
649 extractor.identifier(&depth),
650 Some("ETHUSDT.BINANCE".to_string())
651 );
652 assert_eq!(depth.bids.len(), DEPTH10_LEN);
653 assert_eq!(depth.asks.len(), DEPTH10_LEN);
654 assert_eq!(
655 extractor.timestamps(&depth),
656 Some((
657 UnixNanos::from(1_700_000_000_000_000_700),
658 UnixNanos::from(1_700_000_000_000_000_800),
659 ))
660 );
661 assert_eq!(
662 hex32(&extractor.fingerprint(&depth).expect("fingerprint")),
663 "432e1f3951c660acaefb3a44c429b050e5e8d5b6f9c3b89cc969edaad85e8e4d"
664 );
665 }
666
667 #[rstest]
668 fn deltas_extractor_fields_and_fingerprint() {
669 let deltas = deltas();
670 let registry = DataMarkerExtractorRegistry::default_registry(&[DataClass::BookDeltas]);
671 let extractor = extractor_for(®istry, &deltas);
672
673 assert_eq!(extractor.data_class(), DataClass::BookDeltas);
674 assert_eq!(
675 extractor.identifier(&deltas),
676 Some("ETHUSDT.BINANCE".to_string())
677 );
678 assert_eq!(deltas.deltas.len(), 2);
679 assert_eq!(
680 extractor.timestamps(&deltas),
681 Some((
682 UnixNanos::from(1_700_000_000_000_000_900),
683 UnixNanos::from(1_700_000_000_000_001_000),
684 ))
685 );
686 assert_eq!(
687 hex32(&extractor.fingerprint(&deltas).expect("fingerprint")),
688 "4bb7fedc4454d53e08300e1ae7a59e648747152134a921050a2c80edec6d8f9e"
689 );
690 }
691
692 #[rstest]
693 #[case::bid_price(|q: &mut QuoteTick| q.bid_price = Price::from("3000.13"))]
694 #[case::ask_price(|q: &mut QuoteTick| q.ask_price = Price::from("3000.26"))]
695 #[case::bid_size(|q: &mut QuoteTick| q.bid_size = Quantity::from("1.26"))]
696 #[case::ask_size(|q: &mut QuoteTick| q.ask_size = Quantity::from("2.51"))]
697 #[case::ts_event(|q: &mut QuoteTick| q.ts_event = UnixNanos::from(1))]
698 fn quote_fingerprint_changes_when_hashed_field_changes(#[case] mutate: fn(&mut QuoteTick)) {
699 let base = quote_tick();
700 let mut changed = base;
701 mutate(&mut changed);
702
703 assert_ne!(fingerprint_quote(&base), fingerprint_quote(&changed));
704 }
705
706 #[rstest]
707 #[case::price(|t: &mut TradeTick| t.price = Price::from("3000.19"))]
708 #[case::size(|t: &mut TradeTick| t.size = Quantity::from("0.76"))]
709 #[case::aggressor_side(|t: &mut TradeTick| t.aggressor_side = AggressorSide::Sell)]
710 #[case::trade_id(|t: &mut TradeTick| t.trade_id = TradeId::new("T-ABC-124"))]
711 #[case::ts_event(|t: &mut TradeTick| t.ts_event = UnixNanos::from(1))]
712 fn trade_fingerprint_changes_when_hashed_field_changes(#[case] mutate: fn(&mut TradeTick)) {
713 let base = trade_tick();
714 let mut changed = base;
715 mutate(&mut changed);
716
717 assert_ne!(fingerprint_trade(&base), fingerprint_trade(&changed));
718 }
719
720 #[rstest]
721 #[case::bar_type(|b: &mut Bar| b.bar_type = BarType::from("ETHUSDT.BINANCE-5-MINUTE-LAST-EXTERNAL"))]
722 #[case::open(|b: &mut Bar| b.open = Price::from("3000.01"))]
723 #[case::high(|b: &mut Bar| b.high = Price::from("3010.51"))]
724 #[case::low(|b: &mut Bar| b.low = Price::from("2995.26"))]
725 #[case::close(|b: &mut Bar| b.close = Price::from("3005.76"))]
726 #[case::volume(|b: &mut Bar| b.volume = Quantity::from("42.26"))]
727 #[case::ts_event(|b: &mut Bar| b.ts_event = UnixNanos::from(1))]
728 fn bar_fingerprint_changes_when_hashed_field_changes(#[case] mutate: fn(&mut Bar)) {
729 let base = bar();
730 let mut changed = base;
731 mutate(&mut changed);
732
733 assert_ne!(fingerprint_bar(&base), fingerprint_bar(&changed));
734 }
735
736 #[rstest]
737 #[case::bid_price(|d: &mut OrderBookDepth| d.bids[0].price = price_from_cents(300_001))]
738 #[case::bid_size(|d: &mut OrderBookDepth| d.bids[0].size = quantity_from_cents(10_001))]
739 #[case::bid_count(|d: &mut OrderBookDepth| d.bid_counts[0] = 99)]
740 #[case::ask_price(|d: &mut OrderBookDepth| d.asks[0].price = price_from_cents(300_101))]
741 #[case::ask_size(|d: &mut OrderBookDepth| d.asks[0].size = quantity_from_cents(20_001))]
742 #[case::ask_count(|d: &mut OrderBookDepth| d.ask_counts[0] = 99)]
743 #[case::ts_event(|d: &mut OrderBookDepth| d.ts_event = UnixNanos::from(1))]
744 fn depth_fingerprint_changes_when_hashed_field_changes(
745 #[case] mutate: fn(&mut OrderBookDepth),
746 ) {
747 let base = depth();
748 let mut changed = base.clone();
749 mutate(&mut changed);
750
751 assert_ne!(fingerprint_depth(&base), fingerprint_depth(&changed));
752 }
753
754 #[rstest]
755 #[case::delta_count(|d: &mut OrderBookDeltas| {
756 let instrument_id = d.instrument_id;
757 d.deltas.push(OrderBookDelta::new(
758 instrument_id,
759 BookAction::Delete,
760 BookOrder::new(
761 OrderSide::Buy,
762 Price::from("2999.00"),
763 Quantity::from("3.30"),
764 12,
765 ),
766 0x20,
767 43,
768 UnixNanos::from(1_700_000_000_000_000_900),
769 UnixNanos::from(1_700_000_000_000_001_000),
770 ));
771 })]
772 #[case::action(|d: &mut OrderBookDeltas| d.deltas[0].action = BookAction::Delete)]
773 #[case::side(|d: &mut OrderBookDeltas| d.deltas[0].order.side = OrderSide::Sell.into())]
774 #[case::price(|d: &mut OrderBookDeltas| d.deltas[0].order.price = Price::from("3000.01"))]
775 #[case::size(|d: &mut OrderBookDeltas| d.deltas[0].order.size = Quantity::from("1.11"))]
776 #[case::order_id(|d: &mut OrderBookDeltas| d.deltas[0].order.order_id = 99)]
777 #[case::flags(|d: &mut OrderBookDeltas| d.deltas[0].flags = 0x02)]
778 #[case::order(|d: &mut OrderBookDeltas| d.deltas.reverse())]
779 #[case::ts_event(|d: &mut OrderBookDeltas| d.ts_event = UnixNanos::from(1))]
780 fn deltas_fingerprint_changes_when_hashed_field_changes(
781 #[case] mutate: fn(&mut OrderBookDeltas),
782 ) {
783 let base = deltas();
784 let mut changed = base.clone();
785 mutate(&mut changed);
786
787 assert_ne!(fingerprint_deltas(&base), fingerprint_deltas(&changed));
788 }
789
790 #[rstest]
791 fn fingerprint_domains_are_class_specific() {
792 let fingerprints = [
793 (DataClass::Quote, fingerprint_quote("e_tick())),
794 (DataClass::Trade, fingerprint_trade(&trade_tick())),
795 (DataClass::Bar, fingerprint_bar(&bar())),
796 (DataClass::BookDepth, fingerprint_depth(&depth())),
797 (DataClass::BookDeltas, fingerprint_deltas(&deltas())),
798 ];
799
800 for (index, (left_class, left_fingerprint)) in fingerprints.iter().enumerate() {
801 for (right_class, right_fingerprint) in fingerprints.iter().skip(index + 1) {
802 assert_ne!(
803 left_fingerprint, right_fingerprint,
804 "{left_class:?} and {right_class:?} fingerprints should differ"
805 );
806 }
807 }
808 }
809
810 #[rstest]
811 fn default_registry_installs_only_enabled_builtins() {
812 let quote = quote_tick();
813 let trade = trade_tick();
814 let bar = bar();
815 let depth = depth();
816 let deltas = deltas();
817 let registry = DataMarkerExtractorRegistry::default_registry(&[
818 DataClass::Quote,
819 DataClass::BookDepth,
820 ]);
821
822 assert!(registry.lookup("e).is_some());
823 assert!(registry.lookup(&depth).is_some());
824 assert!(registry.lookup(&trade).is_none());
825 assert!(registry.lookup(&bar).is_none());
826 assert!(registry.lookup(&deltas).is_none());
827 }
828
829 #[rstest]
830 fn registry_returns_none_for_unregistered_type() {
831 let registry = DataMarkerExtractorRegistry::new();
832 let quote = quote_tick();
833 let value = 1_u8;
834
835 assert!(registry.lookup("e).is_none());
836 assert!(registry.lookup(&value).is_none());
837 }
838}