1use nautilus_common::enums::LogColor;
23use nautilus_core::{UUID4, UnixNanos};
24use nautilus_model::{
25 enums::{LiquiditySide, OrderStatus, OrderType},
26 events::{
27 OrderAccepted, OrderCanceled, OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled,
28 OrderRejected, OrderTriggered, OrderUpdated,
29 },
30 identifiers::{AccountId, PositionId, TradeId},
31 instruments::{Instrument, InstrumentAny},
32 orders::{Order, OrderAny, TRIGGERABLE_ORDER_TYPES},
33 reports::{FillReport, OrderStatusReport},
34 types::{Money, Price, Quantity},
35};
36use rust_decimal::Decimal;
37use ustr::Ustr;
38
39use super::{
40 ids::create_inferred_reconciliation_trade_id,
41 positions::{cap_price_at_instrument_max, is_within_single_unit_tolerance},
42};
43
44#[must_use]
55pub fn generate_reconciliation_order_events(
56 order: &OrderAny,
57 report: &OrderStatusReport,
58 instrument: Option<&InstrumentAny>,
59 ts_now: UnixNanos,
60) -> Vec<OrderEventAny> {
61 generate_reconciliation_order_events_inner(
62 order,
63 report,
64 instrument,
65 ts_now,
66 report.order_status == OrderStatus::Voided,
67 None,
68 )
69}
70
71#[must_use]
78pub fn generate_reconciliation_order_snapshot_events(
79 order: &OrderAny,
80 report: &OrderStatusReport,
81 instrument: Option<&InstrumentAny>,
82 ts_now: UnixNanos,
83) -> Vec<OrderEventAny> {
84 generate_reconciliation_order_events_inner(order, report, instrument, ts_now, true, None)
85}
86
87#[must_use]
90pub fn generate_reconciliation_order_snapshot_events_with_commission(
91 order: &OrderAny,
92 report: &OrderStatusReport,
93 instrument: Option<&InstrumentAny>,
94 ts_now: UnixNanos,
95 commission: Option<Money>,
96) -> Vec<OrderEventAny> {
97 generate_reconciliation_order_events_inner(order, report, instrument, ts_now, true, commission)
98}
99
100fn generate_reconciliation_order_events_inner(
101 order: &OrderAny,
102 report: &OrderStatusReport,
103 instrument: Option<&InstrumentAny>,
104 ts_now: UnixNanos,
105 allow_fill_decrease: bool,
106 commission: Option<Money>,
107) -> Vec<OrderEventAny> {
108 if is_superseded_cancel_report(order, report) {
109 let _ = reconcile_order_report(order, report, instrument, ts_now);
110 return Vec::new();
111 }
112
113 if has_material_fill_decrease(order, report) {
114 if !allow_fill_decrease {
115 log::warn!(
116 "Ignoring fill decrease without explicit void evidence for {}: cached={}, venue={}",
117 order.client_order_id(),
118 order.filled_qty(),
119 report.filled_qty,
120 );
121 return reconcile_fill_decrease_terminal(order, report, instrument, ts_now);
122 }
123
124 let mut working = order.clone();
125 let mut events = create_reconciliation_fill_voids(&working, report, ts_now);
126 if events.is_empty() {
127 return reconcile_fill_decrease_terminal(order, report, instrument, ts_now);
128 }
129
130 for event in &events {
131 if let Err(e) = working.apply(event.clone()) {
132 log::warn!(
133 "Cannot project reconciliation fill void for {}: {e}",
134 order.client_order_id()
135 );
136 return reconcile_fill_decrease_terminal(order, report, instrument, ts_now);
137 }
138 }
139
140 if report_is_working(report) && should_reconciliation_update(&working, report) {
141 let updated = create_reconciliation_updated(&working, report, ts_now);
142 if let Err(e) = working.apply(updated.clone()) {
143 log::warn!(
144 "Cannot project reopening reconciliation update for {}: {e}",
145 order.client_order_id()
146 );
147 } else {
148 events.push(updated);
149 }
150 }
151
152 if let Some(terminal) = reconcile_order_report(&working, report, instrument, ts_now) {
153 events.push(terminal);
154 }
155 return events;
156 }
157
158 let (mut working, mut events) = prepare_reconciliation_order(order, report, ts_now);
159
160 if matches!(
161 report.order_status,
162 OrderStatus::Canceled | OrderStatus::Expired,
163 ) && report.filled_qty > working.filled_qty()
164 && let Some(instrument) = instrument
165 && let Some(filled) = create_incremental_inferred_fill(
166 &working,
167 report,
168 &report.account_id,
169 instrument,
170 ts_now,
171 commission,
172 )
173 {
174 if let Err(e) = working.apply(filled.clone()) {
175 log::warn!(
176 "Failed to pre-apply reconciliation fill for {}: {e}",
177 order.client_order_id(),
178 );
179 } else {
180 events.push(filled);
181 }
182 }
183
184 if working.status() == OrderStatus::Filled
185 && matches!(
186 report.order_status,
187 OrderStatus::Canceled | OrderStatus::Expired,
188 )
189 {
190 return events;
191 }
192
193 if let Some(event) =
194 reconcile_order_report_with_commission(&working, report, instrument, ts_now, commission)
195 {
196 events.push(event);
197 }
198
199 events
200}
201
202fn has_material_fill_decrease(order: &OrderAny, report: &OrderStatusReport) -> bool {
203 if report.filled_qty >= order.filled_qty() {
204 return false;
205 }
206
207 let precision = order
208 .filled_qty()
209 .precision
210 .max(report.filled_qty.precision);
211 !is_within_single_unit_tolerance(
212 report.filled_qty.as_decimal(),
213 order.filled_qty().as_decimal(),
214 precision,
215 )
216}
217
218fn reconcile_fill_decrease_terminal(
219 order: &OrderAny,
220 report: &OrderStatusReport,
221 instrument: Option<&InstrumentAny>,
222 ts_now: UnixNanos,
223) -> Vec<OrderEventAny> {
224 let should_reconcile = report.order_status == OrderStatus::Voided
225 || (!order.is_closed()
226 && matches!(
227 report.order_status,
228 OrderStatus::Canceled | OrderStatus::Expired
229 ));
230
231 if !should_reconcile {
232 return Vec::new();
233 }
234
235 reconcile_order_report(order, report, instrument, ts_now)
236 .into_iter()
237 .collect()
238}
239
240#[must_use]
242pub fn generate_reconciliation_order_pre_fill_events(
243 order: &OrderAny,
244 report: &OrderStatusReport,
245 ts_now: UnixNanos,
246) -> Vec<OrderEventAny> {
247 if is_superseded_cancel_report(order, report) {
248 return Vec::new();
249 }
250
251 let (working, mut events) = prepare_reconciliation_order(order, report, ts_now);
252
253 if report.order_status == OrderStatus::Triggered
254 && let Some(triggered) = reconcile_order_report(&working, report, None, ts_now)
255 {
256 events.push(triggered);
257 }
258
259 events
260}
261
262fn prepare_reconciliation_order(
263 order: &OrderAny,
264 report: &OrderStatusReport,
265 ts_now: UnixNanos,
266) -> (OrderAny, Vec<OrderEventAny>) {
267 let mut working = order.clone();
268 let mut events: Vec<OrderEventAny> = Vec::new();
269
270 if should_accept_before_reconciliation(&working, report) {
271 let Some(accepted) = create_reconciliation_accepted(&working, report, ts_now) else {
272 log::warn!(
273 "Cannot create reconciliation acceptance for {}: missing account_id",
274 order.client_order_id(),
275 );
276 return (working, events);
277 };
278
279 if let Err(e) = working.apply(accepted.clone()) {
280 log::warn!(
281 "Failed to pre-apply reconciliation acceptance for {}: {e}",
282 order.client_order_id(),
283 );
284 return (working, events);
285 }
286 events.push(accepted);
287 }
288
289 if report_is_confirmed_state(report)
290 && (local_accepts_amendment(&working)
291 || (matches!(
292 working.status(),
293 OrderStatus::PendingUpdate | OrderStatus::PendingCancel,
294 ) && matches!(
295 report.order_status,
296 OrderStatus::Canceled | OrderStatus::Expired,
297 )))
298 && should_reconciliation_update(&working, report)
299 {
300 let updated = create_reconciliation_updated(&working, report, ts_now);
301 if let Err(e) = working.apply(updated.clone()) {
302 log::warn!(
303 "Failed to pre-apply reconciliation update for {}: {e}",
304 order.client_order_id(),
305 );
306 } else {
307 events.push(updated);
308 }
309 }
310
311 (working, events)
312}
313
314#[must_use]
328pub fn reconcile_order_report(
329 order: &OrderAny,
330 report: &OrderStatusReport,
331 instrument: Option<&InstrumentAny>,
332 ts_now: UnixNanos,
333) -> Option<OrderEventAny> {
334 reconcile_order_report_with_commission(order, report, instrument, ts_now, None)
335}
336
337#[must_use]
339pub fn reconcile_order_report_with_commission(
340 order: &OrderAny,
341 report: &OrderStatusReport,
342 instrument: Option<&InstrumentAny>,
343 ts_now: UnixNanos,
344 commission: Option<Money>,
345) -> Option<OrderEventAny> {
346 if matches!(
347 report.order_status,
348 OrderStatus::PendingUpdate | OrderStatus::PendingCancel
349 ) {
350 log::debug!(
351 "Order {} venue report in pending state: {:?}",
352 order.client_order_id(),
353 report.order_status,
354 );
355 return None;
356 }
357
358 if is_unchanged_accepted_report_during_pending_command(order, report) {
359 log::debug!(
360 "Order {} remains inflight while venue reports an unchanged accepted snapshot",
361 order.client_order_id(),
362 );
363 return None;
364 }
365
366 if is_superseded_cancel_report(order, report) {
367 let cached_venue_order_id = order.venue_order_id().unwrap_or(report.venue_order_id);
368 log::info!(
369 "Suppressing Canceled for {} on previously-promoted venue_order_id {}: \
370 current venue_order_id is {}",
371 order.client_order_id(),
372 report.venue_order_id,
373 cached_venue_order_id,
374 );
375 return None;
376 }
377
378 if order.status() == report.order_status && order.filled_qty() == report.filled_qty {
379 if should_reconciliation_update(order, report) {
380 log::info!(
381 "Order {} has been updated at venue: qty={}->{}, price={:?}->{:?}",
382 order.client_order_id(),
383 order.quantity(),
384 report.quantity,
385 order.price(),
386 report.price
387 );
388 return Some(create_reconciliation_updated(order, report, ts_now));
389 }
390 return None; }
392
393 match report.order_status {
394 OrderStatus::Accepted => {
395 if order.status() == OrderStatus::Accepted
396 && should_reconciliation_update(order, report)
397 {
398 return Some(create_reconciliation_updated(order, report, ts_now));
399 }
400 create_reconciliation_accepted(order, report, ts_now)
401 }
402 OrderStatus::Rejected => {
403 create_reconciliation_rejected(order, report.cancel_reason.as_deref(), ts_now)
404 }
405 OrderStatus::Triggered => {
406 if TRIGGERABLE_ORDER_TYPES.contains(&order.order_type()) {
407 Some(create_reconciliation_triggered(order, report, ts_now))
408 } else {
409 log::debug!(
410 "Skipping OrderTriggered for {} order {}: market-style stops have no TRIGGERED state",
411 order.order_type(),
412 order.client_order_id(),
413 );
414 None
415 }
416 }
417 OrderStatus::Canceled => Some(create_reconciliation_canceled(order, report, ts_now)),
418 OrderStatus::Expired => Some(create_reconciliation_expired(order, report, ts_now)),
419
420 OrderStatus::PartiallyFilled | OrderStatus::Filled => {
421 reconcile_fill_quantity_mismatch(order, report, instrument, ts_now, commission)
422 }
423
424 OrderStatus::Voided => {
425 create_reconciliation_terminal_fill_void(order, report, instrument, ts_now)
426 }
427
428 OrderStatus::PendingUpdate | OrderStatus::PendingCancel => None,
429
430 OrderStatus::Initialized
432 | OrderStatus::Submitted
433 | OrderStatus::Denied
434 | OrderStatus::Emulated
435 | OrderStatus::Released => {
436 log::warn!(
437 "Unexpected order status in venue report for {}: {:?}",
438 order.client_order_id(),
439 report.order_status
440 );
441 None
442 }
443 }
444}
445
446fn is_unchanged_accepted_report_during_pending_command(
447 order: &OrderAny,
448 report: &OrderStatusReport,
449) -> bool {
450 matches!(
451 order.status(),
452 OrderStatus::PendingUpdate | OrderStatus::PendingCancel
453 ) && report.order_status == OrderStatus::Accepted
454 && order.venue_order_id() == Some(report.venue_order_id)
455 && order.filled_qty() == report.filled_qty
456 && !should_reconciliation_update(order, report)
457}
458
459pub fn generate_external_order_status_events(
466 order: &OrderAny,
467 report: &OrderStatusReport,
468 account_id: &AccountId,
469 instrument: &InstrumentAny,
470 ts_now: UnixNanos,
471) -> Vec<OrderEventAny> {
472 generate_external_order_status_events_with_commission(
473 order, report, account_id, instrument, ts_now, None,
474 )
475}
476
477#[must_use]
479pub fn generate_external_order_status_events_with_commission(
480 order: &OrderAny,
481 report: &OrderStatusReport,
482 account_id: &AccountId,
483 instrument: &InstrumentAny,
484 ts_now: UnixNanos,
485 commission: Option<Money>,
486) -> Vec<OrderEventAny> {
487 let accepted = OrderEventAny::Accepted(OrderAccepted::new(
488 order.trader_id(),
489 order.strategy_id(),
490 order.instrument_id(),
491 order.client_order_id(),
492 report.venue_order_id,
493 *account_id,
494 UUID4::new(),
495 report.ts_accepted,
496 ts_now,
497 true, ));
499
500 match report.order_status {
501 OrderStatus::Accepted | OrderStatus::Triggered => vec![accepted],
502 OrderStatus::PartiallyFilled | OrderStatus::Filled => {
503 let mut events = vec![accepted];
504
505 if !report.filled_qty.is_zero()
506 && let Some(filled) =
507 create_inferred_fill(order, report, *account_id, instrument, ts_now, commission)
508 {
509 events.push(filled);
510 }
511
512 events
513 }
514 OrderStatus::Voided => {
515 let mut working = order.clone();
516 if let Err(e) = working.apply(accepted.clone()) {
517 log::warn!(
518 "Cannot project external order acceptance for {}: {e}",
519 order.client_order_id()
520 );
521 return vec![accepted];
522 }
523 let mut events = vec![accepted];
524
525 if !report.filled_qty.is_zero()
526 && let Some(filled) =
527 create_inferred_fill(order, report, *account_id, instrument, ts_now, commission)
528 {
529 if let Err(e) = working.apply(filled.clone()) {
530 log::warn!(
531 "Cannot project external order fill for {}: {e}",
532 order.client_order_id()
533 );
534 } else {
535 events.push(filled);
536 }
537 }
538
539 if let Some(voided) =
540 create_reconciliation_terminal_fill_void(&working, report, Some(instrument), ts_now)
541 {
542 events.push(voided);
543 }
544 events
545 }
546 OrderStatus::Canceled | OrderStatus::Expired => {
547 let terminal = create_external_terminal_event(order, report, *account_id, ts_now);
548 let mut events = vec![accepted];
549
550 let inferred_fill = if report.filled_qty.is_zero() {
551 None
552 } else {
553 create_inferred_fill(order, report, *account_id, instrument, ts_now, commission)
554 };
555 let filled_to_quantity =
556 inferred_fill.is_some() && report.filled_qty >= report.quantity;
557 if let Some(filled) = inferred_fill {
558 events.push(filled);
559 }
560
561 if !filled_to_quantity {
562 events.push(terminal);
563 }
564 events
565 }
566 OrderStatus::Rejected => {
567 let reason = report.cancel_reason.as_deref().unwrap_or("UNKNOWN");
569 vec![OrderEventAny::Rejected(OrderRejected::new(
570 order.trader_id(),
571 order.strategy_id(),
572 order.instrument_id(),
573 order.client_order_id(),
574 *account_id,
575 Ustr::from(reason),
576 UUID4::new(),
577 report.ts_last,
578 ts_now,
579 true, reason_indicates_post_only_rejection(reason),
581 ))]
582 }
583 _ => {
584 log::warn!(
585 "Unhandled order status {} for external order {}",
586 report.order_status,
587 order.client_order_id()
588 );
589 Vec::new()
590 }
591 }
592}
593
594fn create_reconciliation_fill_voids(
595 order: &OrderAny,
596 report: &OrderStatusReport,
597 ts_now: UnixNanos,
598) -> Vec<OrderEventAny> {
599 let mut remaining = order.filled_qty() - report.filled_qty;
600 let order_events = order.events();
601 let mut corrections = Vec::new();
602
603 for candidate in order_events.iter().rev() {
604 let OrderEventAny::Filled(fill) = candidate else {
605 continue;
606 };
607
608 if remaining.is_zero() {
609 break;
610 }
611 let previous = order_events.iter().rev().find_map(|event| match event {
612 OrderEventAny::FillVoided(voided) if voided.trade_id == fill.trade_id => Some(voided),
613 _ => None,
614 });
615 let prior_qty = previous.map_or_else(
616 || Quantity::zero(fill.last_qty.precision),
617 |voided| voided.voided_qty.min(fill.last_qty),
618 );
619 let effective = fill.last_qty - prior_qty;
620 if effective.is_zero() {
621 continue;
622 }
623 let removed = remaining.min(effective);
624 let voided_qty = prior_qty + removed;
625 let commission_voided = fill.commission.and_then(|commission| {
626 let fraction = voided_qty.as_decimal() / fill.last_qty.as_decimal();
627 Money::from_decimal(commission.as_decimal() * fraction, commission.currency).ok()
628 });
629 let mut event = OrderFillVoided::new(
630 fill.trader_id,
631 fill.strategy_id,
632 fill.instrument_id,
633 fill.client_order_id,
634 fill.venue_order_id,
635 fill.account_id,
636 Ustr::from(&format!(
637 "reconciliation-{}-{}",
638 report.report_id, fill.trade_id
639 )),
640 fill.trade_id,
641 voided_qty,
642 commission_voided,
643 fill.order_side,
644 fill.order_type,
645 fill.last_px,
646 fill.currency,
647 fill.liquidity_side,
648 fill.position_id,
649 report.cancel_reason.as_deref().map(Ustr::from),
650 None,
651 UUID4::new(),
652 report.ts_last,
653 ts_now,
654 true,
655 report_is_working(report),
656 );
657 event.causation_id = Some(report.report_id);
658 corrections.push(OrderEventAny::FillVoided(event));
659 remaining = remaining - removed;
660 }
661
662 if !remaining.is_zero() {
663 log::warn!(
664 "Cannot reconcile fill decrease for {}: {} is outside retained fill history",
665 order.client_order_id(),
666 remaining
667 );
668 return Vec::new();
669 }
670 corrections
671}
672
673fn create_reconciliation_terminal_fill_void(
679 order: &OrderAny,
680 report: &OrderStatusReport,
681 instrument: Option<&InstrumentAny>,
682 ts_now: UnixNanos,
683) -> Option<OrderEventAny> {
684 let voided_qty = order.leaves_qty();
685 if voided_qty.is_zero() {
686 return None;
687 }
688 let instrument = instrument?;
689 let order_side = order.order_side();
690
691 let last_px = resolve_fill_price(order, report, instrument)
692 .unwrap_or_else(|| Price::zero(instrument.price_precision()));
693
694 let mut event = OrderFillVoided::new(
695 order.trader_id(),
696 order.strategy_id(),
697 order.instrument_id(),
698 order.client_order_id(),
699 report.venue_order_id,
700 report.account_id,
701 Ustr::from(&format!("reconciliation-{}", report.report_id)),
702 TradeId::new(format!("VOID-{}", report.venue_order_id)),
703 voided_qty,
704 None,
705 order_side,
706 order.order_type(),
707 last_px,
708 instrument.quote_currency(),
709 LiquiditySide::NoLiquiditySide,
710 None,
711 report.cancel_reason.as_deref().map(Ustr::from),
712 None,
713 UUID4::new(),
714 report.ts_last,
715 ts_now,
716 true,
717 false,
718 );
719 event.causation_id = Some(report.report_id);
720 Some(OrderEventAny::FillVoided(event))
721}
722
723fn create_external_terminal_event(
724 order: &OrderAny,
725 report: &OrderStatusReport,
726 account_id: AccountId,
727 ts_now: UnixNanos,
728) -> OrderEventAny {
729 match report.order_status {
730 OrderStatus::Canceled => OrderEventAny::Canceled(OrderCanceled::new(
731 order.trader_id(),
732 order.strategy_id(),
733 order.instrument_id(),
734 order.client_order_id(),
735 UUID4::new(),
736 report.ts_last,
737 ts_now,
738 true, Some(report.venue_order_id),
740 Some(account_id),
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 debug_assert!(
770 !report.last_qty.is_zero(),
771 "fill report last_qty must be non-zero for {}",
772 order.client_order_id(),
773 );
774
775 if order.trade_ids().iter().any(|id| **id == report.trade_id) {
776 log::debug!(
777 "Duplicate fill detected: trade_id {} already exists for order {}",
778 report.trade_id,
779 order.client_order_id()
780 );
781 return None;
782 }
783
784 let potential_filled_qty = order.filled_qty() + report.last_qty;
785 if potential_filled_qty > order.quantity() {
786 if !allow_overfills {
787 log::warn!(
788 "Rejecting fill that would cause overfill for {}: order.quantity={}, order.filled_qty={}, fill.last_qty={}, would result in filled_qty={}",
789 order.client_order_id(),
790 order.quantity(),
791 order.filled_qty(),
792 report.last_qty,
793 potential_filled_qty
794 );
795 return None;
796 }
797 log::warn!(
798 "Allowing overfill during reconciliation for {}: order.quantity={}, order.filled_qty={}, fill.last_qty={}, will result in filled_qty={}",
799 order.client_order_id(),
800 order.quantity(),
801 order.filled_qty(),
802 report.last_qty,
803 potential_filled_qty
804 );
805 }
806
807 let account_id = report.account_id;
808 let venue_order_id = order.venue_order_id().unwrap_or(report.venue_order_id);
809
810 log::info!(
811 color = LogColor::Blue as u8;
812 "Reconciling fill for {}: qty={}, px={}, trade_id={}",
813 order.client_order_id(),
814 report.last_qty,
815 report.last_px,
816 report.trade_id,
817 );
818
819 Some(OrderEventAny::Filled(OrderFilled::new(
820 order.trader_id(),
821 order.strategy_id(),
822 order.instrument_id(),
823 order.client_order_id(),
824 venue_order_id,
825 account_id,
826 report.trade_id,
827 order.order_side(),
828 order.order_type(),
829 report.last_qty,
830 report.last_px,
831 instrument.quote_currency(),
832 report.liquidity_side,
833 UUID4::new(),
834 report.ts_event,
835 ts_now,
836 true, report.venue_position_id,
838 Some(report.commission),
839 None,
840 )))
841}
842
843pub fn should_reconciliation_update(order: &OrderAny, report: &OrderStatusReport) -> bool {
851 if report.quantity != order.quantity() && report.quantity >= order.filled_qty() {
852 return true;
853 }
854
855 let price_drift = report.price.is_some() && report.price != order.price();
856 let trigger_drift =
857 report.trigger_price.is_some() && report.trigger_price != order.trigger_price();
858
859 match order.order_type() {
860 OrderType::Limit => price_drift,
861 OrderType::StopMarket | OrderType::TrailingStopMarket | OrderType::MarketIfTouched => {
862 trigger_drift
863 }
864 OrderType::StopLimit | OrderType::TrailingStopLimit | OrderType::LimitIfTouched => {
865 trigger_drift || price_drift
866 }
867 _ => false,
868 }
869}
870
871#[must_use]
873pub(super) fn create_reconciliation_accepted(
874 order: &OrderAny,
875 report: &OrderStatusReport,
876 ts_now: UnixNanos,
877) -> Option<OrderEventAny> {
878 let account_id = order.account_id()?;
879
880 Some(OrderEventAny::Accepted(OrderAccepted::new(
881 order.trader_id(),
882 order.strategy_id(),
883 order.instrument_id(),
884 order.client_order_id(),
885 order.venue_order_id().unwrap_or(report.venue_order_id),
886 account_id,
887 UUID4::new(),
888 report.ts_accepted,
889 ts_now,
890 true, )))
892}
893
894#[must_use]
896pub fn create_reconciliation_rejected(
897 order: &OrderAny,
898 reason: Option<&str>,
899 ts_now: UnixNanos,
900) -> Option<OrderEventAny> {
901 let account_id = order.account_id()?;
902 let reason = reason.unwrap_or("UNKNOWN");
903
904 Some(OrderEventAny::Rejected(OrderRejected::new(
905 order.trader_id(),
906 order.strategy_id(),
907 order.instrument_id(),
908 order.client_order_id(),
909 account_id,
910 Ustr::from(reason),
911 UUID4::new(),
912 ts_now,
913 ts_now,
914 true, reason_indicates_post_only_rejection(reason),
916 )))
917}
918
919fn reason_indicates_post_only_rejection(reason: &str) -> bool {
920 let normalized: String = reason
921 .chars()
922 .filter_map(|ch| {
923 if ch == '-' || ch == '_' || ch.is_whitespace() {
924 None
925 } else {
926 Some(ch.to_ascii_lowercase())
927 }
928 })
929 .collect();
930
931 normalized.contains("postonly") || normalized.contains("postwouldexecute")
932}
933
934#[must_use]
936pub fn create_reconciliation_triggered(
937 order: &OrderAny,
938 report: &OrderStatusReport,
939 ts_now: UnixNanos,
940) -> OrderEventAny {
941 OrderEventAny::Triggered(OrderTriggered::new(
942 order.trader_id(),
943 order.strategy_id(),
944 order.instrument_id(),
945 order.client_order_id(),
946 UUID4::new(),
947 report.ts_triggered.unwrap_or(ts_now),
948 ts_now,
949 true, order.venue_order_id(),
951 order.account_id(),
952 ))
953}
954
955#[must_use]
957pub(super) fn create_reconciliation_canceled(
958 order: &OrderAny,
959 report: &OrderStatusReport,
960 ts_now: UnixNanos,
961) -> OrderEventAny {
962 OrderEventAny::Canceled(OrderCanceled::new(
963 order.trader_id(),
964 order.strategy_id(),
965 order.instrument_id(),
966 order.client_order_id(),
967 UUID4::new(),
968 report.ts_last,
969 ts_now,
970 true, order.venue_order_id(),
972 order.account_id(),
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}