1use std::sync::Arc;
19
20use arrow::{
21 array::{
22 Float64Builder, StringBuilder, TimestampNanosecondBuilder, UInt8Builder, UInt64Builder,
23 },
24 datatypes::{DataType, Field, Schema},
25 error::ArrowError,
26 record_batch::RecordBatch,
27};
28use nautilus_model::{data::OrderBookDelta, enums::BookAction};
29
30use super::{
31 float64_field, price_to_f64, quantity_to_f64, timestamp_field, unix_nanos_to_i64, utf8_field,
32};
33
34#[must_use]
36pub fn deltas_schema() -> Schema {
37 Schema::new(vec![
38 utf8_field("instrument_id", false),
39 utf8_field("action", false),
40 utf8_field("side", false),
41 float64_field("price", false),
42 float64_field("size", false),
43 utf8_field("order_id", false),
44 Field::new("flags", DataType::UInt8, false),
45 Field::new("sequence", DataType::UInt64, false),
46 timestamp_field("ts_event", false),
47 timestamp_field("ts_init", false),
48 ])
49}
50
51pub fn encode_deltas(data: &[OrderBookDelta]) -> Result<RecordBatch, ArrowError> {
66 let mut instrument_id_builder = StringBuilder::new();
67 let mut action_builder = StringBuilder::new();
68 let mut side_builder = StringBuilder::new();
69 let mut price_builder = Float64Builder::with_capacity(data.len());
70 let mut size_builder = Float64Builder::with_capacity(data.len());
71 let mut order_id_builder = StringBuilder::new();
72 let mut flags_builder = UInt8Builder::with_capacity(data.len());
73 let mut sequence_builder = UInt64Builder::with_capacity(data.len());
74 let mut ts_event_builder = TimestampNanosecondBuilder::with_capacity(data.len());
75 let mut ts_init_builder = TimestampNanosecondBuilder::with_capacity(data.len());
76
77 for delta in data {
78 instrument_id_builder.append_value(delta.instrument_id.to_string());
79 action_builder.append_value(format!("{}", delta.action));
80 side_builder.append_value(
81 delta
82 .order
83 .side
84 .as_ref()
85 .map_or("NO_ORDER_SIDE", AsRef::as_ref),
86 );
87
88 if delta.action == BookAction::Clear {
92 price_builder.append_value(f64::NAN);
93 size_builder.append_value(f64::NAN);
94 } else {
95 price_builder.append_value(price_to_f64(&delta.order.price));
96 size_builder.append_value(quantity_to_f64(&delta.order.size));
97 }
98 order_id_builder.append_value(delta.order.order_id.to_string());
99 flags_builder.append_value(delta.flags);
100 sequence_builder.append_value(delta.sequence);
101 ts_event_builder.append_value(unix_nanos_to_i64(delta.ts_event.as_u64()));
102 ts_init_builder.append_value(unix_nanos_to_i64(delta.ts_init.as_u64()));
103 }
104
105 RecordBatch::try_new(
106 Arc::new(deltas_schema()),
107 vec![
108 Arc::new(instrument_id_builder.finish()),
109 Arc::new(action_builder.finish()),
110 Arc::new(side_builder.finish()),
111 Arc::new(price_builder.finish()),
112 Arc::new(size_builder.finish()),
113 Arc::new(order_id_builder.finish()),
114 Arc::new(flags_builder.finish()),
115 Arc::new(sequence_builder.finish()),
116 Arc::new(ts_event_builder.finish()),
117 Arc::new(ts_init_builder.finish()),
118 ],
119 )
120}
121
122#[cfg(test)]
123mod tests {
124 use arrow::{
125 array::{
126 Array, Float64Array, StringArray, TimestampNanosecondArray, UInt8Array, UInt64Array,
127 },
128 datatypes::TimeUnit,
129 };
130 use nautilus_model::{
131 data::order::BookOrder,
132 enums::{BookAction, OrderSide},
133 identifiers::InstrumentId,
134 types::{Price, Quantity, price::PRICE_UNDEF, quantity::QUANTITY_UNDEF},
135 };
136 use rstest::rstest;
137
138 use super::*;
139
140 fn make_delta(
141 instrument_id: &str,
142 action: BookAction,
143 side: OrderSide,
144 price: &str,
145 order_id: u64,
146 sequence: u64,
147 ts: u64,
148 ) -> OrderBookDelta {
149 OrderBookDelta {
150 instrument_id: InstrumentId::from(instrument_id),
151 action,
152 order: BookOrder {
153 side: Some(side),
154 price: Price::from(price),
155 size: Quantity::from(100),
156 order_id,
157 },
158 flags: 0,
159 sequence,
160 ts_event: ts.into(),
161 ts_init: (ts + 1).into(),
162 }
163 }
164
165 #[rstest]
166 fn test_encode_deltas_schema() {
167 let batch = encode_deltas(&[]).unwrap();
168 let fields = batch.schema().fields().clone();
169 assert_eq!(fields.len(), 10);
170 assert_eq!(fields[0].name(), "instrument_id");
171 assert_eq!(fields[0].data_type(), &DataType::Utf8);
172 assert_eq!(fields[1].name(), "action");
173 assert_eq!(fields[1].data_type(), &DataType::Utf8);
174 assert_eq!(fields[2].name(), "side");
175 assert_eq!(fields[2].data_type(), &DataType::Utf8);
176 assert_eq!(fields[3].name(), "price");
177 assert_eq!(fields[3].data_type(), &DataType::Float64);
178 assert_eq!(fields[4].name(), "size");
179 assert_eq!(fields[4].data_type(), &DataType::Float64);
180 assert_eq!(fields[5].name(), "order_id");
181 assert_eq!(fields[5].data_type(), &DataType::Utf8);
182 assert_eq!(fields[6].name(), "flags");
183 assert_eq!(fields[6].data_type(), &DataType::UInt8);
184 assert_eq!(fields[7].name(), "sequence");
185 assert_eq!(fields[7].data_type(), &DataType::UInt64);
186 assert_eq!(fields[8].name(), "ts_event");
187 assert_eq!(
188 fields[8].data_type(),
189 &DataType::Timestamp(TimeUnit::Nanosecond, None)
190 );
191 assert_eq!(fields[9].name(), "ts_init");
192 }
193
194 #[rstest]
195 fn test_encode_deltas_values() {
196 let deltas = vec![
197 make_delta(
198 "AAPL.XNAS",
199 BookAction::Add,
200 OrderSide::Buy,
201 "100.10",
202 1,
203 10,
204 1_000,
205 ),
206 make_delta(
207 "AAPL.XNAS",
208 BookAction::Update,
209 OrderSide::Sell,
210 "100.20",
211 2,
212 11,
213 2_000,
214 ),
215 ];
216 let batch = encode_deltas(&deltas).unwrap();
217
218 assert_eq!(batch.num_rows(), 2);
219
220 let action_col = batch
221 .column(1)
222 .as_any()
223 .downcast_ref::<StringArray>()
224 .unwrap();
225 let side_col = batch
226 .column(2)
227 .as_any()
228 .downcast_ref::<StringArray>()
229 .unwrap();
230 let price_col = batch
231 .column(3)
232 .as_any()
233 .downcast_ref::<Float64Array>()
234 .unwrap();
235 let size_col = batch
236 .column(4)
237 .as_any()
238 .downcast_ref::<Float64Array>()
239 .unwrap();
240 let order_id_col = batch
241 .column(5)
242 .as_any()
243 .downcast_ref::<StringArray>()
244 .unwrap();
245 let flags_col = batch
246 .column(6)
247 .as_any()
248 .downcast_ref::<UInt8Array>()
249 .unwrap();
250 let sequence_col = batch
251 .column(7)
252 .as_any()
253 .downcast_ref::<UInt64Array>()
254 .unwrap();
255 let ts_event_col = batch
256 .column(8)
257 .as_any()
258 .downcast_ref::<TimestampNanosecondArray>()
259 .unwrap();
260
261 assert_eq!(action_col.value(0), format!("{}", BookAction::Add));
262 assert_eq!(action_col.value(1), format!("{}", BookAction::Update));
263 assert_eq!(side_col.value(0), format!("{}", OrderSide::Buy));
264 assert_eq!(side_col.value(1), format!("{}", OrderSide::Sell));
265 assert!((price_col.value(0) - 100.10).abs() < 1e-9);
266 assert!((price_col.value(1) - 100.20).abs() < 1e-9);
267 assert!((size_col.value(0) - 100.0).abs() < 1e-9);
268 assert_eq!(order_id_col.value(0), "1");
269 assert_eq!(order_id_col.value(1), "2");
270 assert_eq!(flags_col.value(0), 0);
271 assert_eq!(sequence_col.value(0), 10);
272 assert_eq!(sequence_col.value(1), 11);
273 assert_eq!(ts_event_col.value(0), 1_000);
274 }
275
276 #[rstest]
277 fn test_encode_deltas_empty() {
278 let batch = encode_deltas(&[]).unwrap();
279 assert_eq!(batch.num_rows(), 0);
280 }
281
282 #[rstest]
283 fn test_encode_deltas_live_clear_renders_as_nan() {
284 let clear = OrderBookDelta::clear(InstrumentId::from("AAPL.XNAS"), 1, 1.into(), 2.into());
285
286 let batch = encode_deltas(&[clear]).unwrap();
287 let side_col = batch
288 .column(2)
289 .as_any()
290 .downcast_ref::<StringArray>()
291 .unwrap();
292 let price_col = batch
293 .column(3)
294 .as_any()
295 .downcast_ref::<Float64Array>()
296 .unwrap();
297 let size_col = batch
298 .column(4)
299 .as_any()
300 .downcast_ref::<Float64Array>()
301 .unwrap();
302
303 assert_eq!(side_col.value(0), "NO_ORDER_SIDE");
304 assert!(
305 price_col.value(0).is_nan(),
306 "live clear price should be NaN"
307 );
308 assert!(size_col.value(0).is_nan(), "live clear size should be NaN");
309 }
310
311 #[rstest]
312 fn test_encode_deltas_clear_sentinels_render_as_nan() {
313 let clear = OrderBookDelta {
314 instrument_id: InstrumentId::from("AAPL.XNAS"),
315 action: BookAction::Clear,
316 order: BookOrder {
317 side: None,
318 price: Price::from_raw(PRICE_UNDEF, 0),
319 size: Quantity::from_raw(QUANTITY_UNDEF, 0),
320 order_id: 0,
321 },
322 flags: 0,
323 sequence: 1,
324 ts_event: 1.into(),
325 ts_init: 2.into(),
326 };
327
328 let batch = encode_deltas(&[clear]).unwrap();
329 let price_col = batch
330 .column(3)
331 .as_any()
332 .downcast_ref::<Float64Array>()
333 .unwrap();
334 let size_col = batch
335 .column(4)
336 .as_any()
337 .downcast_ref::<Float64Array>()
338 .unwrap();
339
340 assert!(price_col.value(0).is_nan(), "clear price should be NaN");
341 assert!(size_col.value(0).is_nan(), "clear size should be NaN");
342 }
343
344 #[rstest]
345 fn test_encode_deltas_mixed_instruments() {
346 let deltas = vec![
347 make_delta(
348 "AAPL.XNAS",
349 BookAction::Add,
350 OrderSide::Buy,
351 "100.10",
352 1,
353 1,
354 1,
355 ),
356 make_delta(
357 "MSFT.XNAS",
358 BookAction::Add,
359 OrderSide::Sell,
360 "250.00",
361 2,
362 1,
363 2,
364 ),
365 ];
366 let batch = encode_deltas(&deltas).unwrap();
367 let instrument_id_col = batch
368 .column(0)
369 .as_any()
370 .downcast_ref::<StringArray>()
371 .unwrap();
372 assert_eq!(instrument_id_col.value(0), "AAPL.XNAS");
373 assert_eq!(instrument_id_col.value(1), "MSFT.XNAS");
374 }
375}