1use 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
32const PRICE_PRECISION: u8 = 4;
34
35const SIZE_PRECISION: u8 = 0;
37
38#[derive(Debug)]
39struct OrderState {
40 price: Price,
41 size: u32,
42 side: OrderSide,
43}
44
45#[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 #[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 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 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 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 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 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 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 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 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 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 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, self.sequence,
377 ts_event,
378 ts_init,
379 ));
380
381 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 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 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 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 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 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 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; let itch_ts: u64 = 34_200_000_000_000; 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 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()); 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'); body.push(b'N'); body.extend_from_slice(&100u32.to_be_bytes()); body.push(b'Y'); body.push(b'C'); body.extend_from_slice(b"C "); body.push(b'P'); body.push(b'N'); body.push(b'N'); body.push(b'1'); body.push(b'N'); body.extend_from_slice(&0u32.to_be_bytes()); body.push(b'N'); 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'); 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 #[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 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}