1use nautilus_common::enums::LogColor;
22use nautilus_core::{UUID4, UnixNanos};
23use nautilus_model::{
24 enums::{LiquiditySide, OrderStatus, OrderType},
25 events::{
26 OrderAccepted, OrderCanceled, OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled,
27 OrderRejected, OrderTriggered, OrderUpdated,
28 },
29 identifiers::{AccountId, PositionId, TradeId},
30 instruments::{Instrument, InstrumentAny},
31 orders::{Order, OrderAny, TRIGGERABLE_ORDER_TYPES},
32 reports::{FillReport, OrderStatusReport},
33 types::{Money, Price, Quantity},
34};
35use rust_decimal::Decimal;
36use ustr::Ustr;
37
38use super::{
39 ids::create_inferred_reconciliation_trade_id,
40 positions::{cap_price_at_instrument_max, is_within_single_unit_tolerance},
41};
42
43#[must_use]
54pub fn generate_reconciliation_order_events(
55 order: &OrderAny,
56 report: &OrderStatusReport,
57 instrument: Option<&InstrumentAny>,
58 ts_now: UnixNanos,
59) -> Vec<OrderEventAny> {
60 generate_reconciliation_order_events_inner(
61 order,
62 report,
63 instrument,
64 ts_now,
65 report.order_status == OrderStatus::Voided,
66 None,
67 )
68}
69
70#[must_use]
77pub fn generate_reconciliation_order_snapshot_events(
78 order: &OrderAny,
79 report: &OrderStatusReport,
80 instrument: Option<&InstrumentAny>,
81 ts_now: UnixNanos,
82) -> Vec<OrderEventAny> {
83 generate_reconciliation_order_events_inner(order, report, instrument, ts_now, true, None)
84}
85
86#[must_use]
89pub fn generate_reconciliation_order_snapshot_events_with_commission(
90 order: &OrderAny,
91 report: &OrderStatusReport,
92 instrument: Option<&InstrumentAny>,
93 ts_now: UnixNanos,
94 commission: Option<Money>,
95) -> Vec<OrderEventAny> {
96 generate_reconciliation_order_events_inner(order, report, instrument, ts_now, true, commission)
97}
98
99fn generate_reconciliation_order_events_inner(
100 order: &OrderAny,
101 report: &OrderStatusReport,
102 instrument: Option<&InstrumentAny>,
103 ts_now: UnixNanos,
104 allow_fill_decrease: bool,
105 commission: Option<Money>,
106) -> Vec<OrderEventAny> {
107 if is_superseded_cancel_report(order, report) {
108 let _ = reconcile_order_report(order, report, instrument, ts_now);
109 return Vec::new();
110 }
111
112 if has_material_fill_decrease(order, report) {
113 if !allow_fill_decrease {
114 log::warn!(
115 "Ignoring fill decrease without explicit void evidence for {}: cached={}, venue={}",
116 order.client_order_id(),
117 order.filled_qty(),
118 report.filled_qty,
119 );
120 return reconcile_fill_decrease_terminal(order, report, instrument, ts_now);
121 }
122
123 let mut working = order.clone();
124 let mut events = create_reconciliation_fill_voids(&working, report, ts_now);
125 if events.is_empty() {
126 return reconcile_fill_decrease_terminal(order, report, instrument, ts_now);
127 }
128
129 for event in &events {
130 if let Err(e) = working.apply(event.clone()) {
131 log::warn!(
132 "Cannot project reconciliation fill void for {}: {e}",
133 order.client_order_id()
134 );
135 return reconcile_fill_decrease_terminal(order, report, instrument, ts_now);
136 }
137 }
138
139 if report_is_working(report) && should_reconciliation_update(&working, report) {
140 let updated = create_reconciliation_updated(&working, report, ts_now);
141 if let Err(e) = working.apply(updated.clone()) {
142 log::warn!(
143 "Cannot project reopening reconciliation update for {}: {e}",
144 order.client_order_id()
145 );
146 } else {
147 events.push(updated);
148 }
149 }
150
151 if let Some(terminal) = reconcile_order_report(&working, report, instrument, ts_now) {
152 events.push(terminal);
153 }
154 return events;
155 }
156
157 let (mut working, mut events) = prepare_reconciliation_order(order, report, ts_now);
158
159 if matches!(
160 report.order_status,
161 OrderStatus::Canceled | OrderStatus::Expired,
162 ) && report.filled_qty > working.filled_qty()
163 && let Some(instrument) = instrument
164 && let Some(filled) = create_incremental_inferred_fill(
165 &working,
166 report,
167 &report.account_id,
168 instrument,
169 ts_now,
170 commission,
171 )
172 {
173 if let Err(e) = working.apply(filled.clone()) {
174 log::warn!(
175 "Failed to pre-apply reconciliation fill for {}: {e}",
176 order.client_order_id(),
177 );
178 } else {
179 events.push(filled);
180 }
181 }
182
183 if working.status() == OrderStatus::Filled
184 && matches!(
185 report.order_status,
186 OrderStatus::Canceled | OrderStatus::Expired,
187 )
188 {
189 return events;
190 }
191
192 if let Some(event) =
193 reconcile_order_report_with_commission(&working, report, instrument, ts_now, commission)
194 {
195 events.push(event);
196 }
197
198 events
199}
200
201fn has_material_fill_decrease(order: &OrderAny, report: &OrderStatusReport) -> bool {
202 if report.filled_qty >= order.filled_qty() {
203 return false;
204 }
205
206 let precision = order
207 .filled_qty()
208 .precision
209 .max(report.filled_qty.precision);
210 !is_within_single_unit_tolerance(
211 report.filled_qty.as_decimal(),
212 order.filled_qty().as_decimal(),
213 precision,
214 )
215}
216
217fn reconcile_fill_decrease_terminal(
218 order: &OrderAny,
219 report: &OrderStatusReport,
220 instrument: Option<&InstrumentAny>,
221 ts_now: UnixNanos,
222) -> Vec<OrderEventAny> {
223 let should_reconcile = report.order_status == OrderStatus::Voided
224 || (!order.is_closed()
225 && matches!(
226 report.order_status,
227 OrderStatus::Canceled | OrderStatus::Expired
228 ));
229
230 if !should_reconcile {
231 return Vec::new();
232 }
233
234 reconcile_order_report(order, report, instrument, ts_now)
235 .into_iter()
236 .collect()
237}
238
239#[must_use]
241pub fn generate_reconciliation_order_pre_fill_events(
242 order: &OrderAny,
243 report: &OrderStatusReport,
244 ts_now: UnixNanos,
245) -> Vec<OrderEventAny> {
246 if is_superseded_cancel_report(order, report) {
247 return Vec::new();
248 }
249
250 let (working, mut events) = prepare_reconciliation_order(order, report, ts_now);
251
252 if report.order_status == OrderStatus::Triggered
253 && let Some(triggered) = reconcile_order_report(&working, report, None, ts_now)
254 {
255 events.push(triggered);
256 }
257
258 events
259}
260
261fn prepare_reconciliation_order(
262 order: &OrderAny,
263 report: &OrderStatusReport,
264 ts_now: UnixNanos,
265) -> (OrderAny, Vec<OrderEventAny>) {
266 let mut working = order.clone();
267 let mut events: Vec<OrderEventAny> = Vec::new();
268
269 if should_accept_before_reconciliation(&working, report) {
270 let Some(accepted) = create_reconciliation_accepted(&working, report, ts_now) else {
271 log::warn!(
272 "Cannot create reconciliation acceptance for {}: missing account_id",
273 order.client_order_id(),
274 );
275 return (working, events);
276 };
277
278 if let Err(e) = working.apply(accepted.clone()) {
279 log::warn!(
280 "Failed to pre-apply reconciliation acceptance for {}: {e}",
281 order.client_order_id(),
282 );
283 return (working, events);
284 }
285 events.push(accepted);
286 }
287
288 if report_is_confirmed_state(report)
289 && (local_accepts_amendment(&working)
290 || (matches!(
291 working.status(),
292 OrderStatus::PendingUpdate | OrderStatus::PendingCancel,
293 ) && matches!(
294 report.order_status,
295 OrderStatus::Canceled | OrderStatus::Expired,
296 )))
297 && should_reconciliation_update(&working, report)
298 {
299 let updated = create_reconciliation_updated(&working, report, ts_now);
300 if let Err(e) = working.apply(updated.clone()) {
301 log::warn!(
302 "Failed to pre-apply reconciliation update for {}: {e}",
303 order.client_order_id(),
304 );
305 } else {
306 events.push(updated);
307 }
308 }
309
310 (working, events)
311}
312
313#[must_use]
327pub fn reconcile_order_report(
328 order: &OrderAny,
329 report: &OrderStatusReport,
330 instrument: Option<&InstrumentAny>,
331 ts_now: UnixNanos,
332) -> Option<OrderEventAny> {
333 reconcile_order_report_with_commission(order, report, instrument, ts_now, None)
334}
335
336#[must_use]
338pub fn reconcile_order_report_with_commission(
339 order: &OrderAny,
340 report: &OrderStatusReport,
341 instrument: Option<&InstrumentAny>,
342 ts_now: UnixNanos,
343 commission: Option<Money>,
344) -> Option<OrderEventAny> {
345 if matches!(
346 report.order_status,
347 OrderStatus::PendingUpdate | OrderStatus::PendingCancel
348 ) {
349 log::debug!(
350 "Order {} venue report in pending state: {:?}",
351 order.client_order_id(),
352 report.order_status,
353 );
354 return None;
355 }
356
357 if is_unchanged_accepted_report_during_pending_command(order, report) {
358 log::debug!(
359 "Order {} remains inflight while venue reports an unchanged accepted snapshot",
360 order.client_order_id(),
361 );
362 return None;
363 }
364
365 if is_superseded_cancel_report(order, report) {
366 let cached_venue_order_id = order.venue_order_id().unwrap_or(report.venue_order_id);
367 log::info!(
368 "Suppressing Canceled for {} on previously-promoted venue_order_id {}: \
369 current venue_order_id is {}",
370 order.client_order_id(),
371 report.venue_order_id,
372 cached_venue_order_id,
373 );
374 return None;
375 }
376
377 if order.status() == report.order_status && order.filled_qty() == report.filled_qty {
378 if should_reconciliation_update(order, report) {
379 log::info!(
380 "Order {} has been updated at venue: qty={}->{}, price={:?}->{:?}",
381 order.client_order_id(),
382 order.quantity(),
383 report.quantity,
384 order.price(),
385 report.price
386 );
387 return Some(create_reconciliation_updated(order, report, ts_now));
388 }
389 return None; }
391
392 match report.order_status {
393 OrderStatus::Accepted => {
394 if order.status() == OrderStatus::Accepted
395 && should_reconciliation_update(order, report)
396 {
397 return Some(create_reconciliation_updated(order, report, ts_now));
398 }
399 create_reconciliation_accepted(order, report, ts_now)
400 }
401 OrderStatus::Rejected => {
402 create_reconciliation_rejected(order, report.cancel_reason.as_deref(), ts_now)
403 }
404 OrderStatus::Triggered => {
405 if TRIGGERABLE_ORDER_TYPES.contains(&order.order_type()) {
406 Some(create_reconciliation_triggered(order, report, ts_now))
407 } else {
408 log::debug!(
409 "Skipping OrderTriggered for {} order {}: market-style stops have no TRIGGERED state",
410 order.order_type(),
411 order.client_order_id(),
412 );
413 None
414 }
415 }
416 OrderStatus::Canceled => Some(create_reconciliation_canceled(order, report, ts_now)),
417 OrderStatus::Expired => Some(create_reconciliation_expired(order, report, ts_now)),
418
419 OrderStatus::PartiallyFilled | OrderStatus::Filled => {
420 reconcile_fill_quantity_mismatch(order, report, instrument, ts_now, commission)
421 }
422
423 OrderStatus::Voided => {
424 create_reconciliation_terminal_fill_void(order, report, instrument, ts_now)
425 }
426
427 OrderStatus::PendingUpdate | OrderStatus::PendingCancel => None,
428
429 OrderStatus::Initialized
431 | OrderStatus::Submitted
432 | OrderStatus::Denied
433 | OrderStatus::Emulated
434 | OrderStatus::Released => {
435 log::warn!(
436 "Unexpected order status in venue report for {}: {:?}",
437 order.client_order_id(),
438 report.order_status
439 );
440 None
441 }
442 }
443}
444
445fn is_unchanged_accepted_report_during_pending_command(
446 order: &OrderAny,
447 report: &OrderStatusReport,
448) -> bool {
449 matches!(
450 order.status(),
451 OrderStatus::PendingUpdate | OrderStatus::PendingCancel
452 ) && report.order_status == OrderStatus::Accepted
453 && order.venue_order_id() == Some(report.venue_order_id)
454 && order.filled_qty() == report.filled_qty
455 && !should_reconciliation_update(order, report)
456}
457
458pub fn generate_external_order_status_events(
465 order: &OrderAny,
466 report: &OrderStatusReport,
467 account_id: &AccountId,
468 instrument: &InstrumentAny,
469 ts_now: UnixNanos,
470) -> Vec<OrderEventAny> {
471 generate_external_order_status_events_with_commission(
472 order, report, account_id, instrument, ts_now, None,
473 )
474}
475
476#[must_use]
478pub fn generate_external_order_status_events_with_commission(
479 order: &OrderAny,
480 report: &OrderStatusReport,
481 account_id: &AccountId,
482 instrument: &InstrumentAny,
483 ts_now: UnixNanos,
484 commission: Option<Money>,
485) -> Vec<OrderEventAny> {
486 let accepted = OrderEventAny::Accepted(OrderAccepted::new(
487 order.trader_id(),
488 order.strategy_id(),
489 order.instrument_id(),
490 order.client_order_id(),
491 report.venue_order_id,
492 *account_id,
493 UUID4::new(),
494 report.ts_accepted,
495 ts_now,
496 true, ));
498
499 match report.order_status {
500 OrderStatus::Accepted | OrderStatus::Triggered => vec![accepted],
501 OrderStatus::PartiallyFilled | OrderStatus::Filled => {
502 let mut events = vec![accepted];
503
504 if !report.filled_qty.is_zero()
505 && let Some(filled) =
506 create_inferred_fill(order, report, *account_id, instrument, ts_now, commission)
507 {
508 events.push(filled);
509 }
510
511 events
512 }
513 OrderStatus::Voided => {
514 let mut working = order.clone();
515 if let Err(e) = working.apply(accepted.clone()) {
516 log::warn!(
517 "Cannot project external order acceptance for {}: {e}",
518 order.client_order_id()
519 );
520 return vec![accepted];
521 }
522 let mut events = vec![accepted];
523
524 if !report.filled_qty.is_zero()
525 && let Some(filled) =
526 create_inferred_fill(order, report, *account_id, instrument, ts_now, commission)
527 {
528 if let Err(e) = working.apply(filled.clone()) {
529 log::warn!(
530 "Cannot project external order fill for {}: {e}",
531 order.client_order_id()
532 );
533 } else {
534 events.push(filled);
535 }
536 }
537
538 if let Some(voided) =
539 create_reconciliation_terminal_fill_void(&working, report, Some(instrument), ts_now)
540 {
541 events.push(voided);
542 }
543 events
544 }
545 OrderStatus::Canceled | OrderStatus::Expired => {
546 let terminal = create_external_terminal_event(order, report, *account_id, ts_now);
547 let mut events = vec![accepted];
548
549 let inferred_fill = if report.filled_qty.is_zero() {
550 None
551 } else {
552 create_inferred_fill(order, report, *account_id, instrument, ts_now, commission)
553 };
554 let filled_to_quantity =
555 inferred_fill.is_some() && report.filled_qty >= report.quantity;
556 if let Some(filled) = inferred_fill {
557 events.push(filled);
558 }
559
560 if !filled_to_quantity {
561 events.push(terminal);
562 }
563 events
564 }
565 OrderStatus::Rejected => {
566 let reason = report.cancel_reason.as_deref().unwrap_or("UNKNOWN");
568 vec![OrderEventAny::Rejected(OrderRejected::new(
569 order.trader_id(),
570 order.strategy_id(),
571 order.instrument_id(),
572 order.client_order_id(),
573 *account_id,
574 Ustr::from(reason),
575 UUID4::new(),
576 report.ts_last,
577 ts_now,
578 true, reason_indicates_post_only_rejection(reason),
580 ))]
581 }
582 _ => {
583 log::warn!(
584 "Unhandled order status {} for external order {}",
585 report.order_status,
586 order.client_order_id()
587 );
588 Vec::new()
589 }
590 }
591}
592
593fn create_reconciliation_fill_voids(
594 order: &OrderAny,
595 report: &OrderStatusReport,
596 ts_now: UnixNanos,
597) -> Vec<OrderEventAny> {
598 let mut remaining = order.filled_qty() - report.filled_qty;
599 let order_events = order.events();
600 let mut corrections = Vec::new();
601
602 for candidate in order_events.iter().rev() {
603 let OrderEventAny::Filled(fill) = candidate else {
604 continue;
605 };
606
607 if remaining.is_zero() {
608 break;
609 }
610 let previous = order_events.iter().rev().find_map(|event| match event {
611 OrderEventAny::FillVoided(voided) if voided.trade_id == fill.trade_id => Some(voided),
612 _ => None,
613 });
614 let prior_qty = previous.map_or_else(
615 || Quantity::zero(fill.last_qty.precision),
616 |voided| voided.voided_qty.min(fill.last_qty),
617 );
618 let effective = fill.last_qty - prior_qty;
619 if effective.is_zero() {
620 continue;
621 }
622 let removed = remaining.min(effective);
623 let voided_qty = prior_qty + removed;
624 let commission_voided = fill.commission.and_then(|commission| {
625 let fraction = voided_qty.as_decimal() / fill.last_qty.as_decimal();
626 Money::from_decimal(commission.as_decimal() * fraction, commission.currency).ok()
627 });
628 let mut event = OrderFillVoided::new(
629 fill.trader_id,
630 fill.strategy_id,
631 fill.instrument_id,
632 fill.client_order_id,
633 fill.venue_order_id,
634 fill.account_id,
635 Ustr::from(&format!(
636 "reconciliation-{}-{}",
637 report.report_id, fill.trade_id
638 )),
639 fill.trade_id,
640 voided_qty,
641 commission_voided,
642 fill.order_side,
643 fill.order_type,
644 fill.last_px,
645 fill.currency,
646 fill.liquidity_side,
647 fill.position_id,
648 report.cancel_reason.as_deref().map(Ustr::from),
649 None,
650 UUID4::new(),
651 report.ts_last,
652 ts_now,
653 true,
654 report_is_working(report),
655 );
656 event.causation_id = Some(report.report_id);
657 corrections.push(OrderEventAny::FillVoided(event));
658 remaining = remaining - removed;
659 }
660
661 if !remaining.is_zero() {
662 log::warn!(
663 "Cannot reconcile fill decrease for {}: {} is outside retained fill history",
664 order.client_order_id(),
665 remaining
666 );
667 return Vec::new();
668 }
669 corrections
670}
671
672fn create_reconciliation_terminal_fill_void(
678 order: &OrderAny,
679 report: &OrderStatusReport,
680 instrument: Option<&InstrumentAny>,
681 ts_now: UnixNanos,
682) -> Option<OrderEventAny> {
683 let voided_qty = order.leaves_qty();
684 if voided_qty.is_zero() {
685 return None;
686 }
687 let instrument = instrument?;
688 let order_side = order.order_side();
689
690 let last_px = resolve_fill_price(order, report, instrument)
691 .unwrap_or_else(|| Price::zero(instrument.price_precision()));
692
693 let mut event = OrderFillVoided::new(
694 order.trader_id(),
695 order.strategy_id(),
696 order.instrument_id(),
697 order.client_order_id(),
698 report.venue_order_id,
699 report.account_id,
700 Ustr::from(&format!("reconciliation-{}", report.report_id)),
701 TradeId::new(format!("VOID-{}", report.venue_order_id)),
702 voided_qty,
703 None,
704 order_side,
705 order.order_type(),
706 last_px,
707 instrument.quote_currency(),
708 LiquiditySide::NoLiquiditySide,
709 None,
710 report.cancel_reason.as_deref().map(Ustr::from),
711 None,
712 UUID4::new(),
713 report.ts_last,
714 ts_now,
715 true,
716 false,
717 );
718 event.causation_id = Some(report.report_id);
719 Some(OrderEventAny::FillVoided(event))
720}
721
722fn create_external_terminal_event(
723 order: &OrderAny,
724 report: &OrderStatusReport,
725 account_id: AccountId,
726 ts_now: UnixNanos,
727) -> OrderEventAny {
728 match report.order_status {
729 OrderStatus::Canceled => OrderEventAny::Canceled(OrderCanceled::new(
730 order.trader_id(),
731 order.strategy_id(),
732 order.instrument_id(),
733 order.client_order_id(),
734 UUID4::new(),
735 report.ts_last,
736 ts_now,
737 true, Some(report.venue_order_id),
739 Some(account_id),
740 report.cancel_reason.as_deref().map(Ustr::from),
741 )),
742 OrderStatus::Expired => OrderEventAny::Expired(OrderExpired::new(
743 order.trader_id(),
744 order.strategy_id(),
745 order.instrument_id(),
746 order.client_order_id(),
747 UUID4::new(),
748 report.ts_last,
749 ts_now,
750 true, Some(report.venue_order_id),
752 Some(account_id),
753 )),
754 status => unreachable!("cannot create external terminal event for {status}"),
755 }
756}
757
758pub fn reconcile_fill_report(
763 order: &OrderAny,
764 report: &FillReport,
765 instrument: &InstrumentAny,
766 ts_now: UnixNanos,
767 allow_overfills: bool,
768) -> Option<OrderEventAny> {
769 if report.last_qty.is_zero() {
770 log::warn!("Skipping zero-quantity fill report: {report}");
771 return None;
772 }
773
774 if order.trade_ids().iter().any(|id| **id == report.trade_id) {
775 log::debug!(
776 "Duplicate fill detected: trade_id {} already exists for order {}",
777 report.trade_id,
778 order.client_order_id()
779 );
780 return None;
781 }
782
783 let potential_filled_qty = order.filled_qty() + report.last_qty;
784 if potential_filled_qty > order.quantity() {
785 if !allow_overfills {
786 log::warn!(
787 "Rejecting fill that would cause overfill for {}: order.quantity={}, order.filled_qty={}, fill.last_qty={}, would result in filled_qty={}",
788 order.client_order_id(),
789 order.quantity(),
790 order.filled_qty(),
791 report.last_qty,
792 potential_filled_qty
793 );
794 return None;
795 }
796 log::warn!(
797 "Allowing overfill during reconciliation for {}: order.quantity={}, order.filled_qty={}, fill.last_qty={}, will result in filled_qty={}",
798 order.client_order_id(),
799 order.quantity(),
800 order.filled_qty(),
801 report.last_qty,
802 potential_filled_qty
803 );
804 }
805
806 let account_id = report.account_id;
807 let venue_order_id = order.venue_order_id().unwrap_or(report.venue_order_id);
808
809 log::info!(
810 color = LogColor::Blue as u8;
811 "Reconciling fill for {}: qty={}, px={}, trade_id={}",
812 order.client_order_id(),
813 report.last_qty,
814 report.last_px,
815 report.trade_id,
816 );
817
818 Some(OrderEventAny::Filled(OrderFilled::new(
819 order.trader_id(),
820 order.strategy_id(),
821 order.instrument_id(),
822 order.client_order_id(),
823 venue_order_id,
824 account_id,
825 report.trade_id,
826 order.order_side(),
827 order.order_type(),
828 report.last_qty,
829 report.last_px,
830 instrument.quote_currency(),
831 report.liquidity_side,
832 UUID4::new(),
833 report.ts_event,
834 ts_now,
835 true, report.venue_position_id,
837 Some(report.commission),
838 None,
839 )))
840}
841
842pub fn should_reconciliation_update(order: &OrderAny, report: &OrderStatusReport) -> bool {
850 if report.quantity != order.quantity() && report.quantity >= order.filled_qty() {
851 return true;
852 }
853
854 let price_drift = report.price.is_some() && report.price != order.price();
855 let trigger_drift =
856 report.trigger_price.is_some() && report.trigger_price != order.trigger_price();
857
858 match order.order_type() {
859 OrderType::Limit => price_drift,
860 OrderType::StopMarket | OrderType::TrailingStopMarket | OrderType::MarketIfTouched => {
861 trigger_drift
862 }
863 OrderType::StopLimit | OrderType::TrailingStopLimit | OrderType::LimitIfTouched => {
864 trigger_drift || price_drift
865 }
866 _ => false,
867 }
868}
869
870#[must_use]
872pub(super) fn create_reconciliation_accepted(
873 order: &OrderAny,
874 report: &OrderStatusReport,
875 ts_now: UnixNanos,
876) -> Option<OrderEventAny> {
877 let account_id = order.account_id()?;
878
879 Some(OrderEventAny::Accepted(OrderAccepted::new(
880 order.trader_id(),
881 order.strategy_id(),
882 order.instrument_id(),
883 order.client_order_id(),
884 order.venue_order_id().unwrap_or(report.venue_order_id),
885 account_id,
886 UUID4::new(),
887 report.ts_accepted,
888 ts_now,
889 true, )))
891}
892
893#[must_use]
895pub fn create_reconciliation_rejected(
896 order: &OrderAny,
897 reason: Option<&str>,
898 ts_now: UnixNanos,
899) -> Option<OrderEventAny> {
900 let account_id = order.account_id()?;
901 let reason = reason.unwrap_or("UNKNOWN");
902
903 Some(OrderEventAny::Rejected(OrderRejected::new(
904 order.trader_id(),
905 order.strategy_id(),
906 order.instrument_id(),
907 order.client_order_id(),
908 account_id,
909 Ustr::from(reason),
910 UUID4::new(),
911 ts_now,
912 ts_now,
913 true, reason_indicates_post_only_rejection(reason),
915 )))
916}
917
918fn reason_indicates_post_only_rejection(reason: &str) -> bool {
919 let normalized: String = reason
920 .chars()
921 .filter_map(|ch| {
922 if ch == '-' || ch == '_' || ch.is_whitespace() {
923 None
924 } else {
925 Some(ch.to_ascii_lowercase())
926 }
927 })
928 .collect();
929
930 normalized.contains("postonly") || normalized.contains("postwouldexecute")
931}
932
933#[must_use]
935pub fn create_reconciliation_triggered(
936 order: &OrderAny,
937 report: &OrderStatusReport,
938 ts_now: UnixNanos,
939) -> OrderEventAny {
940 OrderEventAny::Triggered(OrderTriggered::new(
941 order.trader_id(),
942 order.strategy_id(),
943 order.instrument_id(),
944 order.client_order_id(),
945 UUID4::new(),
946 report.ts_triggered.unwrap_or(ts_now),
947 ts_now,
948 true, order.venue_order_id(),
950 order.account_id(),
951 ))
952}
953
954#[must_use]
956pub(super) fn create_reconciliation_canceled(
957 order: &OrderAny,
958 report: &OrderStatusReport,
959 ts_now: UnixNanos,
960) -> OrderEventAny {
961 OrderEventAny::Canceled(OrderCanceled::new(
962 order.trader_id(),
963 order.strategy_id(),
964 order.instrument_id(),
965 order.client_order_id(),
966 UUID4::new(),
967 report.ts_last,
968 ts_now,
969 true, order.venue_order_id(),
971 order.account_id(),
972 report.cancel_reason.as_deref().map(Ustr::from),
973 ))
974}
975
976#[must_use]
978pub(super) fn create_reconciliation_expired(
979 order: &OrderAny,
980 report: &OrderStatusReport,
981 ts_now: UnixNanos,
982) -> OrderEventAny {
983 OrderEventAny::Expired(OrderExpired::new(
984 order.trader_id(),
985 order.strategy_id(),
986 order.instrument_id(),
987 order.client_order_id(),
988 UUID4::new(),
989 report.ts_last,
990 ts_now,
991 true, order.venue_order_id(),
993 order.account_id(),
994 ))
995}
996
997#[must_use]
999pub(super) fn create_reconciliation_updated(
1000 order: &OrderAny,
1001 report: &OrderStatusReport,
1002 ts_now: UnixNanos,
1003) -> OrderEventAny {
1004 let trigger_price = match order.order_type() {
1011 OrderType::StopMarket
1012 | OrderType::StopLimit
1013 | OrderType::MarketIfTouched
1014 | OrderType::LimitIfTouched
1015 | OrderType::TrailingStopMarket
1016 | OrderType::TrailingStopLimit => report.trigger_price,
1017 _ => None,
1018 };
1019
1020 OrderEventAny::Updated(OrderUpdated::new(
1021 order.trader_id(),
1022 order.strategy_id(),
1023 order.instrument_id(),
1024 order.client_order_id(),
1025 report.quantity,
1026 UUID4::new(),
1027 report.ts_last,
1028 ts_now,
1029 true, order.venue_order_id(),
1031 order.account_id(),
1032 report.price,
1033 trigger_price,
1034 None, order.is_quote_quantity(),
1036 ))
1037}
1038
1039pub(super) fn create_inferred_fill(
1041 order: &OrderAny,
1042 report: &OrderStatusReport,
1043 account_id: AccountId,
1044 instrument: &InstrumentAny,
1045 ts_now: UnixNanos,
1046 commission: Option<Money>,
1047) -> Option<OrderEventAny> {
1048 let cached_order_side = order.order_side();
1049 let order_side = report.order_side.unwrap_or(cached_order_side);
1050 if order_side != cached_order_side {
1051 log::warn!(
1052 "Order side mismatch for {}: cached={:?}, venue={order_side:?}",
1053 order.client_order_id(),
1054 cached_order_side,
1055 );
1056 }
1057
1058 let liquidity_side = match order.order_type() {
1059 OrderType::Market
1060 | OrderType::StopMarket
1061 | OrderType::MarketToLimit
1062 | OrderType::TrailingStopMarket => LiquiditySide::Taker,
1063 _ if order.is_post_only() => LiquiditySide::Maker,
1064 _ => LiquiditySide::NoLiquiditySide,
1065 };
1066
1067 let Some(last_px) = resolve_fill_price(order, report, instrument) else {
1068 log::warn!(
1069 "Cannot create inferred fill for {}: no avg_px, report price, or order price",
1070 order.client_order_id()
1071 );
1072
1073 return None;
1074 };
1075 let last_px = clamp_inferred_fill_price(last_px, instrument);
1076
1077 let position_id = reconciliation_position_id(report, instrument);
1078 let trade_id = create_inferred_reconciliation_trade_id(
1079 account_id,
1080 order.instrument_id(),
1081 order.client_order_id(),
1082 Some(report.venue_order_id),
1083 order_side,
1084 order.order_type(),
1085 report.filled_qty,
1086 report.filled_qty,
1087 last_px,
1088 position_id,
1089 report.ts_last,
1090 );
1091
1092 log::info!(
1093 "Generated inferred fill for {} ({}) qty={} px={}",
1094 order.client_order_id(),
1095 report.venue_order_id,
1096 report.filled_qty,
1097 last_px,
1098 );
1099
1100 Some(OrderEventAny::Filled(OrderFilled::new(
1101 order.trader_id(),
1102 order.strategy_id(),
1103 order.instrument_id(),
1104 order.client_order_id(),
1105 report.venue_order_id,
1106 account_id,
1107 trade_id,
1108 order_side,
1109 order.order_type(),
1110 report.filled_qty,
1111 last_px,
1112 instrument.quote_currency(),
1113 liquidity_side,
1114 UUID4::new(),
1115 report.ts_last,
1116 ts_now,
1117 true, report.venue_position_id,
1119 commission,
1120 None,
1121 )))
1122}
1123
1124pub fn create_incremental_inferred_fill(
1126 order: &OrderAny,
1127 report: &OrderStatusReport,
1128 account_id: &AccountId,
1129 instrument: &InstrumentAny,
1130 ts_now: UnixNanos,
1131 commission: Option<Money>,
1132) -> Option<OrderEventAny> {
1133 let order_side = order.order_side();
1134
1135 let order_filled_qty = order.filled_qty();
1136 debug_assert!(
1137 report.filled_qty >= order_filled_qty,
1138 "incremental inferred fill requires report.filled_qty ({}) >= order.filled_qty ({}) for {}",
1139 report.filled_qty,
1140 order_filled_qty,
1141 order.client_order_id(),
1142 );
1143 let last_qty = report.filled_qty - order_filled_qty;
1144
1145 if last_qty <= Quantity::zero(instrument.size_precision()) {
1146 return None;
1147 }
1148
1149 let (last_px, liquidity_side) =
1150 incremental_inferred_fill_price_and_liquidity(order, report, instrument)?;
1151
1152 let venue_order_id = order.venue_order_id().unwrap_or(report.venue_order_id);
1153 let position_id = reconciliation_position_id(report, instrument);
1154 let trade_id = create_inferred_reconciliation_trade_id(
1155 *account_id,
1156 order.instrument_id(),
1157 order.client_order_id(),
1158 Some(venue_order_id),
1159 order_side,
1160 order.order_type(),
1161 report.filled_qty,
1162 last_qty,
1163 last_px,
1164 position_id,
1165 report.ts_last,
1166 );
1167
1168 log::info!(
1169 color = LogColor::Blue as u8;
1170 "Generated inferred fill for {}: qty={}, px={}",
1171 order.client_order_id(),
1172 last_qty,
1173 last_px,
1174 );
1175
1176 Some(OrderEventAny::Filled(OrderFilled::new(
1177 order.trader_id(),
1178 order.strategy_id(),
1179 order.instrument_id(),
1180 order.client_order_id(),
1181 venue_order_id,
1182 *account_id,
1183 trade_id,
1184 order_side,
1185 order.order_type(),
1186 last_qty,
1187 last_px,
1188 instrument.quote_currency(),
1189 liquidity_side,
1190 UUID4::new(),
1191 report.ts_last,
1192 ts_now,
1193 true, report.venue_position_id,
1195 commission,
1196 None,
1197 )))
1198}
1199
1200#[must_use]
1208pub fn incremental_inferred_fill_price_and_liquidity(
1209 order: &OrderAny,
1210 report: &OrderStatusReport,
1211 instrument: &InstrumentAny,
1212) -> Option<(Price, LiquiditySide)> {
1213 let last_px = calculate_incremental_fill_price(order, report, instrument)?;
1214
1215 Some((
1216 clamp_inferred_fill_price(last_px, instrument),
1217 inferred_fill_liquidity_side(order),
1218 ))
1219}
1220
1221#[must_use]
1228pub fn inferred_fill_price_and_liquidity(
1229 order: &OrderAny,
1230 report: &OrderStatusReport,
1231 instrument: &InstrumentAny,
1232) -> Option<(Price, LiquiditySide)> {
1233 let last_px = resolve_fill_price(order, report, instrument)?;
1234
1235 Some((
1236 clamp_inferred_fill_price(last_px, instrument),
1237 inferred_fill_liquidity_side(order),
1238 ))
1239}
1240
1241fn inferred_fill_liquidity_side(order: &OrderAny) -> LiquiditySide {
1242 match order.order_type() {
1243 OrderType::Market
1244 | OrderType::StopMarket
1245 | OrderType::MarketToLimit
1246 | OrderType::TrailingStopMarket => LiquiditySide::Taker,
1247 _ if order.is_post_only() => LiquiditySide::Maker,
1248 _ => LiquiditySide::NoLiquiditySide,
1249 }
1250}
1251
1252pub fn create_inferred_fill_for_qty(
1258 order: &OrderAny,
1259 report: &OrderStatusReport,
1260 account_id: &AccountId,
1261 instrument: &InstrumentAny,
1262 fill_qty: Quantity,
1263 ts_now: UnixNanos,
1264 commission: Option<Money>,
1265) -> Option<OrderEventAny> {
1266 if fill_qty.is_zero() {
1267 return None;
1268 }
1269
1270 let order_side = order.order_side();
1271
1272 let Some((last_px, liquidity_side)) =
1273 inferred_fill_price_and_liquidity(order, report, instrument)
1274 else {
1275 log::warn!(
1276 "Cannot determine fill price for {}: no avg_px, report price, or order price",
1277 order.client_order_id()
1278 );
1279
1280 return None;
1281 };
1282
1283 let venue_order_id = order.venue_order_id().unwrap_or(report.venue_order_id);
1284 let position_id = reconciliation_position_id(report, instrument);
1285 let trade_id = create_inferred_reconciliation_trade_id(
1286 *account_id,
1287 order.instrument_id(),
1288 order.client_order_id(),
1289 Some(venue_order_id),
1290 order_side,
1291 order.order_type(),
1292 report.filled_qty,
1293 fill_qty,
1294 last_px,
1295 position_id,
1296 report.ts_last,
1297 );
1298
1299 log::info!(
1300 color = LogColor::Blue as u8;
1301 "Generated inferred fill for {}: qty={}, px={}",
1302 order.client_order_id(),
1303 fill_qty,
1304 last_px,
1305 );
1306
1307 Some(OrderEventAny::Filled(OrderFilled::new(
1308 order.trader_id(),
1309 order.strategy_id(),
1310 order.instrument_id(),
1311 order.client_order_id(),
1312 venue_order_id,
1313 *account_id,
1314 trade_id,
1315 order_side,
1316 order.order_type(),
1317 fill_qty,
1318 last_px,
1319 instrument.quote_currency(),
1320 liquidity_side,
1321 UUID4::new(),
1322 report.ts_last,
1323 ts_now,
1324 true, report.venue_position_id,
1326 commission,
1327 None,
1328 )))
1329}
1330
1331fn report_is_confirmed_state(report: &OrderStatusReport) -> bool {
1332 matches!(
1333 report.order_status,
1334 OrderStatus::Accepted
1335 | OrderStatus::Triggered
1336 | OrderStatus::PartiallyFilled
1337 | OrderStatus::Filled
1338 | OrderStatus::Canceled
1339 | OrderStatus::Expired
1340 )
1341}
1342
1343fn report_is_working(report: &OrderStatusReport) -> bool {
1344 matches!(
1345 report.order_status,
1346 OrderStatus::Accepted | OrderStatus::Triggered | OrderStatus::PartiallyFilled
1347 )
1348}
1349
1350pub(crate) fn is_superseded_cancel_report(order: &OrderAny, report: &OrderStatusReport) -> bool {
1351 if report.order_status != OrderStatus::Canceled {
1352 return false;
1353 }
1354
1355 let Some(cached_venue_order_id) = order.venue_order_id() else {
1356 return false;
1357 };
1358
1359 cached_venue_order_id != report.venue_order_id
1360 && order
1361 .venue_order_ids()
1362 .iter()
1363 .any(|venue_order_id| **venue_order_id == report.venue_order_id)
1364}
1365
1366fn local_accepts_amendment(order: &OrderAny) -> bool {
1367 matches!(
1368 order.status(),
1369 OrderStatus::Accepted | OrderStatus::Triggered | OrderStatus::PartiallyFilled
1370 )
1371}
1372
1373fn should_accept_before_reconciliation(order: &OrderAny, report: &OrderStatusReport) -> bool {
1374 order.status() == OrderStatus::Submitted && report.order_status != OrderStatus::Rejected
1375}
1376
1377fn reconcile_fill_quantity_mismatch(
1381 order: &OrderAny,
1382 report: &OrderStatusReport,
1383 instrument: Option<&InstrumentAny>,
1384 ts_now: UnixNanos,
1385 commission: Option<Money>,
1386) -> Option<OrderEventAny> {
1387 let order_filled_qty = order.filled_qty();
1388 let report_filled_qty = report.filled_qty;
1389
1390 if report_filled_qty < order_filled_qty {
1391 let precision = order_filled_qty.precision.max(report_filled_qty.precision);
1394 if is_within_single_unit_tolerance(
1395 report_filled_qty.as_decimal(),
1396 order_filled_qty.as_decimal(),
1397 precision,
1398 ) {
1399 return None;
1400 }
1401
1402 log::warn!(
1403 "Fill qty mismatch for {} ({}): cached={}, venue={}, order_qty={} (venue < cached)",
1404 order.client_order_id(),
1405 report.venue_order_id,
1406 order_filled_qty,
1407 report_filled_qty,
1408 order.quantity(),
1409 );
1410 return None;
1411 }
1412
1413 if report_filled_qty > order_filled_qty {
1414 if order.is_closed() {
1417 let precision = order_filled_qty.precision.max(report_filled_qty.precision);
1418
1419 if is_within_single_unit_tolerance(
1420 report_filled_qty.as_decimal(),
1421 order_filled_qty.as_decimal(),
1422 precision,
1423 ) {
1424 return None;
1425 }
1426
1427 log::debug!(
1428 "{} {} already closed but reported difference in filled_qty: \
1429 report={}, cached={}, skipping inferred fill generation for closed order",
1430 order.instrument_id(),
1431 order.client_order_id(),
1432 report_filled_qty,
1433 order_filled_qty,
1434 );
1435 return None;
1436 }
1437
1438 let Some(instrument) = instrument else {
1440 log::warn!(
1441 "Cannot generate inferred fill for {}: instrument not available",
1442 order.client_order_id()
1443 );
1444 return None;
1445 };
1446
1447 let account_id = order.account_id()?;
1448 return create_incremental_inferred_fill(
1449 order,
1450 report,
1451 &account_id,
1452 instrument,
1453 ts_now,
1454 commission,
1455 );
1456 }
1457
1458 if order.status() != report.order_status {
1463 if should_reconciliation_update(order, report) {
1464 log::info!(
1465 "Status mismatch with matching fill qty for {}: local={:?}, venue={:?}, \
1466 filled_qty={}, updating quantity {}->{}",
1467 order.client_order_id(),
1468 order.status(),
1469 report.order_status,
1470 report.filled_qty,
1471 order.quantity(),
1472 report.quantity,
1473 );
1474 return Some(create_reconciliation_updated(order, report, ts_now));
1475 }
1476
1477 log::warn!(
1478 "Status mismatch with matching fill qty for {}: local={:?}, venue={:?}, filled_qty={}",
1479 order.client_order_id(),
1480 order.status(),
1481 report.order_status,
1482 report.filled_qty
1483 );
1484 }
1485
1486 None
1487}
1488
1489fn calculate_incremental_fill_price(
1495 order: &OrderAny,
1496 report: &OrderStatusReport,
1497 instrument: &InstrumentAny,
1498) -> Option<Price> {
1499 let order_filled_qty = order.filled_qty();
1500 debug_assert!(
1501 report.filled_qty >= order_filled_qty,
1502 "incremental fill price requires report.filled_qty ({}) >= order.filled_qty ({}) for {}",
1503 report.filled_qty,
1504 order_filled_qty,
1505 order.client_order_id(),
1506 );
1507
1508 if order_filled_qty.is_zero() {
1510 let last_px = resolve_fill_price(order, report, instrument);
1511 if last_px.is_none() {
1512 log::warn!(
1513 "Cannot determine fill price for {}: no avg_px, report price, or order price",
1514 order.client_order_id()
1515 );
1516 }
1517
1518 return last_px;
1519 }
1520
1521 if let Some(report_avg_px) = report.avg_px {
1523 let last_px_decimal = match order.avg_px() {
1524 None => report_avg_px,
1526 Some(order_avg_px) => {
1527 let report_filled_qty = report.filled_qty;
1528 let last_qty = report_filled_qty - order_filled_qty;
1529
1530 let report_notional = report_avg_px * report_filled_qty.as_decimal();
1531 let order_notional = order_avg_px * order_filled_qty.as_decimal();
1532 let last_notional = report_notional - order_notional;
1533 let back_solved = last_notional / last_qty.as_decimal();
1534
1535 if back_solved < Decimal::ZERO && !instrument.allows_negative_price() {
1536 if report_avg_px < Decimal::ZERO {
1537 log::warn!(
1538 "Cannot price inferred fill for {}: back-solved {back_solved} and venue average {report_avg_px} are both negative on an instrument that disallows negative prices",
1539 order.client_order_id(),
1540 );
1541
1542 return None;
1543 }
1544
1545 log::warn!(
1546 "Negative back-solved fill price {back_solved} for {}, using venue average {report_avg_px}",
1547 order.client_order_id(),
1548 );
1549
1550 report_avg_px
1551 } else {
1552 back_solved
1553 }
1554 }
1555 };
1556
1557 return Price::from_decimal_dp(last_px_decimal, instrument.price_precision())
1558 .inspect_err(|e| {
1559 log::warn!(
1560 "Cannot price {} from incremental {last_px_decimal}, falling back: {e}",
1561 order.client_order_id(),
1562 );
1563 })
1564 .ok()
1565 .or_else(|| resolve_fill_price(order, report, instrument));
1566 }
1567
1568 resolve_fill_price(order, report, instrument)
1569}
1570
1571fn resolve_fill_price(
1577 order: &OrderAny,
1578 report: &OrderStatusReport,
1579 instrument: &InstrumentAny,
1580) -> Option<Price> {
1581 report
1582 .avg_px
1583 .and_then(|avg_px| {
1584 Price::from_decimal_dp(avg_px, instrument.price_precision())
1585 .inspect_err(|e| {
1586 log::warn!(
1587 "Cannot price {} from venue average {avg_px}, trying next source: {e}",
1588 order.client_order_id(),
1589 );
1590 })
1591 .ok()
1592 })
1593 .or(report.price)
1594 .or_else(|| order.price())
1595}
1596
1597fn clamp_inferred_fill_price(price: Price, instrument: &InstrumentAny) -> Price {
1602 let px = cap_price_at_instrument_max(price.as_decimal(), instrument);
1603 Price::from_decimal_dp(px, instrument.price_precision()).unwrap_or(price)
1604}
1605
1606fn reconciliation_position_id(
1607 report: &OrderStatusReport,
1608 instrument: &InstrumentAny,
1609) -> PositionId {
1610 report
1611 .venue_position_id
1612 .unwrap_or_else(|| PositionId::new(format!("{}-EXTERNAL", instrument.id())))
1613}