1use 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 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 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 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 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}