Skip to main content

nautilus_serialization/arrow/
quote.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::{Decimal128Array, UInt64Array},
20    datatypes::{DataType, Field, Schema},
21    error::ArrowError,
22    record_batch::RecordBatch,
23};
24use nautilus_model::data::QuoteTick;
25#[cfg(test)]
26use nautilus_model::identifiers::InstrumentId;
27
28use super::{
29    DecodeDataFromRecordBatch, EncodingError, KEY_IDENTIFIER, decode_required_decimal_price,
30    decode_required_decimal_quantity, decode_required_timestamp, extract_column,
31    fixed_decimal_data_type, identifier_array_from_display, parse_metadata,
32    required_price_decimal_array, required_quantity_decimal_array,
33};
34#[cfg(test)]
35use super::{KEY_INSTRUMENT_ID, KEY_PRICE_PRECISION};
36use crate::arrow::{ArrowSchemaProvider, Data, DecodeFromRecordBatch, EncodeToRecordBatch};
37
38impl ArrowSchemaProvider for QuoteTick {
39    fn get_schema(metadata: Option<HashMap<String, String>>) -> Schema {
40        let fields = vec![
41            Field::new("bid_price", fixed_decimal_data_type(), true),
42            Field::new("ask_price", fixed_decimal_data_type(), true),
43            Field::new("bid_size", fixed_decimal_data_type(), true),
44            Field::new("ask_size", fixed_decimal_data_type(), true),
45            Field::new("ts_event", crate::arrow::timestamp_data_type(), false),
46            Field::new("ts_init", crate::arrow::timestamp_data_type(), false),
47            Field::new(KEY_IDENTIFIER, DataType::Utf8, true),
48        ];
49
50        match metadata {
51            Some(metadata) => Schema::new_with_metadata(fields, metadata),
52            None => Schema::new(fields),
53        }
54    }
55}
56
57impl EncodeToRecordBatch for QuoteTick {
58    fn encode_batch<T>(
59        metadata: &HashMap<String, String>,
60        data: &[T],
61    ) -> Result<RecordBatch, ArrowError>
62    where
63        T: std::borrow::Borrow<Self>,
64    {
65        let mut ts_event_builder = UInt64Array::builder(data.len());
66        let mut ts_init_builder = UInt64Array::builder(data.len());
67
68        for quote in data.iter().map(std::borrow::Borrow::borrow) {
69            ts_event_builder.append_value(quote.ts_event.as_u64());
70            ts_init_builder.append_value(quote.ts_init.as_u64());
71        }
72
73        crate::arrow::record_batch_with_timestamps(
74            Self::get_schema(Some(metadata.clone())).into(),
75            vec![
76                Arc::new(required_price_decimal_array(
77                    data.iter()
78                        .map(std::borrow::Borrow::borrow)
79                        .map(|quote| quote.bid_price.raw()),
80                    "bid_price",
81                )?),
82                Arc::new(required_price_decimal_array(
83                    data.iter()
84                        .map(std::borrow::Borrow::borrow)
85                        .map(|quote| quote.ask_price.raw()),
86                    "ask_price",
87                )?),
88                Arc::new(required_quantity_decimal_array(
89                    data.iter()
90                        .map(std::borrow::Borrow::borrow)
91                        .map(|quote| quote.bid_size.raw()),
92                    "bid_size",
93                )?),
94                Arc::new(required_quantity_decimal_array(
95                    data.iter()
96                        .map(std::borrow::Borrow::borrow)
97                        .map(|quote| quote.ask_size.raw()),
98                    "ask_size",
99                )?),
100                Arc::new(ts_event_builder.finish()),
101                Arc::new(ts_init_builder.finish()),
102                Arc::new(identifier_array_from_display(
103                    data.iter()
104                        .map(std::borrow::Borrow::borrow)
105                        .map(|quote| quote.instrument_id),
106                )),
107            ],
108        )
109    }
110
111    fn metadata(&self) -> HashMap<String, String> {
112        Self::get_metadata(
113            &self.instrument_id,
114            self.bid_price.precision,
115            self.bid_size.precision,
116        )
117    }
118}
119
120impl DecodeFromRecordBatch for QuoteTick {
121    fn decode_batch(
122        metadata: &HashMap<String, String>,
123        record_batch: RecordBatch,
124    ) -> Result<Vec<Self>, EncodingError> {
125        let (instrument_id, price_precision, size_precision) = parse_metadata(metadata)?;
126        let record_batch = crate::arrow::record_batch_with_u64_timestamps(&record_batch)?;
127        let record_batch = &record_batch;
128        let cols = record_batch.columns();
129
130        let bid_price_values =
131            extract_column::<Decimal128Array>(cols, "bid_price", 0, fixed_decimal_data_type())?;
132        let ask_price_values =
133            extract_column::<Decimal128Array>(cols, "ask_price", 1, fixed_decimal_data_type())?;
134        let bid_size_values =
135            extract_column::<Decimal128Array>(cols, "bid_size", 2, fixed_decimal_data_type())?;
136        let ask_size_values =
137            extract_column::<Decimal128Array>(cols, "ask_size", 3, fixed_decimal_data_type())?;
138        let ts_event_values = extract_column::<UInt64Array>(cols, "ts_event", 4, DataType::UInt64)?;
139        let ts_init_values = extract_column::<UInt64Array>(cols, "ts_init", 5, DataType::UInt64)?;
140
141        let result: Result<Vec<Self>, EncodingError> = (0..record_batch.num_rows())
142            .map(|row| {
143                let bid_price = decode_required_decimal_price(
144                    bid_price_values,
145                    price_precision,
146                    "bid_price",
147                    row,
148                )?;
149                let ask_price = decode_required_decimal_price(
150                    ask_price_values,
151                    price_precision,
152                    "ask_price",
153                    row,
154                )?;
155                let bid_size = decode_required_decimal_quantity(
156                    bid_size_values,
157                    size_precision,
158                    "bid_size",
159                    row,
160                )?;
161                let ask_size = decode_required_decimal_quantity(
162                    ask_size_values,
163                    size_precision,
164                    "ask_size",
165                    row,
166                )?;
167                Ok(Self {
168                    instrument_id,
169                    bid_price,
170                    ask_price,
171                    bid_size,
172                    ask_size,
173                    ts_event: decode_required_timestamp(ts_event_values, "ts_event", row)?,
174                    ts_init: decode_required_timestamp(ts_init_values, "ts_init", row)?,
175                })
176            })
177            .collect();
178
179        result
180    }
181}
182
183impl DecodeDataFromRecordBatch for QuoteTick {
184    fn decode_data_batch(
185        metadata: &HashMap<String, String>,
186        record_batch: RecordBatch,
187    ) -> Result<Vec<Data>, EncodingError> {
188        let ticks: Vec<Self> = Self::decode_batch(metadata, record_batch)?;
189        Ok(ticks.into_iter().map(Data::from).collect())
190    }
191}
192
193#[cfg(test)]
194mod tests {
195    use std::{collections::HashMap, sync::Arc};
196
197    use arrow::array::{Array, StringArray, TimestampNanosecondArray};
198    use nautilus_model::types::{
199        Price, Quantity, fixed::FIXED_SCALAR, price::PriceRaw, quantity::QuantityRaw,
200    };
201    use rstest::rstest;
202
203    use super::*;
204    use crate::arrow::{KEY_IDENTIFIER, get_raw_price, get_raw_quantity};
205
206    #[rstest]
207    fn test_quote_nanoseconds_round_trip_within_one_microsecond() {
208        let first = QuoteTick {
209            instrument_id: InstrumentId::from("AAPL.XNAS"),
210            bid_price: Price::from("123.45"),
211            ask_price: Price::from("123.67"),
212            bid_size: Quantity::from(17),
213            ask_size: Quantity::from(29),
214            ts_event: 1_788_652_800_123_456_789_u64.into(),
215            ts_init: 1_788_652_800_123_456_799_u64.into(),
216        };
217        let second = QuoteTick {
218            ts_event: 1_788_652_800_123_456_801_u64.into(),
219            ts_init: 1_788_652_800_123_456_899_u64.into(),
220            ..first
221        };
222        let values = vec![first, second];
223        let metadata = QuoteTick::get_metadata(&first.instrument_id, 2, 0);
224        let batch = QuoteTick::encode_batch(&metadata, &values).unwrap();
225
226        for field in ["ts_event", "ts_init"] {
227            assert_eq!(
228                batch.schema().field_with_name(field).unwrap().data_type(),
229                &DataType::Timestamp(arrow::datatypes::TimeUnit::Nanosecond, Some("UTC".into()))
230            );
231        }
232        let decoded = QuoteTick::decode_batch(&metadata, batch).unwrap();
233        assert_eq!(decoded, values);
234    }
235
236    #[rstest]
237    #[case(38, 2)]
238    #[case(38, 9)]
239    #[case(38, 18)]
240    #[case(37, 16)]
241    fn test_decode_rejects_decimal_type(
242        #[case] precision: u8,
243        #[case] scale: i8,
244        #[values(0, 1, 2, 3)] column: usize,
245    ) {
246        let quote = QuoteTick {
247            instrument_id: InstrumentId::from("AAPL.XNAS"),
248            bid_price: Price::from("123.45"),
249            ask_price: Price::from("123.67"),
250            bid_size: Quantity::from(17),
251            ask_size: Quantity::from(29),
252            ts_event: 1.into(),
253            ts_init: 2.into(),
254        };
255        let metadata = quote.metadata();
256        let batch = QuoteTick::encode_batch(&metadata, &[quote]).unwrap();
257        let mut columns = batch.columns().to_vec();
258        let values = columns[column]
259            .as_any()
260            .downcast_ref::<Decimal128Array>()
261            .unwrap()
262            .clone()
263            .with_precision_and_scale(precision, scale)
264            .unwrap();
265        columns[column] = Arc::new(values);
266        let mut fields = batch.schema().fields().to_vec();
267        let name = fields[column].name().clone();
268        fields[column] = Arc::new(
269            fields[column]
270                .as_ref()
271                .clone()
272                .with_data_type(DataType::Decimal128(precision, scale)),
273        );
274        let batch = RecordBatch::try_new(Arc::new(Schema::new(fields)), columns).unwrap();
275
276        let error = QuoteTick::decode_batch(&metadata, batch).unwrap_err();
277
278        match error {
279            EncodingError::InvalidColumnType(field, index, expected, actual) => {
280                assert_eq!(field, name);
281                assert_eq!(index, column);
282                assert_eq!(expected, fixed_decimal_data_type());
283                assert_eq!(actual, DataType::Decimal128(precision, scale));
284            }
285            error => panic!("Unexpected error: {error}"),
286        }
287    }
288
289    #[rstest]
290    fn test_get_schema() {
291        let instrument_id = InstrumentId::from("AAPL.XNAS");
292        let metadata = QuoteTick::get_metadata(&instrument_id, 2, 0);
293        let schema = QuoteTick::get_schema(Some(metadata.clone()));
294
295        let mut expected_fields = Vec::with_capacity(7);
296
297        expected_fields.push(Field::new("bid_price", fixed_decimal_data_type(), true));
298        expected_fields.push(Field::new("ask_price", fixed_decimal_data_type(), true));
299
300        expected_fields.extend(vec![
301            Field::new("bid_size", fixed_decimal_data_type(), true),
302            Field::new("ask_size", fixed_decimal_data_type(), true),
303            Field::new("ts_event", crate::arrow::timestamp_data_type(), false),
304            Field::new("ts_init", crate::arrow::timestamp_data_type(), false),
305            Field::new(KEY_IDENTIFIER, DataType::Utf8, true),
306        ]);
307
308        let expected_schema = Schema::new_with_metadata(expected_fields, metadata);
309        assert_eq!(schema, expected_schema);
310    }
311
312    #[rstest]
313    fn test_get_schema_map() {
314        let arrow_schema = QuoteTick::get_schema_map();
315        let mut expected_map = HashMap::new();
316
317        let fixed_size_binary = "Decimal128(38, 16)".to_string();
318        expected_map.insert("bid_price".to_string(), fixed_size_binary.clone());
319        expected_map.insert("ask_price".to_string(), fixed_size_binary.clone());
320        expected_map.insert("bid_size".to_string(), fixed_size_binary.clone());
321        expected_map.insert("ask_size".to_string(), fixed_size_binary);
322        expected_map.insert(
323            "ts_event".to_string(),
324            "Timestamp(Nanosecond, Some(\"UTC\"))".to_string(),
325        );
326        expected_map.insert(
327            "ts_init".to_string(),
328            "Timestamp(Nanosecond, Some(\"UTC\"))".to_string(),
329        );
330        expected_map.insert(KEY_IDENTIFIER.to_string(), "Utf8".to_string());
331        assert_eq!(arrow_schema, expected_map);
332    }
333
334    #[rstest]
335    fn test_encode_quote_tick() {
336        // Create test data
337        let instrument_id = InstrumentId::from("AAPL.XNAS");
338        let tick1 = QuoteTick {
339            instrument_id,
340            bid_price: Price::from("100.10"),
341            ask_price: Price::from("101.50"),
342            bid_size: Quantity::from(1000),
343            ask_size: Quantity::from(500),
344            ts_event: 1.into(),
345            ts_init: 3.into(),
346        };
347
348        let tick2 = QuoteTick {
349            instrument_id,
350            bid_price: Price::from("100.75"),
351            ask_price: Price::from("100.20"),
352            bid_size: Quantity::from(750),
353            ask_size: Quantity::from(300),
354            ts_event: 2.into(),
355            ts_init: 4.into(),
356        };
357
358        let data = vec![tick1, tick2];
359        let metadata = QuoteTick::get_metadata(&instrument_id, 2, 0);
360        let record_batch = QuoteTick::encode_batch(&metadata, &data).unwrap();
361
362        // Verify the encoded data
363        let columns = record_batch.columns();
364
365        let bid_price_values = columns[0]
366            .as_any()
367            .downcast_ref::<Decimal128Array>()
368            .unwrap();
369        let ask_price_values = columns[1]
370            .as_any()
371            .downcast_ref::<Decimal128Array>()
372            .unwrap();
373        assert_eq!(
374            get_raw_price(bid_price_values.value(0)),
375            (100.10 * FIXED_SCALAR) as PriceRaw
376        );
377        assert_eq!(
378            get_raw_price(bid_price_values.value(1)),
379            (100.75 * FIXED_SCALAR) as PriceRaw
380        );
381        assert_eq!(
382            get_raw_price(ask_price_values.value(0)),
383            (101.50 * FIXED_SCALAR) as PriceRaw
384        );
385        assert_eq!(
386            get_raw_price(ask_price_values.value(1)),
387            (100.20 * FIXED_SCALAR) as PriceRaw
388        );
389
390        let bid_size_values = columns[2]
391            .as_any()
392            .downcast_ref::<Decimal128Array>()
393            .unwrap();
394        let ask_size_values = columns[3]
395            .as_any()
396            .downcast_ref::<Decimal128Array>()
397            .unwrap();
398        let ts_event_values = columns[4]
399            .as_any()
400            .downcast_ref::<TimestampNanosecondArray>()
401            .unwrap();
402        let ts_init_values = columns[5]
403            .as_any()
404            .downcast_ref::<TimestampNanosecondArray>()
405            .unwrap();
406
407        let identifier_values = columns[6].as_any().downcast_ref::<StringArray>().unwrap();
408
409        assert_eq!(columns.len(), 7);
410        assert_eq!(bid_size_values.len(), 2);
411        assert_eq!(
412            get_raw_quantity(bid_size_values.value(0)),
413            (1000.0 * FIXED_SCALAR) as QuantityRaw
414        );
415        assert_eq!(
416            get_raw_quantity(bid_size_values.value(1)),
417            (750.0 * FIXED_SCALAR) as QuantityRaw
418        );
419        assert_eq!(ask_size_values.len(), 2);
420        assert_eq!(
421            get_raw_quantity(ask_size_values.value(0)),
422            (500.0 * FIXED_SCALAR) as QuantityRaw
423        );
424        assert_eq!(
425            get_raw_quantity(ask_size_values.value(1)),
426            (300.0 * FIXED_SCALAR) as QuantityRaw
427        );
428        assert_eq!(ts_event_values.len(), 2);
429        assert_eq!(ts_event_values.value(0), 1);
430        assert_eq!(ts_event_values.value(1), 2);
431        assert_eq!(ts_init_values.len(), 2);
432        assert_eq!(ts_init_values.value(0), 3);
433        assert_eq!(ts_init_values.value(1), 4);
434        assert_eq!(record_batch.schema().field(6).name(), KEY_IDENTIFIER);
435        assert_eq!(identifier_values.len(), 2);
436        assert_eq!(identifier_values.value(0), "AAPL.XNAS");
437        assert_eq!(identifier_values.value(1), "AAPL.XNAS");
438    }
439
440    #[rstest]
441    fn test_decode_batch() {
442        let instrument_id = InstrumentId::from("AAPL.XNAS");
443        let metadata = QuoteTick::get_metadata(&instrument_id, 2, 0);
444
445        let raw_bid1 = (100.00 * FIXED_SCALAR) as PriceRaw;
446        let raw_bid2 = (99.00 * FIXED_SCALAR) as PriceRaw;
447        let raw_ask1 = (101.00 * FIXED_SCALAR) as PriceRaw;
448        let raw_ask2 = (100.00 * FIXED_SCALAR) as PriceRaw;
449
450        let (bid_price, ask_price) = (
451            crate::arrow::test_support::decimal_array_from_bytes(vec![
452                &raw_bid1.to_le_bytes(),
453                &raw_bid2.to_le_bytes(),
454            ]),
455            crate::arrow::test_support::decimal_array_from_bytes(vec![
456                &raw_ask1.to_le_bytes(),
457                &raw_ask2.to_le_bytes(),
458            ]),
459        );
460
461        let bid_size = crate::arrow::test_support::decimal_array_from_bytes(vec![
462            &((100.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
463            &((90.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
464        ]);
465        let ask_size = crate::arrow::test_support::decimal_array_from_bytes(vec![
466            &((110.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
467            &((100.0 * FIXED_SCALAR) as PriceRaw).to_le_bytes(),
468        ]);
469        let ts_event = UInt64Array::from(vec![1, 2]);
470        let ts_init = UInt64Array::from(vec![3, 4]);
471
472        let record_batch = crate::arrow::record_batch_with_timestamps(
473            crate::arrow::schema_without_identifier_column(&QuoteTick::get_schema(Some(
474                metadata.clone(),
475            )))
476            .into(),
477            vec![
478                Arc::new(bid_price),
479                Arc::new(ask_price),
480                Arc::new(bid_size),
481                Arc::new(ask_size),
482                Arc::new(ts_event),
483                Arc::new(ts_init),
484            ],
485        )
486        .unwrap();
487
488        let decoded_data = QuoteTick::decode_batch(&metadata, record_batch).unwrap();
489        assert_eq!(decoded_data.len(), 2);
490
491        // Verify decoded values
492        assert_eq!(decoded_data[0].bid_price, Price::from_raw(raw_bid1, 2));
493        assert_eq!(decoded_data[0].ask_price, Price::from_raw(raw_ask1, 2));
494        assert_eq!(decoded_data[1].bid_price, Price::from_raw(raw_bid2, 2));
495        assert_eq!(decoded_data[1].ask_price, Price::from_raw(raw_ask2, 2));
496    }
497
498    #[rstest]
499    fn test_decode_batch_rejects_null_timestamp_with_field_and_row() {
500        let instrument_id = InstrumentId::from("AAPL.XNAS");
501        let metadata = QuoteTick::get_metadata(&instrument_id, 2, 0);
502        let quote = QuoteTick::new(
503            instrument_id,
504            Price::from("100.00"),
505            Price::from("101.00"),
506            Quantity::from(10),
507            Quantity::from(20),
508            1.into(),
509            2.into(),
510        );
511        let encoded = QuoteTick::encode_batch(&metadata, &[quote]).unwrap();
512        let mut columns = encoded.columns().to_vec();
513        columns[5] = Arc::new(TimestampNanosecondArray::from(vec![None]).with_timezone("UTC"));
514        let fields = encoded
515            .schema()
516            .fields()
517            .iter()
518            .map(|field| {
519                if field.name() == "ts_init" {
520                    Arc::new(field.as_ref().clone().with_nullable(true))
521                } else {
522                    field.clone()
523                }
524            })
525            .collect::<Vec<_>>();
526        let schema = Arc::new(Schema::new_with_metadata(fields, metadata.clone()));
527        let batch = RecordBatch::try_new(schema, columns).unwrap();
528
529        let error = QuoteTick::decode_batch(&metadata, batch).unwrap_err();
530
531        assert!(error.to_string().contains("ts_init"));
532        assert!(error.to_string().contains("row 0"));
533    }
534
535    #[rstest]
536    fn test_decode_batch_invalid_bid_price_returns_error() {
537        let instrument_id = InstrumentId::from("AAPL.XNAS");
538        let metadata = QuoteTick::get_metadata(&instrument_id, 2, 0);
539
540        let invalid_price: PriceRaw = PriceRaw::MAX - 1000;
541        let valid_price = (100.00 * FIXED_SCALAR) as PriceRaw;
542
543        let bid_price = crate::arrow::test_support::decimal_array_from_bytes(vec![
544            &invalid_price.to_le_bytes(),
545        ]);
546        let ask_price =
547            crate::arrow::test_support::decimal_array_from_bytes(vec![&valid_price.to_le_bytes()]);
548        let bid_size = crate::arrow::test_support::decimal_array_from_bytes(vec![
549            &((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes(),
550        ]);
551        let ask_size = crate::arrow::test_support::decimal_array_from_bytes(vec![
552            &((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes(),
553        ]);
554        let ts_event = UInt64Array::from(vec![1]);
555        let ts_init = UInt64Array::from(vec![2]);
556
557        let record_batch = crate::arrow::record_batch_with_timestamps(
558            crate::arrow::schema_without_identifier_column(&QuoteTick::get_schema(Some(
559                metadata.clone(),
560            )))
561            .into(),
562            vec![
563                Arc::new(bid_price),
564                Arc::new(ask_price),
565                Arc::new(bid_size),
566                Arc::new(ask_size),
567                Arc::new(ts_event),
568                Arc::new(ts_init),
569            ],
570        )
571        .unwrap();
572
573        let result = QuoteTick::decode_batch(&metadata, record_batch);
574        assert!(result.is_err());
575        let err = result.unwrap_err();
576        assert!(
577            err.to_string().contains("bid_price") && err.to_string().contains("row 0"),
578            "Expected bid_price error at row 0, was: {err}"
579        );
580    }
581
582    #[rstest]
583    fn test_decode_batch_invalid_ask_size_returns_error() {
584        use nautilus_model::types::{fixed::FIXED_PRECISION, quantity::QUANTITY_RAW_MAX};
585
586        let instrument_id = InstrumentId::from("AAPL.XNAS");
587        // Decode the size at full precision so the out-of-range raw value bypasses the
588        // precision-0 correction, which would otherwise round it back within the bound.
589        let metadata = QuoteTick::get_metadata(&instrument_id, 2, FIXED_PRECISION);
590
591        let valid_price = (100.00 * FIXED_SCALAR) as PriceRaw;
592        let bid_price =
593            crate::arrow::test_support::decimal_array_from_bytes(vec![&valid_price.to_le_bytes()]);
594        let ask_price =
595            crate::arrow::test_support::decimal_array_from_bytes(vec![&valid_price.to_le_bytes()]);
596        let bid_size = crate::arrow::test_support::decimal_array_from_bytes(vec![
597            &((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes(),
598        ]);
599
600        let invalid_size = QUANTITY_RAW_MAX + 1;
601        let ask_size =
602            crate::arrow::test_support::decimal_array_from_bytes(vec![&invalid_size.to_le_bytes()]);
603        let ts_event = UInt64Array::from(vec![1]);
604        let ts_init = UInt64Array::from(vec![2]);
605
606        let record_batch = crate::arrow::record_batch_with_timestamps(
607            crate::arrow::schema_without_identifier_column(&QuoteTick::get_schema(Some(
608                metadata.clone(),
609            )))
610            .into(),
611            vec![
612                Arc::new(bid_price),
613                Arc::new(ask_price),
614                Arc::new(bid_size),
615                Arc::new(ask_size),
616                Arc::new(ts_event),
617                Arc::new(ts_init),
618            ],
619        )
620        .unwrap();
621
622        let result = QuoteTick::decode_batch(&metadata, record_batch);
623        assert!(result.is_err());
624        let err = result.unwrap_err();
625        assert!(
626            err.to_string().contains("ask_size") && err.to_string().contains("row 0"),
627            "Expected ask_size error at row 0, was: {err}"
628        );
629    }
630
631    #[rstest]
632    fn test_decode_batch_missing_instrument_id_returns_error() {
633        let instrument_id = InstrumentId::from("AAPL.XNAS");
634        let mut metadata = QuoteTick::get_metadata(&instrument_id, 2, 0);
635        metadata.remove(KEY_INSTRUMENT_ID);
636
637        let valid_price = (100.00 * FIXED_SCALAR) as PriceRaw;
638        let bid_price =
639            crate::arrow::test_support::decimal_array_from_bytes(vec![&valid_price.to_le_bytes()]);
640        let ask_price =
641            crate::arrow::test_support::decimal_array_from_bytes(vec![&valid_price.to_le_bytes()]);
642        let bid_size = crate::arrow::test_support::decimal_array_from_bytes(vec![
643            &((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes(),
644        ]);
645        let ask_size = crate::arrow::test_support::decimal_array_from_bytes(vec![
646            &((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes(),
647        ]);
648        let ts_event = UInt64Array::from(vec![1]);
649        let ts_init = UInt64Array::from(vec![2]);
650
651        let record_batch = crate::arrow::record_batch_with_timestamps(
652            crate::arrow::schema_without_identifier_column(&QuoteTick::get_schema(Some(
653                metadata.clone(),
654            )))
655            .into(),
656            vec![
657                Arc::new(bid_price),
658                Arc::new(ask_price),
659                Arc::new(bid_size),
660                Arc::new(ask_size),
661                Arc::new(ts_event),
662                Arc::new(ts_init),
663            ],
664        )
665        .unwrap();
666
667        let result = QuoteTick::decode_batch(&metadata, record_batch);
668        assert!(result.is_err());
669        let err = result.unwrap_err();
670        assert!(
671            err.to_string().contains("instrument_id"),
672            "Expected missing instrument_id error, was: {err}"
673        );
674    }
675
676    #[rstest]
677    fn test_decode_batch_missing_price_precision_returns_error() {
678        let instrument_id = InstrumentId::from("AAPL.XNAS");
679        let mut metadata = QuoteTick::get_metadata(&instrument_id, 2, 0);
680        metadata.remove(KEY_PRICE_PRECISION);
681
682        let valid_price = (100.00 * FIXED_SCALAR) as PriceRaw;
683        let bid_price =
684            crate::arrow::test_support::decimal_array_from_bytes(vec![&valid_price.to_le_bytes()]);
685        let ask_price =
686            crate::arrow::test_support::decimal_array_from_bytes(vec![&valid_price.to_le_bytes()]);
687        let bid_size = crate::arrow::test_support::decimal_array_from_bytes(vec![
688            &((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes(),
689        ]);
690        let ask_size = crate::arrow::test_support::decimal_array_from_bytes(vec![
691            &((100.0 * FIXED_SCALAR) as QuantityRaw).to_le_bytes(),
692        ]);
693        let ts_event = UInt64Array::from(vec![1]);
694        let ts_init = UInt64Array::from(vec![2]);
695
696        let record_batch = crate::arrow::record_batch_with_timestamps(
697            crate::arrow::schema_without_identifier_column(&QuoteTick::get_schema(Some(
698                metadata.clone(),
699            )))
700            .into(),
701            vec![
702                Arc::new(bid_price),
703                Arc::new(ask_price),
704                Arc::new(bid_size),
705                Arc::new(ask_size),
706                Arc::new(ts_event),
707                Arc::new(ts_init),
708            ],
709        )
710        .unwrap();
711
712        let result = QuoteTick::decode_batch(&metadata, record_batch);
713        assert!(result.is_err());
714        let err = result.unwrap_err();
715        assert!(
716            err.to_string().contains("price_precision"),
717            "Expected missing price_precision error, was: {err}"
718        );
719    }
720
721    #[rstest]
722    fn test_encode_decode_round_trip() {
723        let instrument_id = InstrumentId::from("AAPL.XNAS");
724        let metadata = QuoteTick::get_metadata(&instrument_id, 2, 0);
725
726        let tick1 = QuoteTick {
727            instrument_id,
728            bid_price: Price::from("100.10"),
729            ask_price: Price::from("100.20"),
730            bid_size: Quantity::from(1000),
731            ask_size: Quantity::from(500),
732            ts_event: 1_000_000_000.into(),
733            ts_init: 1_000_000_001.into(),
734        };
735
736        let tick2 = QuoteTick {
737            instrument_id,
738            bid_price: Price::from("100.15"),
739            ask_price: Price::from("100.25"),
740            bid_size: Quantity::from(750),
741            ask_size: Quantity::from(250),
742            ts_event: 2_000_000_000.into(),
743            ts_init: 2_000_000_001.into(),
744        };
745
746        let original = vec![tick1, tick2];
747        let record_batch = QuoteTick::encode_batch(&metadata, &original).unwrap();
748        let decoded = QuoteTick::decode_batch(&metadata, record_batch).unwrap();
749
750        assert_eq!(decoded.len(), original.len());
751        for (orig, dec) in original.iter().zip(decoded.iter()) {
752            assert_eq!(dec.instrument_id, orig.instrument_id);
753            assert_eq!(dec.bid_price, orig.bid_price);
754            assert_eq!(dec.ask_price, orig.ask_price);
755            assert_eq!(dec.bid_size, orig.bid_size);
756            assert_eq!(dec.ask_size, orig.ask_size);
757            assert_eq!(dec.ts_event, orig.ts_event);
758            assert_eq!(dec.ts_init, orig.ts_init);
759        }
760    }
761}