1use std::{collections::HashMap, sync::Arc};
17
18use arrow::{
19 array::{FixedSizeBinaryArray, FixedSizeBinaryBuilder, UInt8Array, UInt64Array},
20 datatypes::{DataType, Field, Schema},
21 error::ArrowError,
22 record_batch::RecordBatch,
23};
24use nautilus_model::{
25 data::{BookOrder, OrderBookDelta},
26 enums::{BookAction, FromU8, OrderSide},
27 types::fixed::PRECISION_BYTES,
28};
29
30use super::{
31 DecodeDataFromRecordBatch, EncodingError, decode_price_with_sentinel,
32 decode_quantity_with_sentinel, extract_column, parse_price_size_metadata,
33 validate_precision_bytes,
34};
35use crate::arrow::{ArrowSchemaProvider, Data, DecodeFromRecordBatch, EncodeToRecordBatch};
36
37impl ArrowSchemaProvider for OrderBookDelta {
38 fn get_schema(metadata: Option<HashMap<String, String>>) -> Schema {
39 let fields = vec![
40 Field::new("action", DataType::UInt8, false),
41 Field::new("side", DataType::UInt8, false),
42 Field::new("price", DataType::FixedSizeBinary(PRECISION_BYTES), false),
43 Field::new("size", DataType::FixedSizeBinary(PRECISION_BYTES), false),
44 Field::new("order_id", DataType::UInt64, false),
45 Field::new("flags", DataType::UInt8, false),
46 Field::new("sequence", DataType::UInt64, false),
47 Field::new("ts_event", DataType::UInt64, false),
48 Field::new("ts_init", DataType::UInt64, false),
49 ];
50
51 match metadata {
52 Some(metadata) => Schema::new_with_metadata(fields, metadata),
53 None => Schema::new(fields),
54 }
55 }
56}
57
58impl EncodeToRecordBatch for OrderBookDelta {
59 fn encode_batch(
60 metadata: &HashMap<String, String>,
61 data: &[Self],
62 ) -> Result<RecordBatch, ArrowError> {
63 let mut action_builder = UInt8Array::builder(data.len());
64 let mut side_builder = UInt8Array::builder(data.len());
65 let mut price_builder = FixedSizeBinaryBuilder::with_capacity(data.len(), PRECISION_BYTES);
66 let mut size_builder = FixedSizeBinaryBuilder::with_capacity(data.len(), PRECISION_BYTES);
67 let mut order_id_builder = UInt64Array::builder(data.len());
68 let mut flags_builder = UInt8Array::builder(data.len());
69 let mut sequence_builder = UInt64Array::builder(data.len());
70 let mut ts_event_builder = UInt64Array::builder(data.len());
71 let mut ts_init_builder = UInt64Array::builder(data.len());
72
73 for delta in data {
74 action_builder.append_value(delta.action as u8);
75 side_builder.append_value(delta.order.side.map_or(0, |side| side as u8));
76 price_builder
77 .append_value(delta.order.price.raw().to_le_bytes())
78 .unwrap();
79 size_builder
80 .append_value(delta.order.size.raw().to_le_bytes())
81 .unwrap();
82 order_id_builder.append_value(delta.order.order_id);
83 flags_builder.append_value(delta.flags);
84 sequence_builder.append_value(delta.sequence);
85 ts_event_builder.append_value(delta.ts_event.as_u64());
86 ts_init_builder.append_value(delta.ts_init.as_u64());
87 }
88
89 let action_array = action_builder.finish();
90 let side_array = side_builder.finish();
91 let price_array = price_builder.finish();
92 let size_array = size_builder.finish();
93 let order_id_array = order_id_builder.finish();
94 let flags_array = flags_builder.finish();
95 let sequence_array = sequence_builder.finish();
96 let ts_event_array = ts_event_builder.finish();
97 let ts_init_array = ts_init_builder.finish();
98
99 RecordBatch::try_new(
100 Self::get_schema(Some(metadata.clone())).into(),
101 vec![
102 Arc::new(action_array),
103 Arc::new(side_array),
104 Arc::new(price_array),
105 Arc::new(size_array),
106 Arc::new(order_id_array),
107 Arc::new(flags_array),
108 Arc::new(sequence_array),
109 Arc::new(ts_event_array),
110 Arc::new(ts_init_array),
111 ],
112 )
113 }
114
115 fn metadata(&self) -> HashMap<String, String> {
116 Self::get_metadata(
117 &self.instrument_id,
118 self.order.price.precision,
119 self.order.size.precision,
120 )
121 }
122
123 fn chunk_metadata(chunk: &[Self]) -> HashMap<String, String> {
127 chunk
128 .iter()
129 .find(|delta| delta.action != BookAction::Clear)
130 .or_else(|| chunk.first())
131 .map(EncodeToRecordBatch::metadata)
132 .expect("Chunk must have at least one element to encode")
133 }
134}
135
136impl DecodeFromRecordBatch for OrderBookDelta {
137 fn decode_batch(
138 metadata: &HashMap<String, String>,
139 record_batch: RecordBatch,
140 ) -> Result<Vec<Self>, EncodingError> {
141 let (instrument_id, price_precision, size_precision) = parse_price_size_metadata(metadata)?;
142 let cols = record_batch.columns();
143
144 let action_values = extract_column::<UInt8Array>(cols, "action", 0, DataType::UInt8)?;
145 let side_values = extract_column::<UInt8Array>(cols, "side", 1, DataType::UInt8)?;
146 let price_values = extract_column::<FixedSizeBinaryArray>(
147 cols,
148 "price",
149 2,
150 DataType::FixedSizeBinary(PRECISION_BYTES),
151 )?;
152 let size_values = extract_column::<FixedSizeBinaryArray>(
153 cols,
154 "size",
155 3,
156 DataType::FixedSizeBinary(PRECISION_BYTES),
157 )?;
158 let order_id_values = extract_column::<UInt64Array>(cols, "order_id", 4, DataType::UInt64)?;
159 let flags_values = extract_column::<UInt8Array>(cols, "flags", 5, DataType::UInt8)?;
160 let sequence_values = extract_column::<UInt64Array>(cols, "sequence", 6, DataType::UInt64)?;
161 let ts_event_values = extract_column::<UInt64Array>(cols, "ts_event", 7, DataType::UInt64)?;
162 let ts_init_values = extract_column::<UInt64Array>(cols, "ts_init", 8, DataType::UInt64)?;
163
164 validate_precision_bytes(price_values, "price")?;
165 validate_precision_bytes(size_values, "size")?;
166
167 let result: Result<Vec<Self>, EncodingError> = (0..record_batch.num_rows())
168 .map(|i| {
169 let action_value = action_values.value(i);
170 let action = BookAction::from_u8(action_value).ok_or_else(|| {
171 EncodingError::ParseError(
172 stringify!(BookAction),
173 format!("Invalid enum value, was {action_value}"),
174 )
175 })?;
176 let side_value = side_values.value(i);
177 let side = match side_value {
178 0 => None,
179 1 => Some(OrderSide::Buy),
180 2 => Some(OrderSide::Sell),
181 _ => {
182 return Err(EncodingError::ParseError(
183 "Option<OrderSide>",
184 format!("Invalid enum value, was {side_value}"),
185 ));
186 }
187 };
188 let price =
189 decode_price_with_sentinel(price_values.value(i), price_precision, "price", i)?;
190 let size =
191 decode_quantity_with_sentinel(size_values.value(i), size_precision, "size", i)?;
192 let order_id = order_id_values.value(i);
193 let flags = flags_values.value(i);
194 let sequence = sequence_values.value(i);
195 let ts_event = ts_event_values.value(i).into();
196 let ts_init = ts_init_values.value(i).into();
197
198 Ok(Self {
199 instrument_id,
200 action,
201 order: BookOrder {
202 side,
203 price,
204 size,
205 order_id,
206 },
207 flags,
208 sequence,
209 ts_event,
210 ts_init,
211 })
212 })
213 .collect();
214
215 result
216 }
217}
218
219impl DecodeDataFromRecordBatch for OrderBookDelta {
220 fn decode_data_batch(
221 metadata: &HashMap<String, String>,
222 record_batch: RecordBatch,
223 ) -> Result<Vec<Data>, EncodingError> {
224 let deltas: Vec<Self> = Self::decode_batch(metadata, record_batch)?;
225 Ok(deltas.into_iter().map(Data::from).collect())
226 }
227}
228
229#[cfg(test)]
230mod tests {
231 use std::sync::Arc;
232
233 use arrow::{array::Array, record_batch::RecordBatch};
234 use nautilus_model::{
235 enums::OrderSide,
236 identifiers::InstrumentId,
237 types::{
238 Price, Quantity,
239 fixed::FIXED_SCALAR,
240 price::{PRICE_UNDEF, PriceRaw},
241 quantity::{QUANTITY_UNDEF, QuantityRaw},
242 },
243 };
244 use pretty_assertions::assert_eq;
245 use rstest::rstest;
246
247 use super::*;
248 use crate::arrow::{KEY_INSTRUMENT_ID, fixed_size_binary, get_raw_price};
249
250 #[rstest]
251 fn test_get_schema() {
252 let instrument_id = InstrumentId::from("AAPL.XNAS");
253 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
254 let schema = OrderBookDelta::get_schema(Some(metadata.clone()));
255
256 let expected_fields = vec![
257 Field::new("action", DataType::UInt8, false),
258 Field::new("side", DataType::UInt8, false),
259 Field::new("price", DataType::FixedSizeBinary(PRECISION_BYTES), false),
260 Field::new("size", DataType::FixedSizeBinary(PRECISION_BYTES), false),
261 Field::new("order_id", DataType::UInt64, false),
262 Field::new("flags", DataType::UInt8, false),
263 Field::new("sequence", DataType::UInt64, false),
264 Field::new("ts_event", DataType::UInt64, false),
265 Field::new("ts_init", DataType::UInt64, false),
266 ];
267
268 let expected_schema = Schema::new_with_metadata(expected_fields, metadata);
269 assert_eq!(schema, expected_schema);
270 }
271
272 #[rstest]
273 fn test_get_schema_map() {
274 let schema_map = OrderBookDelta::get_schema_map();
275 let fixed_size_binary = format!("FixedSizeBinary({PRECISION_BYTES})");
276
277 assert_eq!(schema_map.get("action").unwrap(), "UInt8");
278 assert_eq!(schema_map.get("side").unwrap(), "UInt8");
279 assert_eq!(*schema_map.get("price").unwrap(), fixed_size_binary);
280 assert_eq!(*schema_map.get("size").unwrap(), fixed_size_binary);
281 assert_eq!(schema_map.get("order_id").unwrap(), "UInt64");
282 assert_eq!(schema_map.get("flags").unwrap(), "UInt8");
283 assert_eq!(schema_map.get("sequence").unwrap(), "UInt64");
284 assert_eq!(schema_map.get("ts_event").unwrap(), "UInt64");
285 assert_eq!(schema_map.get("ts_init").unwrap(), "UInt64");
286 }
287
288 #[rstest]
289 fn test_encode_batch() {
290 let instrument_id = InstrumentId::from("AAPL.XNAS");
291 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
292
293 let delta1 = OrderBookDelta {
294 instrument_id,
295 action: BookAction::Add,
296 order: BookOrder {
297 side: OrderSide::Buy.into(),
298 price: Price::from("100.10"),
299 size: Quantity::from(100),
300 order_id: 1,
301 },
302 flags: 0,
303 sequence: 1,
304 ts_event: 1.into(),
305 ts_init: 3.into(),
306 };
307
308 let delta2 = OrderBookDelta {
309 instrument_id,
310 action: BookAction::Update,
311 order: BookOrder {
312 side: OrderSide::Sell.into(),
313 price: Price::from("101.20"),
314 size: Quantity::from(200),
315 order_id: 2,
316 },
317 flags: 1,
318 sequence: 2,
319 ts_event: 2.into(),
320 ts_init: 4.into(),
321 };
322
323 let data = vec![delta1, delta2];
324 let record_batch = OrderBookDelta::encode_batch(&metadata, &data).unwrap();
325
326 let columns = record_batch.columns();
327 let action_values = columns[0].as_any().downcast_ref::<UInt8Array>().unwrap();
328 let side_values = columns[1].as_any().downcast_ref::<UInt8Array>().unwrap();
329 let price_values = columns[2]
330 .as_any()
331 .downcast_ref::<FixedSizeBinaryArray>()
332 .unwrap();
333 let size_values = columns[3]
334 .as_any()
335 .downcast_ref::<FixedSizeBinaryArray>()
336 .unwrap();
337 let order_id_values = columns[4].as_any().downcast_ref::<UInt64Array>().unwrap();
338 let flags_values = columns[5].as_any().downcast_ref::<UInt8Array>().unwrap();
339 let sequence_values = columns[6].as_any().downcast_ref::<UInt64Array>().unwrap();
340 let ts_event_values = columns[7].as_any().downcast_ref::<UInt64Array>().unwrap();
341 let ts_init_values = columns[8].as_any().downcast_ref::<UInt64Array>().unwrap();
342
343 assert_eq!(columns.len(), 9);
344 assert_eq!(action_values.len(), 2);
345 assert_eq!(action_values.value(0), 1);
346 assert_eq!(action_values.value(1), 2);
347 assert_eq!(side_values.len(), 2);
348 assert_eq!(side_values.value(0), 1);
349 assert_eq!(side_values.value(1), 2);
350
351 assert_eq!(price_values.len(), 2);
352 assert_eq!(
353 get_raw_price(price_values.value(0)),
354 (100.10 * FIXED_SCALAR) as PriceRaw
355 );
356 assert_eq!(
357 get_raw_price(price_values.value(1)),
358 (101.20 * FIXED_SCALAR) as PriceRaw
359 );
360
361 assert_eq!(size_values.len(), 2);
362 assert_eq!(
363 get_raw_price(size_values.value(0)),
364 (100.0 * FIXED_SCALAR) as PriceRaw
365 );
366 assert_eq!(
367 get_raw_price(size_values.value(1)),
368 (200.0 * FIXED_SCALAR) as PriceRaw
369 );
370 assert_eq!(order_id_values.len(), 2);
371 assert_eq!(order_id_values.value(0), 1);
372 assert_eq!(order_id_values.value(1), 2);
373 assert_eq!(flags_values.len(), 2);
374 assert_eq!(flags_values.value(0), 0);
375 assert_eq!(flags_values.value(1), 1);
376 assert_eq!(sequence_values.len(), 2);
377 assert_eq!(sequence_values.value(0), 1);
378 assert_eq!(sequence_values.value(1), 2);
379 assert_eq!(ts_event_values.len(), 2);
380 assert_eq!(ts_event_values.value(0), 1);
381 assert_eq!(ts_event_values.value(1), 2);
382 assert_eq!(ts_init_values.len(), 2);
383 assert_eq!(ts_init_values.value(0), 3);
384 assert_eq!(ts_init_values.value(1), 4);
385 }
386
387 #[rstest]
388 fn test_decode_batch() {
389 let instrument_id = InstrumentId::from("AAPL.XNAS");
390 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
391
392 let action = UInt8Array::from(vec![1, 2]);
393 let side = UInt8Array::from(vec![1, 1]);
394 let price = fixed_size_binary(vec![
395 &((101.10 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
396 &((101.20 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
397 ]);
398 let size = fixed_size_binary(vec![
399 &((10000.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
400 &((9000.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
401 ]);
402 let order_id = UInt64Array::from(vec![1, 2]);
403 let flags = UInt8Array::from(vec![0, 0]);
404 let sequence = UInt64Array::from(vec![1, 2]);
405 let ts_event = UInt64Array::from(vec![1, 2]);
406 let ts_init = UInt64Array::from(vec![3, 4]);
407
408 let record_batch = RecordBatch::try_new(
409 OrderBookDelta::get_schema(Some(metadata.clone())).into(),
410 vec![
411 Arc::new(action),
412 Arc::new(side),
413 Arc::new(price),
414 Arc::new(size),
415 Arc::new(order_id),
416 Arc::new(flags),
417 Arc::new(sequence),
418 Arc::new(ts_event),
419 Arc::new(ts_init),
420 ],
421 )
422 .unwrap();
423
424 let decoded_data = OrderBookDelta::decode_batch(&metadata, record_batch).unwrap();
425 assert_eq!(decoded_data.len(), 2);
426 }
427
428 #[rstest]
429 fn test_decode_batch_with_undef_values() {
430 let instrument_id = InstrumentId::from("PLTR.XNAS");
431 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
432
433 let action = UInt8Array::from(vec![4, 1]); let side = UInt8Array::from(vec![0, 1]); let price = fixed_size_binary(vec![
437 &PRICE_UNDEF.to_le_bytes(),
438 &((100.50 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
439 ]);
440 let size = fixed_size_binary(vec![
441 &QUANTITY_UNDEF.to_le_bytes(),
442 &((1000.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
443 ]);
444 let order_id = UInt64Array::from(vec![0, 1]);
445 let flags = UInt8Array::from(vec![0, 0]);
446 let sequence = UInt64Array::from(vec![1, 2]);
447 let ts_event = UInt64Array::from(vec![1, 2]);
448 let ts_init = UInt64Array::from(vec![3, 4]);
449
450 let record_batch = RecordBatch::try_new(
451 OrderBookDelta::get_schema(Some(metadata.clone())).into(),
452 vec![
453 Arc::new(action),
454 Arc::new(side),
455 Arc::new(price),
456 Arc::new(size),
457 Arc::new(order_id),
458 Arc::new(flags),
459 Arc::new(sequence),
460 Arc::new(ts_event),
461 Arc::new(ts_init),
462 ],
463 )
464 .unwrap();
465
466 let decoded_data = OrderBookDelta::decode_batch(&metadata, record_batch).unwrap();
467 assert_eq!(decoded_data.len(), 2);
468 assert_eq!(decoded_data[0].order.price.raw(), PRICE_UNDEF);
469 assert_eq!(decoded_data[0].order.price.precision, 0);
470 assert_eq!(decoded_data[0].order.size.raw(), QUANTITY_UNDEF);
471 assert_eq!(decoded_data[0].order.size.precision, 0);
472 assert_eq!(decoded_data[1].order.price.precision, 2);
473 assert_eq!(decoded_data[1].order.size.precision, 0);
474 }
475
476 #[rstest]
477 fn test_decode_batch_invalid_price_returns_error() {
478 let instrument_id = InstrumentId::from("AAPL.XNAS");
479 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
480
481 let action = UInt8Array::from(vec![1]);
482 let side = UInt8Array::from(vec![1]);
483
484 let invalid_price: PriceRaw = PriceRaw::MAX - 1000;
485 let price = fixed_size_binary(vec![&invalid_price.to_le_bytes()]);
486 let size = fixed_size_binary(vec![&((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes()]);
487 let order_id = UInt64Array::from(vec![1]);
488 let flags = UInt8Array::from(vec![0]);
489 let sequence = UInt64Array::from(vec![1]);
490 let ts_event = UInt64Array::from(vec![1]);
491 let ts_init = UInt64Array::from(vec![2]);
492
493 let record_batch = RecordBatch::try_new(
494 OrderBookDelta::get_schema(Some(metadata.clone())).into(),
495 vec![
496 Arc::new(action),
497 Arc::new(side),
498 Arc::new(price),
499 Arc::new(size),
500 Arc::new(order_id),
501 Arc::new(flags),
502 Arc::new(sequence),
503 Arc::new(ts_event),
504 Arc::new(ts_init),
505 ],
506 )
507 .unwrap();
508
509 let result = OrderBookDelta::decode_batch(&metadata, record_batch);
510 assert!(result.is_err());
511 let err = result.unwrap_err();
512 assert!(
513 err.to_string().contains("price") && err.to_string().contains("row 0"),
514 "Expected price error at row 0, was: {err}"
515 );
516 }
517
518 #[rstest]
519 fn test_decode_batch_invalid_action_returns_error() {
520 let instrument_id = InstrumentId::from("AAPL.XNAS");
521 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
522
523 let action = UInt8Array::from(vec![99]);
524 let side = UInt8Array::from(vec![1]);
525 let price = fixed_size_binary(vec![&((100.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes()]);
526 let size = fixed_size_binary(vec![&((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes()]);
527 let order_id = UInt64Array::from(vec![1]);
528 let flags = UInt8Array::from(vec![0]);
529 let sequence = UInt64Array::from(vec![1]);
530 let ts_event = UInt64Array::from(vec![1]);
531 let ts_init = UInt64Array::from(vec![2]);
532
533 let record_batch = RecordBatch::try_new(
534 OrderBookDelta::get_schema(Some(metadata.clone())).into(),
535 vec![
536 Arc::new(action),
537 Arc::new(side),
538 Arc::new(price),
539 Arc::new(size),
540 Arc::new(order_id),
541 Arc::new(flags),
542 Arc::new(sequence),
543 Arc::new(ts_event),
544 Arc::new(ts_init),
545 ],
546 )
547 .unwrap();
548
549 let result = OrderBookDelta::decode_batch(&metadata, record_batch);
550 assert!(result.is_err());
551 let err = result.unwrap_err();
552 assert!(
553 err.to_string().contains("BookAction"),
554 "Expected BookAction error, was: {err}"
555 );
556 }
557
558 #[rstest]
559 fn test_decode_batch_missing_instrument_id_returns_error() {
560 let instrument_id = InstrumentId::from("AAPL.XNAS");
561 let mut metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
562 metadata.remove(KEY_INSTRUMENT_ID);
563
564 let action = UInt8Array::from(vec![1]);
565 let side = UInt8Array::from(vec![1]);
566 let price = fixed_size_binary(vec![&((100.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes()]);
567 let size = fixed_size_binary(vec![&((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes()]);
568 let order_id = UInt64Array::from(vec![1]);
569 let flags = UInt8Array::from(vec![0]);
570 let sequence = UInt64Array::from(vec![1]);
571 let ts_event = UInt64Array::from(vec![1]);
572 let ts_init = UInt64Array::from(vec![2]);
573
574 let record_batch = RecordBatch::try_new(
575 OrderBookDelta::get_schema(Some(metadata.clone())).into(),
576 vec![
577 Arc::new(action),
578 Arc::new(side),
579 Arc::new(price),
580 Arc::new(size),
581 Arc::new(order_id),
582 Arc::new(flags),
583 Arc::new(sequence),
584 Arc::new(ts_event),
585 Arc::new(ts_init),
586 ],
587 )
588 .unwrap();
589
590 let result = OrderBookDelta::decode_batch(&metadata, record_batch);
591 assert!(result.is_err());
592 let err = result.unwrap_err();
593 assert!(
594 err.to_string().contains("instrument_id"),
595 "Expected missing instrument_id error, was: {err}"
596 );
597 }
598
599 #[rstest]
600 fn test_encode_decode_round_trip() {
601 let instrument_id = InstrumentId::from("AAPL.XNAS");
602 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
603
604 let delta1 = OrderBookDelta {
605 instrument_id,
606 action: BookAction::Add,
607 order: BookOrder {
608 side: OrderSide::Buy.into(),
609 price: Price::from("100.10"),
610 size: Quantity::from(100),
611 order_id: 1,
612 },
613 flags: 0,
614 sequence: 1,
615 ts_event: 1_000_000_000.into(),
616 ts_init: 1_000_000_001.into(),
617 };
618
619 let delta2 = OrderBookDelta {
620 instrument_id,
621 action: BookAction::Update,
622 order: BookOrder {
623 side: OrderSide::Sell.into(),
624 price: Price::from("101.20"),
625 size: Quantity::from(200),
626 order_id: 2,
627 },
628 flags: 1,
629 sequence: 2,
630 ts_event: 2_000_000_000.into(),
631 ts_init: 2_000_000_001.into(),
632 };
633
634 let original = vec![delta1, delta2];
635 let record_batch = OrderBookDelta::encode_batch(&metadata, &original).unwrap();
636 let decoded = OrderBookDelta::decode_batch(&metadata, record_batch).unwrap();
637
638 assert_eq!(decoded.len(), original.len());
639 for (orig, dec) in original.iter().zip(decoded.iter()) {
640 assert_eq!(dec.instrument_id, orig.instrument_id);
641 assert_eq!(dec.action, orig.action);
642 assert_eq!(dec.order.side, orig.order.side);
643 assert_eq!(dec.order.price, orig.order.price);
644 assert_eq!(dec.order.size, orig.order.size);
645 assert_eq!(dec.order.order_id, orig.order.order_id);
646 assert_eq!(dec.flags, orig.flags);
647 assert_eq!(dec.sequence, orig.sequence);
648 assert_eq!(dec.ts_event, orig.ts_event);
649 assert_eq!(dec.ts_init, orig.ts_init);
650 }
651 }
652}