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