1use std::str::FromStr;
17
18use indexmap::IndexMap;
19use nautilus_core::{UUID4, UnixNanos};
20use nautilus_model::{
21 enums::{
22 ContingencyType, LiquiditySide, OrderSide, OrderStatus, OrderType, TimeInForce, TriggerType,
23 },
24 events::{
25 OrderAccepted, OrderCancelRejected, OrderCanceled, OrderDenied, OrderEmulated,
26 OrderEventAny, OrderExpired, OrderFilled, OrderInitialized, OrderModifyRejected,
27 OrderPendingCancel, OrderPendingUpdate, OrderRejected, OrderReleased, OrderSnapshot,
28 OrderSubmitted, OrderTriggered, OrderUpdated,
29 },
30 identifiers::{
31 AccountId, ClientOrderId, ExecAlgorithmId, InstrumentId, OrderListId, PositionId,
32 StrategyId, TradeId, TraderId, VenueOrderId,
33 },
34 types::{Currency, Money, Price, Quantity},
35};
36use rust_decimal::Decimal;
37use sqlx::{FromRow, Row, postgres::PgRow};
38use ustr::Ustr;
39
40use crate::sql::models::enums::TrailingOffsetTypePg;
41
42#[derive(Debug)]
43pub struct OrderEventAnyRow(pub OrderEventAny);
44
45#[derive(Debug)]
46pub struct OrderAcceptedRow(pub OrderAccepted);
47
48#[derive(Debug)]
49pub struct OrderCancelRejectedRow(pub OrderCancelRejected);
50
51#[derive(Debug)]
52pub struct OrderCanceledRow(pub OrderCanceled);
53
54#[derive(Debug)]
55pub struct OrderDeniedRow(pub OrderDenied);
56
57#[derive(Debug)]
58pub struct OrderEmulatedRow(pub OrderEmulated);
59
60#[derive(Debug)]
61pub struct OrderExpiredRow(pub OrderExpired);
62
63#[derive(Debug)]
64pub struct OrderFilledRow(pub OrderFilled);
65
66#[derive(Debug)]
67pub struct OrderInitializedRow(pub OrderInitialized);
68
69#[derive(Debug)]
70pub struct OrderModifyRejectedRow(pub OrderModifyRejected);
71
72#[derive(Debug)]
73pub struct OrderPendingCancelRow(pub OrderPendingCancel);
74
75#[derive(Debug)]
76pub struct OrderPendingUpdateRow(pub OrderPendingUpdate);
77
78#[derive(Debug)]
79pub struct OrderRejectedRow(pub OrderRejected);
80
81#[derive(Debug)]
82pub struct OrderReleasedRow(pub OrderReleased);
83
84#[derive(Debug)]
85pub struct OrderSubmittedRow(pub OrderSubmitted);
86
87#[derive(Debug)]
88pub struct OrderTriggeredRow(pub OrderTriggered);
89
90#[derive(Debug)]
91pub struct OrderUpdatedRow(pub OrderUpdated);
92
93#[derive(Debug)]
94pub struct OrderSnapshotRow(pub OrderSnapshot);
95
96impl<'r> FromRow<'r, PgRow> for OrderEventAnyRow {
97 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
98 let kind = row.get::<String, _>("kind");
99 if kind == "OrderAccepted" {
100 let row = OrderAcceptedRow::from_row(row)?;
101 Ok(Self(OrderEventAny::Accepted(row.0)))
102 } else if kind == "OrderCancelRejected" {
103 let row = OrderCancelRejectedRow::from_row(row)?;
104 Ok(Self(OrderEventAny::CancelRejected(row.0)))
105 } else if kind == "OrderCanceled" {
106 let row = OrderCanceledRow::from_row(row)?;
107 Ok(Self(OrderEventAny::Canceled(row.0)))
108 } else if kind == "OrderDenied" {
109 let row = OrderDeniedRow::from_row(row)?;
110 Ok(Self(OrderEventAny::Denied(row.0)))
111 } else if kind == "OrderEmulated" {
112 let row = OrderEmulatedRow::from_row(row)?;
113 Ok(Self(OrderEventAny::Emulated(row.0)))
114 } else if kind == "OrderExpired" {
115 let row = OrderExpiredRow::from_row(row)?;
116 Ok(Self(OrderEventAny::Expired(row.0)))
117 } else if kind == "OrderFilled" {
118 let row = OrderFilledRow::from_row(row)?;
119 Ok(Self(OrderEventAny::Filled(row.0)))
120 } else if kind == "OrderInitialized" {
121 let row = OrderInitializedRow::from_row(row)?;
122 Ok(Self(OrderEventAny::Initialized(row.0)))
123 } else if kind == "OrderModifyRejected" {
124 let row = OrderModifyRejectedRow::from_row(row)?;
125 Ok(Self(OrderEventAny::ModifyRejected(row.0)))
126 } else if kind == "OrderPendingCancel" {
127 let row = OrderPendingCancelRow::from_row(row)?;
128 Ok(Self(OrderEventAny::PendingCancel(row.0)))
129 } else if kind == "OrderPendingUpdate" {
130 let row = OrderPendingUpdateRow::from_row(row)?;
131 Ok(Self(OrderEventAny::PendingUpdate(row.0)))
132 } else if kind == "OrderRejected" {
133 let row = OrderRejectedRow::from_row(row)?;
134 Ok(Self(OrderEventAny::Rejected(row.0)))
135 } else if kind == "OrderReleased" {
136 let row = OrderReleasedRow::from_row(row)?;
137 Ok(Self(OrderEventAny::Released(row.0)))
138 } else if kind == "OrderSubmitted" {
139 let row = OrderSubmittedRow::from_row(row)?;
140 Ok(Self(OrderEventAny::Submitted(row.0)))
141 } else if kind == "OrderTriggered" {
142 let row = OrderTriggeredRow::from_row(row)?;
143 Ok(Self(OrderEventAny::Triggered(row.0)))
144 } else if kind == "OrderUpdated" {
145 let row = OrderUpdatedRow::from_row(row)?;
146 Ok(Self(OrderEventAny::Updated(row.0)))
147 } else {
148 Err(sqlx::Error::Decode(
149 format!("Unknown order event kind: {kind} in Postgres transformation").into(),
150 ))
151 }
152 }
153}
154
155impl<'r> FromRow<'r, PgRow> for OrderInitializedRow {
156 #[expect(
157 clippy::too_many_lines,
158 reason = "SQL row mapping mirrors the full order initialized event constructor"
159 )]
160 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
161 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
162 let client_order_id = row
163 .try_get::<&str, _>("client_order_id")
164 .map(ClientOrderId::from)?;
165 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
166 let strategy_id = row
167 .try_get::<&str, _>("strategy_id")
168 .map(StrategyId::from)?;
169 let instrument_id = row
170 .try_get::<&str, _>("instrument_id")
171 .map(InstrumentId::from)?;
172 let order_type = row
173 .try_get::<&str, _>("order_type")
174 .map(|x| OrderType::from_str(x).unwrap())?;
175 let order_side = row
176 .try_get::<&str, _>("order_side")
177 .map(|x| OrderSide::from_str(x).unwrap())?;
178 let quantity = row.try_get::<&str, _>("quantity").map(Quantity::from)?;
179 let time_in_force = row
180 .try_get::<&str, _>("time_in_force")
181 .map(|x| TimeInForce::from_str(x).unwrap())?;
182 let post_only = row.try_get::<bool, _>("post_only")?;
183 let reduce_only = row.try_get::<bool, _>("reduce_only")?;
184 let quote_quantity = row.try_get::<bool, _>("quote_quantity")?;
185 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
186 let ts_event = row.try_get::<String, _>("ts_event").map(UnixNanos::from)?;
187 let ts_init = row.try_get::<String, _>("ts_init").map(UnixNanos::from)?;
188 let price = row
189 .try_get::<Option<&str>, _>("price")
190 .ok()
191 .and_then(|x| x.map(Price::from));
192 let activation_price = row
193 .try_get::<Option<&str>, _>("activation_price")
194 .ok()
195 .and_then(|x| x.map(Price::from));
196 let trigger_price = row
197 .try_get::<Option<&str>, _>("trigger_price")
198 .ok()
199 .and_then(|x| x.map(Price::from));
200 let trigger_type = row
201 .try_get::<Option<&str>, _>("trigger_type")
202 .ok()
203 .and_then(parse_trigger_type);
204 let limit_offset = row
205 .try_get::<Option<&str>, _>("limit_offset")
206 .ok()
207 .and_then(|x| x.and_then(|s| Decimal::from_str(s).ok()));
208 let trailing_offset = row
209 .try_get::<Option<&str>, _>("trailing_offset")
210 .ok()
211 .and_then(|x| x.and_then(|s| Decimal::from_str(s).ok()));
212 let trailing_offset_type = row
213 .try_get::<Option<TrailingOffsetTypePg>, _>("trailing_offset_type")
214 .ok()
215 .flatten()
216 .and_then(|value| value.0);
217 let expire_time = row
218 .try_get::<Option<&str>, _>("expire_time")
219 .ok()
220 .and_then(|x| x.map(UnixNanos::from));
221 let display_qty = row
222 .try_get::<Option<&str>, _>("display_qty")
223 .ok()
224 .and_then(|x| x.map(Quantity::from));
225 let emulation_trigger = row
226 .try_get::<Option<&str>, _>("emulation_trigger")
227 .ok()
228 .and_then(parse_trigger_type);
229 let trigger_instrument_id = row
230 .try_get::<Option<&str>, _>("trigger_instrument_id")
231 .ok()
232 .and_then(|x| x.map(InstrumentId::from));
233 let contingency_type = row
234 .try_get::<Option<&str>, _>("contingency_type")
235 .ok()
236 .and_then(parse_contingency_type);
237 let order_list_id = row
238 .try_get::<Option<&str>, _>("order_list_id")
239 .ok()
240 .and_then(|x| x.map(OrderListId::from));
241 let linked_order_ids = row
242 .try_get::<Vec<String>, _>("linked_order_ids")
243 .ok()
244 .map(|x| x.iter().map(|x| ClientOrderId::from(x.as_str())).collect());
245 let parent_order_id = row
246 .try_get::<Option<&str>, _>("parent_order_id")
247 .ok()
248 .and_then(|x| x.map(ClientOrderId::from));
249 let exec_algorithm_id = row
250 .try_get::<Option<&str>, _>("exec_algorithm_id")
251 .ok()
252 .and_then(|x| x.map(ExecAlgorithmId::from));
253 let exec_algorithm_params: Option<IndexMap<Ustr, Ustr>> = row
254 .try_get::<Option<serde_json::Value>, _>("exec_algorithm_params")
255 .ok()
256 .and_then(|x| x.map(|x| serde_json::from_value::<IndexMap<String, String>>(x).unwrap()))
257 .map(|x| {
258 x.into_iter()
259 .map(|(k, v)| (Ustr::from(k.as_str()), Ustr::from(v.as_str())))
260 .collect()
261 });
262 let exec_spawn_id = row
263 .try_get::<Option<&str>, _>("exec_spawn_id")
264 .ok()
265 .and_then(|x| x.map(ClientOrderId::from));
266 let tags = tags_from_row(row);
267 let order_event = OrderInitialized::new_checked(
268 trader_id,
269 strategy_id,
270 instrument_id,
271 client_order_id,
272 order_side,
273 order_type,
274 quantity,
275 time_in_force,
276 post_only,
277 reduce_only,
278 quote_quantity,
279 reconciliation,
280 event_id,
281 ts_event,
282 ts_init,
283 price,
284 activation_price,
285 trigger_price,
286 trigger_type,
287 limit_offset,
288 trailing_offset,
289 trailing_offset_type,
290 expire_time,
291 display_qty,
292 emulation_trigger,
293 trigger_instrument_id,
294 contingency_type,
295 order_list_id,
296 linked_order_ids,
297 parent_order_id,
298 exec_algorithm_id,
299 exec_algorithm_params,
300 exec_spawn_id,
301 tags,
302 )
303 .map_err(|e| sqlx::Error::Decode(Box::new(e)))?;
304 Ok(Self(order_event))
305 }
306}
307
308impl<'r> FromRow<'r, PgRow> for OrderAcceptedRow {
309 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
310 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
311 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
312 let strategy_id = row
313 .try_get::<&str, _>("strategy_id")
314 .map(StrategyId::from)?;
315 let instrument_id = row
316 .try_get::<&str, _>("instrument_id")
317 .map(InstrumentId::from)?;
318 let client_order_id = row
319 .try_get::<&str, _>("client_order_id")
320 .map(ClientOrderId::from)?;
321 let venue_order_id = row
322 .try_get::<&str, _>("venue_order_id")
323 .map(VenueOrderId::from)?;
324 let account_id = row.try_get::<&str, _>("account_id").map(AccountId::from)?;
325 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
326 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
327 let order_event = OrderAccepted::new(
328 trader_id,
329 strategy_id,
330 instrument_id,
331 client_order_id,
332 venue_order_id,
333 account_id,
334 event_id,
335 ts_event,
336 ts_init,
337 false,
338 );
339 Ok(Self(order_event))
340 }
341}
342
343impl<'r> FromRow<'r, PgRow> for OrderCancelRejectedRow {
344 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
345 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
346 let strategy_id = row
347 .try_get::<&str, _>("strategy_id")
348 .map(StrategyId::from)?;
349 let instrument_id = row
350 .try_get::<&str, _>("instrument_id")
351 .map(InstrumentId::from)?;
352 let client_order_id = row
353 .try_get::<&str, _>("client_order_id")
354 .map(ClientOrderId::from)?;
355 let reason = row.try_get::<&str, _>("reason").map(Ustr::from)?;
356 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
357 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
358 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
359 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
360 let venue_order_id = row
361 .try_get::<Option<&str>, _>("venue_order_id")?
362 .map(Into::into);
363 let account_id = row
364 .try_get::<Option<&str>, _>("account_id")?
365 .map(Into::into);
366 let order_event = OrderCancelRejected::new(
367 trader_id,
368 strategy_id,
369 instrument_id,
370 client_order_id,
371 reason,
372 event_id,
373 ts_event,
374 ts_init,
375 reconciliation,
376 venue_order_id,
377 account_id,
378 );
379 Ok(Self(order_event))
380 }
381}
382
383impl<'r> FromRow<'r, PgRow> for OrderCanceledRow {
384 fn from_row(_row: &'r PgRow) -> Result<Self, sqlx::Error> {
385 todo!()
386 }
387}
388
389impl<'r> FromRow<'r, PgRow> for OrderDeniedRow {
390 fn from_row(_row: &'r PgRow) -> Result<Self, sqlx::Error> {
391 todo!()
392 }
393}
394
395impl<'r> FromRow<'r, PgRow> for OrderEmulatedRow {
396 fn from_row(_row: &'r PgRow) -> Result<Self, sqlx::Error> {
397 todo!()
398 }
399}
400
401impl<'r> FromRow<'r, PgRow> for OrderExpiredRow {
402 fn from_row(_row: &'r PgRow) -> Result<Self, sqlx::Error> {
403 todo!()
404 }
405}
406
407impl<'r> FromRow<'r, PgRow> for OrderFilledRow {
408 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
409 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
410 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
411 let strategy_id = row
412 .try_get::<&str, _>("strategy_id")
413 .map(StrategyId::from)?;
414 let instrument_id = row
415 .try_get::<&str, _>("instrument_id")
416 .map(InstrumentId::from)?;
417 let client_order_id = row
418 .try_get::<&str, _>("client_order_id")
419 .map(ClientOrderId::from)?;
420 let venue_order_id = row
421 .try_get::<&str, _>("venue_order_id")
422 .map(VenueOrderId::from)?;
423 let account_id = row.try_get::<&str, _>("account_id").map(AccountId::from)?;
424 let trade_id = row.try_get::<&str, _>("trade_id").map(TradeId::from)?;
425 let order_side = row
426 .try_get::<&str, _>("order_side")
427 .map(|x| OrderSide::from_str(x).unwrap())?;
428 let order_type = row
429 .try_get::<&str, _>("order_type")
430 .map(|x| OrderType::from_str(x).unwrap())?;
431 let last_px = row.try_get::<&str, _>("last_px").map(Price::from)?;
432 let last_qty = row.try_get::<&str, _>("last_qty").map(Quantity::from)?;
433 let currency = row.try_get::<&str, _>("currency").map(Currency::from)?;
434 let liquidity_side = row
435 .try_get::<&str, _>("liquidity_side")
436 .map(|x| LiquiditySide::from_str(x).unwrap())?;
437 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
438 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
439 let position_id = row
440 .try_get::<Option<&str>, _>("position_id")
441 .map(|x| x.map(PositionId::from))?;
442 let commission = row
443 .try_get::<Option<&str>, _>("commission")
444 .map(|x| x.map(|x| Money::from_str(x).unwrap()))?;
445 let order_event = OrderFilled::new(
446 trader_id,
447 strategy_id,
448 instrument_id,
449 client_order_id,
450 venue_order_id,
451 account_id,
452 trade_id,
453 order_side,
454 order_type,
455 last_qty,
456 last_px,
457 currency,
458 liquidity_side,
459 event_id,
460 ts_event,
461 ts_init,
462 false,
463 position_id,
464 commission,
465 None,
466 );
467 Ok(Self(order_event))
468 }
469}
470
471impl<'r> FromRow<'r, PgRow> for OrderModifyRejectedRow {
472 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
473 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
474 let strategy_id = row
475 .try_get::<&str, _>("strategy_id")
476 .map(StrategyId::from)?;
477 let instrument_id = row
478 .try_get::<&str, _>("instrument_id")
479 .map(InstrumentId::from)?;
480 let client_order_id = row
481 .try_get::<&str, _>("client_order_id")
482 .map(ClientOrderId::from)?;
483 let reason = row.try_get::<&str, _>("reason").map(Ustr::from)?;
484 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
485 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
486 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
487 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
488 let venue_order_id = row
489 .try_get::<Option<&str>, _>("venue_order_id")?
490 .map(Into::into);
491 let account_id = row
492 .try_get::<Option<&str>, _>("account_id")?
493 .map(Into::into);
494 let order_event = OrderModifyRejected::new(
495 trader_id,
496 strategy_id,
497 instrument_id,
498 client_order_id,
499 reason,
500 event_id,
501 ts_event,
502 ts_init,
503 reconciliation,
504 venue_order_id,
505 account_id,
506 );
507 Ok(Self(order_event))
508 }
509}
510
511impl<'r> FromRow<'r, PgRow> for OrderPendingCancelRow {
512 fn from_row(_row: &'r PgRow) -> Result<Self, sqlx::Error> {
513 todo!()
514 }
515}
516
517impl<'r> FromRow<'r, PgRow> for OrderPendingUpdateRow {
518 fn from_row(_row: &'r PgRow) -> Result<Self, sqlx::Error> {
519 todo!()
520 }
521}
522
523impl<'r> FromRow<'r, PgRow> for OrderRejectedRow {
524 fn from_row(_row: &'r PgRow) -> Result<Self, sqlx::Error> {
525 todo!()
526 }
527}
528
529impl<'r> FromRow<'r, PgRow> for OrderReleasedRow {
530 fn from_row(_row: &'r PgRow) -> Result<Self, sqlx::Error> {
531 todo!()
532 }
533}
534
535impl<'r> FromRow<'r, PgRow> for OrderSubmittedRow {
536 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
537 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
538 let strategy_id = row
539 .try_get::<&str, _>("strategy_id")
540 .map(StrategyId::from)?;
541 let instrument_id = row
542 .try_get::<&str, _>("instrument_id")
543 .map(InstrumentId::from)?;
544 let client_order_id = row
545 .try_get::<&str, _>("client_order_id")
546 .map(ClientOrderId::from)?;
547 let account_id = row.try_get::<&str, _>("account_id").map(AccountId::from)?;
548 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
549 let ts_event = row
550 .try_get::<String, _>("ts_event")
551 .map(|res| UnixNanos::from(res.as_str()))?;
552 let ts_init = row
553 .try_get::<String, _>("ts_init")
554 .map(|res| UnixNanos::from(res.as_str()))?;
555 let order_event = OrderSubmitted::new(
556 trader_id,
557 strategy_id,
558 instrument_id,
559 client_order_id,
560 account_id,
561 event_id,
562 ts_event,
563 ts_init,
564 );
565 Ok(Self(order_event))
566 }
567}
568
569impl<'r> FromRow<'r, PgRow> for OrderTriggeredRow {
570 fn from_row(_row: &'r PgRow) -> Result<Self, sqlx::Error> {
571 todo!()
572 }
573}
574
575impl<'r> FromRow<'r, PgRow> for OrderUpdatedRow {
576 fn from_row(_row: &'r PgRow) -> Result<Self, sqlx::Error> {
577 todo!()
578 }
579}
580
581impl<'r> FromRow<'r, PgRow> for OrderSnapshotRow {
582 #[expect(
583 clippy::too_many_lines,
584 reason = "SQL row mapping mirrors the full order snapshot constructor"
585 )]
586 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
587 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
588 let strategy_id = row
589 .try_get::<&str, _>("strategy_id")
590 .map(StrategyId::from)?;
591 let instrument_id = row
592 .try_get::<&str, _>("instrument_id")
593 .map(InstrumentId::from)?;
594 let client_order_id = row
595 .try_get::<&str, _>("client_order_id")
596 .map(ClientOrderId::from)?;
597 let venue_order_id = row
598 .try_get::<Option<&str>, _>("venue_order_id")
599 .ok()
600 .and_then(|x| x.map(VenueOrderId::from));
601 let position_id = row
602 .try_get::<Option<&str>, _>("position_id")
603 .ok()
604 .and_then(|x| x.map(PositionId::from));
605 let account_id = row
606 .try_get::<Option<&str>, _>("account_id")
607 .ok()
608 .and_then(|x| x.map(AccountId::from));
609 let last_trade_id = row
610 .try_get::<Option<&str>, _>("last_trade_id")
611 .ok()
612 .and_then(|x| x.map(TradeId::from));
613 let order_type = row
614 .try_get::<&str, _>("order_type")
615 .map(|x| OrderType::from_str(x).expect("Invalid `OrderType`"))?;
616 let order_side = row
617 .try_get::<&str, _>("order_side")
618 .map(|x| OrderSide::from_str(x).expect("Invalid `OrderSide`"))?;
619 let quantity = row.try_get::<&str, _>("quantity").map(Quantity::from)?;
620 let price = row
621 .try_get::<Option<&str>, _>("price")
622 .ok()
623 .and_then(|x| x.map(Price::from));
624 let activation_price = row
625 .try_get::<Option<&str>, _>("activation_price")
626 .ok()
627 .and_then(|x| x.map(Price::from));
628 let trigger_price = row
629 .try_get::<Option<&str>, _>("trigger_price")
630 .ok()
631 .and_then(|x| x.map(Price::from));
632 let trigger_type = row
633 .try_get::<Option<&str>, _>("trigger_type")
634 .ok()
635 .and_then(parse_trigger_type);
636 let limit_offset = row
637 .try_get::<Option<&str>, _>("limit_offset")
638 .ok()
639 .and_then(|x| x.and_then(|s| Decimal::from_str(s).ok()));
640 let trailing_offset = row
641 .try_get::<Option<&str>, _>("trailing_offset")
642 .ok()
643 .and_then(|x| x.and_then(|s| Decimal::from_str(s).ok()));
644 let trailing_offset_type = row
645 .try_get::<Option<TrailingOffsetTypePg>, _>("trailing_offset_type")
646 .ok()
647 .flatten()
648 .and_then(|value| value.0);
649 let time_in_force = row
650 .try_get::<&str, _>("time_in_force")
651 .map(|x| TimeInForce::from_str(x).expect("Invalid `TimeInForce`"))?;
652 let expire_time = row
653 .try_get::<Option<&str>, _>("expire_time")
654 .ok()
655 .and_then(|x| x.map(UnixNanos::from));
656 let filled_qty = row.try_get::<&str, _>("filled_qty").map(Quantity::from)?;
657 let liquidity_side = row
658 .try_get::<Option<&str>, _>("liquidity_side")
659 .ok()
660 .and_then(|x| x.map(|x| LiquiditySide::from_str(x).expect("Invalid `LiquiditySide`")));
661 let avg_px = row.try_get::<Option<Decimal>, _>("avg_px").ok().flatten();
662 let slippage = row.try_get::<Option<Decimal>, _>("slippage").ok().flatten();
663 let commissions = row
664 .try_get::<Option<Vec<String>>, _>("commissions")?
665 .map_or_else(Vec::new, |c| {
666 c.into_iter().map(|s| Money::from(&s)).collect()
667 });
668 let status = row
669 .try_get::<&str, _>("status")
670 .map(|x| OrderStatus::from_str(x).expect("Invalid `OrderStatus`"))?;
671 let is_post_only = row.try_get::<bool, _>("is_post_only")?;
672 let is_reduce_only = row.try_get::<bool, _>("is_reduce_only")?;
673 let is_quote_quantity = row.try_get::<bool, _>("is_quote_quantity")?;
674 let display_qty = row
675 .try_get::<Option<&str>, _>("display_qty")
676 .ok()
677 .and_then(|x| x.map(Quantity::from));
678 let emulation_trigger = row
679 .try_get::<Option<&str>, _>("emulation_trigger")
680 .ok()
681 .and_then(parse_trigger_type);
682 let trigger_instrument_id = row
683 .try_get::<Option<&str>, _>("trigger_instrument_id")
684 .ok()
685 .and_then(|x| x.map(InstrumentId::from));
686 let contingency_type = row
687 .try_get::<Option<&str>, _>("contingency_type")
688 .ok()
689 .and_then(parse_contingency_type);
690 let order_list_id = row
691 .try_get::<Option<&str>, _>("order_list_id")
692 .ok()
693 .and_then(|x| x.map(OrderListId::from));
694 let linked_order_ids = row
695 .try_get::<Option<Vec<String>>, _>("linked_order_ids")
696 .ok()
697 .and_then(|ids| ids.map(|ids| ids.into_iter().map(ClientOrderId::from).collect()));
698 let parent_order_id = row
699 .try_get::<Option<&str>, _>("parent_order_id")
700 .ok()
701 .and_then(|x| x.map(ClientOrderId::from));
702 let exec_algorithm_id = row
703 .try_get::<Option<&str>, _>("exec_algorithm_id")
704 .ok()
705 .and_then(|x| x.map(ExecAlgorithmId::from));
706 let exec_algorithm_params: Option<IndexMap<Ustr, Ustr>> = row
707 .try_get::<Option<serde_json::Value>, _>("exec_algorithm_params")
708 .ok()
709 .and_then(|x| {
710 x.map(|x| {
711 serde_json::from_value::<IndexMap<String, String>>(x)
712 .expect("Invalid exec algorithm params")
713 })
714 })
715 .map(|x| {
716 x.into_iter()
717 .map(|(k, v)| (Ustr::from(k.as_str()), Ustr::from(v.as_str())))
718 .collect()
719 });
720 let exec_spawn_id = row
721 .try_get::<Option<&str>, _>("exec_spawn_id")
722 .ok()
723 .and_then(|x| x.map(ClientOrderId::from));
724 let tags = tags_from_row(row);
725 let init_id = row.try_get::<&str, _>("init_id").map(UUID4::from)?;
726 let ts_init = row.try_get::<String, _>("ts_init").map(UnixNanos::from)?;
727 let ts_last = row.try_get::<String, _>("ts_last").map(UnixNanos::from)?;
728
729 let snapshot = OrderSnapshot {
730 trader_id,
731 strategy_id,
732 instrument_id,
733 client_order_id,
734 venue_order_id,
735 position_id,
736 account_id,
737 last_trade_id,
738 order_type,
739 order_side,
740 quantity,
741 price,
742 activation_price,
743 trigger_price,
744 trigger_type,
745 limit_offset,
746 trailing_offset,
747 trailing_offset_type,
748 time_in_force,
749 expire_time,
750 filled_qty,
751 liquidity_side,
752 avg_px,
753 slippage,
754 commissions,
755 status,
756 is_post_only,
757 is_reduce_only,
758 is_quote_quantity,
759 display_qty,
760 emulation_trigger,
761 trigger_instrument_id,
762 contingency_type,
763 order_list_id,
764 linked_order_ids,
765 parent_order_id,
766 exec_algorithm_id,
767 exec_algorithm_params,
768 exec_spawn_id,
769 tags,
770 init_id,
771 ts_init,
772 ts_last,
773 causation_id: None,
774 };
775
776 Ok(Self(snapshot))
777 }
778}
779
780fn tags_from_row(row: &PgRow) -> Option<Vec<Ustr>> {
781 row.try_get::<Vec<String>, _>("tags")
782 .ok()
783 .map(|tags| tags.iter().map(|tag| Ustr::from(tag.as_str())).collect())
784}
785
786fn parse_trigger_type(value: Option<&str>) -> Option<TriggerType> {
787 value.and_then(|value| {
788 if value.eq_ignore_ascii_case("NO_TRIGGER") {
789 None
790 } else {
791 Some(TriggerType::from_str(value).expect("Invalid `TriggerType`"))
792 }
793 })
794}
795
796fn parse_contingency_type(value: Option<&str>) -> Option<ContingencyType> {
797 value.and_then(|value| {
798 if value.eq_ignore_ascii_case("NO_CONTINGENCY") {
799 None
800 } else {
801 Some(ContingencyType::from_str(value).expect("Invalid `ContingencyType`"))
802 }
803 })
804}
805
806#[cfg(test)]
807mod tests {
808 use rstest::rstest;
809
810 use super::*;
811
812 #[rstest]
813 #[case(None, None)]
814 #[case(Some("NO_TRIGGER"), None)]
815 #[case(Some("LAST_PRICE"), Some(TriggerType::LastPrice))]
816 fn test_parse_trigger_type_accepts_legacy_absence(
817 #[case] value: Option<&str>,
818 #[case] expected: Option<TriggerType>,
819 ) {
820 assert_eq!(parse_trigger_type(value), expected);
821 }
822
823 #[rstest]
824 #[case(None, None)]
825 #[case(Some("NO_CONTINGENCY"), None)]
826 #[case(Some("OCO"), Some(ContingencyType::Oco))]
827 fn test_parse_contingency_type_accepts_legacy_absence(
828 #[case] value: Option<&str>,
829 #[case] expected: Option<ContingencyType>,
830 ) {
831 assert_eq!(parse_contingency_type(value), expected);
832 }
833}