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, 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    /// Extracts metadata from the first non-clear delta, falling back to the first clear.
149    ///
150    /// Clear deltas use sentinel values whose precision does not describe the following book data.
151    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        // Create test data with 'R' (clear) action which has PRICE_UNDEF and QUANTITY_UNDEF
450        let action = UInt8Array::from(vec![4, 1]); // 4 = Clear, 1 = Add
451        let side = UInt8Array::from(vec![0, 1]); // NoOrderSide for Clear, Buy for Add
452        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}