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