Skip to main content

nautilus_event_store/markers/
extractor.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Per-class market-data extractors for marker sidecar capture.
17
18use 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
39/// Extracts marker sidecar fields from a concrete market-data message.
40///
41/// The bus tap passes messages as `&dyn Any`, so implementations downcast to the type they were
42/// registered for and return `None` when the registry is miswired.
43pub trait DataMarkerExtractor: Send + Sync {
44    /// Returns the data class handled by this extractor.
45    fn data_class(&self) -> DataClass;
46
47    /// Returns the stable stream identifier for the message.
48    fn identifier(&self, msg: &dyn Any) -> Option<String>;
49
50    /// Returns `(ts_event, ts_init)` for the message.
51    fn timestamps(&self, msg: &dyn Any) -> Option<(UnixNanos, UnixNanos)>;
52
53    /// Returns the class-specific canonical content fingerprint for the message.
54    fn fingerprint(&self, msg: &dyn Any) -> Option<[u8; 32]>;
55}
56
57/// Registry of marker extractors keyed by concrete message [`TypeId`].
58///
59/// Registration happens before the capture tap is installed. Lookups use the concrete type behind
60/// `&dyn Any`, so the hot path avoids trying every extractor.
61pub struct DataMarkerExtractorRegistry {
62    by_type: AHashMap<TypeId, Box<dyn DataMarkerExtractor>>,
63}
64
65impl DataMarkerExtractorRegistry {
66    /// Creates an empty extractor registry.
67    #[must_use]
68    pub fn new() -> Self {
69        Self {
70            by_type: AHashMap::new(),
71        }
72    }
73
74    /// Registers `ex` as the marker extractor for `T`.
75    ///
76    /// Replaces any previous extractor for `T`; callers should finish registration before sharing
77    /// the registry with the capture path.
78    pub fn register<T: 'static>(&mut self, ex: Box<dyn DataMarkerExtractor>) {
79        self.by_type.insert(TypeId::of::<T>(), ex);
80    }
81
82    /// Creates a registry containing builtin extractors for the enabled `classes`.
83    #[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    /// Returns the extractor registered for the concrete type behind `msg`.
111    #[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(&registry, &quote);
568
569        assert_eq!(extractor.data_class(), DataClass::Quote);
570        assert_eq!(
571            extractor.identifier(&quote),
572            Some("ETHUSDT.BINANCE".to_string())
573        );
574        assert_eq!(
575            extractor.timestamps(&quote),
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(&quote).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(&registry, &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(&registry, &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(&registry, &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(&registry, &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(&quote_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(&quote).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(&quote).is_none());
836        assert!(registry.lookup(&value).is_none());
837    }
838}