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, 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
38/// Extracts marker sidecar fields from a concrete market-data message.
39///
40/// The bus tap passes messages as `&dyn Any`, so implementations downcast to the type they were
41/// registered for and return `None` when the registry is miswired.
42pub trait DataMarkerExtractor: Send + Sync {
43    /// Returns the data class handled by this extractor.
44    fn data_class(&self) -> DataClass;
45
46    /// Returns the stable stream identifier for the message.
47    fn identifier(&self, msg: &dyn Any) -> Option<String>;
48
49    /// Returns `(ts_event, ts_init)` for the message.
50    fn timestamps(&self, msg: &dyn Any) -> Option<(UnixNanos, UnixNanos)>;
51
52    /// Returns the class-specific canonical content fingerprint for the message.
53    fn fingerprint(&self, msg: &dyn Any) -> Option<[u8; 32]>;
54}
55
56/// Registry of marker extractors keyed by concrete message [`TypeId`].
57///
58/// Registration happens before the capture tap is installed. Lookups use the concrete type behind
59/// `&dyn Any`, so the hot path avoids trying every extractor.
60pub struct DataMarkerExtractorRegistry {
61    by_type: AHashMap<TypeId, Box<dyn DataMarkerExtractor>>,
62}
63
64impl DataMarkerExtractorRegistry {
65    /// Creates an empty extractor registry.
66    #[must_use]
67    pub fn new() -> Self {
68        Self {
69            by_type: AHashMap::new(),
70        }
71    }
72
73    /// Registers `ex` as the marker extractor for `T`.
74    ///
75    /// Replaces any previous extractor for `T`; callers should finish registration before sharing
76    /// the registry with the capture path.
77    pub fn register<T: 'static>(&mut self, ex: Box<dyn DataMarkerExtractor>) {
78        self.by_type.insert(TypeId::of::<T>(), ex);
79    }
80
81    /// Creates a registry containing builtin extractors for the enabled `classes`.
82    #[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    /// Returns the extractor registered for the concrete type behind `msg`.
110    #[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(&registry, &quote);
532
533        assert_eq!(extractor.data_class(), DataClass::Quote);
534        assert_eq!(
535            extractor.identifier(&quote),
536            Some("ETHUSDT.BINANCE".to_string())
537        );
538        assert_eq!(
539            extractor.timestamps(&quote),
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(&quote).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(&registry, &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(&registry, &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(&registry, &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(&registry, &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(&quote_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(&quote).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(&quote).is_none());
800        assert!(registry.lookup(&value).is_none());
801    }
802}