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.map_or(0, |side| 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 = match side_value {
203 0 => None,
204 1 => Some(OrderSide::Buy),
205 2 => Some(OrderSide::Sell),
206 _ => {
207 return Err(EncodingError::ParseError(
208 "Option<OrderSide>",
209 format!("Invalid enum value, was {side_value}"),
210 ));
211 }
212 };
213 let price =
214 decode_price_with_sentinel(price_values.value(i), price_precision, "price", i)?;
215 let size =
216 decode_quantity_with_sentinel(size_values.value(i), size_precision, "size", i)?;
217 let order_id = order_id_values.value(i);
218 let flags = flags_values.value(i);
219 let sequence = sequence_values.value(i);
220 let ts_event = ts_event_values.value(i).into();
221 let ts_init = ts_init_values.value(i).into();
222
223 Ok(Self {
224 instrument_id,
225 action,
226 order: BookOrder {
227 side,
228 price,
229 size,
230 order_id,
231 },
232 flags,
233 sequence,
234 ts_event,
235 ts_init,
236 })
237 })
238 .collect();
239
240 result
241 }
242}
243
244impl DecodeDataFromRecordBatch for OrderBookDelta {
245 fn decode_data_batch(
246 metadata: &HashMap<String, String>,
247 record_batch: RecordBatch,
248 ) -> Result<Vec<Data>, EncodingError> {
249 let deltas: Vec<Self> = Self::decode_batch(metadata, record_batch)?;
250 Ok(deltas.into_iter().map(Data::from).collect())
251 }
252}
253
254#[cfg(test)]
255mod tests {
256 use std::sync::Arc;
257
258 use arrow::{array::Array, record_batch::RecordBatch};
259 use nautilus_model::{
260 enums::OrderSide,
261 types::{
262 Price, Quantity,
263 fixed::FIXED_SCALAR,
264 price::{PRICE_UNDEF, PriceRaw},
265 quantity::{QUANTITY_UNDEF, QuantityRaw},
266 },
267 };
268 use pretty_assertions::assert_eq;
269 use rstest::rstest;
270
271 use super::*;
272 use crate::arrow::{fixed_size_binary, get_raw_price};
273
274 #[rstest]
275 fn test_get_schema() {
276 let instrument_id = InstrumentId::from("AAPL.XNAS");
277 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
278 let schema = OrderBookDelta::get_schema(Some(metadata.clone()));
279
280 let expected_fields = vec![
281 Field::new("action", DataType::UInt8, false),
282 Field::new("side", DataType::UInt8, false),
283 Field::new("price", DataType::FixedSizeBinary(PRECISION_BYTES), false),
284 Field::new("size", DataType::FixedSizeBinary(PRECISION_BYTES), false),
285 Field::new("order_id", DataType::UInt64, false),
286 Field::new("flags", DataType::UInt8, false),
287 Field::new("sequence", DataType::UInt64, false),
288 Field::new("ts_event", DataType::UInt64, false),
289 Field::new("ts_init", DataType::UInt64, false),
290 ];
291
292 let expected_schema = Schema::new_with_metadata(expected_fields, metadata);
293 assert_eq!(schema, expected_schema);
294 }
295
296 #[rstest]
297 fn test_get_schema_map() {
298 let schema_map = OrderBookDelta::get_schema_map();
299 let fixed_size_binary = format!("FixedSizeBinary({PRECISION_BYTES})");
300
301 assert_eq!(schema_map.get("action").unwrap(), "UInt8");
302 assert_eq!(schema_map.get("side").unwrap(), "UInt8");
303 assert_eq!(*schema_map.get("price").unwrap(), fixed_size_binary);
304 assert_eq!(*schema_map.get("size").unwrap(), fixed_size_binary);
305 assert_eq!(schema_map.get("order_id").unwrap(), "UInt64");
306 assert_eq!(schema_map.get("flags").unwrap(), "UInt8");
307 assert_eq!(schema_map.get("sequence").unwrap(), "UInt64");
308 assert_eq!(schema_map.get("ts_event").unwrap(), "UInt64");
309 assert_eq!(schema_map.get("ts_init").unwrap(), "UInt64");
310 }
311
312 #[rstest]
313 fn test_encode_batch() {
314 let instrument_id = InstrumentId::from("AAPL.XNAS");
315 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
316
317 let delta1 = OrderBookDelta {
318 instrument_id,
319 action: BookAction::Add,
320 order: BookOrder {
321 side: OrderSide::Buy.into(),
322 price: Price::from("100.10"),
323 size: Quantity::from(100),
324 order_id: 1,
325 },
326 flags: 0,
327 sequence: 1,
328 ts_event: 1.into(),
329 ts_init: 3.into(),
330 };
331
332 let delta2 = OrderBookDelta {
333 instrument_id,
334 action: BookAction::Update,
335 order: BookOrder {
336 side: OrderSide::Sell.into(),
337 price: Price::from("101.20"),
338 size: Quantity::from(200),
339 order_id: 2,
340 },
341 flags: 1,
342 sequence: 2,
343 ts_event: 2.into(),
344 ts_init: 4.into(),
345 };
346
347 let data = vec![delta1, delta2];
348 let record_batch = OrderBookDelta::encode_batch(&metadata, &data).unwrap();
349
350 let columns = record_batch.columns();
351 let action_values = columns[0].as_any().downcast_ref::<UInt8Array>().unwrap();
352 let side_values = columns[1].as_any().downcast_ref::<UInt8Array>().unwrap();
353 let price_values = columns[2]
354 .as_any()
355 .downcast_ref::<FixedSizeBinaryArray>()
356 .unwrap();
357 let size_values = columns[3]
358 .as_any()
359 .downcast_ref::<FixedSizeBinaryArray>()
360 .unwrap();
361 let order_id_values = columns[4].as_any().downcast_ref::<UInt64Array>().unwrap();
362 let flags_values = columns[5].as_any().downcast_ref::<UInt8Array>().unwrap();
363 let sequence_values = columns[6].as_any().downcast_ref::<UInt64Array>().unwrap();
364 let ts_event_values = columns[7].as_any().downcast_ref::<UInt64Array>().unwrap();
365 let ts_init_values = columns[8].as_any().downcast_ref::<UInt64Array>().unwrap();
366
367 assert_eq!(columns.len(), 9);
368 assert_eq!(action_values.len(), 2);
369 assert_eq!(action_values.value(0), 1);
370 assert_eq!(action_values.value(1), 2);
371 assert_eq!(side_values.len(), 2);
372 assert_eq!(side_values.value(0), 1);
373 assert_eq!(side_values.value(1), 2);
374
375 assert_eq!(price_values.len(), 2);
376 assert_eq!(
377 get_raw_price(price_values.value(0)),
378 (100.10 * FIXED_SCALAR) as PriceRaw
379 );
380 assert_eq!(
381 get_raw_price(price_values.value(1)),
382 (101.20 * FIXED_SCALAR) as PriceRaw
383 );
384
385 assert_eq!(size_values.len(), 2);
386 assert_eq!(
387 get_raw_price(size_values.value(0)),
388 (100.0 * FIXED_SCALAR) as PriceRaw
389 );
390 assert_eq!(
391 get_raw_price(size_values.value(1)),
392 (200.0 * FIXED_SCALAR) as PriceRaw
393 );
394 assert_eq!(order_id_values.len(), 2);
395 assert_eq!(order_id_values.value(0), 1);
396 assert_eq!(order_id_values.value(1), 2);
397 assert_eq!(flags_values.len(), 2);
398 assert_eq!(flags_values.value(0), 0);
399 assert_eq!(flags_values.value(1), 1);
400 assert_eq!(sequence_values.len(), 2);
401 assert_eq!(sequence_values.value(0), 1);
402 assert_eq!(sequence_values.value(1), 2);
403 assert_eq!(ts_event_values.len(), 2);
404 assert_eq!(ts_event_values.value(0), 1);
405 assert_eq!(ts_event_values.value(1), 2);
406 assert_eq!(ts_init_values.len(), 2);
407 assert_eq!(ts_init_values.value(0), 3);
408 assert_eq!(ts_init_values.value(1), 4);
409 }
410
411 #[rstest]
412 fn test_decode_batch() {
413 let instrument_id = InstrumentId::from("AAPL.XNAS");
414 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
415
416 let action = UInt8Array::from(vec![1, 2]);
417 let side = UInt8Array::from(vec![1, 1]);
418 let price = fixed_size_binary(vec![
419 &((101.10 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
420 &((101.20 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
421 ]);
422 let size = fixed_size_binary(vec![
423 &((10000.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
424 &((9000.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
425 ]);
426 let order_id = UInt64Array::from(vec![1, 2]);
427 let flags = UInt8Array::from(vec![0, 0]);
428 let sequence = UInt64Array::from(vec![1, 2]);
429 let ts_event = UInt64Array::from(vec![1, 2]);
430 let ts_init = UInt64Array::from(vec![3, 4]);
431
432 let record_batch = RecordBatch::try_new(
433 OrderBookDelta::get_schema(Some(metadata.clone())).into(),
434 vec![
435 Arc::new(action),
436 Arc::new(side),
437 Arc::new(price),
438 Arc::new(size),
439 Arc::new(order_id),
440 Arc::new(flags),
441 Arc::new(sequence),
442 Arc::new(ts_event),
443 Arc::new(ts_init),
444 ],
445 )
446 .unwrap();
447
448 let decoded_data = OrderBookDelta::decode_batch(&metadata, record_batch).unwrap();
449 assert_eq!(decoded_data.len(), 2);
450 }
451
452 #[rstest]
453 fn test_decode_batch_with_undef_values() {
454 let instrument_id = InstrumentId::from("PLTR.XNAS");
455 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
456
457 let action = UInt8Array::from(vec![4, 1]); let side = UInt8Array::from(vec![0, 1]); let price = fixed_size_binary(vec![
461 &PRICE_UNDEF.to_le_bytes(),
462 &((100.50 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
463 ]);
464 let size = fixed_size_binary(vec![
465 &QUANTITY_UNDEF.to_le_bytes(),
466 &((1000.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
467 ]);
468 let order_id = UInt64Array::from(vec![0, 1]);
469 let flags = UInt8Array::from(vec![0, 0]);
470 let sequence = UInt64Array::from(vec![1, 2]);
471 let ts_event = UInt64Array::from(vec![1, 2]);
472 let ts_init = UInt64Array::from(vec![3, 4]);
473
474 let record_batch = RecordBatch::try_new(
475 OrderBookDelta::get_schema(Some(metadata.clone())).into(),
476 vec![
477 Arc::new(action),
478 Arc::new(side),
479 Arc::new(price),
480 Arc::new(size),
481 Arc::new(order_id),
482 Arc::new(flags),
483 Arc::new(sequence),
484 Arc::new(ts_event),
485 Arc::new(ts_init),
486 ],
487 )
488 .unwrap();
489
490 let decoded_data = OrderBookDelta::decode_batch(&metadata, record_batch).unwrap();
491 assert_eq!(decoded_data.len(), 2);
492 assert_eq!(decoded_data[0].order.price.raw, PRICE_UNDEF);
493 assert_eq!(decoded_data[0].order.price.precision, 0);
494 assert_eq!(decoded_data[0].order.size.raw, QUANTITY_UNDEF);
495 assert_eq!(decoded_data[0].order.size.precision, 0);
496 assert_eq!(decoded_data[1].order.price.precision, 2);
497 assert_eq!(decoded_data[1].order.size.precision, 0);
498 }
499
500 #[rstest]
501 fn test_decode_batch_invalid_price_returns_error() {
502 let instrument_id = InstrumentId::from("AAPL.XNAS");
503 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
504
505 let action = UInt8Array::from(vec![1]);
506 let side = UInt8Array::from(vec![1]);
507
508 let invalid_price: PriceRaw = PriceRaw::MAX - 1000;
509 let price = fixed_size_binary(vec![&invalid_price.to_le_bytes()]);
510 let size = fixed_size_binary(vec![&((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes()]);
511 let order_id = UInt64Array::from(vec![1]);
512 let flags = UInt8Array::from(vec![0]);
513 let sequence = UInt64Array::from(vec![1]);
514 let ts_event = UInt64Array::from(vec![1]);
515 let ts_init = UInt64Array::from(vec![2]);
516
517 let record_batch = RecordBatch::try_new(
518 OrderBookDelta::get_schema(Some(metadata.clone())).into(),
519 vec![
520 Arc::new(action),
521 Arc::new(side),
522 Arc::new(price),
523 Arc::new(size),
524 Arc::new(order_id),
525 Arc::new(flags),
526 Arc::new(sequence),
527 Arc::new(ts_event),
528 Arc::new(ts_init),
529 ],
530 )
531 .unwrap();
532
533 let result = OrderBookDelta::decode_batch(&metadata, record_batch);
534 assert!(result.is_err());
535 let err = result.unwrap_err();
536 assert!(
537 err.to_string().contains("price") && err.to_string().contains("row 0"),
538 "Expected price error at row 0, was: {err}"
539 );
540 }
541
542 #[rstest]
543 fn test_decode_batch_invalid_action_returns_error() {
544 let instrument_id = InstrumentId::from("AAPL.XNAS");
545 let metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
546
547 let action = UInt8Array::from(vec![99]);
548 let side = UInt8Array::from(vec![1]);
549 let price = fixed_size_binary(vec![&((100.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes()]);
550 let size = fixed_size_binary(vec![&((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes()]);
551 let order_id = UInt64Array::from(vec![1]);
552 let flags = UInt8Array::from(vec![0]);
553 let sequence = UInt64Array::from(vec![1]);
554 let ts_event = UInt64Array::from(vec![1]);
555 let ts_init = UInt64Array::from(vec![2]);
556
557 let record_batch = RecordBatch::try_new(
558 OrderBookDelta::get_schema(Some(metadata.clone())).into(),
559 vec![
560 Arc::new(action),
561 Arc::new(side),
562 Arc::new(price),
563 Arc::new(size),
564 Arc::new(order_id),
565 Arc::new(flags),
566 Arc::new(sequence),
567 Arc::new(ts_event),
568 Arc::new(ts_init),
569 ],
570 )
571 .unwrap();
572
573 let result = OrderBookDelta::decode_batch(&metadata, record_batch);
574 assert!(result.is_err());
575 let err = result.unwrap_err();
576 assert!(
577 err.to_string().contains("BookAction"),
578 "Expected BookAction error, was: {err}"
579 );
580 }
581
582 #[rstest]
583 fn test_decode_batch_missing_instrument_id_returns_error() {
584 let instrument_id = InstrumentId::from("AAPL.XNAS");
585 let mut metadata = OrderBookDelta::get_metadata(&instrument_id, 2, 0);
586 metadata.remove(KEY_INSTRUMENT_ID);
587
588 let action = UInt8Array::from(vec![1]);
589 let side = UInt8Array::from(vec![1]);
590 let price = fixed_size_binary(vec![&((100.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes()]);
591 let size = fixed_size_binary(vec![&((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes()]);
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.into(),
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.into(),
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}