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, OrderFillVoided, OrderFilled, OrderInitialized,
27 OrderModifyRejected, OrderPendingCancel, OrderPendingUpdate, OrderRejected, OrderReleased,
28 OrderSnapshot, 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 OrderFillVoidedRow(pub OrderFillVoided);
68
69#[derive(Debug)]
70pub struct OrderInitializedRow(pub OrderInitialized);
71
72#[derive(Debug)]
73pub struct OrderModifyRejectedRow(pub OrderModifyRejected);
74
75#[derive(Debug)]
76pub struct OrderPendingCancelRow(pub OrderPendingCancel);
77
78#[derive(Debug)]
79pub struct OrderPendingUpdateRow(pub OrderPendingUpdate);
80
81#[derive(Debug)]
82pub struct OrderRejectedRow(pub OrderRejected);
83
84#[derive(Debug)]
85pub struct OrderReleasedRow(pub OrderReleased);
86
87#[derive(Debug)]
88pub struct OrderSubmittedRow(pub OrderSubmitted);
89
90#[derive(Debug)]
91pub struct OrderTriggeredRow(pub OrderTriggered);
92
93#[derive(Debug)]
94pub struct OrderUpdatedRow(pub OrderUpdated);
95
96#[derive(Debug)]
97pub struct OrderSnapshotRow(pub OrderSnapshot);
98
99impl<'r> FromRow<'r, PgRow> for OrderEventAnyRow {
100 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
101 let kind = row.get::<String, _>("kind");
102 if kind == "OrderAccepted" {
103 let row = OrderAcceptedRow::from_row(row)?;
104 Ok(Self(OrderEventAny::Accepted(row.0)))
105 } else if kind == "OrderCancelRejected" {
106 let row = OrderCancelRejectedRow::from_row(row)?;
107 Ok(Self(OrderEventAny::CancelRejected(row.0)))
108 } else if kind == "OrderCanceled" {
109 let row = OrderCanceledRow::from_row(row)?;
110 Ok(Self(OrderEventAny::Canceled(row.0)))
111 } else if kind == "OrderDenied" {
112 let row = OrderDeniedRow::from_row(row)?;
113 Ok(Self(OrderEventAny::Denied(row.0)))
114 } else if kind == "OrderEmulated" {
115 let row = OrderEmulatedRow::from_row(row)?;
116 Ok(Self(OrderEventAny::Emulated(row.0)))
117 } else if kind == "OrderExpired" {
118 let row = OrderExpiredRow::from_row(row)?;
119 Ok(Self(OrderEventAny::Expired(row.0)))
120 } else if kind == "OrderFillVoided" {
121 let row = OrderFillVoidedRow::from_row(row)?;
122 Ok(Self(OrderEventAny::FillVoided(row.0)))
123 } else if kind == "OrderFilled" {
124 let row = OrderFilledRow::from_row(row)?;
125 Ok(Self(OrderEventAny::Filled(row.0)))
126 } else if kind == "OrderInitialized" {
127 let row = OrderInitializedRow::from_row(row)?;
128 Ok(Self(OrderEventAny::Initialized(row.0)))
129 } else if kind == "OrderModifyRejected" {
130 let row = OrderModifyRejectedRow::from_row(row)?;
131 Ok(Self(OrderEventAny::ModifyRejected(row.0)))
132 } else if kind == "OrderPendingCancel" {
133 let row = OrderPendingCancelRow::from_row(row)?;
134 Ok(Self(OrderEventAny::PendingCancel(row.0)))
135 } else if kind == "OrderPendingUpdate" {
136 let row = OrderPendingUpdateRow::from_row(row)?;
137 Ok(Self(OrderEventAny::PendingUpdate(row.0)))
138 } else if kind == "OrderRejected" {
139 let row = OrderRejectedRow::from_row(row)?;
140 Ok(Self(OrderEventAny::Rejected(row.0)))
141 } else if kind == "OrderReleased" {
142 let row = OrderReleasedRow::from_row(row)?;
143 Ok(Self(OrderEventAny::Released(row.0)))
144 } else if kind == "OrderSubmitted" {
145 let row = OrderSubmittedRow::from_row(row)?;
146 Ok(Self(OrderEventAny::Submitted(row.0)))
147 } else if kind == "OrderTriggered" {
148 let row = OrderTriggeredRow::from_row(row)?;
149 Ok(Self(OrderEventAny::Triggered(row.0)))
150 } else if kind == "OrderUpdated" {
151 let row = OrderUpdatedRow::from_row(row)?;
152 Ok(Self(OrderEventAny::Updated(row.0)))
153 } else {
154 Err(sqlx::Error::Decode(
155 format!("Unknown order event kind: {kind} in Postgres transformation").into(),
156 ))
157 }
158 }
159}
160
161impl<'r> FromRow<'r, PgRow> for OrderInitializedRow {
162 #[expect(
163 clippy::too_many_lines,
164 reason = "SQL row mapping mirrors the full order initialized event constructor"
165 )]
166 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
167 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
168 let client_order_id = row
169 .try_get::<&str, _>("client_order_id")
170 .map(ClientOrderId::from)?;
171 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
172 let strategy_id = row
173 .try_get::<&str, _>("strategy_id")
174 .map(StrategyId::from)?;
175 let instrument_id = row
176 .try_get::<&str, _>("instrument_id")
177 .map(InstrumentId::from)?;
178 let order_type = row
179 .try_get::<&str, _>("order_type")
180 .map(|x| OrderType::from_str(x).unwrap())?;
181 let order_side = row
182 .try_get::<&str, _>("order_side")
183 .map(|x| OrderSide::from_str(x).unwrap())?;
184 let quantity = row.try_get::<&str, _>("quantity").map(Quantity::from)?;
185 let time_in_force = row
186 .try_get::<&str, _>("time_in_force")
187 .map(|x| TimeInForce::from_str(x).unwrap())?;
188 let post_only = row.try_get::<bool, _>("post_only")?;
189 let reduce_only = row.try_get::<bool, _>("reduce_only")?;
190 let quote_quantity = row.try_get::<bool, _>("quote_quantity")?;
191 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
192 let ts_event = row.try_get::<String, _>("ts_event").map(UnixNanos::from)?;
193 let ts_init = row.try_get::<String, _>("ts_init").map(UnixNanos::from)?;
194 let price = row
195 .try_get::<Option<&str>, _>("price")
196 .ok()
197 .and_then(|x| x.map(Price::from));
198 let activation_price = row
199 .try_get::<Option<&str>, _>("activation_price")
200 .ok()
201 .and_then(|x| x.map(Price::from));
202 let trigger_price = row
203 .try_get::<Option<&str>, _>("trigger_price")
204 .ok()
205 .and_then(|x| x.map(Price::from));
206 let trigger_type = row
207 .try_get::<Option<&str>, _>("trigger_type")
208 .ok()
209 .and_then(parse_trigger_type);
210 let limit_offset = row
211 .try_get::<Option<&str>, _>("limit_offset")
212 .ok()
213 .and_then(|x| x.and_then(|s| Decimal::from_str(s).ok()));
214 let trailing_offset = row
215 .try_get::<Option<&str>, _>("trailing_offset")
216 .ok()
217 .and_then(|x| x.and_then(|s| Decimal::from_str(s).ok()));
218 let trailing_offset_type = row
219 .try_get::<Option<TrailingOffsetTypePg>, _>("trailing_offset_type")
220 .ok()
221 .flatten()
222 .and_then(|value| value.0);
223 let expire_time = row
224 .try_get::<Option<&str>, _>("expire_time")
225 .ok()
226 .and_then(|x| x.map(UnixNanos::from));
227 let display_qty = row
228 .try_get::<Option<&str>, _>("display_qty")
229 .ok()
230 .and_then(|x| x.map(Quantity::from));
231 let emulation_trigger = row
232 .try_get::<Option<&str>, _>("emulation_trigger")
233 .ok()
234 .and_then(parse_trigger_type);
235 let trigger_instrument_id = row
236 .try_get::<Option<&str>, _>("trigger_instrument_id")
237 .ok()
238 .and_then(|x| x.map(InstrumentId::from));
239 let contingency_type = row
240 .try_get::<Option<&str>, _>("contingency_type")
241 .ok()
242 .and_then(parse_contingency_type);
243 let order_list_id = row
244 .try_get::<Option<&str>, _>("order_list_id")
245 .ok()
246 .and_then(|x| x.map(OrderListId::from));
247 let linked_order_ids = row
248 .try_get::<Vec<String>, _>("linked_order_ids")
249 .ok()
250 .map(|x| x.iter().map(|x| ClientOrderId::from(x.as_str())).collect());
251 let parent_order_id = row
252 .try_get::<Option<&str>, _>("parent_order_id")
253 .ok()
254 .and_then(|x| x.map(ClientOrderId::from));
255 let exec_algorithm_id = row
256 .try_get::<Option<&str>, _>("exec_algorithm_id")
257 .ok()
258 .and_then(|x| x.map(ExecAlgorithmId::from));
259 let exec_algorithm_params: Option<IndexMap<Ustr, Ustr>> = row
260 .try_get::<Option<serde_json::Value>, _>("exec_algorithm_params")
261 .ok()
262 .and_then(|x| x.map(|x| serde_json::from_value::<IndexMap<String, String>>(x).unwrap()))
263 .map(|x| {
264 x.into_iter()
265 .map(|(k, v)| (Ustr::from(k.as_str()), Ustr::from(v.as_str())))
266 .collect()
267 });
268 let exec_spawn_id = row
269 .try_get::<Option<&str>, _>("exec_spawn_id")
270 .ok()
271 .and_then(|x| x.map(ClientOrderId::from));
272 let tags = tags_from_row(row);
273 let mut order_event = OrderInitialized::new_checked(
274 trader_id,
275 strategy_id,
276 instrument_id,
277 client_order_id,
278 order_side,
279 order_type,
280 quantity,
281 time_in_force,
282 post_only,
283 reduce_only,
284 quote_quantity,
285 reconciliation,
286 event_id,
287 ts_event,
288 ts_init,
289 price,
290 activation_price,
291 trigger_price,
292 trigger_type,
293 limit_offset,
294 trailing_offset,
295 trailing_offset_type,
296 expire_time,
297 display_qty,
298 emulation_trigger,
299 trigger_instrument_id,
300 contingency_type,
301 order_list_id,
302 linked_order_ids,
303 parent_order_id,
304 exec_algorithm_id,
305 exec_algorithm_params,
306 exec_spawn_id,
307 tags,
308 )
309 .map_err(|e| sqlx::Error::Decode(Box::new(e)))?;
310 order_event.causation_id = causation_id_from_row(row)?;
311 Ok(Self(order_event))
312 }
313}
314
315impl<'r> FromRow<'r, PgRow> for OrderAcceptedRow {
316 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
317 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
318 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
319 let strategy_id = row
320 .try_get::<&str, _>("strategy_id")
321 .map(StrategyId::from)?;
322 let instrument_id = row
323 .try_get::<&str, _>("instrument_id")
324 .map(InstrumentId::from)?;
325 let client_order_id = row
326 .try_get::<&str, _>("client_order_id")
327 .map(ClientOrderId::from)?;
328 let venue_order_id = row
329 .try_get::<&str, _>("venue_order_id")
330 .map(VenueOrderId::from)?;
331 let account_id = row.try_get::<&str, _>("account_id").map(AccountId::from)?;
332 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
333 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
334 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
335 let causation_id = causation_id_from_row(row)?;
336 let mut order_event = OrderAccepted::new(
337 trader_id,
338 strategy_id,
339 instrument_id,
340 client_order_id,
341 venue_order_id,
342 account_id,
343 event_id,
344 ts_event,
345 ts_init,
346 reconciliation,
347 );
348 order_event.causation_id = causation_id;
349 Ok(Self(order_event))
350 }
351}
352
353impl<'r> FromRow<'r, PgRow> for OrderCancelRejectedRow {
354 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
355 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
356 let strategy_id = row
357 .try_get::<&str, _>("strategy_id")
358 .map(StrategyId::from)?;
359 let instrument_id = row
360 .try_get::<&str, _>("instrument_id")
361 .map(InstrumentId::from)?;
362 let client_order_id = row
363 .try_get::<&str, _>("client_order_id")
364 .map(ClientOrderId::from)?;
365 let reason = row.try_get::<&str, _>("reason").map(Ustr::from)?;
366 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
367 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
368 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
369 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
370 let venue_order_id = row
371 .try_get::<Option<&str>, _>("venue_order_id")?
372 .map(Into::into);
373 let account_id = row
374 .try_get::<Option<&str>, _>("account_id")?
375 .map(Into::into);
376 let causation_id = causation_id_from_row(row)?;
377 let mut order_event = OrderCancelRejected::new(
378 trader_id,
379 strategy_id,
380 instrument_id,
381 client_order_id,
382 reason,
383 event_id,
384 ts_event,
385 ts_init,
386 reconciliation,
387 venue_order_id,
388 account_id,
389 );
390 order_event.causation_id = causation_id;
391 Ok(Self(order_event))
392 }
393}
394
395impl<'r> FromRow<'r, PgRow> for OrderCanceledRow {
396 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
397 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
398 let strategy_id = row
399 .try_get::<&str, _>("strategy_id")
400 .map(StrategyId::from)?;
401 let instrument_id = row
402 .try_get::<&str, _>("instrument_id")
403 .map(InstrumentId::from)?;
404 let client_order_id = row
405 .try_get::<&str, _>("client_order_id")
406 .map(ClientOrderId::from)?;
407 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
408 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
409 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
410 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
411 let venue_order_id = row
412 .try_get::<Option<&str>, _>("venue_order_id")?
413 .map(Into::into);
414 let account_id = row
415 .try_get::<Option<&str>, _>("account_id")?
416 .map(Into::into);
417 let reason = row.try_get::<Option<&str>, _>("reason")?.map(Ustr::from);
418 let causation_id = causation_id_from_row(row)?;
419 let mut order_event = OrderCanceled::new(
420 trader_id,
421 strategy_id,
422 instrument_id,
423 client_order_id,
424 event_id,
425 ts_event,
426 ts_init,
427 reconciliation,
428 venue_order_id,
429 account_id,
430 reason,
431 );
432 order_event.causation_id = causation_id;
433 Ok(Self(order_event))
434 }
435}
436
437impl<'r> FromRow<'r, PgRow> for OrderDeniedRow {
438 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
439 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
440 let strategy_id = row
441 .try_get::<&str, _>("strategy_id")
442 .map(StrategyId::from)?;
443 let instrument_id = row
444 .try_get::<&str, _>("instrument_id")
445 .map(InstrumentId::from)?;
446 let client_order_id = row
447 .try_get::<&str, _>("client_order_id")
448 .map(ClientOrderId::from)?;
449 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
450 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
451 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
452 let reason = row.try_get::<&str, _>("reason").map(Ustr::from)?;
453 let causation_id = causation_id_from_row(row)?;
454 let mut order_event = OrderDenied::new(
455 trader_id,
456 strategy_id,
457 instrument_id,
458 client_order_id,
459 reason,
460 event_id,
461 ts_event,
462 ts_init,
463 );
464 order_event.causation_id = causation_id;
465 Ok(Self(order_event))
466 }
467}
468
469impl<'r> FromRow<'r, PgRow> for OrderEmulatedRow {
470 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
471 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
472 let strategy_id = row
473 .try_get::<&str, _>("strategy_id")
474 .map(StrategyId::from)?;
475 let instrument_id = row
476 .try_get::<&str, _>("instrument_id")
477 .map(InstrumentId::from)?;
478 let client_order_id = row
479 .try_get::<&str, _>("client_order_id")
480 .map(ClientOrderId::from)?;
481 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
482 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
483 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
484 let causation_id = causation_id_from_row(row)?;
485 let mut order_event = OrderEmulated::new(
486 trader_id,
487 strategy_id,
488 instrument_id,
489 client_order_id,
490 event_id,
491 ts_event,
492 ts_init,
493 );
494 order_event.causation_id = causation_id;
495 Ok(Self(order_event))
496 }
497}
498
499impl<'r> FromRow<'r, PgRow> for OrderExpiredRow {
500 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
501 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
502 let strategy_id = row
503 .try_get::<&str, _>("strategy_id")
504 .map(StrategyId::from)?;
505 let instrument_id = row
506 .try_get::<&str, _>("instrument_id")
507 .map(InstrumentId::from)?;
508 let client_order_id = row
509 .try_get::<&str, _>("client_order_id")
510 .map(ClientOrderId::from)?;
511 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
512 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
513 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
514 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
515 let venue_order_id = row
516 .try_get::<Option<&str>, _>("venue_order_id")?
517 .map(Into::into);
518 let account_id = row
519 .try_get::<Option<&str>, _>("account_id")?
520 .map(Into::into);
521 let causation_id = causation_id_from_row(row)?;
522 let mut order_event = OrderExpired::new(
523 trader_id,
524 strategy_id,
525 instrument_id,
526 client_order_id,
527 event_id,
528 ts_event,
529 ts_init,
530 reconciliation,
531 venue_order_id,
532 account_id,
533 );
534 order_event.causation_id = causation_id;
535 Ok(Self(order_event))
536 }
537}
538
539impl<'r> FromRow<'r, PgRow> for OrderFilledRow {
540 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
541 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
542 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
543 let strategy_id = row
544 .try_get::<&str, _>("strategy_id")
545 .map(StrategyId::from)?;
546 let instrument_id = row
547 .try_get::<&str, _>("instrument_id")
548 .map(InstrumentId::from)?;
549 let client_order_id = row
550 .try_get::<&str, _>("client_order_id")
551 .map(ClientOrderId::from)?;
552 let venue_order_id = row
553 .try_get::<&str, _>("venue_order_id")
554 .map(VenueOrderId::from)?;
555 let account_id = row.try_get::<&str, _>("account_id").map(AccountId::from)?;
556 let trade_id = row.try_get::<&str, _>("trade_id").map(TradeId::from)?;
557 let order_side = row
558 .try_get::<&str, _>("order_side")
559 .map(|x| OrderSide::from_str(x).unwrap())?;
560 let order_type = row
561 .try_get::<&str, _>("order_type")
562 .map(|x| OrderType::from_str(x).unwrap())?;
563 let last_px = row.try_get::<&str, _>("last_px").map(Price::from)?;
564 let last_qty = row.try_get::<&str, _>("last_qty").map(Quantity::from)?;
565 let currency = row.try_get::<&str, _>("currency").map(Currency::from)?;
566 let liquidity_side = row
567 .try_get::<&str, _>("liquidity_side")
568 .map(|x| LiquiditySide::from_str(x).unwrap())?;
569 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
570 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
571 let position_id = row
572 .try_get::<Option<&str>, _>("position_id")
573 .map(|x| x.map(PositionId::from))?;
574 let commission = row
575 .try_get::<Option<&str>, _>("commission")
576 .map(|x| x.map(|x| Money::from_str(x).unwrap()))?;
577 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
578 let info = decode_info(row)?;
579 let causation_id = causation_id_from_row(row)?;
580 let mut order_event = OrderFilled::new(
581 trader_id,
582 strategy_id,
583 instrument_id,
584 client_order_id,
585 venue_order_id,
586 account_id,
587 trade_id,
588 order_side,
589 order_type,
590 last_qty,
591 last_px,
592 currency,
593 liquidity_side,
594 event_id,
595 ts_event,
596 ts_init,
597 reconciliation,
598 position_id,
599 commission,
600 info,
601 );
602 order_event.causation_id = causation_id;
603 Ok(Self(order_event))
604 }
605}
606
607impl<'r> FromRow<'r, PgRow> for OrderFillVoidedRow {
608 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
609 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
610 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
611 let strategy_id = row
612 .try_get::<&str, _>("strategy_id")
613 .map(StrategyId::from)?;
614 let instrument_id = row
615 .try_get::<&str, _>("instrument_id")
616 .map(InstrumentId::from)?;
617 let client_order_id = row
618 .try_get::<&str, _>("client_order_id")
619 .map(ClientOrderId::from)?;
620 let venue_order_id = row
621 .try_get::<&str, _>("venue_order_id")
622 .map(VenueOrderId::from)?;
623 let account_id = row.try_get::<&str, _>("account_id").map(AccountId::from)?;
624 let correction_id = row
625 .try_get::<Option<&str>, _>("correction_id")?
626 .map(Ustr::from)
627 .ok_or_else(|| {
628 sqlx::Error::Decode(
629 "OrderFillVoided row has no correction_id; it predates the column and \
630 the value cannot be recovered"
631 .into(),
632 )
633 })?;
634 let trade_id = row.try_get::<&str, _>("trade_id").map(TradeId::from)?;
635 let voided_qty = row.try_get::<&str, _>("quantity").map(Quantity::from)?;
636 let commission_voided = row
637 .try_get::<Option<&str>, _>("commission")?
638 .map(|x| Money::from_str(x).map_err(|e| sqlx::Error::Decode(e.into())))
639 .transpose()?;
640 let order_side = OrderSide::from_str(row.try_get::<&str, _>("order_side")?)
641 .map_err(|e| sqlx::Error::Decode(Box::new(e)))?;
642 let order_type = OrderType::from_str(row.try_get::<&str, _>("order_type")?)
643 .map_err(|e| sqlx::Error::Decode(Box::new(e)))?;
644 let last_px = row.try_get::<&str, _>("last_px").map(Price::from)?;
645 let currency = row.try_get::<&str, _>("currency").map(Currency::from)?;
646 let liquidity_side = LiquiditySide::from_str(row.try_get::<&str, _>("liquidity_side")?)
647 .map_err(|e| sqlx::Error::Decode(Box::new(e)))?;
648 let position_id = row
649 .try_get::<Option<&str>, _>("position_id")?
650 .map(PositionId::from);
651 let reason = row.try_get::<Option<&str>, _>("reason")?.map(Ustr::from);
652 let info = decode_info(row)?;
653 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
654 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
655 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
656 let is_reopened = row
657 .try_get::<Option<bool>, _>("is_reopened")?
658 .unwrap_or(false);
659 let causation_id = causation_id_from_row(row)?;
660 let mut order_event = OrderFillVoided::new(
661 trader_id,
662 strategy_id,
663 instrument_id,
664 client_order_id,
665 venue_order_id,
666 account_id,
667 correction_id,
668 trade_id,
669 voided_qty,
670 commission_voided,
671 order_side,
672 order_type,
673 last_px,
674 currency,
675 liquidity_side,
676 position_id,
677 reason,
678 info,
679 event_id,
680 ts_event,
681 ts_init,
682 reconciliation,
683 is_reopened,
684 );
685 order_event.causation_id = causation_id;
686 Ok(Self(order_event))
687 }
688}
689
690impl<'r> FromRow<'r, PgRow> for OrderModifyRejectedRow {
691 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
692 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
693 let strategy_id = row
694 .try_get::<&str, _>("strategy_id")
695 .map(StrategyId::from)?;
696 let instrument_id = row
697 .try_get::<&str, _>("instrument_id")
698 .map(InstrumentId::from)?;
699 let client_order_id = row
700 .try_get::<&str, _>("client_order_id")
701 .map(ClientOrderId::from)?;
702 let reason = row.try_get::<&str, _>("reason").map(Ustr::from)?;
703 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
704 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
705 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
706 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
707 let venue_order_id = row
708 .try_get::<Option<&str>, _>("venue_order_id")?
709 .map(Into::into);
710 let account_id = row
711 .try_get::<Option<&str>, _>("account_id")?
712 .map(Into::into);
713 let causation_id = causation_id_from_row(row)?;
714 let mut order_event = OrderModifyRejected::new(
715 trader_id,
716 strategy_id,
717 instrument_id,
718 client_order_id,
719 reason,
720 event_id,
721 ts_event,
722 ts_init,
723 reconciliation,
724 venue_order_id,
725 account_id,
726 );
727 order_event.causation_id = causation_id;
728 Ok(Self(order_event))
729 }
730}
731
732impl<'r> FromRow<'r, PgRow> for OrderPendingCancelRow {
733 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
734 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
735 let strategy_id = row
736 .try_get::<&str, _>("strategy_id")
737 .map(StrategyId::from)?;
738 let instrument_id = row
739 .try_get::<&str, _>("instrument_id")
740 .map(InstrumentId::from)?;
741 let client_order_id = row
742 .try_get::<&str, _>("client_order_id")
743 .map(ClientOrderId::from)?;
744 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
745 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
746 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
747 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
748 let venue_order_id = row
749 .try_get::<Option<&str>, _>("venue_order_id")?
750 .map(Into::into);
751 let account_id = row
752 .try_get::<Option<&str>, _>("account_id")?
753 .map(Into::into);
754 let causation_id = causation_id_from_row(row)?;
755 let mut order_event = OrderPendingCancel::new(
756 trader_id,
757 strategy_id,
758 instrument_id,
759 client_order_id,
760 account_id,
761 event_id,
762 ts_event,
763 ts_init,
764 reconciliation,
765 venue_order_id,
766 );
767 order_event.causation_id = causation_id;
768 Ok(Self(order_event))
769 }
770}
771
772impl<'r> FromRow<'r, PgRow> for OrderPendingUpdateRow {
773 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
774 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
775 let strategy_id = row
776 .try_get::<&str, _>("strategy_id")
777 .map(StrategyId::from)?;
778 let instrument_id = row
779 .try_get::<&str, _>("instrument_id")
780 .map(InstrumentId::from)?;
781 let client_order_id = row
782 .try_get::<&str, _>("client_order_id")
783 .map(ClientOrderId::from)?;
784 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
785 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
786 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
787 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
788 let venue_order_id = row
789 .try_get::<Option<&str>, _>("venue_order_id")?
790 .map(Into::into);
791 let account_id = row
792 .try_get::<Option<&str>, _>("account_id")?
793 .map(Into::into);
794 let causation_id = causation_id_from_row(row)?;
795 let mut order_event = OrderPendingUpdate::new(
796 trader_id,
797 strategy_id,
798 instrument_id,
799 client_order_id,
800 account_id,
801 event_id,
802 ts_event,
803 ts_init,
804 reconciliation,
805 venue_order_id,
806 );
807 order_event.causation_id = causation_id;
808 Ok(Self(order_event))
809 }
810}
811
812impl<'r> FromRow<'r, PgRow> for OrderRejectedRow {
813 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
814 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
815 let strategy_id = row
816 .try_get::<&str, _>("strategy_id")
817 .map(StrategyId::from)?;
818 let instrument_id = row
819 .try_get::<&str, _>("instrument_id")
820 .map(InstrumentId::from)?;
821 let client_order_id = row
822 .try_get::<&str, _>("client_order_id")
823 .map(ClientOrderId::from)?;
824 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
825 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
826 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
827 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
828 let account_id = row.try_get::<&str, _>("account_id").map(AccountId::from)?;
829 let reason = row.try_get::<&str, _>("reason").map(Ustr::from)?;
830 let due_post_only = row
832 .try_get::<Option<bool>, _>("due_post_only")?
833 .unwrap_or(false);
834 let causation_id = causation_id_from_row(row)?;
835 let mut order_event = OrderRejected::new(
836 trader_id,
837 strategy_id,
838 instrument_id,
839 client_order_id,
840 account_id,
841 reason,
842 event_id,
843 ts_event,
844 ts_init,
845 reconciliation,
846 due_post_only,
847 );
848 order_event.causation_id = causation_id;
849 Ok(Self(order_event))
850 }
851}
852
853impl<'r> FromRow<'r, PgRow> for OrderReleasedRow {
854 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
855 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
856 let strategy_id = row
857 .try_get::<&str, _>("strategy_id")
858 .map(StrategyId::from)?;
859 let instrument_id = row
860 .try_get::<&str, _>("instrument_id")
861 .map(InstrumentId::from)?;
862 let client_order_id = row
863 .try_get::<&str, _>("client_order_id")
864 .map(ClientOrderId::from)?;
865 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
866 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
867 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
868 let released_price = row
869 .try_get::<Option<&str>, _>("released_price")?
870 .map(Price::from)
871 .ok_or_else(|| {
872 sqlx::Error::Decode(
873 "OrderReleased row has no released_price; it predates the column and \
874 the value cannot be recovered"
875 .into(),
876 )
877 })?;
878 let causation_id = causation_id_from_row(row)?;
879 let mut order_event = OrderReleased::new(
880 trader_id,
881 strategy_id,
882 instrument_id,
883 client_order_id,
884 released_price,
885 event_id,
886 ts_event,
887 ts_init,
888 );
889 order_event.causation_id = causation_id;
890 Ok(Self(order_event))
891 }
892}
893
894impl<'r> FromRow<'r, PgRow> for OrderSubmittedRow {
895 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
896 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
897 let strategy_id = row
898 .try_get::<&str, _>("strategy_id")
899 .map(StrategyId::from)?;
900 let instrument_id = row
901 .try_get::<&str, _>("instrument_id")
902 .map(InstrumentId::from)?;
903 let client_order_id = row
904 .try_get::<&str, _>("client_order_id")
905 .map(ClientOrderId::from)?;
906 let account_id = row.try_get::<&str, _>("account_id").map(AccountId::from)?;
907 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
908 let ts_event = row
909 .try_get::<String, _>("ts_event")
910 .map(|res| UnixNanos::from(res.as_str()))?;
911 let ts_init = row
912 .try_get::<String, _>("ts_init")
913 .map(|res| UnixNanos::from(res.as_str()))?;
914 let causation_id = causation_id_from_row(row)?;
915 let mut order_event = OrderSubmitted::new(
916 trader_id,
917 strategy_id,
918 instrument_id,
919 client_order_id,
920 account_id,
921 event_id,
922 ts_event,
923 ts_init,
924 );
925 order_event.causation_id = causation_id;
926 Ok(Self(order_event))
927 }
928}
929
930impl<'r> FromRow<'r, PgRow> for OrderTriggeredRow {
931 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
932 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
933 let strategy_id = row
934 .try_get::<&str, _>("strategy_id")
935 .map(StrategyId::from)?;
936 let instrument_id = row
937 .try_get::<&str, _>("instrument_id")
938 .map(InstrumentId::from)?;
939 let client_order_id = row
940 .try_get::<&str, _>("client_order_id")
941 .map(ClientOrderId::from)?;
942 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
943 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
944 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
945 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
946 let venue_order_id = row
947 .try_get::<Option<&str>, _>("venue_order_id")?
948 .map(Into::into);
949 let account_id = row
950 .try_get::<Option<&str>, _>("account_id")?
951 .map(Into::into);
952 let causation_id = causation_id_from_row(row)?;
953 let mut order_event = OrderTriggered::new(
954 trader_id,
955 strategy_id,
956 instrument_id,
957 client_order_id,
958 event_id,
959 ts_event,
960 ts_init,
961 reconciliation,
962 venue_order_id,
963 account_id,
964 );
965 order_event.causation_id = causation_id;
966 Ok(Self(order_event))
967 }
968}
969
970impl<'r> FromRow<'r, PgRow> for OrderUpdatedRow {
971 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
972 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
973 let strategy_id = row
974 .try_get::<&str, _>("strategy_id")
975 .map(StrategyId::from)?;
976 let instrument_id = row
977 .try_get::<&str, _>("instrument_id")
978 .map(InstrumentId::from)?;
979 let client_order_id = row
980 .try_get::<&str, _>("client_order_id")
981 .map(ClientOrderId::from)?;
982 let event_id = row.try_get::<&str, _>("id").map(UUID4::from)?;
983 let ts_event = row.try_get::<&str, _>("ts_event").map(UnixNanos::from)?;
984 let ts_init = row.try_get::<&str, _>("ts_init").map(UnixNanos::from)?;
985 let reconciliation = row.try_get::<bool, _>("reconciliation")?;
986 let venue_order_id = row
987 .try_get::<Option<&str>, _>("venue_order_id")?
988 .map(Into::into);
989 let account_id = row
990 .try_get::<Option<&str>, _>("account_id")?
991 .map(Into::into);
992 let quantity = row.try_get::<&str, _>("quantity").map(Quantity::from)?;
993 let price = row.try_get::<Option<&str>, _>("price")?.map(Price::from);
994 let trigger_price = row
995 .try_get::<Option<&str>, _>("trigger_price")?
996 .map(Price::from);
997 let protection_price = row
998 .try_get::<Option<&str>, _>("protection_price")?
999 .map(Price::from);
1000 let is_quote_quantity = row.try_get::<bool, _>("quote_quantity")?;
1001 let causation_id = causation_id_from_row(row)?;
1002 let mut order_event = OrderUpdated::new(
1003 trader_id,
1004 strategy_id,
1005 instrument_id,
1006 client_order_id,
1007 quantity,
1008 event_id,
1009 ts_event,
1010 ts_init,
1011 reconciliation,
1012 venue_order_id,
1013 account_id,
1014 price,
1015 trigger_price,
1016 protection_price,
1017 is_quote_quantity,
1018 );
1019 order_event.causation_id = causation_id;
1020 Ok(Self(order_event))
1021 }
1022}
1023
1024impl<'r> FromRow<'r, PgRow> for OrderSnapshotRow {
1025 #[expect(
1026 clippy::too_many_lines,
1027 reason = "SQL row mapping mirrors the full order snapshot constructor"
1028 )]
1029 fn from_row(row: &'r PgRow) -> Result<Self, sqlx::Error> {
1030 let trader_id = row.try_get::<&str, _>("trader_id").map(TraderId::from)?;
1031 let strategy_id = row
1032 .try_get::<&str, _>("strategy_id")
1033 .map(StrategyId::from)?;
1034 let instrument_id = row
1035 .try_get::<&str, _>("instrument_id")
1036 .map(InstrumentId::from)?;
1037 let client_order_id = row
1038 .try_get::<&str, _>("client_order_id")
1039 .map(ClientOrderId::from)?;
1040 let venue_order_id = row
1041 .try_get::<Option<&str>, _>("venue_order_id")
1042 .ok()
1043 .and_then(|x| x.map(VenueOrderId::from));
1044 let position_id = row
1045 .try_get::<Option<&str>, _>("position_id")
1046 .ok()
1047 .and_then(|x| x.map(PositionId::from));
1048 let account_id = row
1049 .try_get::<Option<&str>, _>("account_id")
1050 .ok()
1051 .and_then(|x| x.map(AccountId::from));
1052 let last_trade_id = row
1053 .try_get::<Option<&str>, _>("last_trade_id")
1054 .ok()
1055 .and_then(|x| x.map(TradeId::from));
1056 let order_type = row
1057 .try_get::<&str, _>("order_type")
1058 .map(|x| OrderType::from_str(x).expect("Invalid `OrderType`"))?;
1059 let order_side = row
1060 .try_get::<&str, _>("order_side")
1061 .map(|x| OrderSide::from_str(x).expect("Invalid `OrderSide`"))?;
1062 let quantity = row.try_get::<&str, _>("quantity").map(Quantity::from)?;
1063 let price = row
1064 .try_get::<Option<&str>, _>("price")
1065 .ok()
1066 .and_then(|x| x.map(Price::from));
1067 let activation_price = row
1068 .try_get::<Option<&str>, _>("activation_price")
1069 .ok()
1070 .and_then(|x| x.map(Price::from));
1071 let trigger_price = row
1072 .try_get::<Option<&str>, _>("trigger_price")
1073 .ok()
1074 .and_then(|x| x.map(Price::from));
1075 let trigger_type = row
1076 .try_get::<Option<&str>, _>("trigger_type")
1077 .ok()
1078 .and_then(parse_trigger_type);
1079 let limit_offset = row
1080 .try_get::<Option<&str>, _>("limit_offset")
1081 .ok()
1082 .and_then(|x| x.and_then(|s| Decimal::from_str(s).ok()));
1083 let trailing_offset = row
1084 .try_get::<Option<&str>, _>("trailing_offset")
1085 .ok()
1086 .and_then(|x| x.and_then(|s| Decimal::from_str(s).ok()));
1087 let trailing_offset_type = row
1088 .try_get::<Option<TrailingOffsetTypePg>, _>("trailing_offset_type")
1089 .ok()
1090 .flatten()
1091 .and_then(|value| value.0);
1092 let time_in_force = row
1093 .try_get::<&str, _>("time_in_force")
1094 .map(|x| TimeInForce::from_str(x).expect("Invalid `TimeInForce`"))?;
1095 let expire_time = row
1096 .try_get::<Option<&str>, _>("expire_time")
1097 .ok()
1098 .and_then(|x| x.map(UnixNanos::from));
1099 let filled_qty = row.try_get::<&str, _>("filled_qty").map(Quantity::from)?;
1100 let liquidity_side = row
1101 .try_get::<Option<&str>, _>("liquidity_side")
1102 .ok()
1103 .and_then(|x| x.map(|x| LiquiditySide::from_str(x).expect("Invalid `LiquiditySide`")));
1104 let avg_px = row.try_get::<Option<Decimal>, _>("avg_px").ok().flatten();
1105 let slippage = row.try_get::<Option<Decimal>, _>("slippage").ok().flatten();
1106 let commissions = row
1107 .try_get::<Option<Vec<String>>, _>("commissions")?
1108 .map_or_else(Vec::new, |c| {
1109 c.into_iter().map(|s| Money::from(&s)).collect()
1110 });
1111 let status = row
1112 .try_get::<&str, _>("status")
1113 .map(|x| OrderStatus::from_str(x).expect("Invalid `OrderStatus`"))?;
1114 let is_post_only = row.try_get::<bool, _>("is_post_only")?;
1115 let is_reduce_only = row.try_get::<bool, _>("is_reduce_only")?;
1116 let is_quote_quantity = row.try_get::<bool, _>("is_quote_quantity")?;
1117 let display_qty = row
1118 .try_get::<Option<&str>, _>("display_qty")
1119 .ok()
1120 .and_then(|x| x.map(Quantity::from));
1121 let emulation_trigger = row
1122 .try_get::<Option<&str>, _>("emulation_trigger")
1123 .ok()
1124 .and_then(parse_trigger_type);
1125 let trigger_instrument_id = row
1126 .try_get::<Option<&str>, _>("trigger_instrument_id")
1127 .ok()
1128 .and_then(|x| x.map(InstrumentId::from));
1129 let contingency_type = row
1130 .try_get::<Option<&str>, _>("contingency_type")
1131 .ok()
1132 .and_then(parse_contingency_type);
1133 let order_list_id = row
1134 .try_get::<Option<&str>, _>("order_list_id")
1135 .ok()
1136 .and_then(|x| x.map(OrderListId::from));
1137 let linked_order_ids = row
1138 .try_get::<Option<Vec<String>>, _>("linked_order_ids")
1139 .ok()
1140 .and_then(|ids| ids.map(|ids| ids.into_iter().map(ClientOrderId::from).collect()));
1141 let parent_order_id = row
1142 .try_get::<Option<&str>, _>("parent_order_id")
1143 .ok()
1144 .and_then(|x| x.map(ClientOrderId::from));
1145 let exec_algorithm_id = row
1146 .try_get::<Option<&str>, _>("exec_algorithm_id")
1147 .ok()
1148 .and_then(|x| x.map(ExecAlgorithmId::from));
1149 let exec_algorithm_params: Option<IndexMap<Ustr, Ustr>> = row
1150 .try_get::<Option<serde_json::Value>, _>("exec_algorithm_params")
1151 .ok()
1152 .and_then(|x| {
1153 x.map(|x| {
1154 serde_json::from_value::<IndexMap<String, String>>(x)
1155 .expect("Invalid exec algorithm params")
1156 })
1157 })
1158 .map(|x| {
1159 x.into_iter()
1160 .map(|(k, v)| (Ustr::from(k.as_str()), Ustr::from(v.as_str())))
1161 .collect()
1162 });
1163 let exec_spawn_id = row
1164 .try_get::<Option<&str>, _>("exec_spawn_id")
1165 .ok()
1166 .and_then(|x| x.map(ClientOrderId::from));
1167 let tags = tags_from_row(row);
1168 let init_id = row.try_get::<&str, _>("init_id").map(UUID4::from)?;
1169 let ts_init = row.try_get::<String, _>("ts_init").map(UnixNanos::from)?;
1170 let ts_last = row.try_get::<String, _>("ts_last").map(UnixNanos::from)?;
1171
1172 let snapshot = OrderSnapshot {
1173 trader_id,
1174 strategy_id,
1175 instrument_id,
1176 client_order_id,
1177 venue_order_id,
1178 position_id,
1179 account_id,
1180 last_trade_id,
1181 order_type,
1182 order_side,
1183 quantity,
1184 price,
1185 activation_price,
1186 trigger_price,
1187 trigger_type,
1188 limit_offset,
1189 trailing_offset,
1190 trailing_offset_type,
1191 time_in_force,
1192 expire_time,
1193 filled_qty,
1194 liquidity_side,
1195 avg_px,
1196 slippage,
1197 commissions,
1198 status,
1199 is_post_only,
1200 is_reduce_only,
1201 is_quote_quantity,
1202 display_qty,
1203 emulation_trigger,
1204 trigger_instrument_id,
1205 contingency_type,
1206 order_list_id,
1207 linked_order_ids,
1208 parent_order_id,
1209 exec_algorithm_id,
1210 exec_algorithm_params,
1211 exec_spawn_id,
1212 tags,
1213 init_id,
1214 ts_init,
1215 ts_last,
1216 causation_id: None,
1217 };
1218
1219 Ok(Self(snapshot))
1220 }
1221}
1222
1223fn causation_id_from_row(row: &PgRow) -> Result<Option<UUID4>, sqlx::Error> {
1224 row.try_get::<Option<&str>, _>("causation_id")?
1225 .map(|value| UUID4::from_str(value).map_err(|e| sqlx::Error::Decode(e.into())))
1226 .transpose()
1227}
1228
1229fn decode_info(row: &PgRow) -> Result<Option<IndexMap<Ustr, Ustr>>, sqlx::Error> {
1230 let value = row.try_get::<Option<serde_json::Value>, _>("info")?;
1231 let Some(value) = value else {
1232 return Ok(None);
1233 };
1234
1235 let decoded: IndexMap<String, String> =
1236 serde_json::from_value(value).map_err(|e| sqlx::Error::Decode(Box::new(e)))?;
1237
1238 Ok(Some(
1239 decoded
1240 .into_iter()
1241 .map(|(k, v)| (Ustr::from(k.as_str()), Ustr::from(v.as_str())))
1242 .collect(),
1243 ))
1244}
1245
1246fn tags_from_row(row: &PgRow) -> Option<Vec<Ustr>> {
1247 row.try_get::<Vec<String>, _>("tags")
1248 .ok()
1249 .map(|tags| tags.iter().map(|tag| Ustr::from(tag.as_str())).collect())
1250}
1251
1252fn parse_trigger_type(value: Option<&str>) -> Option<TriggerType> {
1253 value.and_then(|value| {
1254 if value.eq_ignore_ascii_case("NO_TRIGGER") {
1255 None
1256 } else {
1257 Some(TriggerType::from_str(value).expect("Invalid `TriggerType`"))
1258 }
1259 })
1260}
1261
1262fn parse_contingency_type(value: Option<&str>) -> Option<ContingencyType> {
1263 value.and_then(|value| {
1264 if value.eq_ignore_ascii_case("NO_CONTINGENCY") {
1265 None
1266 } else {
1267 Some(ContingencyType::from_str(value).expect("Invalid `ContingencyType`"))
1268 }
1269 })
1270}
1271
1272#[cfg(test)]
1273mod tests {
1274 use rstest::rstest;
1275
1276 use super::*;
1277
1278 #[rstest]
1279 #[case(None, None)]
1280 #[case(Some("NO_TRIGGER"), None)]
1281 #[case(Some("LAST_PRICE"), Some(TriggerType::LastPrice))]
1282 fn test_parse_trigger_type_accepts_legacy_absence(
1283 #[case] value: Option<&str>,
1284 #[case] expected: Option<TriggerType>,
1285 ) {
1286 assert_eq!(parse_trigger_type(value), expected);
1287 }
1288
1289 #[rstest]
1290 #[case(None, None)]
1291 #[case(Some("NO_CONTINGENCY"), None)]
1292 #[case(Some("OCO"), Some(ContingencyType::Oco))]
1293 fn test_parse_contingency_type_accepts_legacy_absence(
1294 #[case] value: Option<&str>,
1295 #[case] expected: Option<ContingencyType>,
1296 ) {
1297 assert_eq!(parse_contingency_type(value), expected);
1298 }
1299}