1use std::collections::HashMap;
17
18use arrow::{datatypes::Schema, error::ArrowError, record_batch::RecordBatch};
19use nautilus_model::reports::{
20 ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport,
21};
22
23use super::{
24 ArrowSchemaProvider, DecodeTypedFromRecordBatch, EncodeToRecordBatch, EncodingError,
25 KEY_INSTRUMENT_ID,
26 json::{JsonFieldSpec, decode_batch, encode_batch, metadata_for_type, schema_for_type},
27};
28
29const ORDER_STATUS_REPORT_FIELDS: &[JsonFieldSpec] = &[
30 JsonFieldSpec::utf8("account_id", false),
31 JsonFieldSpec::utf8("instrument_id", false),
32 JsonFieldSpec::utf8("client_order_id", true),
33 JsonFieldSpec::utf8("venue_order_id", false),
34 JsonFieldSpec::utf8("order_side", false),
35 JsonFieldSpec::utf8("order_type", false),
36 JsonFieldSpec::utf8("time_in_force", false),
37 JsonFieldSpec::utf8("order_status", false),
38 JsonFieldSpec::utf8("quantity", false),
39 JsonFieldSpec::utf8("filled_qty", false),
40 JsonFieldSpec::utf8("report_id", false),
41 JsonFieldSpec::u64("ts_accepted", false),
42 JsonFieldSpec::u64("ts_last", false),
43 JsonFieldSpec::u64("ts_init", false),
44 JsonFieldSpec::utf8("order_list_id", true),
45 JsonFieldSpec::utf8("venue_position_id", true),
46 JsonFieldSpec::utf8_json("linked_order_ids", true),
47 JsonFieldSpec::utf8("parent_order_id", true),
48 JsonFieldSpec::utf8("contingency_type", false),
49 JsonFieldSpec::u64("expire_time", true),
50 JsonFieldSpec::utf8("price", true),
51 JsonFieldSpec::utf8("activation_price", true),
52 JsonFieldSpec::utf8("trigger_price", true),
53 JsonFieldSpec::utf8("trigger_type", true),
54 JsonFieldSpec::utf8("limit_offset", true),
55 JsonFieldSpec::utf8("trailing_offset", true),
56 JsonFieldSpec::utf8("trailing_offset_type", false),
57 JsonFieldSpec::utf8("avg_px", true),
58 JsonFieldSpec::utf8("display_qty", true),
59 JsonFieldSpec::boolean("post_only", false),
60 JsonFieldSpec::boolean("reduce_only", false),
61 JsonFieldSpec::utf8("cancel_reason", true),
62 JsonFieldSpec::u64("ts_triggered", true),
63];
64
65const FILL_REPORT_FIELDS: &[JsonFieldSpec] = &[
66 JsonFieldSpec::utf8("account_id", false),
67 JsonFieldSpec::utf8("instrument_id", false),
68 JsonFieldSpec::utf8("venue_order_id", false),
69 JsonFieldSpec::utf8("trade_id", false),
70 JsonFieldSpec::utf8("order_side", false),
71 JsonFieldSpec::utf8("last_qty", false),
72 JsonFieldSpec::utf8("last_px", false),
73 JsonFieldSpec::utf8("commission", false),
74 JsonFieldSpec::utf8("liquidity_side", false),
75 JsonFieldSpec::utf8("report_id", false),
76 JsonFieldSpec::u64("ts_event", false),
77 JsonFieldSpec::u64("ts_init", false),
78 JsonFieldSpec::utf8("client_order_id", true),
79 JsonFieldSpec::utf8("venue_position_id", true),
80];
81
82const POSITION_STATUS_REPORT_FIELDS: &[JsonFieldSpec] = &[
83 JsonFieldSpec::utf8("account_id", false),
84 JsonFieldSpec::utf8("instrument_id", false),
85 JsonFieldSpec::utf8("position_side", false),
86 JsonFieldSpec::utf8("quantity", false),
87 JsonFieldSpec::utf8("signed_decimal_qty", false),
88 JsonFieldSpec::utf8("report_id", false),
89 JsonFieldSpec::u64("ts_last", false),
90 JsonFieldSpec::u64("ts_init", false),
91 JsonFieldSpec::utf8("venue_position_id", true),
92 JsonFieldSpec::utf8("avg_px_open", true),
93];
94
95const EXECUTION_MASS_STATUS_FIELDS: &[JsonFieldSpec] = &[
96 JsonFieldSpec::utf8("client_id", false),
97 JsonFieldSpec::utf8("account_id", false),
98 JsonFieldSpec::utf8("venue", false),
99 JsonFieldSpec::utf8("report_id", false),
100 JsonFieldSpec::u64("ts_init", false),
101 JsonFieldSpec::utf8_json("order_reports", false),
102 JsonFieldSpec::utf8_json("fill_reports", false),
103 JsonFieldSpec::utf8_json("position_reports", false),
104];
105
106fn instrument_metadata(type_name: &'static str, instrument_id: &str) -> HashMap<String, String> {
107 let mut metadata = metadata_for_type(type_name);
108 metadata.insert(KEY_INSTRUMENT_ID.to_string(), instrument_id.to_string());
109 metadata
110}
111
112macro_rules! impl_report_arrow {
113 ($type:ty, $type_name:expr, $fields:expr) => {
114 impl ArrowSchemaProvider for $type {
115 fn get_schema(metadata: Option<HashMap<String, String>>) -> Schema {
116 schema_for_type($type_name, metadata, $fields)
117 }
118 }
119
120 impl EncodeToRecordBatch for $type {
121 fn encode_batch(
122 metadata: &HashMap<String, String>,
123 data: &[Self],
124 ) -> Result<RecordBatch, ArrowError> {
125 encode_batch($type_name, metadata, data, $fields)
126 }
127
128 fn metadata(&self) -> HashMap<String, String> {
129 instrument_metadata($type_name, &self.instrument_id.to_string())
130 }
131 }
132
133 impl DecodeTypedFromRecordBatch for $type {
134 fn decode_typed_batch(
135 metadata: &HashMap<String, String>,
136 record_batch: RecordBatch,
137 ) -> Result<Vec<Self>, EncodingError> {
138 decode_batch(metadata, &record_batch, $fields, Some($type_name))
139 }
140 }
141 };
142}
143
144impl_report_arrow!(
145 OrderStatusReport,
146 "OrderStatusReport",
147 ORDER_STATUS_REPORT_FIELDS
148);
149impl_report_arrow!(FillReport, "FillReport", FILL_REPORT_FIELDS);
150impl_report_arrow!(
151 PositionStatusReport,
152 "PositionStatusReport",
153 POSITION_STATUS_REPORT_FIELDS
154);
155
156impl ArrowSchemaProvider for ExecutionMassStatus {
157 fn get_schema(metadata: Option<HashMap<String, String>>) -> Schema {
158 schema_for_type(
159 "ExecutionMassStatus",
160 metadata,
161 EXECUTION_MASS_STATUS_FIELDS,
162 )
163 }
164}
165
166impl EncodeToRecordBatch for ExecutionMassStatus {
167 fn encode_batch(
168 metadata: &HashMap<String, String>,
169 data: &[Self],
170 ) -> Result<RecordBatch, ArrowError> {
171 encode_batch(
172 "ExecutionMassStatus",
173 metadata,
174 data,
175 EXECUTION_MASS_STATUS_FIELDS,
176 )
177 }
178
179 fn metadata(&self) -> HashMap<String, String> {
180 metadata_for_type("ExecutionMassStatus")
181 }
182}
183
184impl DecodeTypedFromRecordBatch for ExecutionMassStatus {
185 fn decode_typed_batch(
186 metadata: &HashMap<String, String>,
187 record_batch: RecordBatch,
188 ) -> Result<Vec<Self>, EncodingError> {
189 decode_batch(
190 metadata,
191 &record_batch,
192 EXECUTION_MASS_STATUS_FIELDS,
193 Some("ExecutionMassStatus"),
194 )
195 }
196}
197
198#[cfg(test)]
199mod tests {
200 use std::str::FromStr;
201
202 use nautilus_core::{UUID4, UnixNanos};
203 use nautilus_model::{
204 enums::{OrderSide, OrderStatus, OrderType, PositionSideSpecified, TimeInForce},
205 identifiers::{AccountId, ClientOrderId, InstrumentId, PositionId, VenueOrderId},
206 reports::{OrderStatusReport, PositionStatusReport},
207 types::{Price, Quantity},
208 };
209 use rstest::rstest;
210 use rust_decimal::Decimal;
211
212 use super::*;
213
214 #[rstest]
215 fn test_order_status_report_round_trip() {
216 let report = OrderStatusReport::new(
217 AccountId::from("SIM-001"),
218 InstrumentId::from("AUDUSD.SIM"),
219 Some(ClientOrderId::from("O-19700101-000000-001-001-1")),
220 VenueOrderId::from("1"),
221 OrderSide::Buy,
222 OrderType::Limit,
223 TimeInForce::Gtc,
224 OrderStatus::Accepted,
225 Quantity::from("100"),
226 Quantity::from("25"),
227 UnixNanos::from(1_000_000_000),
228 UnixNanos::from(2_000_000_000),
229 UnixNanos::from(3_000_000_000),
230 None,
231 )
232 .with_linked_order_ids([ClientOrderId::from("O-19700101-000000-001-001-2")]);
233 let report = OrderStatusReport {
234 activation_price: Some(Price::from("1.05000")),
235 limit_offset: Some(Decimal::from_str("0.123456789123456789").unwrap()),
236 trailing_offset: Some(Decimal::from_str("0.987654321987654321").unwrap()),
237 avg_px: Some(Decimal::from_str("1.23456789123456789").unwrap()),
238 ..report
239 };
240
241 let metadata = report.metadata();
242 let batch =
243 OrderStatusReport::encode_batch(&metadata, std::slice::from_ref(&report)).unwrap();
244 let decoded =
245 OrderStatusReport::decode_typed_batch(batch.schema().metadata(), batch).unwrap();
246
247 assert_eq!(decoded, vec![report]);
248 }
249
250 #[rstest]
251 fn test_position_status_report_round_trip_preserves_decimal_precision() {
252 let report = PositionStatusReport {
253 account_id: AccountId::from("SIM-001"),
254 instrument_id: InstrumentId::from("AUDUSD.SIM"),
255 position_side: PositionSideSpecified::Long,
256 quantity: Quantity::from("100.25"),
257 signed_decimal_qty: Decimal::from_str("100.250000000123456789").unwrap(),
258 report_id: UUID4::default(),
259 ts_last: UnixNanos::from(1_000_000_000),
260 ts_init: UnixNanos::from(2_000_000_000),
261 venue_position_id: Some(PositionId::from("P-001")),
262 avg_px_open: Some(Decimal::from_str("1.23456789123456789").unwrap()),
263 };
264 let metadata = report.metadata();
265 let batch =
266 PositionStatusReport::encode_batch(&metadata, std::slice::from_ref(&report)).unwrap();
267 let decoded =
268 PositionStatusReport::decode_typed_batch(batch.schema().metadata(), batch).unwrap();
269
270 assert_eq!(decoded, vec![report]);
271 }
272}