Skip to main content

nautilus_serialization/arrow/
delta.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
16use 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    /// Extracts metadata from the first non-clear delta, falling back to the first clear.
124    ///
125    /// Clear deltas use sentinel values whose precision does not describe the following book data.
126    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        // Create test data with 'R' (clear) action which has PRICE_UNDEF and QUANTITY_UNDEF
434        let action = UInt8Array::from(vec![4, 1]); // 4 = Clear, 1 = Add
435        let side = UInt8Array::from(vec![0, 1]); // No side for Clear, Buy for Add
436        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}