Skip to main content

nautilus_testkit/itch/
parse.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//! ITCH 5.0 message to [`OrderBookDelta`] conversion.
17
18use std::{io::Read, path::Path};
19
20use ahash::AHashMap;
21use nautilus_core::{
22    UnixNanos,
23    time::{AtomicTime, get_atomic_clock_realtime},
24};
25use nautilus_model::{
26    data::{delta::OrderBookDelta, order::BookOrder},
27    enums::{BookAction, OrderSide, RecordFlag},
28    identifiers::InstrumentId,
29    types::{Price, Quantity},
30};
31
32/// Price precision for US equities (4 decimal places, $0.0001 increments).
33const PRICE_PRECISION: u8 = 4;
34
35/// Size precision for US equities (whole shares).
36const SIZE_PRECISION: u8 = 0;
37
38#[derive(Debug)]
39struct OrderState {
40    price: Price,
41    size: u32,
42    side: OrderSide,
43}
44
45/// Converts a stream of ITCH 5.0 messages into [`OrderBookDelta`] events
46/// for a single instrument.
47///
48/// Maintains internal order state to compute remaining sizes after partial
49/// executions and cancellations. Each delta carries the ITCH message time as
50/// `ts_event` and the conversion wall-clock time as `ts_init`.
51#[derive(Debug)]
52pub struct ItchParser {
53    clock: &'static AtomicTime,
54    instrument_id: InstrumentId,
55    target_locate: Option<u16>,
56    target_stock: String,
57    base_ns: u64,
58    orders: AHashMap<u64, OrderState>,
59    sequence: u64,
60}
61
62impl ItchParser {
63    /// Creates a new [`ItchParser`] for the given instrument.
64    ///
65    /// # Arguments
66    ///
67    /// - `instrument_id` - The NautilusTrader instrument ID for output deltas.
68    /// - `stock` - The ITCH stock symbol to filter for (e.g., "AAPL").
69    /// - `base_ns` - Base UNIX nanoseconds for midnight of the trading day
70    ///   (ITCH timestamps are nanoseconds since midnight).
71    #[must_use]
72    pub fn new(instrument_id: InstrumentId, stock: &str, base_ns: u64) -> Self {
73        Self {
74            clock: get_atomic_clock_realtime(),
75            instrument_id,
76            target_locate: None,
77            target_stock: stock.to_string(),
78            base_ns,
79            orders: AHashMap::new(),
80            sequence: 0,
81        }
82    }
83
84    /// Parses all ITCH messages from a gzip-compressed file and returns
85    /// the filtered [`OrderBookDelta`] events.
86    ///
87    /// # Errors
88    ///
89    /// Returns an error if the file cannot be opened or contains invalid data.
90    pub fn parse_gzip_file(&mut self, path: &Path) -> anyhow::Result<Vec<OrderBookDelta>> {
91        let stream = itchy::MessageStream::from_gzip(path)
92            .map_err(|e| anyhow::anyhow!("Failed to open ITCH gzip: {e}"))?;
93        self.parse_stream(stream)
94    }
95
96    /// Parses all ITCH messages from a reader and returns filtered deltas.
97    ///
98    /// # Errors
99    ///
100    /// Returns an error if the stream contains invalid data.
101    pub fn parse_reader<R: Read>(&mut self, reader: R) -> anyhow::Result<Vec<OrderBookDelta>> {
102        let stream = itchy::MessageStream::from_reader(reader);
103        self.parse_stream(stream)
104    }
105
106    fn parse_stream<R: Read>(
107        &mut self,
108        stream: itchy::MessageStream<R>,
109    ) -> anyhow::Result<Vec<OrderBookDelta>> {
110        let mut deltas = Vec::new();
111
112        for result in stream {
113            let msg = result.map_err(|e| anyhow::anyhow!("ITCH parse error: {e}"))?;
114
115            // Handle feed-level messages before stock locate filtering
116            match msg.body {
117                itchy::Body::StockDirectory(ref dir) => {
118                    let symbol = dir.stock.trim();
119                    if symbol == self.target_stock {
120                        self.target_locate = Some(msg.stock_locate);
121                    }
122                    continue;
123                }
124                itchy::Body::SystemEvent {
125                    event: itchy::EventCode::EndOfMessages,
126                } => {
127                    let ts_event = UnixNanos::from(self.base_ns + msg.timestamp);
128                    let ts_init = self.clock.get_time_ns();
129                    self.handle_end_of_messages(ts_event, ts_init, &mut deltas);
130                    continue;
131                }
132                _ => {}
133            }
134
135            // Filter by target stock locate
136            let Some(locate) = self.target_locate else {
137                continue;
138            };
139
140            if msg.stock_locate != locate {
141                continue;
142            }
143
144            let ts_event = UnixNanos::from(self.base_ns + msg.timestamp);
145            let ts_init = self.clock.get_time_ns();
146
147            match msg.body {
148                itchy::Body::AddOrder(ref add) => {
149                    self.handle_add_order(add, ts_event, ts_init, &mut deltas);
150                }
151                itchy::Body::DeleteOrder { reference } => {
152                    self.handle_delete_order(reference, ts_event, ts_init, &mut deltas);
153                }
154                itchy::Body::OrderCancelled {
155                    reference,
156                    cancelled,
157                } => {
158                    self.handle_cancel(reference, cancelled, ts_event, ts_init, &mut deltas);
159                }
160                itchy::Body::OrderExecuted {
161                    reference,
162                    executed,
163                    ..
164                }
165                | itchy::Body::OrderExecutedWithPrice {
166                    reference,
167                    executed,
168                    ..
169                } => {
170                    self.handle_execution(reference, executed, ts_event, ts_init, &mut deltas);
171                }
172                itchy::Body::ReplaceOrder(ref replace) => {
173                    self.handle_replace(replace, ts_event, ts_init, &mut deltas);
174                }
175                _ => {}
176            }
177        }
178
179        // Set F_LAST on the final delta
180        if let Some(last) = deltas.last_mut() {
181            last.flags |= RecordFlag::F_LAST as u8;
182        }
183
184        Ok(deltas)
185    }
186
187    fn handle_add_order(
188        &mut self,
189        add: &itchy::AddOrder,
190        ts_event: UnixNanos,
191        ts_init: UnixNanos,
192        deltas: &mut Vec<OrderBookDelta>,
193    ) {
194        let side = convert_side(add.side);
195        let price = convert_price(add.price);
196
197        self.orders.insert(
198            add.reference,
199            OrderState {
200                price,
201                size: add.shares,
202                side,
203            },
204        );
205
206        self.sequence += 1;
207        let order = BookOrder::new(
208            side,
209            price,
210            Quantity::new(f64::from(add.shares), SIZE_PRECISION),
211            add.reference,
212        );
213        deltas.push(OrderBookDelta::new(
214            self.instrument_id,
215            BookAction::Add,
216            order,
217            RecordFlag::F_LAST as u8,
218            self.sequence,
219            ts_event,
220            ts_init,
221        ));
222    }
223
224    fn handle_delete_order(
225        &mut self,
226        reference: u64,
227        ts_event: UnixNanos,
228        ts_init: UnixNanos,
229        deltas: &mut Vec<OrderBookDelta>,
230    ) {
231        if let Some(state) = self.orders.remove(&reference) {
232            self.sequence += 1;
233            let order = BookOrder::new(
234                state.side,
235                state.price,
236                Quantity::new(0.0, SIZE_PRECISION),
237                reference,
238            );
239            deltas.push(OrderBookDelta::new(
240                self.instrument_id,
241                BookAction::Delete,
242                order,
243                RecordFlag::F_LAST as u8,
244                self.sequence,
245                ts_event,
246                ts_init,
247            ));
248        }
249    }
250
251    fn handle_cancel(
252        &mut self,
253        reference: u64,
254        cancelled: u32,
255        ts_event: UnixNanos,
256        ts_init: UnixNanos,
257        deltas: &mut Vec<OrderBookDelta>,
258    ) {
259        if let Some(state) = self.orders.get_mut(&reference) {
260            state.size = state.size.saturating_sub(cancelled);
261
262            if state.size == 0 {
263                // Full cancel
264                let state = self.orders.remove(&reference).unwrap();
265                self.sequence += 1;
266                let order = BookOrder::new(
267                    state.side,
268                    state.price,
269                    Quantity::new(0.0, SIZE_PRECISION),
270                    reference,
271                );
272                deltas.push(OrderBookDelta::new(
273                    self.instrument_id,
274                    BookAction::Delete,
275                    order,
276                    RecordFlag::F_LAST as u8,
277                    self.sequence,
278                    ts_event,
279                    ts_init,
280                ));
281            } else {
282                // Partial cancel
283                self.sequence += 1;
284                let order = BookOrder::new(
285                    state.side,
286                    state.price,
287                    Quantity::new(f64::from(state.size), SIZE_PRECISION),
288                    reference,
289                );
290                deltas.push(OrderBookDelta::new(
291                    self.instrument_id,
292                    BookAction::Update,
293                    order,
294                    RecordFlag::F_LAST as u8,
295                    self.sequence,
296                    ts_event,
297                    ts_init,
298                ));
299            }
300        }
301    }
302
303    fn handle_execution(
304        &mut self,
305        reference: u64,
306        executed: u32,
307        ts_event: UnixNanos,
308        ts_init: UnixNanos,
309        deltas: &mut Vec<OrderBookDelta>,
310    ) {
311        if let Some(state) = self.orders.get_mut(&reference) {
312            state.size = state.size.saturating_sub(executed);
313
314            if state.size == 0 {
315                // Fully consumed
316                let state = self.orders.remove(&reference).unwrap();
317                self.sequence += 1;
318                let order = BookOrder::new(
319                    state.side,
320                    state.price,
321                    Quantity::new(0.0, SIZE_PRECISION),
322                    reference,
323                );
324                deltas.push(OrderBookDelta::new(
325                    self.instrument_id,
326                    BookAction::Delete,
327                    order,
328                    RecordFlag::F_LAST as u8,
329                    self.sequence,
330                    ts_event,
331                    ts_init,
332                ));
333            } else {
334                // Partial execution
335                self.sequence += 1;
336                let order = BookOrder::new(
337                    state.side,
338                    state.price,
339                    Quantity::new(f64::from(state.size), SIZE_PRECISION),
340                    reference,
341                );
342                deltas.push(OrderBookDelta::new(
343                    self.instrument_id,
344                    BookAction::Update,
345                    order,
346                    RecordFlag::F_LAST as u8,
347                    self.sequence,
348                    ts_event,
349                    ts_init,
350                ));
351            }
352        }
353    }
354
355    fn handle_replace(
356        &mut self,
357        replace: &itchy::ReplaceOrder,
358        ts_event: UnixNanos,
359        ts_init: UnixNanos,
360        deltas: &mut Vec<OrderBookDelta>,
361    ) {
362        // Delete old order
363        if let Some(old_state) = self.orders.remove(&replace.old_reference) {
364            self.sequence += 1;
365            let old_order = BookOrder::new(
366                old_state.side,
367                old_state.price,
368                Quantity::new(0.0, SIZE_PRECISION),
369                replace.old_reference,
370            );
371            deltas.push(OrderBookDelta::new(
372                self.instrument_id,
373                BookAction::Delete,
374                old_order,
375                0, // Not the last in this event group
376                self.sequence,
377                ts_event,
378                ts_init,
379            ));
380
381            // Add new order (inherits side from old order)
382            let new_price = convert_price(replace.price);
383            self.orders.insert(
384                replace.new_reference,
385                OrderState {
386                    price: new_price,
387                    size: replace.shares,
388                    side: old_state.side,
389                },
390            );
391
392            self.sequence += 1;
393            let new_order = BookOrder::new(
394                old_state.side,
395                new_price,
396                Quantity::new(f64::from(replace.shares), SIZE_PRECISION),
397                replace.new_reference,
398            );
399            deltas.push(OrderBookDelta::new(
400                self.instrument_id,
401                BookAction::Add,
402                new_order,
403                RecordFlag::F_LAST as u8,
404                self.sequence,
405                ts_event,
406                ts_init,
407            ));
408        }
409    }
410
411    fn handle_end_of_messages(
412        &mut self,
413        ts_event: UnixNanos,
414        ts_init: UnixNanos,
415        deltas: &mut Vec<OrderBookDelta>,
416    ) {
417        self.sequence += 1;
418        deltas.push(OrderBookDelta::clear(
419            self.instrument_id,
420            self.sequence,
421            ts_event,
422            ts_init,
423        ));
424    }
425}
426
427fn convert_side(side: itchy::Side) -> OrderSide {
428    match side {
429        itchy::Side::Buy => OrderSide::Buy,
430        itchy::Side::Sell => OrderSide::Sell,
431    }
432}
433
434fn convert_price(price: itchy::Price4) -> Price {
435    Price::new(f64::from(price.raw()) / 10_000.0, PRICE_PRECISION)
436}
437
438#[cfg(test)]
439mod tests {
440    use std::{fs, fs::File, path::PathBuf, sync::Arc};
441
442    use nautilus_model::{data::OrderBookDelta, enums::OrderSide};
443    use nautilus_serialization::arrow::{ArrowSchemaProvider, EncodeToRecordBatch};
444    use parquet::{arrow::ArrowWriter, file::properties::WriterProperties};
445    use rstest::rstest;
446
447    use super::*;
448
449    const AAPL_ID: &str = "AAPL.XNAS";
450
451    fn setup_parser(base_ns: u64) -> ItchParser {
452        ItchParser::new(InstrumentId::from(AAPL_ID), "AAPL", base_ns)
453    }
454
455    fn aapl_stream_with(messages: &[Vec<u8>]) -> Vec<u8> {
456        let mut buf = build_stock_directory_msg(1, b"AAPL    ");
457        for msg in messages {
458            buf.extend_from_slice(msg);
459        }
460        buf
461    }
462
463    #[rstest]
464    fn test_convert_side() {
465        assert_eq!(convert_side(itchy::Side::Buy), OrderSide::Buy);
466        assert_eq!(convert_side(itchy::Side::Sell), OrderSide::Sell);
467    }
468
469    #[rstest]
470    fn test_convert_price() {
471        let price = convert_price(itchy::Price4::from(1_2345));
472        assert_eq!(price.as_f64(), 1.2345);
473        assert_eq!(price.precision, PRICE_PRECISION);
474    }
475
476    #[rstest]
477    fn test_convert_price_whole_dollar() {
478        let price = convert_price(itchy::Price4::from(100_0000));
479        assert_eq!(price.as_f64(), 100.0);
480    }
481
482    #[rstest]
483    fn test_convert_price_sub_penny() {
484        let price = convert_price(itchy::Price4::from(150_2501));
485        assert_eq!(price.as_f64(), 150.2501);
486    }
487
488    #[rstest]
489    fn test_add_order() {
490        let buf = aapl_stream_with(&[build_add_order_msg(1, 42, b'B', 100, 1_502_500)]);
491        let mut parser = setup_parser(0);
492        let deltas = parser.parse_reader(&buf[..]).unwrap();
493
494        assert_eq!(deltas.len(), 1);
495        assert_eq!(deltas[0].action, BookAction::Add);
496        assert_eq!(deltas[0].order.side, Some(OrderSide::Buy));
497        assert_eq!(deltas[0].order.price.as_f64(), 150.25);
498        assert_eq!(deltas[0].order.size.as_f64(), 100.0);
499        assert_eq!(deltas[0].order.order_id, 42);
500    }
501
502    #[rstest]
503    fn test_delete_order() {
504        let buf = aapl_stream_with(&[
505            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
506            build_delete_order_msg(1, 42),
507        ]);
508        let mut parser = setup_parser(0);
509        let deltas = parser.parse_reader(&buf[..]).unwrap();
510
511        assert_eq!(deltas.len(), 2);
512        assert_eq!(deltas[1].action, BookAction::Delete);
513        assert_eq!(deltas[1].order.order_id, 42);
514        assert_eq!(deltas[1].order.size.as_f64(), 0.0);
515    }
516
517    #[rstest]
518    fn test_delete_unknown_order_is_ignored() {
519        let buf = aapl_stream_with(&[build_delete_order_msg(1, 999)]);
520        let mut parser = setup_parser(0);
521        let deltas = parser.parse_reader(&buf[..]).unwrap();
522
523        assert_eq!(deltas.len(), 0);
524    }
525
526    #[rstest]
527    fn test_partial_cancel_reduces_size() {
528        let buf = aapl_stream_with(&[
529            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
530            build_order_cancelled_msg(1, 42, 30),
531        ]);
532        let mut parser = setup_parser(0);
533        let deltas = parser.parse_reader(&buf[..]).unwrap();
534
535        assert_eq!(deltas.len(), 2);
536        assert_eq!(deltas[1].action, BookAction::Update);
537        assert_eq!(deltas[1].order.size.as_f64(), 70.0);
538        assert_eq!(deltas[1].order.price.as_f64(), 150.0);
539        assert_eq!(deltas[1].order.order_id, 42);
540    }
541
542    #[rstest]
543    fn test_full_cancel_deletes_order() {
544        let buf = aapl_stream_with(&[
545            build_add_order_msg(1, 42, b'S', 100, 1_500_000),
546            build_order_cancelled_msg(1, 42, 100),
547        ]);
548        let mut parser = setup_parser(0);
549        let deltas = parser.parse_reader(&buf[..]).unwrap();
550
551        assert_eq!(deltas.len(), 2);
552        assert_eq!(deltas[1].action, BookAction::Delete);
553        assert_eq!(deltas[1].order.size.as_f64(), 0.0);
554    }
555
556    #[rstest]
557    fn test_cancel_unknown_order_is_ignored() {
558        let buf = aapl_stream_with(&[build_order_cancelled_msg(1, 999, 50)]);
559        let mut parser = setup_parser(0);
560        let deltas = parser.parse_reader(&buf[..]).unwrap();
561
562        assert_eq!(deltas.len(), 0);
563    }
564
565    #[rstest]
566    fn test_partial_execution_updates_size() {
567        let buf = aapl_stream_with(&[
568            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
569            build_order_executed_msg(1, 42, 40, 1001),
570        ]);
571        let mut parser = setup_parser(0);
572        let deltas = parser.parse_reader(&buf[..]).unwrap();
573
574        assert_eq!(deltas.len(), 2);
575        assert_eq!(deltas[1].action, BookAction::Update);
576        assert_eq!(deltas[1].order.size.as_f64(), 60.0);
577        assert_eq!(deltas[1].order.order_id, 42);
578    }
579
580    #[rstest]
581    fn test_full_execution_deletes_order() {
582        let buf = aapl_stream_with(&[
583            build_add_order_msg(1, 42, b'S', 100, 1_500_000),
584            build_order_executed_msg(1, 42, 100, 1001),
585        ]);
586        let mut parser = setup_parser(0);
587        let deltas = parser.parse_reader(&buf[..]).unwrap();
588
589        assert_eq!(deltas.len(), 2);
590        assert_eq!(deltas[1].action, BookAction::Delete);
591        assert_eq!(deltas[1].order.size.as_f64(), 0.0);
592    }
593
594    #[rstest]
595    fn test_multiple_partial_executions_then_full() {
596        let buf = aapl_stream_with(&[
597            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
598            build_order_executed_msg(1, 42, 30, 1001),
599            build_order_executed_msg(1, 42, 30, 1002),
600            build_order_executed_msg(1, 42, 40, 1003),
601        ]);
602        let mut parser = setup_parser(0);
603        let deltas = parser.parse_reader(&buf[..]).unwrap();
604
605        assert_eq!(deltas.len(), 4);
606        assert_eq!(deltas[1].action, BookAction::Update);
607        assert_eq!(deltas[1].order.size.as_f64(), 70.0);
608        assert_eq!(deltas[2].action, BookAction::Update);
609        assert_eq!(deltas[2].order.size.as_f64(), 40.0);
610        assert_eq!(deltas[3].action, BookAction::Delete);
611        assert_eq!(deltas[3].order.size.as_f64(), 0.0);
612    }
613
614    #[rstest]
615    fn test_executed_with_price_partial() {
616        let buf = aapl_stream_with(&[
617            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
618            build_order_executed_with_price_msg(1, 42, 25, 2001, 1_505_000),
619        ]);
620        let mut parser = setup_parser(0);
621        let deltas = parser.parse_reader(&buf[..]).unwrap();
622
623        assert_eq!(deltas.len(), 2);
624        assert_eq!(deltas[1].action, BookAction::Update);
625        assert_eq!(deltas[1].order.size.as_f64(), 75.0);
626        // Book order retains the resting price, not the execution price
627        assert_eq!(deltas[1].order.price.as_f64(), 150.0);
628    }
629
630    #[rstest]
631    fn test_execution_unknown_order_is_ignored() {
632        let buf = aapl_stream_with(&[build_order_executed_msg(1, 999, 50, 1001)]);
633        let mut parser = setup_parser(0);
634        let deltas = parser.parse_reader(&buf[..]).unwrap();
635
636        assert_eq!(deltas.len(), 0);
637    }
638
639    #[rstest]
640    fn test_replace_order() {
641        let buf = aapl_stream_with(&[
642            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
643            build_replace_order_msg(1, 42, 43, 150, 1_510_000),
644        ]);
645        let mut parser = setup_parser(0);
646        let deltas = parser.parse_reader(&buf[..]).unwrap();
647
648        assert_eq!(deltas.len(), 3);
649        assert_eq!(deltas[1].action, BookAction::Delete);
650        assert_eq!(deltas[1].order.order_id, 42);
651        assert_eq!(deltas[2].action, BookAction::Add);
652        assert_eq!(deltas[2].order.order_id, 43);
653        assert_eq!(deltas[2].order.price.as_f64(), 151.0);
654        assert_eq!(deltas[2].order.size.as_f64(), 150.0);
655        assert_eq!(deltas[2].order.side, Some(OrderSide::Buy));
656    }
657
658    #[rstest]
659    fn test_replace_inherits_side_from_original() {
660        let buf = aapl_stream_with(&[
661            build_add_order_msg(1, 42, b'S', 100, 1_500_000),
662            build_replace_order_msg(1, 42, 43, 200, 1_490_000),
663        ]);
664        let mut parser = setup_parser(0);
665        let deltas = parser.parse_reader(&buf[..]).unwrap();
666
667        assert_eq!(deltas[2].order.side, Some(OrderSide::Sell));
668    }
669
670    #[rstest]
671    fn test_replace_delete_has_no_f_last_flag() {
672        // The delete in a replace pair is not the last event in the group
673        let buf = aapl_stream_with(&[
674            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
675            build_replace_order_msg(1, 42, 43, 100, 1_510_000),
676        ]);
677        let mut parser = setup_parser(0);
678        let deltas = parser.parse_reader(&buf[..]).unwrap();
679
680        assert_eq!(deltas[1].flags, 0);
681        assert_ne!(deltas[2].flags & RecordFlag::F_LAST as u8, 0);
682    }
683
684    #[rstest]
685    fn test_replace_unknown_order_is_ignored() {
686        let buf = aapl_stream_with(&[build_replace_order_msg(1, 999, 1000, 100, 1_500_000)]);
687        let mut parser = setup_parser(0);
688        let deltas = parser.parse_reader(&buf[..]).unwrap();
689
690        assert_eq!(deltas.len(), 0);
691    }
692
693    #[rstest]
694    fn test_end_of_messages_emits_clear() {
695        let buf = aapl_stream_with(&[
696            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
697            // EndOfMessages uses locate=0 (feed-level, not stock-specific)
698            build_system_event_msg(0, b'C'),
699        ]);
700        let mut parser = setup_parser(0);
701        let deltas = parser.parse_reader(&buf[..]).unwrap();
702
703        assert_eq!(deltas.len(), 2);
704        assert_eq!(deltas[1].action, BookAction::Clear);
705    }
706
707    #[rstest]
708    fn test_end_of_messages_with_different_locate() {
709        // Regression: EndOfMessages must be processed regardless of locate code
710        let buf = aapl_stream_with(&[
711            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
712            build_system_event_msg(99, b'C'),
713        ]);
714        let mut parser = setup_parser(0);
715        let deltas = parser.parse_reader(&buf[..]).unwrap();
716
717        assert_eq!(deltas.len(), 2);
718        assert_eq!(deltas[1].action, BookAction::Clear);
719    }
720
721    #[rstest]
722    fn test_filters_by_stock_locate() {
723        let mut buf = build_stock_directory_msg(1, b"AAPL    ");
724        buf.extend_from_slice(&build_stock_directory_msg(2, b"MSFT    "));
725        buf.extend_from_slice(&build_add_order_msg(2, 10, b'B', 50, 3_000_000));
726        buf.extend_from_slice(&build_add_order_msg(1, 11, b'S', 200, 1_500_000));
727
728        let mut parser = setup_parser(0);
729        let deltas = parser.parse_reader(&buf[..]).unwrap();
730
731        assert_eq!(deltas.len(), 1);
732        assert_eq!(deltas[0].order.order_id, 11);
733    }
734
735    #[rstest]
736    fn test_messages_before_directory_are_ignored() {
737        let mut buf = Vec::new();
738        // AddOrder arrives before any StockDirectory
739        buf.extend_from_slice(&build_add_order_msg(1, 42, b'B', 100, 1_500_000));
740        buf.extend_from_slice(&build_stock_directory_msg(1, b"AAPL    "));
741        buf.extend_from_slice(&build_add_order_msg(1, 43, b'B', 100, 1_500_000));
742
743        let mut parser = setup_parser(0);
744        let deltas = parser.parse_reader(&buf[..]).unwrap();
745
746        // Only the second add (after directory) should be captured
747        assert_eq!(deltas.len(), 1);
748        assert_eq!(deltas[0].order.order_id, 43);
749    }
750
751    #[rstest]
752    fn test_timestamp_offset_from_midnight() {
753        let base_ns: u64 = 1_548_806_400_000_000_000; // 2019-01-30 midnight UTC
754        let itch_ts: u64 = 34_200_000_000_000; // 9:30 AM (ns since midnight)
755        let buf = aapl_stream_with(&[build_add_order_msg_with_ts(
756            1, 42, b'B', 100, 1_500_000, itch_ts,
757        )]);
758        let mut parser = setup_parser(base_ns);
759        let deltas = parser.parse_reader(&buf[..]).unwrap();
760
761        assert_eq!(deltas[0].ts_event, UnixNanos::from(base_ns + itch_ts));
762        // ts_init is the conversion wall-clock stamp, later than the 2019 event time
763        assert!(deltas[0].ts_init > deltas[0].ts_event);
764    }
765
766    #[rstest]
767    fn test_f_last_set_on_final_delta() {
768        let buf = aapl_stream_with(&[
769            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
770            build_add_order_msg(1, 43, b'S', 200, 1_510_000),
771        ]);
772        let mut parser = setup_parser(0);
773        let deltas = parser.parse_reader(&buf[..]).unwrap();
774
775        assert_eq!(deltas.len(), 2);
776        assert_ne!(deltas[1].flags & RecordFlag::F_LAST as u8, 0);
777    }
778
779    #[rstest]
780    fn test_sequence_numbers_are_monotonic() {
781        let buf = aapl_stream_with(&[
782            build_add_order_msg(1, 42, b'B', 100, 1_500_000),
783            build_order_executed_msg(1, 42, 50, 1001),
784            build_add_order_msg(1, 43, b'S', 200, 1_510_000),
785            build_delete_order_msg(1, 43),
786        ]);
787        let mut parser = setup_parser(0);
788        let deltas = parser.parse_reader(&buf[..]).unwrap();
789
790        for i in 1..deltas.len() {
791            assert!(deltas[i].sequence > deltas[i - 1].sequence);
792        }
793    }
794
795    #[rstest]
796    fn test_empty_stream() {
797        let buf: &[u8] = &[];
798        let mut parser = setup_parser(0);
799        let deltas = parser.parse_reader(buf).unwrap();
800
801        assert_eq!(deltas.len(), 0);
802    }
803
804    fn build_msg(tag: u8, stock_locate: u16, timestamp: u64, body: &[u8]) -> Vec<u8> {
805        let msg_len = (1 + 2 + 2 + 6 + body.len()) as u16;
806        let mut buf = Vec::new();
807        buf.extend_from_slice(&msg_len.to_be_bytes());
808        buf.push(tag);
809        buf.extend_from_slice(&stock_locate.to_be_bytes());
810        buf.extend_from_slice(&0u16.to_be_bytes()); // tracking_number
811        // 6-byte timestamp (big-endian u48)
812        buf.push((timestamp >> 40) as u8);
813        buf.push((timestamp >> 32) as u8);
814        buf.push((timestamp >> 24) as u8);
815        buf.push((timestamp >> 16) as u8);
816        buf.push((timestamp >> 8) as u8);
817        buf.push(timestamp as u8);
818        buf.extend_from_slice(body);
819        buf
820    }
821
822    fn build_stock_directory_msg(locate: u16, stock: &[u8; 8]) -> Vec<u8> {
823        let mut body = Vec::new();
824        body.extend_from_slice(stock);
825        body.push(b'Q'); // market_category
826        body.push(b'N'); // financial_status
827        body.extend_from_slice(&100u32.to_be_bytes()); // round_lot_size
828        body.push(b'Y'); // round_lots_only
829        body.push(b'C'); // issue_classification
830        body.extend_from_slice(b"C "); // issue_subtype
831        body.push(b'P'); // authenticity
832        body.push(b'N'); // short_sale_threshold
833        body.push(b'N'); // ipo_flag
834        body.push(b'1'); // luld_ref_price_tier
835        body.push(b'N'); // etp_flag
836        body.extend_from_slice(&0u32.to_be_bytes()); // etp_leverage_factor
837        body.push(b'N'); // inverse_indicator
838        build_msg(b'R', locate, 0, &body)
839    }
840
841    fn build_add_order_msg(
842        locate: u16,
843        reference: u64,
844        side: u8,
845        shares: u32,
846        price: u32,
847    ) -> Vec<u8> {
848        build_add_order_msg_with_ts(locate, reference, side, shares, price, 0)
849    }
850
851    fn build_add_order_msg_with_ts(
852        locate: u16,
853        reference: u64,
854        side: u8,
855        shares: u32,
856        price: u32,
857        timestamp: u64,
858    ) -> Vec<u8> {
859        let mut body = Vec::new();
860        body.extend_from_slice(&reference.to_be_bytes());
861        body.push(side);
862        body.extend_from_slice(&shares.to_be_bytes());
863        body.extend_from_slice(b"AAPL    ");
864        body.extend_from_slice(&price.to_be_bytes());
865        build_msg(b'A', locate, timestamp, &body)
866    }
867
868    fn build_delete_order_msg(locate: u16, reference: u64) -> Vec<u8> {
869        build_msg(b'D', locate, 0, &reference.to_be_bytes())
870    }
871
872    fn build_order_cancelled_msg(locate: u16, reference: u64, cancelled: u32) -> Vec<u8> {
873        let mut body = Vec::new();
874        body.extend_from_slice(&reference.to_be_bytes());
875        body.extend_from_slice(&cancelled.to_be_bytes());
876        build_msg(b'X', locate, 0, &body)
877    }
878
879    fn build_order_executed_msg(
880        locate: u16,
881        reference: u64,
882        executed: u32,
883        match_number: u64,
884    ) -> Vec<u8> {
885        let mut body = Vec::new();
886        body.extend_from_slice(&reference.to_be_bytes());
887        body.extend_from_slice(&executed.to_be_bytes());
888        body.extend_from_slice(&match_number.to_be_bytes());
889        build_msg(b'E', locate, 0, &body)
890    }
891
892    fn build_order_executed_with_price_msg(
893        locate: u16,
894        reference: u64,
895        executed: u32,
896        match_number: u64,
897        price: u32,
898    ) -> Vec<u8> {
899        let mut body = Vec::new();
900        body.extend_from_slice(&reference.to_be_bytes());
901        body.extend_from_slice(&executed.to_be_bytes());
902        body.extend_from_slice(&match_number.to_be_bytes());
903        body.push(b'Y'); // printable
904        body.extend_from_slice(&price.to_be_bytes());
905        build_msg(b'C', locate, 0, &body)
906    }
907
908    fn build_replace_order_msg(
909        locate: u16,
910        old_reference: u64,
911        new_reference: u64,
912        shares: u32,
913        price: u32,
914    ) -> Vec<u8> {
915        let mut body = Vec::new();
916        body.extend_from_slice(&old_reference.to_be_bytes());
917        body.extend_from_slice(&new_reference.to_be_bytes());
918        body.extend_from_slice(&shares.to_be_bytes());
919        body.extend_from_slice(&price.to_be_bytes());
920        build_msg(b'U', locate, 0, &body)
921    }
922
923    fn build_system_event_msg(locate: u16, event_code: u8) -> Vec<u8> {
924        build_msg(b'S', locate, 0, &[event_code])
925    }
926
927    // Curates AAPL L3 deltas from NASDAQ ITCH 5.0 binary into NautilusTrader Parquet.
928    // Download source: https://emi.nasdaq.com/ITCH/Nasdaq%20ITCH/01302019.NASDAQ_ITCH50.gz
929    // Run: cargo test -p nautilus-testkit --lib test_curate_aapl_itch -- --ignored --nocapture
930    #[rstest]
931    #[ignore = "one-time dataset curation, not for routine CI"]
932    fn test_curate_aapl_itch() {
933        let itch_path = PathBuf::from("/tmp/01302019.NASDAQ_ITCH50.gz");
934        let instrument_id = InstrumentId::from("AAPL.XNAS");
935
936        // 2019-01-30 midnight EST (UTC-5) as Unix nanoseconds
937        let base_ns: u64 = 1_548_824_400_000_000_000;
938        let parquet_path = "/tmp/itch_AAPL.XNAS_2019-01-30_deltas.parquet";
939
940        println!("Parsing ITCH from {}", itch_path.display());
941        let mut parser = ItchParser::new(instrument_id, "AAPL", base_ns);
942        let deltas = parser.parse_gzip_file(&itch_path).unwrap();
943        let count = deltas.len();
944        println!("Parsed {count} deltas for AAPL");
945
946        let metadata =
947            OrderBookDelta::get_metadata(&instrument_id, PRICE_PRECISION, SIZE_PRECISION);
948        let schema = OrderBookDelta::get_schema(Some(metadata.clone()));
949
950        println!("Writing Parquet to {parquet_path}");
951        let file = File::create(parquet_path).unwrap();
952        let zstd_level = parquet::basic::ZstdLevel::try_new(3).unwrap();
953        let props = WriterProperties::builder()
954            .set_compression(parquet::basic::Compression::ZSTD(zstd_level))
955            .set_max_row_group_row_count(Some(1_000_000))
956            .build();
957        let mut writer = ArrowWriter::try_new(file, Arc::new(schema), Some(props)).unwrap();
958
959        let chunk_size = 1_000_000;
960        for (i, chunk) in deltas.chunks(chunk_size).enumerate() {
961            println!("  Encoding chunk {} ({} records)...", i + 1, chunk.len());
962            let batch = OrderBookDelta::encode_batch(&metadata, chunk).unwrap();
963            writer.write(&batch).unwrap();
964        }
965        writer.close().unwrap();
966
967        let file_size = fs::metadata(parquet_path).unwrap().len();
968        println!("\nRecords: {count}");
969        println!("Price precision: {PRICE_PRECISION}");
970        println!("Size precision: {SIZE_PRECISION}");
971        println!(
972            "File size: {} bytes ({:.1} MB)",
973            file_size,
974            file_size as f64 / 1_048_576.0
975        );
976        println!("Output: {parquet_path}");
977        println!("\nNext steps:");
978        println!("  sha256sum {parquet_path}");
979    }
980}