1use std::str::FromStr;
29
30use anyhow::Context;
31use indexmap::IndexMap;
32use nautilus_common::cache::fifo::{FifoCache, FifoCacheMap};
33use nautilus_core::{UUID4, UnixNanos, collections::AtomicMap, time::AtomicTime};
34use nautilus_live::ExecutionEventEmitter;
35use nautilus_model::{
36 enums::{LiquiditySide, OrderSide, OrderStatus, OrderType, TimeInForce},
37 events::{
38 OrderAccepted, OrderCanceled, OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled,
39 OrderRejected, OrderUpdated,
40 },
41 identifiers::{AccountId, TradeId, VenueOrderId},
42 instruments::{Instrument, InstrumentAny},
43 reports::{FillReport, OrderStatusReport},
44 types::{Money, Price, Quantity},
45};
46use rust_decimal::Decimal;
47use ustr::Ustr;
48
49use super::{
50 messages::{
51 PolymarketUserOrder, PolymarketUserOrderStatus, PolymarketUserTrade, UserWsMessage,
52 },
53 parse::parse_timestamp_ms,
54};
55use crate::{
56 common::{
57 enums::{
58 PolymarketLiquiditySide, PolymarketOrderSide, PolymarketOrderStatus,
59 PolymarketOrderType, PolymarketTradeStatus,
60 },
61 models::PolymarketMakerOrder,
62 },
63 execution::{
64 get_pusd_currency,
65 identity::{OrderIdentity, OrderIdentityRegistry},
66 is_post_only_crossing,
67 order_fill_tracker::{BufferedFill, FillCorrectionMetadata, OrderFillTrackerMap},
68 parse::{
69 build_maker_fill_report, compute_commission, determine_order_side,
70 instrument_fee_exponent, instrument_taker_fee, parse_liquidity_side,
71 },
72 pending::PendingSubmitTracker,
73 },
74 http::error::sanitize_error_text,
75};
76
77#[derive(Debug)]
79pub(crate) struct AccountRefreshRequest;
80
81#[derive(Debug, Default)]
83pub(crate) struct WsDispatchState {
84 pub processed_fills: FifoCache<String, 10_000>,
85 matched_fills: FifoCacheMap<String, Vec<OrderFilled>, 10_000>,
86 voided_trades: FifoCache<String, 10_000>,
87 confirmed_trades: FifoCache<String, 10_000>,
88 pending_terminal_orders: FifoCacheMap<VenueOrderId, PendingTerminalOrder, 10_000>,
89 terminal_cancel_reports: FifoCacheMap<VenueOrderId, OrderStatusReport, 10_000>,
93}
94
95impl WsDispatchState {
96 pub(crate) fn restore_matched_trade(&mut self, key: String, fills: Vec<OrderFilled>) {
97 self.processed_fills.add(key.clone());
98 self.matched_fills.insert(key, fills);
99 }
100
101 pub(crate) fn restore_voided_trade(&mut self, key: String) {
102 self.processed_fills.add(key.clone());
103 self.matched_fills.remove(&key);
104 self.voided_trades.add(key);
105 }
106}
107
108#[cfg(test)]
109impl WsDispatchState {
110 pub(crate) fn matched_fill_count(&self, key: &str) -> usize {
111 self.matched_fills.get(&key.to_string()).map_or(0, Vec::len)
112 }
113
114 pub(crate) fn is_voided_trade(&self, key: &str) -> bool {
115 self.voided_trades.contains(&key.to_string())
116 }
117}
118
119#[derive(Clone, Debug)]
120struct PendingTerminalOrder {
121 trade_ids: Vec<String>,
122 ts_event: UnixNanos,
123}
124
125#[derive(Debug)]
127pub(crate) struct WsDispatchContext<'a> {
128 pub token_instruments: &'a AtomicMap<Ustr, InstrumentAny>,
129 pub fill_tracker: &'a OrderFillTrackerMap,
130 pub pending_submits: &'a PendingSubmitTracker,
131 pub order_identities: &'a OrderIdentityRegistry,
132 pub emitter: &'a ExecutionEventEmitter,
133 pub account_id: AccountId,
134 pub clock: &'static AtomicTime,
135 pub user_address: &'a str,
136 pub user_api_key: &'a str,
137}
138
139pub(crate) fn dispatch_user_message(
141 message: &UserWsMessage,
142 ctx: &WsDispatchContext<'_>,
143 state: &mut WsDispatchState,
144) -> Option<AccountRefreshRequest> {
145 match message {
146 UserWsMessage::Order(order) => {
147 dispatch_order_update(order, ctx, state);
148 None
149 }
150 UserWsMessage::Trade(trade) => dispatch_trade_update(trade, ctx, state),
151 }
152}
153
154fn dispatch_order_update(
155 order: &PolymarketUserOrder,
156 ctx: &WsDispatchContext<'_>,
157 state: &mut WsDispatchState,
158) {
159 let Some(status) = order.status.as_ref() else {
160 log::warn!("Ignoring order update without status: {}", order.id);
161 return;
162 };
163
164 let Some(order_type) = order.order_type else {
165 log::warn!("Ignoring order update without order_type: {}", order.id);
166 return;
167 };
168
169 let instruments = ctx.token_instruments.load();
170 let instrument = match instruments.get(&order.asset_id) {
171 Some(i) => i,
172 None => {
173 log::warn!("Unknown asset_id in order update: {}", order.asset_id);
174 return;
175 }
176 };
177
178 let ts_event = parse_timestamp_ms(&order.timestamp).unwrap_or_else(|_| ctx.clock.get_time_ns());
179 let venue_order_id = VenueOrderId::from(order.id.as_str());
180
181 let ts_init = ctx.clock.get_time_ns();
182 let mut report = build_ws_order_status_report(
183 order,
184 status,
185 order_type,
186 instrument,
187 ctx.account_id,
188 ts_event,
189 ts_init,
190 );
191 let local_client_order_id = ctx.pending_submits.client_order_id(&venue_order_id);
192 let mut is_accepted = ctx.fill_tracker.contains(&venue_order_id);
193 report.client_order_id = local_client_order_id;
194
195 let buffered_fills = if local_client_order_id.is_some()
197 && !is_accepted
198 && report.order_status != OrderStatus::Rejected
199 {
200 is_accepted = true;
201 ctx.fill_tracker.register_and_take_pending_fills(
202 venue_order_id,
203 local_client_order_id,
204 report.quantity,
205 report
206 .order_side
207 .expect("WebSocket order report side must be Buy or Sell"),
208 )
209 } else if is_accepted {
210 ctx.fill_tracker
211 .take_pending_fills(venue_order_id, local_client_order_id)
212 } else {
213 Vec::new()
214 };
215
216 if let Some(tracked_filled) = ctx.fill_tracker.get_cumulative_filled(&venue_order_id)
219 && report.filled_qty > tracked_filled
220 {
221 log::debug!(
222 "Capping filled_qty for {venue_order_id} from {} to {} (awaiting trade messages)",
223 report.filled_qty,
224 tracked_filled,
225 );
226 report.filled_qty = tracked_filled;
227 }
228
229 if report.order_status == OrderStatus::Canceled {
233 state
234 .terminal_cancel_reports
235 .insert(venue_order_id, report.clone());
236 }
237
238 let identity = ctx.order_identities.get(&venue_order_id);
241
242 for fill in buffered_fills {
244 match identity {
245 Some(identity) => {
246 emit_buffered_order_filled(&identity, &fill, ctx);
247 }
248 None => ctx.emitter.send_fill_report(fill.report),
249 }
250 }
251
252 if is_accepted || local_client_order_id.is_some() {
253 match identity {
254 Some(identity) => emit_tracked_order_status(&report, &identity, ts_event, ctx),
255 None => ctx.emitter.send_order_status_report(report),
256 }
257 } else if let Some(report) = ctx
258 .fill_tracker
259 .accept_or_buffer_report(venue_order_id, report)
260 {
261 match ctx.order_identities.get(&venue_order_id) {
263 Some(identity) => emit_tracked_order_status(&report, &identity, ts_event, ctx),
264 None => ctx.emitter.send_order_status_report(report),
265 }
266 }
267
268 if status.status == PolymarketOrderStatus::Matched
269 && let Some(trade_ids) = order.associate_trades.clone().filter(|ids| !ids.is_empty())
270 {
271 state.pending_terminal_orders.insert(
272 venue_order_id,
273 PendingTerminalOrder {
274 trade_ids,
275 ts_event,
276 },
277 );
278 emit_quantity_normalization_if_ready(venue_order_id, ctx, state);
279 }
280}
281
282fn emit_buffered_order_filled(
283 identity: &OrderIdentity,
284 buffered: &BufferedFill,
285 ctx: &WsDispatchContext<'_>,
286) {
287 let fill = &buffered.report;
288 ensure_accepted(identity, fill.venue_order_id, fill.ts_event, ctx);
289
290 let info = buffered
291 .correction
292 .as_ref()
293 .and_then(|correction| correction.info.clone());
294 let filled = build_order_filled(identity, fill, info, ctx);
295 ctx.fill_tracker
296 .emit_buffered_fill(filled, buffered.correction.as_ref(), |filled, new_qty| {
297 if let Some(new_qty) = new_qty {
298 emit_buy_overfill_update(
299 identity,
300 fill.venue_order_id,
301 new_qty,
302 fill.ts_event,
303 ctx,
304 );
305 }
306 ctx.emitter.send_order_event(OrderEventAny::Filled(filled));
307 });
308}
309
310fn emit_quantity_normalization_if_ready(
311 venue_order_id: VenueOrderId,
312 ctx: &WsDispatchContext<'_>,
313 state: &mut WsDispatchState,
314) {
315 let is_ready = state
316 .pending_terminal_orders
317 .get(&venue_order_id)
318 .is_some_and(|pending| {
319 pending
320 .trade_ids
321 .iter()
322 .all(|trade_id| state.confirmed_trades.contains(trade_id))
323 });
324
325 if !is_ready {
326 return;
327 }
328
329 let Some(pending) = state.pending_terminal_orders.remove(&venue_order_id) else {
330 return;
331 };
332
333 let Some(identity) = ctx.order_identities.get(&venue_order_id) else {
334 log::warn!("Cannot normalize terminal order {venue_order_id} without a local identity");
335 return;
336 };
337
338 if let Some(quantity) = ctx
339 .fill_tracker
340 .check_terminal_quantity_normalization(&venue_order_id)
341 {
342 emit_terminal_quantity_update(&identity, venue_order_id, quantity, pending.ts_event, ctx);
343 }
344}
345
346fn emit_taker_terminal_status(
352 trade: &PolymarketUserTrade,
353 ctx: &WsDispatchContext<'_>,
354 ts_event: UnixNanos,
355) {
356 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
357
358 let Some(identity) = ctx.order_identities.get(&venue_order_id) else {
359 return;
360 };
361
362 if identity.requires_terminal_quantity_normalization() {
363 if let Some(quantity) = ctx
364 .fill_tracker
365 .check_terminal_quantity_normalization(&venue_order_id)
366 {
367 emit_terminal_quantity_update(&identity, venue_order_id, quantity, ts_event, ctx);
368 }
369 return;
370 }
371
372 if identity.time_in_force == TimeInForce::Ioc
373 && let Some(remainder) = ctx
374 .fill_tracker
375 .take_terminal_ioc_remainder(&venue_order_id)
376 {
377 log::debug!(
378 "Closing terminal IOC order {venue_order_id} as Canceled (unfilled remainder={remainder})"
379 );
380 emit_order_canceled(&identity, venue_order_id, ts_event, ctx);
381 }
382}
383
384fn dispatch_trade_update(
385 trade: &PolymarketUserTrade,
386 ctx: &WsDispatchContext<'_>,
387 state: &mut WsDispatchState,
388) -> Option<AccountRefreshRequest> {
389 let dedup_key = format!("{}-{}", trade.id, trade.taker_order_id);
390 if trade.status == PolymarketTradeStatus::Failed {
391 void_failed_trade(trade, dedup_key, ctx, state);
392 return Some(AccountRefreshRequest);
393 }
394
395 if matches!(
396 trade.status,
397 PolymarketTradeStatus::Mined | PolymarketTradeStatus::Retrying
398 ) {
399 log::debug!("Waiting for terminal trade status: {}", trade.id);
400 return None;
401 }
402
403 if has_unknown_trade_instrument(trade, ctx) {
404 log::warn!(
405 "Deferring trade {} until its instrument is available",
406 trade.id
407 );
408 return None;
409 }
410
411 let is_confirmed = trade.status == PolymarketTradeStatus::Confirmed;
412 if !dispatch_trade_fills(trade, &dedup_key, is_confirmed, ctx, state) {
413 return None;
414 }
415
416 if !is_confirmed {
417 return None;
418 }
419
420 confirm_trade(trade, &dedup_key, ctx, state);
421 Some(AccountRefreshRequest)
422}
423
424fn void_failed_trade(
425 trade: &PolymarketUserTrade,
426 dedup_key: String,
427 ctx: &WsDispatchContext<'_>,
428 state: &mut WsDispatchState,
429) {
430 if state.voided_trades.contains(&dedup_key) {
431 return;
432 }
433
434 let direct_fills = state.matched_fills.remove(&dedup_key).unwrap_or_default();
435 for fill in &direct_fills {
436 ctx.fill_tracker
437 .reverse_fill(&fill.venue_order_id, fill.last_qty);
438 }
439
440 let mut fills = direct_fills;
441 fills.extend(ctx.fill_tracker.void_buffered_trade(&dedup_key));
442 for fill in fills {
443 emit_order_fill_voided(&fill, trade, Some(fill.event_id), ctx);
444 }
445
446 state.processed_fills.add(dedup_key.clone());
447 state.voided_trades.add(dedup_key);
448 state.confirmed_trades.remove(&trade.id);
449}
450
451fn has_unknown_trade_instrument(trade: &PolymarketUserTrade, ctx: &WsDispatchContext<'_>) -> bool {
452 let instruments = ctx.token_instruments.load();
453
454 if trade.trader_side == PolymarketLiquiditySide::Maker {
455 trade
456 .maker_orders
457 .iter()
458 .filter(|order| is_user_maker_order(order, ctx))
459 .any(|order| !instruments.contains_key(&order.asset_id))
460 } else {
461 !instruments.contains_key(&trade.asset_id)
462 }
463}
464
465fn dispatch_trade_fills(
466 trade: &PolymarketUserTrade,
467 dedup_key: &String,
468 is_confirmed: bool,
469 ctx: &WsDispatchContext<'_>,
470 state: &mut WsDispatchState,
471) -> bool {
472 if state.processed_fills.contains(dedup_key) {
473 log::debug!("Duplicate fill skipped: {dedup_key}");
474 return true;
475 }
476
477 let fills = if trade.trader_side == PolymarketLiquiditySide::Maker {
478 let reports = match build_ws_maker_fill_reports(trade, ctx) {
479 Ok(reports) => reports,
480 Err(e) => {
481 log::error!("Cannot build maker fills for trade {}: {e}", trade.id);
482 return false;
483 }
484 };
485 dispatch_maker_fill_reports(reports, trade, dedup_key, is_confirmed, ctx, state)
486 } else {
487 let report = match build_ws_taker_fill_report_for_trade(trade, ctx) {
488 Ok(report) => report,
489 Err(e) => {
490 log::error!("Cannot build taker fill for trade {}: {e}", trade.id);
491 return false;
492 }
493 };
494 dispatch_taker_fill_report(report, trade, dedup_key, is_confirmed, ctx, state)
495 };
496
497 if !fills.is_empty() {
498 state.matched_fills.insert(dedup_key.clone(), fills);
499 }
500 state.processed_fills.add(dedup_key.clone());
501 true
502}
503
504fn confirm_trade(
505 trade: &PolymarketUserTrade,
506 dedup_key: &str,
507 ctx: &WsDispatchContext<'_>,
508 state: &mut WsDispatchState,
509) {
510 let ts_event = parse_timestamp_ms(&trade.timestamp).unwrap_or_else(|_| ctx.clock.get_time_ns());
511 ctx.fill_tracker.mark_trade_confirmed(dedup_key);
512 state.confirmed_trades.add(trade.id.clone());
513 if trade.trader_side == PolymarketLiquiditySide::Maker {
514 for order in trade
515 .maker_orders
516 .iter()
517 .filter(|order| is_user_maker_order(order, ctx))
518 {
519 emit_quantity_normalization_if_ready(
520 VenueOrderId::from(order.order_id.as_str()),
521 ctx,
522 state,
523 );
524 }
525 } else {
526 emit_quantity_normalization_if_ready(
527 VenueOrderId::from(trade.taker_order_id.as_str()),
528 ctx,
529 state,
530 );
531 emit_taker_terminal_status(trade, ctx, ts_event);
532 }
533}
534
535fn build_ws_maker_fill_reports(
536 trade: &PolymarketUserTrade,
537 ctx: &WsDispatchContext<'_>,
538) -> anyhow::Result<Vec<FillReport>> {
539 let user_orders: Vec<_> = trade
540 .maker_orders
541 .iter()
542 .filter(|order| is_user_maker_order(order, ctx))
543 .collect();
544
545 if user_orders.is_empty() {
546 log::warn!("No matching maker orders for user in trade: {}", trade.id);
547 return Ok(Vec::new());
548 }
549
550 let instruments = ctx.token_instruments.load();
551 let liquidity_side = parse_liquidity_side(trade.trader_side);
552 let ts_event = parse_timestamp_ms(&trade.timestamp).unwrap_or_else(|_| ctx.clock.get_time_ns());
553 let ts_init = ctx.clock.get_time_ns();
554 let mut reports = Vec::with_capacity(user_orders.len());
555
556 for mo in user_orders {
557 let asset_id = mo.asset_id;
558 let instrument = instruments
559 .get(&asset_id)
560 .with_context(|| format!("unknown asset_id in maker order: {asset_id}"))?;
561 let mut report = build_maker_fill_report(
562 mo,
563 &trade.id,
564 trade.trader_side,
565 trade.side,
566 trade.asset_id.as_str(),
567 ctx.account_id,
568 instrument.id(),
569 instrument.price_precision(),
570 instrument.size_precision(),
571 crate::execution::get_pusd_currency(),
572 liquidity_side,
573 ts_event,
574 ts_init,
575 )
576 .with_context(|| format!("failed to build maker fill for asset {asset_id}"))?;
577
578 let maker_venue_order_id = report.venue_order_id;
579 report.client_order_id = ctx.pending_submits.client_order_id(&maker_venue_order_id);
580 report.last_qty = ctx
581 .fill_tracker
582 .snap_fill_qty(&maker_venue_order_id, report.last_qty);
583 reports.push(report);
584 }
585
586 Ok(reports)
587}
588
589fn dispatch_maker_fill_reports(
590 reports: Vec<FillReport>,
591 trade: &PolymarketUserTrade,
592 correction_key: &str,
593 is_confirmed: bool,
594 ctx: &WsDispatchContext<'_>,
595 state: &WsDispatchState,
596) -> Vec<OrderFilled> {
597 let fill_info = trade_fill_info(trade);
598 let mut fills = Vec::new();
599
600 for report in reports {
601 let maker_venue_order_id = report.venue_order_id;
602
603 if let Some(report) = ctx.fill_tracker.accept_or_buffer_fill(
604 maker_venue_order_id,
605 report,
606 FillCorrectionMetadata {
607 correction_key: correction_key.to_string(),
608 info: fill_info.clone(),
609 is_confirmed,
610 },
611 ) {
612 match ctx.order_identities.get(&maker_venue_order_id) {
613 Some(identity) => {
614 fills.push(emit_order_filled(
615 &identity,
616 &report,
617 fill_info.clone(),
618 ctx,
619 ));
620 }
621 None => ctx.emitter.send_fill_report(report),
622 }
623 reemit_terminal_cancel(maker_venue_order_id, state, ctx);
624 }
625 }
626 fills
627}
628
629fn is_user_maker_order(order: &PolymarketMakerOrder, ctx: &WsDispatchContext<'_>) -> bool {
630 order.is_owned_by(ctx.user_address, ctx.user_api_key)
631}
632
633fn build_ws_taker_fill_report_for_trade(
634 trade: &PolymarketUserTrade,
635 ctx: &WsDispatchContext<'_>,
636) -> anyhow::Result<FillReport> {
637 let instruments = ctx.token_instruments.load();
638 let instrument = instruments
639 .get(&trade.asset_id)
640 .with_context(|| format!("unknown asset_id in trade: {}", trade.asset_id))?;
641 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
642 let liquidity_side = parse_liquidity_side(trade.trader_side);
643 let ts_event = parse_timestamp_ms(&trade.timestamp).unwrap_or_else(|_| ctx.clock.get_time_ns());
644 let ts_init = ctx.clock.get_time_ns();
645
646 let mut report = build_ws_taker_fill_report(
647 trade,
648 instrument,
649 ctx.account_id,
650 liquidity_side,
651 ts_event,
652 ts_init,
653 )?;
654 report.client_order_id = ctx.pending_submits.client_order_id(&venue_order_id);
655 report.last_qty = ctx
656 .fill_tracker
657 .snap_fill_qty(&venue_order_id, report.last_qty);
658 Ok(report)
659}
660
661fn dispatch_taker_fill_report(
662 report: FillReport,
663 trade: &PolymarketUserTrade,
664 correction_key: &str,
665 is_confirmed: bool,
666 ctx: &WsDispatchContext<'_>,
667 state: &WsDispatchState,
668) -> Vec<OrderFilled> {
669 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
670
671 if let Some(report) = ctx.fill_tracker.accept_or_buffer_fill(
672 venue_order_id,
673 report,
674 FillCorrectionMetadata {
675 correction_key: correction_key.to_string(),
676 info: trade_fill_info(trade),
677 is_confirmed,
678 },
679 ) {
680 match ctx.order_identities.get(&venue_order_id) {
681 Some(identity) => {
682 let fill = emit_order_filled(&identity, &report, trade_fill_info(trade), ctx);
683 reemit_terminal_cancel(venue_order_id, state, ctx);
684 return vec![fill];
685 }
686 None => ctx.emitter.send_fill_report(report),
687 }
688 reemit_terminal_cancel(venue_order_id, state, ctx);
689 }
690 Vec::new()
691}
692
693fn reemit_terminal_cancel(
703 venue_order_id: VenueOrderId,
704 state: &WsDispatchState,
705 ctx: &WsDispatchContext<'_>,
706) {
707 if ctx.fill_tracker.is_fully_filled(&venue_order_id) {
708 return;
709 }
710
711 if let Some(cancel_report) = state.terminal_cancel_reports.get(&venue_order_id) {
712 log::debug!("Re-emitting cancel for {venue_order_id} after fill to restore terminal state");
713 match ctx.order_identities.get(&venue_order_id) {
714 Some(identity) => {
715 emit_order_canceled(&identity, venue_order_id, cancel_report.ts_last, ctx);
716 }
717 None => ctx.emitter.send_order_status_report(cancel_report.clone()),
718 }
719 }
720}
721
722fn build_ws_order_status_report(
723 order: &PolymarketUserOrder,
724 status: &PolymarketUserOrderStatus,
725 order_type: PolymarketOrderType,
726 instrument: &InstrumentAny,
727 account_id: AccountId,
728 ts_event: UnixNanos,
729 ts_init: UnixNanos,
730) -> OrderStatusReport {
731 let venue_order_id = VenueOrderId::from(order.id.as_str());
732 let order_status =
733 crate::execution::parse::resolve_order_status(status.status, order.event_type);
734 let order_side = OrderSide::from(order.side);
735 let time_in_force = TimeInForce::from(order_type);
736 let size_precision = instrument.size_precision();
737 let price_precision = instrument.price_precision();
738 let price_dec = Decimal::from_str(&order.price).unwrap_or_default();
739 let quantity = Decimal::from_str(&order.original_size)
740 .ok()
741 .map(|size| original_size_to_shares(size, price_dec, order.side, order_type))
742 .and_then(|d| Quantity::from_decimal_dp(d, size_precision).ok())
743 .unwrap_or_else(|| Quantity::zero(size_precision));
744 let filled_qty = Decimal::from_str(&order.size_matched)
745 .ok()
746 .and_then(|d| Quantity::from_decimal_dp(d, size_precision).ok())
747 .unwrap_or_else(|| Quantity::zero(size_precision));
748 let price = Price::from_decimal_dp(price_dec, price_precision)
749 .unwrap_or_else(|_| Price::zero(price_precision));
750
751 let mut report = OrderStatusReport::new(
752 account_id,
753 instrument.id(),
754 None,
755 venue_order_id,
756 order_side.into(),
757 OrderType::Limit,
758 time_in_force,
759 order_status,
760 quantity,
761 filled_qty,
762 ts_event,
763 ts_event,
764 ts_init,
765 None,
766 );
767 report.price = Some(price);
768
769 if order_status == OrderStatus::Rejected {
770 report.cancel_reason.clone_from(&status.reason);
771 }
772
773 report
774}
775
776fn original_size_to_shares(
787 original_size: Decimal,
788 price: Decimal,
789 side: PolymarketOrderSide,
790 order_type: PolymarketOrderType,
791) -> Decimal {
792 if side != PolymarketOrderSide::Buy
793 || !matches!(
794 order_type,
795 PolymarketOrderType::FAK | PolymarketOrderType::FOK
796 )
797 {
798 return original_size;
799 }
800
801 if price <= Decimal::ZERO {
802 log::warn!(
803 "Cannot convert {order_type} BUY size {original_size} pUSD to shares \
804 without a positive price, reporting the venue amount"
805 );
806 return original_size;
807 }
808
809 original_size / price
810}
811
812fn build_ws_taker_fill_report(
813 trade: &PolymarketUserTrade,
814 instrument: &InstrumentAny,
815 account_id: AccountId,
816 liquidity_side: LiquiditySide,
817 ts_event: UnixNanos,
818 ts_init: UnixNanos,
819) -> anyhow::Result<FillReport> {
820 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
821 let trade_id = TradeId::from(trade.id.as_str());
822 let order_side = determine_order_side(
823 trade.trader_side,
824 trade.side,
825 trade.asset_id.as_str(),
826 trade.asset_id.as_str(),
827 );
828
829 let size_precision = instrument.size_precision();
830 let price_precision = instrument.price_precision();
831 let size_dec = Decimal::from_str(&trade.size).unwrap_or_default();
832 let price_dec = Decimal::from_str(&trade.price).unwrap_or_default();
833 let last_qty = Quantity::from_decimal_dp(size_dec, size_precision)
834 .unwrap_or_else(|_| Quantity::zero(size_precision));
835 let last_px = Price::from_decimal_dp(price_dec, price_precision)
836 .unwrap_or_else(|_| Price::zero(price_precision));
837
838 let fee_rate = instrument_taker_fee(instrument);
839 let commission_value = compute_commission(
840 fee_rate,
841 instrument_fee_exponent(instrument),
842 size_dec,
843 price_dec,
844 liquidity_side,
845 );
846 let pusd = crate::execution::get_pusd_currency();
847
848 Ok(FillReport {
849 account_id,
850 instrument_id: instrument.id(),
851 venue_order_id,
852 trade_id,
853 order_side,
854 last_qty,
855 last_px,
856 commission: Money::from_decimal(commission_value, pusd)
857 .context("commission is not representable as Money")?,
858 liquidity_side,
859 avg_px: None,
860 report_id: UUID4::new(),
861 ts_event,
862 ts_init,
863 client_order_id: None,
864 venue_position_id: None,
865 })
866}
867
868fn emit_tracked_order_status(
874 report: &OrderStatusReport,
875 identity: &OrderIdentity,
876 ts_event: UnixNanos,
877 ctx: &WsDispatchContext<'_>,
878) {
879 let venue_order_id = report.venue_order_id;
880 match report.order_status {
881 OrderStatus::Accepted => ensure_accepted(identity, venue_order_id, ts_event, ctx),
882 OrderStatus::PartiallyFilled | OrderStatus::Filled => {
883 ensure_accepted(identity, venue_order_id, ts_event, ctx);
884 }
885 OrderStatus::Canceled => {
886 ensure_accepted(identity, venue_order_id, ts_event, ctx);
887 emit_order_canceled(identity, venue_order_id, ts_event, ctx);
888 }
889 OrderStatus::Expired => {
890 ensure_accepted(identity, venue_order_id, ts_event, ctx);
891 emit_order_expired(identity, venue_order_id, ts_event, ctx);
892 }
893 OrderStatus::Rejected => {
894 let reason = report
895 .cancel_reason
896 .clone()
897 .unwrap_or_else(|| "REJECTED".to_string());
898
899 emit_order_rejected(identity, &reason, ts_event, ctx);
900 }
901 other => log::debug!("No order event for status {other:?} on {venue_order_id}"),
902 }
903}
904
905fn ensure_accepted(
911 identity: &OrderIdentity,
912 venue_order_id: VenueOrderId,
913 ts_event: UnixNanos,
914 ctx: &WsDispatchContext<'_>,
915) {
916 if !ctx.order_identities.mark_accepted(venue_order_id) {
917 return;
918 }
919 let accepted = OrderAccepted::new(
920 ctx.emitter.trader_id(),
921 identity.strategy_id,
922 identity.instrument_id,
923 identity.client_order_id,
924 venue_order_id,
925 ctx.account_id,
926 UUID4::new(),
927 ts_event,
928 ctx.clock.get_time_ns(),
929 false,
930 );
931 ctx.emitter
932 .send_order_event(OrderEventAny::Accepted(accepted));
933}
934
935fn emit_order_filled(
940 identity: &OrderIdentity,
941 fill: &FillReport,
942 info: Option<IndexMap<Ustr, Ustr>>,
943 ctx: &WsDispatchContext<'_>,
944) -> OrderFilled {
945 ensure_accepted(identity, fill.venue_order_id, fill.ts_event, ctx);
946
947 if let Some(new_qty) = ctx.fill_tracker.buy_overfill_bump(&fill.venue_order_id) {
948 emit_buy_overfill_update(identity, fill.venue_order_id, new_qty, fill.ts_event, ctx);
949 }
950
951 let filled = build_order_filled(identity, fill, info, ctx);
952 ctx.emitter
953 .send_order_event(OrderEventAny::Filled(filled.clone()));
954 filled
955}
956
957fn build_order_filled(
958 identity: &OrderIdentity,
959 fill: &FillReport,
960 info: Option<IndexMap<Ustr, Ustr>>,
961 ctx: &WsDispatchContext<'_>,
962) -> OrderFilled {
963 OrderFilled::new(
964 ctx.emitter.trader_id(),
965 identity.strategy_id,
966 identity.instrument_id,
967 identity.client_order_id,
968 fill.venue_order_id,
969 ctx.account_id,
970 fill.trade_id,
971 identity.order_side,
972 identity.order_type,
973 fill.last_qty,
974 fill.last_px,
975 get_pusd_currency(),
976 fill.liquidity_side,
977 UUID4::new(),
978 fill.ts_event,
979 fill.ts_init,
980 false,
981 fill.venue_position_id,
982 Some(fill.commission),
983 info,
984 )
985}
986
987fn emit_order_fill_voided(
988 fill: &OrderFilled,
989 trade: &PolymarketUserTrade,
990 causation_id: Option<UUID4>,
991 ctx: &WsDispatchContext<'_>,
992) {
993 let ts_event = parse_timestamp_ms(&trade.timestamp).unwrap_or_else(|_| ctx.clock.get_time_ns());
994 let mut voided = OrderFillVoided::new(
995 fill.trader_id,
996 fill.strategy_id,
997 fill.instrument_id,
998 fill.client_order_id,
999 fill.venue_order_id,
1000 fill.account_id,
1001 Ustr::from(&format!("{}-FAILED-{}", trade.id, fill.client_order_id)),
1002 fill.trade_id,
1003 fill.last_qty,
1004 fill.commission,
1005 fill.order_side,
1006 fill.order_type,
1007 fill.last_px,
1008 fill.currency,
1009 fill.liquidity_side,
1010 fill.position_id,
1011 Some(Ustr::from("FAILED")),
1012 trade_fill_info(trade),
1013 UUID4::new(),
1014 ts_event,
1015 ctx.clock.get_time_ns(),
1016 false,
1017 false,
1018 );
1019 voided.causation_id = causation_id;
1020 ctx.emitter
1021 .send_order_event(OrderEventAny::FillVoided(voided));
1022}
1023
1024fn trade_fill_info(trade: &PolymarketUserTrade) -> Option<IndexMap<Ustr, Ustr>> {
1029 let value = serde_json::to_value(trade).ok()?;
1030 let object = value.as_object()?;
1031 let mut info = IndexMap::with_capacity(object.len());
1032 for (key, val) in object {
1033 let val_str = match val {
1034 serde_json::Value::String(s) => s.clone(),
1035 other => other.to_string(),
1036 };
1037 info.insert(Ustr::from(key.as_str()), Ustr::from(val_str.as_str()));
1038 }
1039 Some(info)
1040}
1041
1042fn emit_buy_overfill_update(
1048 identity: &OrderIdentity,
1049 venue_order_id: VenueOrderId,
1050 new_qty: Quantity,
1051 ts_event: UnixNanos,
1052 ctx: &WsDispatchContext<'_>,
1053) {
1054 let updated = OrderUpdated::new(
1055 ctx.emitter.trader_id(),
1056 identity.strategy_id,
1057 identity.instrument_id,
1058 identity.client_order_id,
1059 new_qty,
1060 UUID4::new(),
1061 ts_event,
1062 ctx.clock.get_time_ns(),
1063 false,
1064 Some(venue_order_id),
1065 Some(ctx.account_id),
1066 None,
1067 None,
1068 None,
1069 false,
1070 );
1071 ctx.emitter
1072 .send_order_event(OrderEventAny::Updated(updated));
1073}
1074
1075fn emit_terminal_quantity_update(
1077 identity: &OrderIdentity,
1078 venue_order_id: VenueOrderId,
1079 quantity: Quantity,
1080 ts_event: UnixNanos,
1081 ctx: &WsDispatchContext<'_>,
1082) {
1083 let updated = OrderUpdated::new(
1084 ctx.emitter.trader_id(),
1085 identity.strategy_id,
1086 identity.instrument_id,
1087 identity.client_order_id,
1088 quantity,
1089 UUID4::new(),
1090 ts_event,
1091 ctx.clock.get_time_ns(),
1092 true,
1093 Some(venue_order_id),
1094 Some(ctx.account_id),
1095 None,
1096 None,
1097 None,
1098 false,
1099 );
1100 ctx.emitter
1101 .send_order_event(OrderEventAny::Updated(updated));
1102}
1103
1104fn emit_order_canceled(
1105 identity: &OrderIdentity,
1106 venue_order_id: VenueOrderId,
1107 ts_event: UnixNanos,
1108 ctx: &WsDispatchContext<'_>,
1109) {
1110 let canceled = OrderCanceled::new(
1111 ctx.emitter.trader_id(),
1112 identity.strategy_id,
1113 identity.instrument_id,
1114 identity.client_order_id,
1115 UUID4::new(),
1116 ts_event,
1117 ctx.clock.get_time_ns(),
1118 false,
1119 Some(venue_order_id),
1120 Some(ctx.account_id),
1121 );
1122 ctx.emitter
1123 .send_order_event(OrderEventAny::Canceled(canceled));
1124}
1125
1126fn emit_order_expired(
1127 identity: &OrderIdentity,
1128 venue_order_id: VenueOrderId,
1129 ts_event: UnixNanos,
1130 ctx: &WsDispatchContext<'_>,
1131) {
1132 let expired = OrderExpired::new(
1133 ctx.emitter.trader_id(),
1134 identity.strategy_id,
1135 identity.instrument_id,
1136 identity.client_order_id,
1137 UUID4::new(),
1138 ts_event,
1139 ctx.clock.get_time_ns(),
1140 false,
1141 Some(venue_order_id),
1142 Some(ctx.account_id),
1143 );
1144 ctx.emitter
1145 .send_order_event(OrderEventAny::Expired(expired));
1146}
1147
1148fn emit_order_rejected(
1149 identity: &OrderIdentity,
1150 reason: &str,
1151 ts_event: UnixNanos,
1152 ctx: &WsDispatchContext<'_>,
1153) {
1154 let reason = sanitize_error_text(reason);
1155
1156 let rejected = OrderRejected::new(
1157 ctx.emitter.trader_id(),
1158 identity.strategy_id,
1159 identity.instrument_id,
1160 identity.client_order_id,
1161 ctx.account_id,
1162 Ustr::from(&reason),
1163 UUID4::new(),
1164 ts_event,
1165 ctx.clock.get_time_ns(),
1166 false,
1167 is_post_only_crossing(&reason),
1168 );
1169 ctx.emitter
1170 .send_order_event(OrderEventAny::Rejected(rejected));
1171}
1172
1173#[cfg(test)]
1174mod tests {
1175 use nautilus_common::messages::{ExecutionEvent, ExecutionReport};
1176 use nautilus_core::time::AtomicTime;
1177 use nautilus_model::{
1178 enums::{AccountType, OrderSide, OrderStatus},
1179 events::OrderEventAny,
1180 identifiers::{ClientOrderId, InstrumentId, StrategyId, TraderId},
1181 orders::{Order, builder::OrderTestBuilder},
1182 types::Currency,
1183 };
1184 use rstest::rstest;
1185 use rust_decimal_macros::dec;
1186
1187 use super::*;
1188 use crate::http::{
1189 models::GammaMarket,
1190 parse::{create_instrument_from_def, parse_gamma_market},
1191 };
1192
1193 fn register_identity(
1195 order_identities: &OrderIdentityRegistry,
1196 venue_order_id: VenueOrderId,
1197 instrument_id: InstrumentId,
1198 client_order_id: &str,
1199 ) {
1200 order_identities.register_order_identity(
1201 venue_order_id,
1202 OrderIdentity {
1203 client_order_id: ClientOrderId::from(client_order_id),
1204 strategy_id: StrategyId::from("S-001"),
1205 instrument_id,
1206 order_side: OrderSide::Buy,
1207 order_type: OrderType::Limit,
1208 time_in_force: TimeInForce::Gtc,
1209 },
1210 );
1211 }
1212
1213 fn load<T: serde::de::DeserializeOwned>(filename: &str) -> T {
1214 let path = format!("test_data/{filename}");
1215 let content = std::fs::read_to_string(path).expect("Failed to read test data");
1216 serde_json::from_str(&content).expect("Failed to parse test data")
1217 }
1218
1219 fn test_instrument() -> InstrumentAny {
1220 let market: GammaMarket = load("gamma_market.json");
1221 let defs = parse_gamma_market(&market).unwrap();
1222 create_instrument_from_def(&defs[0], UnixNanos::from(1_000_000_000u64)).unwrap()
1223 }
1224
1225 fn test_emitter() -> ExecutionEventEmitter {
1226 ExecutionEventEmitter::new(
1227 nautilus_core::time::get_atomic_clock_realtime(),
1228 TraderId::from("TESTER-001"),
1229 AccountId::from("POLY-001"),
1230 AccountType::Cash,
1231 Some(Currency::pUSD()),
1232 )
1233 }
1234
1235 #[rstest]
1236 fn test_emit_order_rejected_uses_bounded_clean_reason() {
1237 let instrument = test_instrument();
1238 let token_instruments = AtomicMap::new();
1239 let fill_tracker = OrderFillTrackerMap::new();
1240 let pending_submits = PendingSubmitTracker::default();
1241 let order_identities = OrderIdentityRegistry::default();
1242 let mut emitter = test_emitter();
1243 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
1244
1245 emitter.set_sender(sender);
1246
1247 let ctx = WsDispatchContext {
1248 token_instruments: &token_instruments,
1249 fill_tracker: &fill_tracker,
1250 pending_submits: &pending_submits,
1251 order_identities: &order_identities,
1252 emitter: &emitter,
1253 account_id: AccountId::from("POLY-001"),
1254 clock: nautilus_core::time::get_atomic_clock_realtime(),
1255 user_address: "0xtest",
1256 user_api_key: "test-key",
1257 };
1258 let identity = OrderIdentity {
1259 client_order_id: ClientOrderId::from("O-WS-REJECT"),
1260 strategy_id: StrategyId::from("S-001"),
1261 instrument_id: instrument.id(),
1262 order_side: OrderSide::Buy,
1263 order_type: OrderType::Limit,
1264 time_in_force: TimeInForce::Gtc,
1265 };
1266
1267 emit_order_rejected(
1268 &identity,
1269 " invalid post-only order:\norder crosses book ",
1270 UnixNanos::from(1_000_000_000),
1271 &ctx,
1272 );
1273
1274 match receiver.try_recv().expect("expected rejected event") {
1275 ExecutionEvent::Order(OrderEventAny::Rejected(event)) => {
1276 assert_eq!(
1277 event.reason.as_str(),
1278 "invalid post-only order: order crosses book"
1279 );
1280 assert!(event.due_post_only);
1281 }
1282 other => panic!("expected rejected event, was {other:?}"),
1283 }
1284 }
1285
1286 #[rstest]
1287 fn test_build_ws_order_status_report() {
1288 let order: PolymarketUserOrder = load("ws_user_order_placement.json");
1289 let instrument = test_instrument();
1290 let ts_event = UnixNanos::from(1_000_000_000u64);
1291 let ts_init = UnixNanos::from(2_000_000_000u64);
1292
1293 let report = build_ws_order_status_report(
1294 &order,
1295 order.status.as_ref().unwrap(),
1296 order.order_type.unwrap(),
1297 &instrument,
1298 AccountId::from("POLY-001"),
1299 ts_event,
1300 ts_init,
1301 );
1302
1303 assert_eq!(report.order_side, Some(OrderSide::Buy));
1304 assert_eq!(report.order_type, OrderType::Limit);
1305 assert_eq!(report.quantity.as_decimal(), dec!(100));
1307 assert_eq!(
1308 report.price.map(|price| price.as_decimal()),
1309 Some(dec!(0.5))
1310 );
1311 assert_eq!(report.ts_accepted, ts_event);
1312 assert_eq!(report.ts_init, ts_init);
1313 }
1314
1315 #[rstest]
1316 fn test_build_ws_order_status_report_venue_cancel_maps_to_canceled() {
1317 let order: PolymarketUserOrder = load("ws_user_order_venue_cancel.json");
1318 let instrument = test_instrument();
1319 let ts_event = UnixNanos::from(1_000_000_000u64);
1320 let ts_init = UnixNanos::from(2_000_000_000u64);
1321
1322 let report = build_ws_order_status_report(
1323 &order,
1324 order.status.as_ref().unwrap(),
1325 order.order_type.unwrap(),
1326 &instrument,
1327 AccountId::from("POLY-001"),
1328 ts_event,
1329 ts_init,
1330 );
1331
1332 assert_eq!(report.order_status, OrderStatus::Canceled);
1333 }
1334
1335 #[rstest]
1338 #[case(
1339 PolymarketOrderSide::Buy,
1340 PolymarketOrderType::FOK,
1341 dec!(1.01),
1342 dec!(0.01),
1343 dec!(101)
1344 )]
1345 #[case(
1346 PolymarketOrderSide::Buy,
1347 PolymarketOrderType::FOK,
1348 dec!(12),
1349 dec!(0.6),
1350 dec!(20)
1351 )]
1352 #[case(
1353 PolymarketOrderSide::Buy,
1354 PolymarketOrderType::FAK,
1355 dec!(1),
1356 dec!(0.01),
1357 dec!(100)
1358 )]
1359 #[case(
1360 PolymarketOrderSide::Buy,
1361 PolymarketOrderType::GTC,
1362 dec!(20),
1363 dec!(0.18),
1364 dec!(20)
1365 )]
1366 #[case(
1367 PolymarketOrderSide::Buy,
1368 PolymarketOrderType::GTD,
1369 dec!(20),
1370 dec!(0.18),
1371 dec!(20)
1372 )]
1373 #[case(
1374 PolymarketOrderSide::Sell,
1375 PolymarketOrderType::FOK,
1376 dec!(20),
1377 dec!(0.6),
1378 dec!(20)
1379 )]
1380 #[case(
1381 PolymarketOrderSide::Buy,
1382 PolymarketOrderType::FOK,
1383 dec!(1.01),
1384 dec!(0),
1385 dec!(1.01)
1386 )]
1387 fn test_original_size_to_shares(
1388 #[case] side: PolymarketOrderSide,
1389 #[case] order_type: PolymarketOrderType,
1390 #[case] original_size: Decimal,
1391 #[case] price: Decimal,
1392 #[case] expected: Decimal,
1393 ) {
1394 let shares = original_size_to_shares(original_size, price, side, order_type);
1395
1396 assert_eq!(shares, expected);
1397 }
1398
1399 #[rstest]
1402 #[case("1", "0.03", "33.333333", "0.03")]
1403 #[case("1.01", "", "1.01", "0")]
1404 fn test_build_ws_order_status_report_fok_buy_quantity(
1405 #[case] original_size: &str,
1406 #[case] price: &str,
1407 #[case] expected_quantity: &str,
1408 #[case] expected_price: &str,
1409 ) {
1410 let mut order: PolymarketUserOrder = load("ws_user_order_fok_buy_pusd_size.json");
1411 order.original_size = original_size.to_string();
1412 order.price = price.to_string();
1413 let instrument = test_instrument();
1414
1415 let report = build_ws_order_status_report(
1416 &order,
1417 order.status.as_ref().unwrap(),
1418 order.order_type.unwrap(),
1419 &instrument,
1420 AccountId::from("POLY-001"),
1421 UnixNanos::from(1_000_000_000u64),
1422 UnixNanos::from(2_000_000_000u64),
1423 );
1424
1425 assert_eq!(
1426 report.quantity.as_decimal(),
1427 Decimal::from_str_exact(expected_quantity).unwrap()
1428 );
1429 assert_eq!(
1430 report.price.map(|price| price.as_decimal()),
1431 Some(Decimal::from_str_exact(expected_price).unwrap())
1432 );
1433 }
1434
1435 #[rstest]
1436 fn test_dispatch_fok_buy_registers_share_quantity_for_in_flight_submit() {
1437 let order: PolymarketUserOrder = load("ws_user_order_fok_buy_pusd_size.json");
1438 let instrument = test_instrument();
1439
1440 let token_instruments = AtomicMap::new();
1441 token_instruments.insert(order.asset_id, instrument.clone());
1442
1443 let fill_tracker = OrderFillTrackerMap::new();
1445 let pending_submits = PendingSubmitTracker::default();
1446 let order_identities = OrderIdentityRegistry::default();
1447 let emitter = test_emitter();
1448
1449 let venue_order_id = VenueOrderId::from(order.id.as_str());
1450 let client_order_id = ClientOrderId::from("O-FOK-IN-FLIGHT");
1451 pending_submits.insert(venue_order_id, client_order_id);
1452 register_identity(
1453 &order_identities,
1454 venue_order_id,
1455 instrument.id(),
1456 client_order_id.as_str(),
1457 );
1458
1459 let ctx = WsDispatchContext {
1460 token_instruments: &token_instruments,
1461 fill_tracker: &fill_tracker,
1462 pending_submits: &pending_submits,
1463 order_identities: &order_identities,
1464 emitter: &emitter,
1465 account_id: AccountId::from("POLY-001"),
1466 clock: nautilus_core::time::get_atomic_clock_realtime(),
1467 user_address: "0xtest",
1468 user_api_key: "test-key",
1469 };
1470 let mut state = WsDispatchState::default();
1471
1472 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
1473
1474 assert_eq!(
1476 fill_tracker
1477 .submitted_qty(&venue_order_id)
1478 .map(|qty| qty.as_decimal()),
1479 Some(dec!(101)),
1480 );
1481 }
1482
1483 #[rstest]
1484 fn test_dispatch_fok_buy_report_quantity_is_shares_without_identity() {
1485 let order: PolymarketUserOrder = load("ws_user_order_fok_buy_pusd_size.json");
1486 let instrument = test_instrument();
1487
1488 let token_instruments = AtomicMap::new();
1489 token_instruments.insert(order.asset_id, instrument.clone());
1490
1491 let fill_tracker = OrderFillTrackerMap::new();
1492 let venue_order_id = VenueOrderId::from(order.id.as_str());
1493 fill_tracker.register(
1494 venue_order_id,
1495 Quantity::from("101"),
1496 OrderSide::Buy,
1497 instrument.id(),
1498 instrument.size_precision(),
1499 instrument.price_precision(),
1500 );
1501
1502 let pending_submits = PendingSubmitTracker::default();
1503 let order_identities = OrderIdentityRegistry::default();
1505 let mut emitter = test_emitter();
1506 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
1507 emitter.set_sender(sender);
1508
1509 let ctx = WsDispatchContext {
1510 token_instruments: &token_instruments,
1511 fill_tracker: &fill_tracker,
1512 pending_submits: &pending_submits,
1513 order_identities: &order_identities,
1514 emitter: &emitter,
1515 account_id: AccountId::from("POLY-001"),
1516 clock: nautilus_core::time::get_atomic_clock_realtime(),
1517 user_address: "0xtest",
1518 user_api_key: "test-key",
1519 };
1520 let mut state = WsDispatchState::default();
1521
1522 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
1523
1524 let event = receiver.try_recv().expect("expected order report");
1525 let ExecutionEvent::Report(ExecutionReport::Order(report)) = event else {
1526 panic!("expected an order report, was {event:?}");
1527 };
1528
1529 assert_eq!(report.venue_order_id, venue_order_id);
1530 assert_eq!(report.order_side, Some(OrderSide::Buy));
1531 assert_eq!(report.time_in_force, TimeInForce::Fok);
1532 assert_eq!(report.order_status, OrderStatus::Canceled);
1533 assert_eq!(report.quantity.as_decimal(), dec!(101));
1534 assert_eq!(report.filled_qty.as_decimal(), dec!(0));
1535 assert_eq!(
1536 report.price.map(|price| price.as_decimal()),
1537 Some(dec!(0.01))
1538 );
1539 }
1540
1541 #[rstest]
1542 fn test_build_ws_taker_fill_report() {
1543 let trade: PolymarketUserTrade = load("ws_user_trade.json");
1544 let instrument = test_instrument();
1545 let ts_event = UnixNanos::from(1_000_000_000u64);
1546 let ts_init = UnixNanos::from(2_000_000_000u64);
1547
1548 let report = build_ws_taker_fill_report(
1549 &trade,
1550 &instrument,
1551 AccountId::from("POLY-001"),
1552 LiquiditySide::Taker,
1553 ts_event,
1554 ts_init,
1555 )
1556 .expect("representable commission builds a fill report");
1557
1558 assert_eq!(report.order_side, OrderSide::Buy);
1559 assert_eq!(report.liquidity_side, LiquiditySide::Taker);
1560 assert_eq!(report.trade_id.as_str(), trade.id);
1561 assert_eq!(report.ts_event, ts_event);
1562 assert_eq!(report.ts_init, ts_init);
1563 }
1564
1565 #[rstest]
1566 fn test_trade_fill_info_flattens_raw_trade() {
1567 let trade: PolymarketUserTrade = load("ws_user_trade.json");
1568
1569 let info = trade_fill_info(&trade).expect("info should be present");
1570
1571 assert_eq!(info.len(), 21);
1573 assert_eq!(info[&Ustr::from("id")], Ustr::from("trade-0xabcdef1234"));
1574 assert_eq!(info[&Ustr::from("fee_rate_bps")], Ustr::from("0"));
1575 assert_eq!(
1576 info[&Ustr::from("transaction_hash")],
1577 Ustr::from("0xabcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890ab")
1578 );
1579 assert_eq!(info[&Ustr::from("bucket_index")], Ustr::from("1"));
1581 assert_eq!(info[&Ustr::from("size")], Ustr::from("25.0"));
1582 assert_eq!(
1583 info[&Ustr::from("taker_order_id")],
1584 Ustr::from("0x1234567890abcdef1234567890abcdef1234567890abcdef1234567890abcdef12")
1585 );
1586 assert_eq!(info[&Ustr::from("type")], Ustr::from("TRADE"));
1588 let maker_orders = info[&Ustr::from("maker_orders")].as_str();
1590 assert!(maker_orders.starts_with('['));
1591 assert!(maker_orders.contains("order_id"));
1592
1593 let empty_hash_trade: PolymarketUserTrade = load("ws_user_trade_msg.json");
1594 let empty_hash_info =
1595 trade_fill_info(&empty_hash_trade).expect("empty hash info should be present");
1596 assert!(!empty_hash_info.contains_key(&Ustr::from("transaction_hash")));
1597 }
1598
1599 #[rstest]
1600 fn test_dispatch_order_message_buffers_when_not_accepted() {
1601 let order: PolymarketUserOrder = load("ws_user_order_placement.json");
1602 let instrument = test_instrument();
1603
1604 let token_instruments = AtomicMap::new();
1605 token_instruments.insert(order.asset_id, instrument);
1606
1607 let fill_tracker = OrderFillTrackerMap::new();
1608 let pending_submits = PendingSubmitTracker::default();
1609 let order_identities = OrderIdentityRegistry::default();
1610 let emitter = test_emitter();
1611
1612 let ctx = WsDispatchContext {
1613 token_instruments: &token_instruments,
1614 fill_tracker: &fill_tracker,
1615 pending_submits: &pending_submits,
1616 order_identities: &order_identities,
1617 emitter: &emitter,
1618 account_id: AccountId::from("POLY-001"),
1619 clock: nautilus_core::time::get_atomic_clock_realtime(),
1620 user_address: "0xtest",
1621 user_api_key: "test-key",
1622 };
1623 let mut state = WsDispatchState::default();
1624
1625 let result = dispatch_user_message(&UserWsMessage::Order(order.clone()), &ctx, &mut state);
1626 assert!(result.is_none());
1627
1628 let venue_order_id = VenueOrderId::from(order.id.as_str());
1630 assert!(fill_tracker.has_pending_report(&venue_order_id));
1631 }
1632
1633 #[rstest]
1634 fn test_dispatch_order_message_ignores_missing_lifecycle_fields() {
1635 let order: PolymarketUserOrder = load("ws_user_order_placement.json");
1636 let instrument = test_instrument();
1637 let token_instruments = AtomicMap::new();
1638 token_instruments.insert(order.asset_id, instrument);
1639 let fill_tracker = OrderFillTrackerMap::new();
1640 let pending_submits = PendingSubmitTracker::default();
1641 let order_identities = OrderIdentityRegistry::default();
1642 let emitter = test_emitter();
1643 let ctx = WsDispatchContext {
1644 token_instruments: &token_instruments,
1645 fill_tracker: &fill_tracker,
1646 pending_submits: &pending_submits,
1647 order_identities: &order_identities,
1648 emitter: &emitter,
1649 account_id: AccountId::from("POLY-001"),
1650 clock: nautilus_core::time::get_atomic_clock_realtime(),
1651 user_address: "0xtest",
1652 user_api_key: "test-key",
1653 };
1654 let venue_order_id = VenueOrderId::from(order.id.as_str());
1655
1656 let mut missing_status = order.clone();
1657 missing_status.status = None;
1658 dispatch_user_message(
1659 &UserWsMessage::Order(missing_status),
1660 &ctx,
1661 &mut WsDispatchState::default(),
1662 );
1663 assert!(!fill_tracker.has_pending_report(&venue_order_id));
1664
1665 let mut missing_order_type = order;
1666 missing_order_type.order_type = None;
1667 dispatch_user_message(
1668 &UserWsMessage::Order(missing_order_type),
1669 &ctx,
1670 &mut WsDispatchState::default(),
1671 );
1672 assert!(!fill_tracker.has_pending_report(&venue_order_id));
1673 }
1674
1675 #[rstest]
1676 fn test_dispatch_order_message_uses_pending_submit_client_order_id() {
1677 let order: PolymarketUserOrder = load("ws_user_order_placement.json");
1678 let instrument = test_instrument();
1679
1680 let token_instruments = AtomicMap::new();
1681 token_instruments.insert(order.asset_id, instrument);
1682
1683 let fill_tracker = OrderFillTrackerMap::new();
1684 let pending_submits = PendingSubmitTracker::default();
1685 let order_identities = OrderIdentityRegistry::default();
1686 let mut emitter = test_emitter();
1687 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
1688 emitter.set_sender(sender);
1689
1690 let venue_order_id = VenueOrderId::from(order.id.as_str());
1691 let client_order_id = ClientOrderId::from("O-UNKNOWN-SUBMIT");
1692 pending_submits.insert(venue_order_id, client_order_id);
1693 register_identity(
1694 &order_identities,
1695 venue_order_id,
1696 test_instrument().id(),
1697 "O-UNKNOWN-SUBMIT",
1698 );
1699
1700 let ctx = WsDispatchContext {
1701 token_instruments: &token_instruments,
1702 fill_tracker: &fill_tracker,
1703 pending_submits: &pending_submits,
1704 order_identities: &order_identities,
1705 emitter: &emitter,
1706 account_id: AccountId::from("POLY-001"),
1707 clock: nautilus_core::time::get_atomic_clock_realtime(),
1708 user_address: "0xtest",
1709 user_api_key: "test-key",
1710 };
1711 let mut state = WsDispatchState::default();
1712
1713 let _ = dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
1714
1715 let event = receiver.try_recv().expect("expected accepted event");
1717 match event {
1718 ExecutionEvent::Order(OrderEventAny::Accepted(accepted)) => {
1719 assert_eq!(accepted.client_order_id, client_order_id);
1720 }
1721 other => panic!("Expected accepted event, was {other:?}"),
1722 }
1723
1724 assert!(!fill_tracker.has_pending_report(&venue_order_id));
1725 }
1726
1727 #[rstest]
1728 fn test_dispatch_maker_fill_owned_by_case_variant_address() {
1729 let mut trade: PolymarketUserTrade = load("ws_user_trade.json");
1730 trade.trader_side = PolymarketLiquiditySide::Maker;
1731 let configured_address = trade.maker_orders[0].maker_address.clone();
1732 let case_variant_address = configured_address
1733 .to_ascii_uppercase()
1734 .replacen("0X", "0x", 1);
1735 assert_ne!(case_variant_address, configured_address);
1736 trade.maker_orders[0].maker_address = case_variant_address;
1737 let foreign_api_key = "ffffffff-ffff-ffff-ffff-ffffffffffff";
1738 assert_ne!(trade.maker_orders[0].owner, foreign_api_key);
1739
1740 let venue_order_id = VenueOrderId::from(trade.maker_orders[0].order_id.as_str());
1741 let token_instruments = AtomicMap::new();
1742 token_instruments.insert(trade.maker_orders[0].asset_id, test_instrument());
1743 let fill_tracker = OrderFillTrackerMap::new();
1744 let pending_submits = PendingSubmitTracker::default();
1745 let order_identities = OrderIdentityRegistry::default();
1746 let emitter = test_emitter();
1747 let ctx = WsDispatchContext {
1748 token_instruments: &token_instruments,
1749 fill_tracker: &fill_tracker,
1750 pending_submits: &pending_submits,
1751 order_identities: &order_identities,
1752 emitter: &emitter,
1753 account_id: AccountId::from("POLY-001"),
1754 clock: nautilus_core::time::get_atomic_clock_realtime(),
1755 user_address: &configured_address,
1756 user_api_key: foreign_api_key,
1757 };
1758 let mut state = WsDispatchState::default();
1759
1760 let _ = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
1761
1762 let fills = fill_tracker.pending_fills_for(&venue_order_id);
1763 assert_eq!(fills.len(), 1);
1764 assert_eq!(fills[0].venue_order_id, venue_order_id);
1765 }
1766
1767 #[rstest]
1768 fn test_dispatch_trade_dedup() {
1769 let trade: PolymarketUserTrade = load("ws_user_trade.json");
1770 let instrument = test_instrument();
1771
1772 let token_instruments = AtomicMap::new();
1773 token_instruments.insert(trade.asset_id, instrument);
1774
1775 let fill_tracker = OrderFillTrackerMap::new();
1776 let pending_submits = PendingSubmitTracker::default();
1777 let order_identities = OrderIdentityRegistry::default();
1778 let emitter = test_emitter();
1779
1780 let ctx = WsDispatchContext {
1781 token_instruments: &token_instruments,
1782 fill_tracker: &fill_tracker,
1783 pending_submits: &pending_submits,
1784 order_identities: &order_identities,
1785 emitter: &emitter,
1786 account_id: AccountId::from("POLY-001"),
1787 clock: nautilus_core::time::get_atomic_clock_realtime(),
1788 user_address: "0xtest",
1789 user_api_key: "test-key",
1790 };
1791 let mut state = WsDispatchState::default();
1792
1793 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
1794
1795 let _ = dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
1797 assert_eq!(fill_tracker.pending_fills_for(&venue_order_id).len(), 1);
1798
1799 let _ = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
1801 assert_eq!(fill_tracker.pending_fills_for(&venue_order_id).len(), 1);
1802 }
1803
1804 #[rstest]
1805 fn test_dispatch_taker_commission_failure_preserves_replay_state() {
1806 let trade: PolymarketUserTrade = load("ws_user_trade.json");
1807 let valid_instrument = test_instrument();
1808 let mut invalid_instrument = valid_instrument.clone();
1809 let InstrumentAny::BinaryOption(binary_option) = &mut invalid_instrument else {
1810 panic!("expected binary option test instrument");
1811 };
1812 binary_option.taker_fee =
1813 Decimal::from_i128_with_scale(100_000_000_000_000_000_000_000_000i128, 0);
1814
1815 let token_instruments = AtomicMap::new();
1816 token_instruments.insert(trade.asset_id, invalid_instrument);
1817 let fill_tracker = OrderFillTrackerMap::new();
1818 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
1819 fill_tracker.register(
1820 venue_order_id,
1821 Quantity::from("100"),
1822 OrderSide::Buy,
1823 valid_instrument.id(),
1824 valid_instrument.size_precision(),
1825 valid_instrument.price_precision(),
1826 );
1827 let pending_submits = PendingSubmitTracker::default();
1828 let order_identities = OrderIdentityRegistry::default();
1829 register_identity(
1830 &order_identities,
1831 venue_order_id,
1832 valid_instrument.id(),
1833 "O-COMMISSION-REPLAY",
1834 );
1835 order_identities.mark_accepted(venue_order_id);
1836 let mut emitter = test_emitter();
1837 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
1838 emitter.set_sender(sender);
1839 let ctx = WsDispatchContext {
1840 token_instruments: &token_instruments,
1841 fill_tracker: &fill_tracker,
1842 pending_submits: &pending_submits,
1843 order_identities: &order_identities,
1844 emitter: &emitter,
1845 account_id: AccountId::from("POLY-001"),
1846 clock: nautilus_core::time::get_atomic_clock_realtime(),
1847 user_address: "0xtest",
1848 user_api_key: "test-key",
1849 };
1850 let mut state = WsDispatchState::default();
1851 let dedup_key = format!("{}-{}", trade.id, trade.taker_order_id);
1852
1853 let failed = dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
1854
1855 assert!(failed.is_none());
1856 assert!(!state.processed_fills.contains(&dedup_key));
1857 assert!(!state.confirmed_trades.contains(&trade.id));
1858 assert!(!fill_tracker.is_trade_confirmed(&dedup_key));
1859 assert_eq!(
1860 fill_tracker.get_cumulative_filled(&venue_order_id),
1861 Some(Quantity::zero(valid_instrument.size_precision()))
1862 );
1863 assert!(receiver.try_recv().is_err());
1864
1865 token_instruments.insert(trade.asset_id, valid_instrument);
1866 let replay = dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
1867 let emitted = receiver.try_recv().expect("valid replay emits one fill");
1868 let duplicate = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
1869
1870 assert!(replay.is_some());
1871 assert!(duplicate.is_some());
1872 assert!(state.processed_fills.contains(&dedup_key));
1873 assert!(
1874 state
1875 .confirmed_trades
1876 .contains(&"trade-0xabcdef1234".to_string())
1877 );
1878 assert!(fill_tracker.is_trade_confirmed(&dedup_key));
1879 assert!(matches!(
1880 emitted,
1881 ExecutionEvent::Order(OrderEventAny::Filled(_))
1882 ));
1883 assert_eq!(
1884 fill_tracker.get_cumulative_filled(&venue_order_id),
1885 Some(Quantity::from("25.0"))
1886 );
1887 assert!(receiver.try_recv().is_err());
1888 }
1889
1890 #[rstest]
1891 fn test_dispatch_trade_replays_after_instrument_becomes_available() {
1892 let trade: PolymarketUserTrade = load("ws_user_trade.json");
1893 let instrument = test_instrument();
1894 let token_instruments = AtomicMap::new();
1895 let fill_tracker = OrderFillTrackerMap::new();
1896 let pending_submits = PendingSubmitTracker::default();
1897 let order_identities = OrderIdentityRegistry::default();
1898 let emitter = test_emitter();
1899 let ctx = WsDispatchContext {
1900 token_instruments: &token_instruments,
1901 fill_tracker: &fill_tracker,
1902 pending_submits: &pending_submits,
1903 order_identities: &order_identities,
1904 emitter: &emitter,
1905 account_id: AccountId::from("POLY-001"),
1906 clock: nautilus_core::time::get_atomic_clock_realtime(),
1907 user_address: "0xtest",
1908 user_api_key: "test-key",
1909 };
1910 let mut state = WsDispatchState::default();
1911 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
1912
1913 let first_result =
1914 dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
1915 token_instruments.insert(trade.asset_id, instrument);
1916 let replay_result = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
1917
1918 assert!(first_result.is_none());
1919 assert!(replay_result.is_some());
1920 assert_eq!(fill_tracker.pending_fills_for(&venue_order_id).len(), 1);
1921 }
1922
1923 #[rstest]
1924 #[case(crate::common::enums::PolymarketTradeStatus::Mined)]
1925 #[case(crate::common::enums::PolymarketTradeStatus::Retrying)]
1926 fn test_dispatch_trade_ignores_pending_settlement_status(
1927 #[case] status: crate::common::enums::PolymarketTradeStatus,
1928 ) {
1929 let mut trade: PolymarketUserTrade = load("ws_user_trade.json");
1930 trade.status = status;
1931 let instrument = test_instrument();
1932 let token_instruments = AtomicMap::new();
1933 token_instruments.insert(trade.asset_id, instrument);
1934 let fill_tracker = OrderFillTrackerMap::new();
1935 let pending_submits = PendingSubmitTracker::default();
1936 let order_identities = OrderIdentityRegistry::default();
1937 let emitter = test_emitter();
1938 let ctx = WsDispatchContext {
1939 token_instruments: &token_instruments,
1940 fill_tracker: &fill_tracker,
1941 pending_submits: &pending_submits,
1942 order_identities: &order_identities,
1943 emitter: &emitter,
1944 account_id: AccountId::from("POLY-001"),
1945 clock: nautilus_core::time::get_atomic_clock_realtime(),
1946 user_address: "0xtest",
1947 user_api_key: "test-key",
1948 };
1949 let mut state = WsDispatchState::default();
1950 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
1951
1952 let result = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
1953
1954 assert!(result.is_none());
1955 assert!(fill_tracker.pending_fills_for(&venue_order_id).is_empty());
1956 }
1957
1958 #[rstest]
1959 fn test_dispatch_matched_trade_emits_fill_and_failed_trade_voids_it() {
1960 let mut trade: PolymarketUserTrade = load("ws_user_trade.json");
1961 trade.status = crate::common::enums::PolymarketTradeStatus::Matched;
1962 let instrument = test_instrument();
1963 let token_instruments = AtomicMap::new();
1964 token_instruments.insert(trade.asset_id, instrument.clone());
1965 let fill_tracker = OrderFillTrackerMap::new();
1966 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
1967 fill_tracker.register(
1968 venue_order_id,
1969 Quantity::from("100"),
1970 OrderSide::Buy,
1971 instrument.id(),
1972 instrument.size_precision(),
1973 instrument.price_precision(),
1974 );
1975 let pending_submits = PendingSubmitTracker::default();
1976 let order_identities = OrderIdentityRegistry::default();
1977 register_identity(
1978 &order_identities,
1979 venue_order_id,
1980 instrument.id(),
1981 "O-MATCHED-FAILED",
1982 );
1983 order_identities.mark_accepted(venue_order_id);
1984 let mut emitter = test_emitter();
1985 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
1986 emitter.set_sender(sender);
1987 let ctx = WsDispatchContext {
1988 token_instruments: &token_instruments,
1989 fill_tracker: &fill_tracker,
1990 pending_submits: &pending_submits,
1991 order_identities: &order_identities,
1992 emitter: &emitter,
1993 account_id: AccountId::from("POLY-001"),
1994 clock: nautilus_core::time::get_atomic_clock_realtime(),
1995 user_address: "0xtest",
1996 user_api_key: "test-key",
1997 };
1998 let mut state = WsDispatchState::default();
1999
2000 let matched = dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
2001 let filled = match receiver.try_recv().unwrap() {
2002 ExecutionEvent::Order(OrderEventAny::Filled(event)) => event,
2003 other => panic!("expected matched fill, was {other:?}"),
2004 };
2005 trade.status = crate::common::enums::PolymarketTradeStatus::Failed;
2006 let failed = dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
2007 let voided = match receiver.try_recv().unwrap() {
2008 ExecutionEvent::Order(OrderEventAny::FillVoided(event)) => event,
2009 other => panic!("expected failed fill correction, was {other:?}"),
2010 };
2011
2012 let mut failed_first_state = WsDispatchState::default();
2013 let failed_first = dispatch_user_message(
2014 &UserWsMessage::Trade(trade.clone()),
2015 &ctx,
2016 &mut failed_first_state,
2017 );
2018 let dedup_key = format!("{}-{}", trade.id, trade.taker_order_id);
2019 trade.status = crate::common::enums::PolymarketTradeStatus::Matched;
2020 let matched_after_failure =
2021 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut failed_first_state);
2022
2023 assert!(matched.is_none());
2024 assert!(failed.is_some());
2025 assert!(failed_first.is_some());
2026 assert!(matched_after_failure.is_none());
2027 assert_eq!(voided.trade_id, filled.trade_id);
2028 assert_eq!(voided.voided_qty, filled.last_qty);
2029 assert_eq!(voided.commission_voided, filled.commission);
2030 assert_eq!(voided.last_px, filled.last_px);
2031 assert!(!voided.is_reopened);
2032 assert_eq!(voided.causation_id, Some(filled.event_id));
2033 assert_eq!(
2034 fill_tracker.get_cumulative_filled(&venue_order_id),
2035 Some(Quantity::zero(instrument.size_precision()))
2036 );
2037 assert!(failed_first_state.processed_fills.contains(&dedup_key));
2038 assert!(failed_first_state.is_voided_trade(&dedup_key));
2039 assert!(receiver.try_recv().is_err());
2040 }
2041
2042 #[rstest]
2043 fn test_dispatch_trade_uses_pending_submit_client_order_id() {
2044 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2045 let instrument = test_instrument();
2046
2047 let token_instruments = AtomicMap::new();
2048 token_instruments.insert(trade.asset_id, instrument);
2049
2050 let fill_tracker = OrderFillTrackerMap::new();
2051 let pending_submits = PendingSubmitTracker::default();
2052 let order_identities = OrderIdentityRegistry::default();
2053 let emitter = test_emitter();
2054
2055 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
2056 let client_order_id = ClientOrderId::from("O-UNKNOWN-FILL");
2057 pending_submits.insert(venue_order_id, client_order_id);
2058
2059 let ctx = WsDispatchContext {
2060 token_instruments: &token_instruments,
2061 fill_tracker: &fill_tracker,
2062 pending_submits: &pending_submits,
2063 order_identities: &order_identities,
2064 emitter: &emitter,
2065 account_id: AccountId::from("POLY-001"),
2066 clock: nautilus_core::time::get_atomic_clock_realtime(),
2067 user_address: "0xtest",
2068 user_api_key: "test-key",
2069 };
2070 let mut state = WsDispatchState::default();
2071
2072 let _ = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2073
2074 let fills = fill_tracker.pending_fills_for(&venue_order_id);
2075 assert_eq!(fills[0].client_order_id, Some(client_order_id));
2076 }
2077
2078 #[rstest]
2079 fn test_dispatch_late_fill_stays_tracked_after_later_registrations() {
2080 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2081 let market: GammaMarket = load("gamma_market_sports_market_money_line.json");
2082 let defs = parse_gamma_market(&market).unwrap();
2083 let instrument =
2084 create_instrument_from_def(&defs[0], UnixNanos::from(1_000_000_000u64)).unwrap();
2085 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
2086
2087 let token_instruments = AtomicMap::new();
2088 token_instruments.insert(trade.asset_id, instrument.clone());
2089
2090 let fill_tracker = OrderFillTrackerMap::new();
2091 fill_tracker.register(
2092 venue_order_id,
2093 Quantity::from("100"),
2094 OrderSide::Buy,
2095 instrument.id(),
2096 instrument.size_precision(),
2097 instrument.price_precision(),
2098 );
2099
2100 let pending_submits = PendingSubmitTracker::default();
2101 let order_identities = OrderIdentityRegistry::default();
2102 register_identity(
2103 &order_identities,
2104 venue_order_id,
2105 instrument.id(),
2106 "O-LATE-FILL",
2107 );
2108 order_identities.mark_accepted(venue_order_id);
2109 assert!(order_identities.get(&venue_order_id).is_some());
2110
2111 for index in 0..10_000 {
2112 let later_venue_order_id = VenueOrderId::from(format!("V-LATER-{index}").as_str());
2113 let later_client_order_id = format!("O-LATER-{index}");
2114 register_identity(
2115 &order_identities,
2116 later_venue_order_id,
2117 instrument.id(),
2118 &later_client_order_id,
2119 );
2120 order_identities.mark_accepted(later_venue_order_id);
2121 fill_tracker.register(
2122 later_venue_order_id,
2123 Quantity::from("1"),
2124 OrderSide::Sell,
2125 instrument.id(),
2126 instrument.size_precision(),
2127 instrument.price_precision(),
2128 );
2129 }
2130 assert!(order_identities.get(&venue_order_id).is_some());
2131 assert!(fill_tracker.contains(&venue_order_id));
2132
2133 let mut emitter = test_emitter();
2134 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2135 emitter.set_sender(sender);
2136
2137 let ctx = WsDispatchContext {
2138 token_instruments: &token_instruments,
2139 fill_tracker: &fill_tracker,
2140 pending_submits: &pending_submits,
2141 order_identities: &order_identities,
2142 emitter: &emitter,
2143 account_id: AccountId::from("POLY-001"),
2144 clock: nautilus_core::time::get_atomic_clock_realtime(),
2145 user_address: "0xtest",
2146 user_api_key: "test-key",
2147 };
2148 let mut state = WsDispatchState::default();
2149
2150 dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
2151
2152 let event = receiver.try_recv().expect("expected tracked late fill");
2153 let ExecutionEvent::Order(OrderEventAny::Filled(filled)) = event else {
2154 panic!("expected tracked OrderFilled after later registrations, was {event:?}");
2155 };
2156
2157 assert_eq!(filled.client_order_id, ClientOrderId::from("O-LATE-FILL"));
2158 assert_eq!(filled.venue_order_id, venue_order_id);
2159 assert_eq!(filled.trade_id, TradeId::from(trade.id.as_str()));
2160 assert_eq!(filled.instrument_id, instrument.id());
2161 assert_eq!(
2162 filled.last_qty.as_decimal(),
2163 Decimal::from_str_exact(&trade.size).unwrap()
2164 );
2165 assert_eq!(
2166 filled.last_px.as_decimal(),
2167 Decimal::from_str_exact(&trade.price).unwrap()
2168 );
2169 assert_eq!(filled.order_side, OrderSide::Buy);
2170 assert_eq!(filled.liquidity_side, LiquiditySide::Taker);
2171 let commission = filled.commission.expect("tracked fill has commission");
2172 assert_eq!(commission.as_decimal(), dec!(0.1875));
2173 assert_eq!(commission.currency, Currency::pUSD());
2174 assert!(receiver.try_recv().is_err());
2175 }
2176
2177 #[rstest]
2178 fn test_dispatch_order_matched_caps_filled_qty_when_no_trades_tracked() {
2179 let order: PolymarketUserOrder = load("ws_user_order_matched.json");
2180 let instrument = test_instrument();
2181
2182 let token_instruments = AtomicMap::new();
2183 token_instruments.insert(order.asset_id, instrument.clone());
2184
2185 let fill_tracker = OrderFillTrackerMap::new();
2186 let venue_order_id = VenueOrderId::from(order.id.as_str());
2187
2188 fill_tracker.register(
2190 venue_order_id,
2191 Quantity::from("100"),
2192 OrderSide::Buy,
2193 instrument.id(),
2194 instrument.size_precision(),
2195 instrument.price_precision(),
2196 );
2197
2198 let pending_submits = PendingSubmitTracker::default();
2199 let order_identities = OrderIdentityRegistry::default();
2202 let mut emitter = test_emitter();
2203 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2204 emitter.set_sender(sender);
2205
2206 let ctx = WsDispatchContext {
2207 token_instruments: &token_instruments,
2208 fill_tracker: &fill_tracker,
2209 pending_submits: &pending_submits,
2210 order_identities: &order_identities,
2211 emitter: &emitter,
2212 account_id: AccountId::from("POLY-001"),
2213 clock: nautilus_core::time::get_atomic_clock_realtime(),
2214 user_address: "0xtest",
2215 user_api_key: "test-key",
2216 };
2217 let mut state = WsDispatchState::default();
2218
2219 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
2220
2221 let event = receiver.try_recv().expect("Expected report");
2222 match event {
2223 ExecutionEvent::Report(report) => match report {
2224 ExecutionReport::Order(order_report) => {
2225 assert_eq!(order_report.filled_qty, Quantity::from("0"));
2226 }
2227 other => panic!("Expected order report, was {other:?}"),
2228 },
2229 other => panic!("Expected report event, was {other:?}"),
2230 }
2231 }
2232
2233 #[rstest]
2234 fn test_dispatch_order_matched_uses_tracked_fills_for_filled_qty() {
2235 let order: PolymarketUserOrder = load("ws_user_order_matched.json");
2236 let instrument = test_instrument();
2237
2238 let token_instruments = AtomicMap::new();
2239 token_instruments.insert(order.asset_id, instrument.clone());
2240
2241 let fill_tracker = OrderFillTrackerMap::new();
2242 let venue_order_id = VenueOrderId::from(order.id.as_str());
2243
2244 fill_tracker.register(
2246 venue_order_id,
2247 Quantity::from("100"),
2248 OrderSide::Buy,
2249 instrument.id(),
2250 instrument.size_precision(),
2251 instrument.price_precision(),
2252 );
2253 fill_tracker.record_fill(&venue_order_id, Quantity::new(50.0, 6));
2254
2255 let pending_submits = PendingSubmitTracker::default();
2256 let order_identities = OrderIdentityRegistry::default();
2259 let mut emitter = test_emitter();
2260 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2261 emitter.set_sender(sender);
2262
2263 let ctx = WsDispatchContext {
2264 token_instruments: &token_instruments,
2265 fill_tracker: &fill_tracker,
2266 pending_submits: &pending_submits,
2267 order_identities: &order_identities,
2268 emitter: &emitter,
2269 account_id: AccountId::from("POLY-001"),
2270 clock: nautilus_core::time::get_atomic_clock_realtime(),
2271 user_address: "0xtest",
2272 user_api_key: "test-key",
2273 };
2274 let mut state = WsDispatchState::default();
2275
2276 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
2277
2278 let event = receiver.try_recv().expect("Expected report");
2279 match event {
2280 ExecutionEvent::Report(report) => match report {
2281 ExecutionReport::Order(order_report) => {
2282 assert_eq!(order_report.filled_qty, Quantity::from("50"));
2283 }
2284 other => panic!("Expected order report, was {other:?}"),
2285 },
2286 other => panic!("Expected report event, was {other:?}"),
2287 }
2288 }
2289
2290 #[rstest]
2291 fn test_dispatch_order_matched_normalizes_quantity_without_fill() {
2292 let order: PolymarketUserOrder = load("ws_user_order_matched.json");
2293 let instrument = test_instrument();
2294
2295 let token_instruments = AtomicMap::new();
2296 token_instruments.insert(order.asset_id, instrument.clone());
2297
2298 let fill_tracker = OrderFillTrackerMap::new();
2299 let venue_order_id = VenueOrderId::from(order.id.as_str());
2300 fill_tracker.register(
2301 venue_order_id,
2302 Quantity::from("100"),
2303 OrderSide::Buy,
2304 instrument.id(),
2305 instrument.size_precision(),
2306 instrument.price_precision(),
2307 );
2308 fill_tracker.record_fill(&venue_order_id, Quantity::new(99.995, 6));
2309
2310 let pending_submits = PendingSubmitTracker::default();
2311 let order_identities = OrderIdentityRegistry::default();
2312 register_identity(
2313 &order_identities,
2314 venue_order_id,
2315 instrument.id(),
2316 "O-MATCHED",
2317 );
2318 order_identities.mark_accepted(venue_order_id);
2319 let mut emitter = test_emitter();
2320 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2321 emitter.set_sender(sender);
2322
2323 let clock = Box::leak(Box::new(AtomicTime::new(
2324 false,
2325 UnixNanos::from(2_000_000_000u64),
2326 )));
2327
2328 let ctx = WsDispatchContext {
2329 token_instruments: &token_instruments,
2330 fill_tracker: &fill_tracker,
2331 pending_submits: &pending_submits,
2332 order_identities: &order_identities,
2333 emitter: &emitter,
2334 account_id: AccountId::from("POLY-001"),
2335 clock,
2336 user_address: "0xtest",
2337 user_api_key: "test-key",
2338 };
2339 let mut state = WsDispatchState::default();
2340 state.confirmed_trades.add("trade-0xfill1".to_string());
2341
2342 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
2343
2344 let event = receiver.try_recv().expect("expected quantity update");
2345 match event {
2346 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
2347 assert_eq!(
2348 updated.ts_event,
2349 UnixNanos::from(1_703_875_201_000_000_000u64)
2350 );
2351 assert_eq!(updated.ts_init, UnixNanos::from(2_000_000_000u64));
2352 assert_eq!(updated.quantity, Quantity::new(99.995, 6));
2353 assert!(updated.reconciliation);
2354 }
2355 other => panic!("expected updated event, was {other:?}"),
2356 }
2357 assert!(receiver.try_recv().is_err());
2358 }
2359
2360 #[rstest]
2361 fn test_confirmed_trade_normalizes_pending_matched_quantity() {
2362 let mut order: PolymarketUserOrder = load("ws_user_order_matched.json");
2363 let mut trade: PolymarketUserTrade = load("ws_user_trade.json");
2364 let instrument = test_instrument();
2365 order.associate_trades = Some(vec![trade.id.clone()]);
2366 trade.size = "99.995".to_string();
2367 trade.price = order.price.clone();
2368
2369 let token_instruments = AtomicMap::new();
2370 token_instruments.insert(order.asset_id, instrument.clone());
2371 let fill_tracker = OrderFillTrackerMap::new();
2372 let venue_order_id = VenueOrderId::from(order.id.as_str());
2373 fill_tracker.register(
2374 venue_order_id,
2375 Quantity::from("100"),
2376 OrderSide::Buy,
2377 instrument.id(),
2378 instrument.size_precision(),
2379 instrument.price_precision(),
2380 );
2381 let pending_submits = PendingSubmitTracker::default();
2382 let order_identities = OrderIdentityRegistry::default();
2383 register_identity(
2384 &order_identities,
2385 venue_order_id,
2386 instrument.id(),
2387 "O-CONFIRMED-DUST",
2388 );
2389 order_identities.mark_accepted(venue_order_id);
2390 let mut emitter = test_emitter();
2391 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2392 emitter.set_sender(sender);
2393 let ctx = WsDispatchContext {
2394 token_instruments: &token_instruments,
2395 fill_tracker: &fill_tracker,
2396 pending_submits: &pending_submits,
2397 order_identities: &order_identities,
2398 emitter: &emitter,
2399 account_id: AccountId::from("POLY-001"),
2400 clock: nautilus_core::time::get_atomic_clock_realtime(),
2401 user_address: "0xtest",
2402 user_api_key: "test-key",
2403 };
2404 let mut state = WsDispatchState::default();
2405
2406 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
2407 assert!(receiver.try_recv().is_err());
2408
2409 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2410
2411 let real_fill = receiver.try_recv().expect("expected confirmed venue fill");
2412 let normalized = receiver
2413 .try_recv()
2414 .expect("expected quantity normalization");
2415
2416 match (real_fill, normalized) {
2417 (
2418 ExecutionEvent::Order(OrderEventAny::Filled(real)),
2419 ExecutionEvent::Order(OrderEventAny::Updated(updated)),
2420 ) => {
2421 assert_eq!(real.last_qty, Quantity::from("99.995"));
2422 assert_eq!(updated.quantity, Quantity::from("99.995"));
2423 assert!(updated.reconciliation);
2424 }
2425 other => panic!("expected fill then quantity update, was {other:?}"),
2426 }
2427 assert!(receiver.try_recv().is_err());
2428 }
2429
2430 #[rstest]
2431 fn test_cancel_reemitted_after_fill_for_canceled_order() {
2432 let cancel_order: PolymarketUserOrder = load("ws_user_order_cancellation.json");
2433 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2434 let instrument = test_instrument();
2435
2436 let token_instruments = AtomicMap::new();
2437 token_instruments.insert(cancel_order.asset_id, instrument.clone());
2438
2439 let fill_tracker = OrderFillTrackerMap::new();
2440 let venue_order_id = VenueOrderId::from(cancel_order.id.as_str());
2441
2442 fill_tracker.register(
2444 venue_order_id,
2445 Quantity::from("100"),
2446 OrderSide::Buy,
2447 instrument.id(),
2448 instrument.size_precision(),
2449 instrument.price_precision(),
2450 );
2451
2452 let pending_submits = PendingSubmitTracker::default();
2453 let order_identities = OrderIdentityRegistry::default();
2454 register_identity(
2455 &order_identities,
2456 venue_order_id,
2457 instrument.id(),
2458 "O-CANCEL",
2459 );
2460 order_identities.mark_accepted(venue_order_id);
2461 let mut emitter = test_emitter();
2462 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2463 emitter.set_sender(sender);
2464
2465 let ctx = WsDispatchContext {
2466 token_instruments: &token_instruments,
2467 fill_tracker: &fill_tracker,
2468 pending_submits: &pending_submits,
2469 order_identities: &order_identities,
2470 emitter: &emitter,
2471 account_id: AccountId::from("POLY-001"),
2472 clock: nautilus_core::time::get_atomic_clock_realtime(),
2473 user_address: "0xtest",
2474 user_api_key: "test-key",
2475 };
2476 let mut state = WsDispatchState::default();
2477
2478 dispatch_user_message(&UserWsMessage::Order(cancel_order), &ctx, &mut state);
2480 let cancel_event = receiver.try_recv().expect("Expected canceled event");
2481 match &cancel_event {
2482 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
2483 assert_eq!(c.venue_order_id, Some(venue_order_id));
2484 }
2485 other => panic!("Expected canceled event, was {other:?}"),
2486 }
2487
2488 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2490
2491 let fill_event = receiver.try_recv().expect("Expected filled event");
2493 match &fill_event {
2494 ExecutionEvent::Order(OrderEventAny::Filled(f)) => {
2495 assert_eq!(f.venue_order_id, venue_order_id);
2496 }
2497 other => panic!("Expected filled event, was {other:?}"),
2498 }
2499
2500 let reemitted_cancel = receiver
2501 .try_recv()
2502 .expect("Expected re-emitted canceled event");
2503
2504 match &reemitted_cancel {
2505 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
2506 assert_eq!(c.venue_order_id, Some(venue_order_id));
2507 }
2508 other => panic!("Expected canceled event, was {other:?}"),
2509 }
2510 }
2511
2512 #[rstest]
2513 fn test_cancel_not_reemitted_when_fill_completes_order() {
2514 let cancel_order: PolymarketUserOrder = load("ws_user_order_cancellation.json");
2515 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2516 let instrument = test_instrument();
2517
2518 let token_instruments = AtomicMap::new();
2519 token_instruments.insert(cancel_order.asset_id, instrument.clone());
2520
2521 let fill_tracker = OrderFillTrackerMap::new();
2522 let venue_order_id = VenueOrderId::from(cancel_order.id.as_str());
2523
2524 fill_tracker.register(
2526 venue_order_id,
2527 Quantity::from("25"),
2528 OrderSide::Buy,
2529 instrument.id(),
2530 instrument.size_precision(),
2531 instrument.price_precision(),
2532 );
2533
2534 let pending_submits = PendingSubmitTracker::default();
2535 let order_identities = OrderIdentityRegistry::default();
2536 register_identity(
2537 &order_identities,
2538 venue_order_id,
2539 instrument.id(),
2540 "O-CANCEL-FULL",
2541 );
2542 order_identities.mark_accepted(venue_order_id);
2543 let mut emitter = test_emitter();
2544 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2545 emitter.set_sender(sender);
2546
2547 let ctx = WsDispatchContext {
2548 token_instruments: &token_instruments,
2549 fill_tracker: &fill_tracker,
2550 pending_submits: &pending_submits,
2551 order_identities: &order_identities,
2552 emitter: &emitter,
2553 account_id: AccountId::from("POLY-001"),
2554 clock: nautilus_core::time::get_atomic_clock_realtime(),
2555 user_address: "0xtest",
2556 user_api_key: "test-key",
2557 };
2558 let mut state = WsDispatchState::default();
2559
2560 dispatch_user_message(&UserWsMessage::Order(cancel_order), &ctx, &mut state);
2562 let _cancel = receiver.try_recv().expect("Expected canceled event");
2563
2564 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2565 let _fill = receiver.try_recv().expect("Expected filled event");
2566
2567 assert!(
2569 receiver.try_recv().is_err(),
2570 "Should not re-emit cancel when fill completes the order"
2571 );
2572 }
2573
2574 #[rstest]
2575 fn test_cancel_saved_before_acceptance() {
2576 let cancel_order: PolymarketUserOrder = load("ws_user_order_cancellation.json");
2577 let instrument = test_instrument();
2578
2579 let token_instruments = AtomicMap::new();
2580 token_instruments.insert(cancel_order.asset_id, instrument);
2581
2582 let fill_tracker = OrderFillTrackerMap::new();
2584 let venue_order_id = VenueOrderId::from(cancel_order.id.as_str());
2585
2586 let pending_submits = PendingSubmitTracker::default();
2587 let order_identities = OrderIdentityRegistry::default();
2588 let emitter = test_emitter();
2589
2590 let ctx = WsDispatchContext {
2591 token_instruments: &token_instruments,
2592 fill_tracker: &fill_tracker,
2593 pending_submits: &pending_submits,
2594 order_identities: &order_identities,
2595 emitter: &emitter,
2596 account_id: AccountId::from("POLY-001"),
2597 clock: nautilus_core::time::get_atomic_clock_realtime(),
2598 user_address: "0xtest",
2599 user_api_key: "test-key",
2600 };
2601 let mut state = WsDispatchState::default();
2602
2603 dispatch_user_message(&UserWsMessage::Order(cancel_order), &ctx, &mut state);
2605
2606 assert!(fill_tracker.has_pending_report(&venue_order_id));
2608 assert!(state.terminal_cancel_reports.get(&venue_order_id).is_some());
2609 }
2610
2611 #[rstest]
2614 #[case(PolymarketOrderStatus::Canceled, "Canceled", OrderStatus::Canceled)]
2615 #[case(
2616 PolymarketOrderStatus::CanceledMarketResolved,
2617 "Expired",
2618 OrderStatus::Expired
2619 )]
2620 fn test_buffered_fill_emitted_before_terminal_status(
2621 #[case] status: PolymarketOrderStatus,
2622 #[case] expected_terminal: &str,
2623 #[case] expected_order_status: OrderStatus,
2624 ) {
2625 let mut terminal_order: PolymarketUserOrder = load("ws_user_order_cancellation.json");
2626 terminal_order.status = Some(status.into());
2627 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2628 let instrument = test_instrument();
2629
2630 let token_instruments = AtomicMap::new();
2631 token_instruments.insert(terminal_order.asset_id, instrument.clone());
2632
2633 let fill_tracker = OrderFillTrackerMap::new();
2635 let venue_order_id = VenueOrderId::from(terminal_order.id.as_str());
2636 let client_order_id = ClientOrderId::from("O-BUFFERED");
2637
2638 let pending_submits = PendingSubmitTracker::default();
2639 pending_submits.insert(venue_order_id, client_order_id);
2640 let order_identities = OrderIdentityRegistry::default();
2641 register_identity(
2642 &order_identities,
2643 venue_order_id,
2644 instrument.id(),
2645 client_order_id.as_str(),
2646 );
2647 let mut emitter = test_emitter();
2648 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2649 emitter.set_sender(sender);
2650
2651 let ctx = WsDispatchContext {
2652 token_instruments: &token_instruments,
2653 fill_tracker: &fill_tracker,
2654 pending_submits: &pending_submits,
2655 order_identities: &order_identities,
2656 emitter: &emitter,
2657 account_id: AccountId::from("POLY-001"),
2658 clock: nautilus_core::time::get_atomic_clock_realtime(),
2659 user_address: "0xtest",
2660 user_api_key: "test-key",
2661 };
2662 let mut state = WsDispatchState::default();
2663
2664 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2666 assert!(
2667 receiver.try_recv().is_err(),
2668 "a buffered fill must emit no event before the order is registered",
2669 );
2670
2671 dispatch_user_message(&UserWsMessage::Order(terminal_order), &ctx, &mut state);
2673
2674 let mut emitted = Vec::new();
2675
2676 while let Ok(event) = receiver.try_recv() {
2677 match event {
2678 ExecutionEvent::Order(order_event) => emitted.push(order_event),
2679 other => panic!("expected only order events, was {other:?}"),
2680 }
2681 }
2682
2683 assert_eq!(emitted.len(), 3, "emitted sequence was {emitted:?}");
2684 match &emitted[0] {
2685 OrderEventAny::Accepted(accepted) => {
2686 assert_eq!(accepted.client_order_id, client_order_id);
2687 assert_eq!(accepted.venue_order_id, venue_order_id);
2688 }
2689 other => panic!("expected accepted event first, was {other:?}"),
2690 }
2691
2692 match &emitted[1] {
2693 OrderEventAny::Filled(filled) => {
2694 assert_eq!(filled.client_order_id, client_order_id);
2695 assert_eq!(filled.venue_order_id, venue_order_id);
2696 assert_eq!(filled.last_qty.as_decimal(), dec!(25));
2697 }
2698 other => panic!("expected filled event before the terminal status, was {other:?}"),
2699 }
2700 let terminal = match &emitted[2] {
2701 OrderEventAny::Canceled(canceled) => {
2702 assert_eq!(canceled.client_order_id, client_order_id);
2703 assert_eq!(canceled.venue_order_id, Some(venue_order_id));
2704 "Canceled"
2705 }
2706 OrderEventAny::Expired(expired) => {
2707 assert_eq!(expired.client_order_id, client_order_id);
2708 assert_eq!(expired.venue_order_id, Some(venue_order_id));
2709 "Expired"
2710 }
2711 other => panic!("expected a terminal order event last, was {other:?}"),
2712 };
2713 assert_eq!(terminal, expected_terminal);
2714
2715 let mut order = OrderTestBuilder::new(OrderType::Limit)
2717 .instrument_id(instrument.id())
2718 .client_order_id(client_order_id)
2719 .strategy_id(StrategyId::from("S-001"))
2720 .side(OrderSide::Buy)
2721 .price(Price::from("0.5"))
2722 .quantity(Quantity::from("100"))
2723 .build();
2724
2725 for event in emitted {
2726 order.apply(event).expect("emitted sequence must be valid");
2727 }
2728
2729 assert_eq!(order.status(), expected_order_status);
2730 assert_eq!(order.filled_qty().as_decimal(), dec!(25));
2731 }
2732
2733 #[rstest]
2745 fn test_issue_3797_interleaved_cancel_fill_sequence() {
2746 use crate::common::{
2747 enums::{
2748 PolymarketEventType, PolymarketLiquiditySide, PolymarketOrderSide,
2749 PolymarketOrderStatus, PolymarketOrderType, PolymarketOutcome,
2750 PolymarketTradeStatus,
2751 },
2752 models::PolymarketMakerOrder,
2753 };
2754
2755 let instrument = test_instrument();
2756 let asset_id = instrument.id().symbol.inner();
2757
2758 let order_id =
2759 "0xe743f6c823ecdfa9ddaaf08673b2441d15a38d89e14dcb25b3b70c284be4f6ad".to_string();
2760 let venue_order_id = VenueOrderId::from(order_id.as_str());
2761
2762 let token_instruments = AtomicMap::new();
2763 token_instruments.insert(asset_id, instrument.clone());
2764
2765 let fill_tracker = OrderFillTrackerMap::new();
2766 fill_tracker.register(
2767 venue_order_id,
2768 Quantity::from("20"),
2769 OrderSide::Buy,
2770 instrument.id(),
2771 instrument.size_precision(),
2772 instrument.price_precision(),
2773 );
2774
2775 let pending_submits = PendingSubmitTracker::default();
2776 let order_identities = OrderIdentityRegistry::default();
2777 register_identity(&order_identities, venue_order_id, instrument.id(), "O-3797");
2778 order_identities.mark_accepted(venue_order_id);
2779 let mut emitter = test_emitter();
2780 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2781 emitter.set_sender(sender);
2782
2783 let ctx = WsDispatchContext {
2784 token_instruments: &token_instruments,
2785 fill_tracker: &fill_tracker,
2786 pending_submits: &pending_submits,
2787 order_identities: &order_identities,
2788 emitter: &emitter,
2789 account_id: AccountId::from("POLY-001"),
2790 clock: nautilus_core::time::get_atomic_clock_realtime(),
2791 user_address: "0xabc",
2792 user_api_key: "xxx",
2793 };
2794 let mut state = WsDispatchState::default();
2795
2796 let make_order =
2798 |size_matched: &str, ts: &str, event_type: PolymarketEventType| PolymarketUserOrder {
2799 asset_id,
2800 associate_trades: None,
2801 created_at: Some("1775074735".to_string()),
2802 expiration: Some("0".to_string()),
2803 id: order_id.clone(),
2804 maker_address: Some(Ustr::from("0xabc")),
2805 market: Ustr::from("0x4134"),
2806 order_owner: Some(Ustr::from("xxx")),
2807 order_type: Some(PolymarketOrderType::GTC),
2808 original_size: "20".to_string(),
2809 outcome: Some(PolymarketOutcome::yes()),
2810 owner: Ustr::from("xxx"),
2811 price: "0.18".to_string(),
2812 side: PolymarketOrderSide::Buy,
2813 size_matched: size_matched.to_string(),
2814 status: Some(PolymarketOrderStatus::Canceled.into()),
2815 timestamp: ts.to_string(),
2816 event_type,
2817 };
2818
2819 let make_trade = |trade_id: &str, matched_amount: f64, ts: &str| PolymarketUserTrade {
2821 asset_id,
2822 bucket_index: 0,
2823 fee_rate_bps: "1000".to_string(),
2824 id: trade_id.to_string(),
2825 last_update: "1775074738".to_string(),
2826 maker_address: Ustr::from("0xother"),
2827 maker_orders: vec![PolymarketMakerOrder {
2828 asset_id,
2829 maker_address: "0xabc".to_string(),
2830 matched_amount: Decimal::from_f64_retain(matched_amount).unwrap_or(Decimal::ZERO),
2831 order_id: order_id.clone(),
2832 outcome: PolymarketOutcome::yes(),
2833 owner: "xxx".to_string(),
2834 price: Decimal::from_f64_retain(0.18).unwrap_or(Decimal::ZERO),
2835 side: None,
2836 }],
2837 market: Ustr::from("0x4134"),
2838 match_time: "1775074735".to_string(),
2839 outcome: PolymarketOutcome::yes(),
2840 owner: Ustr::from("other-owner"),
2841 price: "0.82".to_string(),
2842 side: PolymarketOrderSide::Buy,
2843 size: "1.219511".to_string(),
2844 status: PolymarketTradeStatus::Confirmed,
2845 taker_order_id: "0xtaker01".to_string(),
2846 timestamp: ts.to_string(),
2847 trade_owner: Ustr::from("other-owner"),
2848 transaction_hash: None,
2849 trader_side: PolymarketLiquiditySide::Maker,
2850 event_type: PolymarketEventType::Trade,
2851 };
2852
2853 let msg_a = make_order("0", "1775074738031", PolymarketEventType::Cancellation);
2855 dispatch_user_message(&UserWsMessage::Order(msg_a), &ctx, &mut state);
2856
2857 let evt = receiver.try_recv().expect("(A) canceled event");
2858 match &evt {
2859 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
2860 assert_eq!(c.venue_order_id, Some(venue_order_id));
2861 }
2862 other => panic!("(A) expected canceled event, was {other:?}"),
2863 }
2864
2865 let msg_b = make_trade("trade-b", 1.219511, "1775074738032");
2867 dispatch_user_message(&UserWsMessage::Trade(msg_b), &ctx, &mut state);
2868
2869 let evt = receiver.try_recv().expect("(B) filled event");
2870 match &evt {
2871 ExecutionEvent::Order(OrderEventAny::Filled(f)) => {
2872 assert_eq!(f.venue_order_id, venue_order_id);
2873 }
2874 other => panic!("(B) expected filled event, was {other:?}"),
2875 }
2876 let evt = receiver.try_recv().expect("(B) re-emitted cancel");
2878 match &evt {
2879 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
2880 assert_eq!(c.venue_order_id, Some(venue_order_id));
2881 }
2882 other => panic!("(B) expected re-emitted cancel, was {other:?}"),
2883 }
2884
2885 let msg_c = make_order("1.219511", "1775074738034", PolymarketEventType::Update);
2887 dispatch_user_message(&UserWsMessage::Order(msg_c), &ctx, &mut state);
2888
2889 let evt = receiver.try_recv().expect("(C) canceled event");
2890 match &evt {
2891 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
2892 assert_eq!(c.venue_order_id, Some(venue_order_id));
2893 }
2894 other => panic!("(C) expected canceled event, was {other:?}"),
2895 }
2896
2897 let msg_d = make_order("2.560972", "1775074738038", PolymarketEventType::Update);
2899 dispatch_user_message(&UserWsMessage::Order(msg_d), &ctx, &mut state);
2900
2901 let evt = receiver.try_recv().expect("(D) canceled event");
2902 match &evt {
2903 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
2904 assert_eq!(c.venue_order_id, Some(venue_order_id));
2905 }
2906 other => panic!("(D) expected canceled event, was {other:?}"),
2907 }
2908
2909 let msg_e = make_trade("trade-e", 1.341461, "1775074738036");
2911 dispatch_user_message(&UserWsMessage::Trade(msg_e), &ctx, &mut state);
2912
2913 let evt = receiver.try_recv().expect("(E) filled event");
2914 match &evt {
2915 ExecutionEvent::Order(OrderEventAny::Filled(f)) => {
2916 assert_eq!(f.venue_order_id, venue_order_id);
2917 }
2918 other => panic!("(E) expected filled event, was {other:?}"),
2919 }
2920
2921 let evt = receiver.try_recv().expect("(E) re-emitted cancel");
2923 match &evt {
2924 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
2925 assert_eq!(c.venue_order_id, Some(venue_order_id));
2926 }
2927 other => panic!("(E) expected re-emitted cancel, was {other:?}"),
2928 }
2929
2930 assert!(
2932 receiver.try_recv().is_err(),
2933 "No further events expected after the sequence"
2934 );
2935 }
2936
2937 #[rstest]
2938 fn test_dispatch_taker_fill_snaps_overfill_to_submitted_qty() {
2939 use crate::common::enums::{
2944 PolymarketEventType, PolymarketOrderSide, PolymarketOutcome, PolymarketTradeStatus,
2945 };
2946
2947 let instrument = test_instrument();
2948 let asset_id = instrument.id().symbol.inner();
2949 let token_instruments = AtomicMap::new();
2950 token_instruments.insert(asset_id, instrument.clone());
2951
2952 let fill_tracker = OrderFillTrackerMap::new();
2953 let venue_order_id = VenueOrderId::from("0xtaker-overfill");
2954 let submitted = Quantity::new(714.285710, instrument.size_precision());
2956 fill_tracker.register(
2957 venue_order_id,
2958 submitted,
2959 OrderSide::Buy,
2960 instrument.id(),
2961 instrument.size_precision(),
2962 instrument.price_precision(),
2963 );
2964
2965 let pending_submits = PendingSubmitTracker::default();
2966 let order_identities = OrderIdentityRegistry::default();
2967 register_identity(
2968 &order_identities,
2969 venue_order_id,
2970 instrument.id(),
2971 "O-OVERFILL",
2972 );
2973 order_identities.mark_accepted(venue_order_id);
2974 let mut emitter = test_emitter();
2975 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2976 emitter.set_sender(sender);
2977
2978 let ctx = WsDispatchContext {
2979 token_instruments: &token_instruments,
2980 fill_tracker: &fill_tracker,
2981 pending_submits: &pending_submits,
2982 order_identities: &order_identities,
2983 emitter: &emitter,
2984 account_id: AccountId::from("POLY-001"),
2985 clock: nautilus_core::time::get_atomic_clock_realtime(),
2986 user_address: "0xtest",
2987 user_api_key: "test-key",
2988 };
2989 let mut state = WsDispatchState::default();
2990
2991 let trade = PolymarketUserTrade {
2992 asset_id,
2993 bucket_index: 0,
2994 fee_rate_bps: "0".to_string(),
2995 id: "trade-overfill".to_string(),
2996 last_update: "1700000001".to_string(),
2997 maker_address: Ustr::from("0xmaker"),
2998 maker_orders: vec![],
2999 market: Ustr::from("0xmarket"),
3000 match_time: "1700000000".to_string(),
3001 outcome: PolymarketOutcome::yes(),
3002 owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
3003 price: "0.014".to_string(),
3004 side: PolymarketOrderSide::Buy,
3005 size: "714.285714".to_string(),
3008 status: PolymarketTradeStatus::Confirmed,
3009 taker_order_id: venue_order_id.as_str().to_string(),
3010 timestamp: "1700000000000".to_string(),
3011 trade_owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
3012 transaction_hash: None,
3013 trader_side: PolymarketLiquiditySide::Taker,
3014 event_type: PolymarketEventType::Trade,
3015 };
3016
3017 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
3018
3019 let cumulative = fill_tracker
3023 .get_cumulative_filled(&venue_order_id)
3024 .expect("order must be registered");
3025 assert_eq!(cumulative, submitted);
3026
3027 let event = receiver.try_recv().expect("expected a filled event");
3030 match event {
3031 ExecutionEvent::Order(OrderEventAny::Filled(filled)) => {
3032 assert_eq!(
3033 filled.last_qty, submitted,
3034 "filled qty must be snapped to submitted",
3035 );
3036 assert_eq!(filled.venue_order_id, venue_order_id);
3037 }
3038 other => panic!("expected filled event, was {other:?}"),
3039 }
3040 }
3041
3042 #[rstest]
3043 #[case(
3044 TimeInForce::Ioc,
3045 OrderType::Market,
3046 OrderSide::Buy,
3047 "5.202910",
3048 "5.202897",
3049 false,
3050 true
3051 )]
3052 #[case(
3053 TimeInForce::Fok,
3054 OrderType::Limit,
3055 OrderSide::Buy,
3056 "5.202910",
3057 "5.202897",
3058 true,
3059 false
3060 )]
3061 #[case(
3062 TimeInForce::Ioc,
3063 OrderType::Limit,
3064 OrderSide::Buy,
3065 "30",
3066 "20",
3067 false,
3068 true
3069 )]
3070 #[case(
3071 TimeInForce::Ioc,
3072 OrderType::Market,
3073 OrderSide::Sell,
3074 "5.202910",
3075 "5.202897",
3076 false,
3077 true
3078 )]
3079 #[case(
3080 TimeInForce::Gtc,
3081 OrderType::Limit,
3082 OrderSide::Buy,
3083 "5.202910",
3084 "5.202897",
3085 false,
3086 false
3087 )]
3088 fn test_taker_terminal_status_on_trade_confirm(
3089 #[case] time_in_force: TimeInForce,
3090 #[case] order_type: OrderType,
3091 #[case] order_side: OrderSide,
3092 #[case] submitted_qty: &str,
3093 #[case] fill_qty: &str,
3094 #[case] expect_normalization: bool,
3095 #[case] expect_cancel: bool,
3096 ) {
3097 use crate::common::enums::{
3101 PolymarketEventType, PolymarketOrderSide, PolymarketOutcome, PolymarketTradeStatus,
3102 };
3103
3104 let instrument = test_instrument();
3105 let asset_id = instrument.id().symbol.inner();
3106 let token_instruments = AtomicMap::new();
3107 token_instruments.insert(asset_id, instrument.clone());
3108
3109 let fill_tracker = OrderFillTrackerMap::new();
3110 let venue_order_id = VenueOrderId::from("0xtaker-one-shot-dust");
3111 let submitted = Quantity::from_decimal_dp(
3112 Decimal::from_str_exact(submitted_qty).unwrap(),
3113 instrument.size_precision(),
3114 )
3115 .unwrap();
3116 fill_tracker.register(
3117 venue_order_id,
3118 submitted,
3119 order_side,
3120 instrument.id(),
3121 instrument.size_precision(),
3122 instrument.price_precision(),
3123 );
3124
3125 let pending_submits = PendingSubmitTracker::default();
3126 let order_identities = OrderIdentityRegistry::default();
3127 order_identities.register_order_identity(
3128 venue_order_id,
3129 OrderIdentity {
3130 client_order_id: ClientOrderId::from("O-ONE-SHOT"),
3131 strategy_id: StrategyId::from("S-001"),
3132 instrument_id: instrument.id(),
3133 order_side,
3134 order_type,
3135 time_in_force,
3136 },
3137 );
3138 order_identities.mark_accepted(venue_order_id);
3139 let mut emitter = test_emitter();
3140 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3141 emitter.set_sender(sender);
3142
3143 let ctx = WsDispatchContext {
3144 token_instruments: &token_instruments,
3145 fill_tracker: &fill_tracker,
3146 pending_submits: &pending_submits,
3147 order_identities: &order_identities,
3148 emitter: &emitter,
3149 account_id: AccountId::from("POLY-001"),
3150 clock: nautilus_core::time::get_atomic_clock_realtime(),
3151 user_address: "0xtest",
3152 user_api_key: "test-key",
3153 };
3154 let mut state = WsDispatchState::default();
3155
3156 let trade = PolymarketUserTrade {
3157 asset_id,
3158 bucket_index: 0,
3159 fee_rate_bps: "0".to_string(),
3160 id: "trade-one-shot-dust".to_string(),
3161 last_update: "1700000001".to_string(),
3162 maker_address: Ustr::from("0xmaker"),
3163 maker_orders: vec![],
3164 market: Ustr::from("0xmarket"),
3165 match_time: "1700000000".to_string(),
3166 outcome: PolymarketOutcome::yes(),
3167 owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
3168 price: "0.963".to_string(),
3169 side: if order_side == OrderSide::Buy {
3170 PolymarketOrderSide::Buy
3171 } else {
3172 PolymarketOrderSide::Sell
3173 },
3174 size: fill_qty.to_string(),
3175 status: PolymarketTradeStatus::Confirmed,
3176 taker_order_id: venue_order_id.as_str().to_string(),
3177 timestamp: "1700000000000".to_string(),
3178 trade_owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
3179 transaction_hash: None,
3180 trader_side: PolymarketLiquiditySide::Taker,
3181 event_type: PolymarketEventType::Trade,
3182 };
3183
3184 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
3185
3186 let event = receiver.try_recv().expect("expected the venue fill event");
3187 match event {
3188 ExecutionEvent::Order(OrderEventAny::Filled(filled)) => {
3189 assert_eq!(
3190 filled.last_qty,
3191 Quantity::from_decimal_dp(
3192 Decimal::from_str_exact(fill_qty).unwrap(),
3193 instrument.size_precision(),
3194 )
3195 .unwrap(),
3196 );
3197 }
3198 other => panic!("expected filled event, was {other:?}"),
3199 }
3200
3201 if expect_normalization {
3202 let event = receiver.try_recv().expect("expected quantity update");
3203 match event {
3204 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
3205 assert_eq!(
3206 updated.quantity,
3207 Quantity::new(5.202897, instrument.size_precision()),
3208 );
3209 assert_eq!(updated.venue_order_id, Some(venue_order_id));
3210 assert!(updated.reconciliation);
3211 }
3212 other => panic!("expected updated event, was {other:?}"),
3213 }
3214 assert!(
3215 fill_tracker
3216 .get_cumulative_filled(&venue_order_id)
3217 .is_none(),
3218 "order must be settled and removed from the tracker",
3219 );
3220 } else if expect_cancel {
3221 let event = receiver.try_recv().expect("expected IOC cancellation");
3222 match event {
3223 ExecutionEvent::Order(OrderEventAny::Canceled(canceled)) => {
3224 assert_eq!(canceled.venue_order_id, Some(venue_order_id));
3225 }
3226 other => panic!("expected canceled event, was {other:?}"),
3227 }
3228 assert!(
3229 fill_tracker
3230 .get_cumulative_filled(&venue_order_id)
3231 .is_none(),
3232 "canceled IOC must be settled and removed from the tracker",
3233 );
3234 } else {
3235 assert!(
3236 receiver.try_recv().is_err(),
3237 "resting order must not receive a terminal event",
3238 );
3239 assert!(
3240 fill_tracker
3241 .get_cumulative_filled(&venue_order_id)
3242 .is_some(),
3243 "ineligible order must stay tracked with open leaves",
3244 );
3245 }
3246 }
3247
3248 #[rstest]
3249 fn test_dispatch_taker_fill_gross_overfill_raises_qty_then_fills() {
3250 use crate::common::enums::{
3254 PolymarketEventType, PolymarketOrderSide, PolymarketOutcome, PolymarketTradeStatus,
3255 };
3256
3257 let instrument = test_instrument();
3258 let asset_id = instrument.id().symbol.inner();
3259 let size_precision = instrument.size_precision();
3260 let token_instruments = AtomicMap::new();
3261 token_instruments.insert(asset_id, instrument.clone());
3262
3263 let fill_tracker = OrderFillTrackerMap::new();
3264 let venue_order_id = VenueOrderId::from("0xtaker-gross-overfill");
3265 let submitted = Quantity::new(30.0, size_precision);
3266 fill_tracker.register(
3267 venue_order_id,
3268 submitted,
3269 OrderSide::Buy,
3270 instrument.id(),
3271 size_precision,
3272 instrument.price_precision(),
3273 );
3274
3275 let pending_submits = PendingSubmitTracker::default();
3276 let order_identities = OrderIdentityRegistry::default();
3277 register_identity(
3278 &order_identities,
3279 venue_order_id,
3280 instrument.id(),
3281 "O-GROSS-OVERFILL",
3282 );
3283 order_identities.mark_accepted(venue_order_id);
3284 let mut emitter = test_emitter();
3285 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3286 emitter.set_sender(sender);
3287
3288 let ctx = WsDispatchContext {
3289 token_instruments: &token_instruments,
3290 fill_tracker: &fill_tracker,
3291 pending_submits: &pending_submits,
3292 order_identities: &order_identities,
3293 emitter: &emitter,
3294 account_id: AccountId::from("POLY-001"),
3295 clock: nautilus_core::time::get_atomic_clock_realtime(),
3296 user_address: "0xtest",
3297 user_api_key: "test-key",
3298 };
3299 let mut state = WsDispatchState::default();
3300
3301 let trade = PolymarketUserTrade {
3303 asset_id,
3304 bucket_index: 0,
3305 fee_rate_bps: "0".to_string(),
3306 id: "trade-gross-overfill".to_string(),
3307 last_update: "1700000001".to_string(),
3308 maker_address: Ustr::from("0xmaker"),
3309 maker_orders: vec![],
3310 market: Ustr::from("0xmarket"),
3311 match_time: "1700000000".to_string(),
3312 outcome: PolymarketOutcome::yes(),
3313 owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
3314 price: "0.014".to_string(),
3315 side: PolymarketOrderSide::Buy,
3316 size: "33.846152".to_string(),
3317 status: PolymarketTradeStatus::Confirmed,
3318 taker_order_id: venue_order_id.as_str().to_string(),
3319 timestamp: "1700000000000".to_string(),
3320 trade_owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
3321 transaction_hash: None,
3322 trader_side: PolymarketLiquiditySide::Taker,
3323 event_type: PolymarketEventType::Trade,
3324 };
3325
3326 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
3327
3328 let expected_qty = Quantity::new(33.846152, size_precision);
3329
3330 match receiver.try_recv().expect("expected an updated event") {
3332 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
3333 assert_eq!(updated.quantity, expected_qty);
3334 assert_eq!(updated.venue_order_id, Some(venue_order_id));
3335 }
3336 other => panic!("expected updated event raising qty to the fill, was {other:?}"),
3337 }
3338
3339 match receiver.try_recv().expect("expected a filled event") {
3340 ExecutionEvent::Order(OrderEventAny::Filled(filled)) => {
3341 assert_eq!(filled.last_qty, expected_qty);
3342 assert_eq!(filled.venue_order_id, venue_order_id);
3343 }
3344 other => panic!("expected filled event, was {other:?}"),
3345 }
3346 }
3347
3348 #[rstest]
3351 #[case(
3352 crate::common::enums::PolymarketOrderStatus::Unmatched,
3353 Some("invalid post-only order: order crosses book"),
3354 "Rejected"
3355 )]
3356 #[case(
3357 crate::common::enums::PolymarketOrderStatus::CanceledMarketResolved,
3358 None,
3359 "Expired"
3360 )]
3361 fn test_dispatch_order_terminal_status_emits_event(
3362 #[case] status: crate::common::enums::PolymarketOrderStatus,
3363 #[case] reason: Option<&str>,
3364 #[case] expected: &str,
3365 ) {
3366 use crate::common::enums::{
3367 PolymarketEventType, PolymarketOrderSide, PolymarketOrderType, PolymarketOutcome,
3368 };
3369
3370 let instrument = test_instrument();
3371 let asset_id = instrument.id().symbol.inner();
3372 let order_id = "0xterminal-order".to_string();
3373 let venue_order_id = VenueOrderId::from(order_id.as_str());
3374
3375 let token_instruments = AtomicMap::new();
3376 token_instruments.insert(asset_id, instrument.clone());
3377
3378 let fill_tracker = OrderFillTrackerMap::new();
3379 fill_tracker.register(
3380 venue_order_id,
3381 Quantity::from("10"),
3382 OrderSide::Buy,
3383 instrument.id(),
3384 instrument.size_precision(),
3385 instrument.price_precision(),
3386 );
3387
3388 let pending_submits = PendingSubmitTracker::default();
3389 let order_identities = OrderIdentityRegistry::default();
3390 register_identity(
3391 &order_identities,
3392 venue_order_id,
3393 instrument.id(),
3394 "O-TERMINAL",
3395 );
3396 order_identities.mark_accepted(venue_order_id);
3397 let mut emitter = test_emitter();
3398 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3399 emitter.set_sender(sender);
3400
3401 let ctx = WsDispatchContext {
3402 token_instruments: &token_instruments,
3403 fill_tracker: &fill_tracker,
3404 pending_submits: &pending_submits,
3405 order_identities: &order_identities,
3406 emitter: &emitter,
3407 account_id: AccountId::from("POLY-001"),
3408 clock: nautilus_core::time::get_atomic_clock_realtime(),
3409 user_address: "0xabc",
3410 user_api_key: "xxx",
3411 };
3412 let mut state = WsDispatchState::default();
3413
3414 let order = PolymarketUserOrder {
3415 asset_id,
3416 associate_trades: None,
3417 created_at: Some("1775074735".to_string()),
3418 expiration: Some("0".to_string()),
3419 id: order_id,
3420 maker_address: Some(Ustr::from("0xabc")),
3421 market: Ustr::from("0x4134"),
3422 order_owner: Some(Ustr::from("xxx")),
3423 order_type: Some(PolymarketOrderType::FOK),
3424 original_size: "10".to_string(),
3425 outcome: Some(PolymarketOutcome::yes()),
3426 owner: Ustr::from("xxx"),
3427 price: "0.50".to_string(),
3428 side: PolymarketOrderSide::Buy,
3429 size_matched: "0".to_string(),
3430 status: Some(PolymarketUserOrderStatus::new(status, reason)),
3431 timestamp: "1775074738031".to_string(),
3432 event_type: PolymarketEventType::Placement,
3433 };
3434
3435 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
3436
3437 let event = receiver.try_recv().expect("expected terminal order event");
3438 match event {
3439 ExecutionEvent::Order(order_event) => {
3440 assert!(
3441 format!("{order_event:?}").starts_with(expected),
3442 "expected {expected}, was {order_event:?}"
3443 );
3444 assert_eq!(
3445 order_event.client_order_id(),
3446 ClientOrderId::from("O-TERMINAL")
3447 );
3448
3449 if let OrderEventAny::Rejected(rejected) = order_event {
3450 assert_eq!(
3451 rejected.reason.as_str(),
3452 "invalid post-only order: order crosses book"
3453 );
3454 assert!(rejected.due_post_only);
3455 }
3456 }
3457 other => panic!("expected order event, was {other:?}"),
3458 }
3459 }
3460}