1use std::fmt::Debug;
29
30use ahash::{AHashMap, AHashSet};
31use anyhow::Context;
32use indexmap::IndexMap;
33use nautilus_common::cache::fifo::{FifoCache, FifoCacheMap};
34use nautilus_core::{
35 UUID4, UnixNanos, collections::AtomicMap, string::secret::REDACTED, time::AtomicTime,
36};
37use nautilus_live::{ExecutionEventEmitter, execution::context::OrderContext};
38use nautilus_model::{
39 enums::{LiquiditySide, OrderSide, OrderStatus, OrderType, TimeInForce},
40 events::{
41 OrderAccepted, OrderCanceled, OrderEventAny, OrderExpired, OrderFillVoided, OrderFilled,
42 OrderRejected, OrderUpdated,
43 },
44 identifiers::{AccountId, ClientOrderId, InstrumentId, TradeId, VenueOrderId},
45 instruments::{Instrument, InstrumentAny},
46 reports::{FillReport, OrderStatusReport},
47 types::{Money, Price, Quantity},
48};
49use rust_decimal::Decimal;
50use ustr::Ustr;
51
52use super::{
53 messages::{
54 PolymarketUserOrder, PolymarketUserOrderStatus, PolymarketUserTrade, UserWsMessage,
55 },
56 parse::parse_timestamp_ms,
57};
58use crate::{
59 common::{
60 enums::{
61 PolymarketLiquiditySide, PolymarketOrderSide, PolymarketOrderStatus,
62 PolymarketOrderType, PolymarketSignerType, PolymarketTradeStatus,
63 },
64 models::PolymarketMakerOrder,
65 parse::parse_decimal_exact,
66 },
67 execution::{
68 context::OrderContextRegistry,
69 get_pusd_currency, is_post_only_crossing,
70 order_fill_tracker::{BufferedFill, FillCorrectionMetadata, OrderFillTrackerMap},
71 parse::{
72 build_maker_fill_report, compute_commission, determine_order_side,
73 instrument_fee_exponent, instrument_taker_fee, parse_liquidity_side,
74 },
75 pending::PendingSubmitTracker,
76 },
77 http::error::sanitize_error_text,
78};
79
80#[derive(Debug)]
82pub(crate) struct AccountRefreshRequest;
83
84#[derive(Debug, Default)]
89pub(crate) struct WsDispatchState {
90 pub processed_fills: FifoCache<String, 10_000>,
91 matched_fills: FifoCacheMap<String, Vec<OrderFilled>, 10_000>,
92 voided_trades: FifoCache<String, 10_000>,
93 confirmed_trades: FifoCache<String, 10_000>,
94 reconciled_fills: FifoCache<(TradeId, VenueOrderId), 10_000>,
95
96 pending_terminal_orders: FifoCacheMap<VenueOrderId, PendingTerminalOrder, 10_000>,
97 terminal_cancel_reports: FifoCacheMap<VenueOrderId, OrderStatusReport, 10_000>,
98
99 pending_commands: AHashMap<ClientOrderId, PendingCommand>,
100 inflight_cancel_markets: AHashSet<InstrumentId>,
101 replaced_venue_order_ids: FifoCache<VenueOrderId, 10_000>,
102 closed_modify_venue_order_ids: FifoCacheMap<VenueOrderId, UnixNanos, 10_000>,
103}
104
105impl WsDispatchState {
106 pub(crate) fn restore_matched_trade(&mut self, key: String, fills: Vec<OrderFilled>) {
107 self.processed_fills.add(key.clone());
108 self.matched_fills.insert(key, fills);
109 }
110
111 pub(crate) fn restore_voided_trade(&mut self, key: String) {
112 self.processed_fills.add(key.clone());
113 self.matched_fills.remove(&key);
114 self.voided_trades.add(key);
115 }
116
117 pub(crate) fn begin_modify(
118 &mut self,
119 client_order_id: ClientOrderId,
120 old_venue_order_id: VenueOrderId,
121 instrument_id: InstrumentId,
122 ) -> bool {
123 if self.pending_commands.contains_key(&client_order_id)
124 || self.replaced_venue_order_ids.contains(&old_venue_order_id)
125 || self.inflight_cancel_markets.contains(&instrument_id)
126 {
127 return false;
128 }
129
130 self.pending_commands.insert(
131 client_order_id,
132 PendingCommand::Modify {
133 old_venue_order_id,
134 instrument_id,
135 cancel_ts: None,
136 replacement: None,
137 },
138 );
139 true
140 }
141
142 pub(crate) fn is_modifying(&self, client_order_id: &ClientOrderId) -> bool {
143 matches!(
144 self.pending_commands.get(client_order_id),
145 Some(PendingCommand::Modify { .. })
146 )
147 }
148
149 pub(crate) fn confirm_modify_cancel(
150 &mut self,
151 client_order_id: ClientOrderId,
152 expected_old_venue_order_id: VenueOrderId,
153 ts_event: UnixNanos,
154 ) -> bool {
155 let Some(PendingCommand::Modify {
156 old_venue_order_id,
157 cancel_ts,
158 ..
159 }) = self.pending_commands.get_mut(&client_order_id)
160 else {
161 return false;
162 };
163
164 if *old_venue_order_id != expected_old_venue_order_id {
165 return false;
166 }
167
168 cancel_ts.get_or_insert(ts_event);
169 true
170 }
171
172 pub(crate) fn set_modify_replacement(
173 &mut self,
174 client_order_id: ClientOrderId,
175 venue_order_id: VenueOrderId,
176 quantity: Quantity,
177 leg_quantity: Quantity,
178 price: Price,
179 ) -> bool {
180 let Some(PendingCommand::Modify { replacement, .. }) =
181 self.pending_commands.get_mut(&client_order_id)
182 else {
183 return false;
184 };
185
186 *replacement = Some(PendingModifyReplacement {
187 venue_order_id,
188 quantity,
189 leg_quantity,
190 price,
191 });
192 true
193 }
194
195 pub(crate) fn claim_modify_replacement(
196 &mut self,
197 venue_order_id: VenueOrderId,
198 ) -> Option<ModifyPromotion> {
199 let promotion = self.pending_modify_promotion(venue_order_id)?;
200 self.pending_commands.remove(&promotion.client_order_id)?;
201 self.replaced_venue_order_ids
202 .add(promotion.old_venue_order_id);
203 Some(promotion)
204 }
205
206 pub(crate) fn finish_modify_without_replacement(
207 &mut self,
208 client_order_id: ClientOrderId,
209 expected_old_venue_order_id: VenueOrderId,
210 cancellation_proven: bool,
211 ts_event: UnixNanos,
212 ) -> Option<(VenueOrderId, Option<UnixNanos>)> {
213 let Some(PendingCommand::Modify {
214 old_venue_order_id, ..
215 }) = self.pending_commands.get(&client_order_id)
216 else {
217 return None;
218 };
219
220 if *old_venue_order_id != expected_old_venue_order_id {
221 return None;
222 }
223
224 let PendingCommand::Modify {
225 old_venue_order_id,
226 cancel_ts,
227 ..
228 } = self.pending_commands.remove(&client_order_id)?
229 else {
230 return None;
231 };
232
233 let cancel_ts = self
234 .terminal_cancel_reports
235 .get(&old_venue_order_id)
236 .map(|report| report.ts_last)
237 .or(cancel_ts)
238 .or(cancellation_proven.then_some(ts_event));
239
240 if let Some(cancel_ts) = cancel_ts {
241 self.closed_modify_venue_order_ids
242 .insert(old_venue_order_id, cancel_ts);
243 }
244
245 Some((old_venue_order_id, cancel_ts))
246 }
247
248 pub(crate) fn finish_unsubmitted_modifies(
249 &mut self,
250 ) -> Vec<(ClientOrderId, VenueOrderId, Option<UnixNanos>)> {
251 let terminal_cancel_reports = &self.terminal_cancel_reports;
252 let closed_modify_venue_order_ids = &mut self.closed_modify_venue_order_ids;
253 let mut finished = Vec::new();
254
255 self.pending_commands
256 .retain(|client_order_id, command| match command {
257 PendingCommand::Modify {
258 old_venue_order_id,
259 cancel_ts,
260 replacement: None,
261 ..
262 } => {
263 let cancel_ts = terminal_cancel_reports
264 .get(old_venue_order_id)
265 .map(|report| report.ts_last)
266 .or(*cancel_ts);
267 if let Some(cancel_ts) = cancel_ts {
268 closed_modify_venue_order_ids.insert(*old_venue_order_id, cancel_ts);
269 }
270
271 finished.push((*client_order_id, *old_venue_order_id, cancel_ts));
272 false
273 }
274 _ => true,
275 });
276
277 finished
278 }
279
280 pub(crate) fn pending_modify_promotion(
281 &self,
282 venue_order_id: VenueOrderId,
283 ) -> Option<ModifyPromotion> {
284 self.pending_commands
285 .iter()
286 .find_map(|(client_order_id, command)| match command {
287 PendingCommand::Modify {
288 old_venue_order_id,
289 replacement: Some(replacement),
290 ..
291 } if replacement.venue_order_id == venue_order_id => Some(ModifyPromotion {
292 client_order_id: *client_order_id,
293 old_venue_order_id: *old_venue_order_id,
294 venue_order_id,
295 quantity: replacement.quantity,
296 leg_quantity: replacement.leg_quantity,
297 price: replacement.price,
298 }),
299 _ => None,
300 })
301 }
302
303 pub(crate) fn pending_modify_promotions(&self) -> Vec<ModifyPromotion> {
304 self.pending_commands
305 .iter()
306 .filter_map(|(client_order_id, command)| match command {
307 PendingCommand::Modify {
308 old_venue_order_id,
309 replacement: Some(replacement),
310 ..
311 } => Some(ModifyPromotion {
312 client_order_id: *client_order_id,
313 old_venue_order_id: *old_venue_order_id,
314 venue_order_id: replacement.venue_order_id,
315 quantity: replacement.quantity,
316 leg_quantity: replacement.leg_quantity,
317 price: replacement.price,
318 }),
319 _ => None,
320 })
321 .collect()
322 }
323
324 pub(crate) fn begin_cancels(&mut self, orders: &[(ClientOrderId, InstrumentId)]) -> bool {
325 if orders.iter().any(|(client_order_id, instrument_id)| {
326 self.pending_commands.contains_key(client_order_id)
327 || self.inflight_cancel_markets.contains(instrument_id)
328 }) {
329 return false;
330 }
331
332 self.pending_commands.extend(
333 orders
334 .iter()
335 .map(|(client_order_id, _)| (*client_order_id, PendingCommand::Cancel)),
336 );
337 true
338 }
339
340 pub(crate) fn begin_available_cancels(
341 &mut self,
342 orders: &[(ClientOrderId, InstrumentId)],
343 ) -> Option<Vec<ClientOrderId>> {
344 if orders.iter().any(|(client_order_id, instrument_id)| {
345 matches!(
346 self.pending_commands.get(client_order_id),
347 Some(PendingCommand::Modify { .. })
348 ) || self.inflight_cancel_markets.contains(instrument_id)
349 }) {
350 return None;
351 }
352
353 let client_order_ids = orders
354 .iter()
355 .filter_map(|(client_order_id, _)| {
356 if self.pending_commands.contains_key(client_order_id) {
357 return None;
358 }
359
360 self.pending_commands
361 .insert(*client_order_id, PendingCommand::Cancel);
362 Some(*client_order_id)
363 })
364 .collect();
365 Some(client_order_ids)
366 }
367
368 pub(crate) fn finish_cancels(&mut self, client_order_ids: &[ClientOrderId]) {
369 for client_order_id in client_order_ids {
370 if matches!(
371 self.pending_commands.get(client_order_id),
372 Some(PendingCommand::Cancel)
373 ) {
374 self.pending_commands.remove(client_order_id);
375 }
376 }
377 }
378
379 pub(crate) fn begin_market_cancel(&mut self, instrument_id: InstrumentId) -> bool {
380 if self.pending_commands.values().any(|command| match command {
381 PendingCommand::Modify {
382 instrument_id: pending_instrument_id,
383 ..
384 } => *pending_instrument_id == instrument_id,
385 PendingCommand::Cancel => false,
386 }) {
387 return false;
388 }
389
390 self.inflight_cancel_markets.insert(instrument_id)
391 }
392
393 pub(crate) fn finish_market_cancel(&mut self, instrument_id: InstrumentId) {
394 self.inflight_cancel_markets.remove(&instrument_id);
395 }
396
397 pub(crate) fn record_terminal_cancel_report(&mut self, report: OrderStatusReport) {
398 self.terminal_cancel_reports
399 .insert(report.venue_order_id, report);
400 }
401
402 pub(crate) fn record_reconciled_fill(
403 &mut self,
404 trade_id: TradeId,
405 venue_order_id: VenueOrderId,
406 ) {
407 self.reconciled_fills.add((trade_id, venue_order_id));
408 }
409
410 pub(crate) fn replaced_venue_order_id(&self, venue_order_id: VenueOrderId) -> bool {
411 self.replaced_venue_order_ids.contains(&venue_order_id)
412 }
413
414 pub(crate) fn reset_session(&mut self) {
415 let retained_cancel_reports = self
416 .pending_commands
417 .values()
418 .filter_map(|command| match command {
419 PendingCommand::Modify {
420 old_venue_order_id, ..
421 } => self
422 .terminal_cancel_reports
423 .get(old_venue_order_id)
424 .cloned(),
425 PendingCommand::Cancel => None,
426 })
427 .collect::<Vec<_>>();
428
429 self.processed_fills.clear();
430 self.matched_fills.clear();
431 self.voided_trades.clear();
432 self.confirmed_trades.clear();
433 self.pending_terminal_orders.clear();
434 self.terminal_cancel_reports.clear();
435 for report in retained_cancel_reports {
436 self.terminal_cancel_reports
437 .insert(report.venue_order_id, report);
438 }
439
440 self.pending_commands
441 .retain(|_, command| matches!(command, PendingCommand::Modify { .. }));
442 self.inflight_cancel_markets.clear();
443 }
444
445 fn suppress_modify_cancel(&self, venue_order_id: VenueOrderId) -> bool {
446 self.closed_modify_venue_order_ids
447 .contains_key(&venue_order_id)
448 || self.suppress_modify_cancel_reemit(venue_order_id)
449 }
450
451 pub(crate) fn suppress_modify_cancel_reemit(&self, venue_order_id: VenueOrderId) -> bool {
452 self.replaced_venue_order_ids.contains(&venue_order_id)
453 || self.pending_commands.values().any(|command| {
454 matches!(
455 command,
456 PendingCommand::Modify {
457 old_venue_order_id,
458 ..
459 } if *old_venue_order_id == venue_order_id
460 )
461 })
462 }
463}
464
465#[derive(Clone, Copy, Debug)]
466enum PendingCommand {
467 Modify {
468 old_venue_order_id: VenueOrderId,
469 instrument_id: InstrumentId,
470 cancel_ts: Option<UnixNanos>,
471 replacement: Option<PendingModifyReplacement>,
472 },
473 Cancel,
474}
475
476#[derive(Clone, Copy, Debug)]
477struct PendingModifyReplacement {
478 venue_order_id: VenueOrderId,
479 quantity: Quantity,
480 leg_quantity: Quantity,
481 price: Price,
482}
483
484#[derive(Clone, Copy, Debug)]
485pub(crate) struct ModifyPromotion {
486 pub(crate) client_order_id: ClientOrderId,
487 pub(crate) old_venue_order_id: VenueOrderId,
488 pub(crate) venue_order_id: VenueOrderId,
489 pub(crate) quantity: Quantity,
490 pub(crate) leg_quantity: Quantity,
491 pub(crate) price: Price,
492}
493
494#[cfg(test)]
495impl WsDispatchState {
496 pub(crate) fn matched_fill_count(&self, key: &str) -> usize {
497 self.matched_fills.get(&key.to_string()).map_or(0, Vec::len)
498 }
499
500 pub(crate) fn is_voided_trade(&self, key: &str) -> bool {
501 self.voided_trades.contains(&key.to_string())
502 }
503}
504
505#[derive(Clone, Debug)]
506struct PendingTerminalOrder {
507 trade_ids: Vec<String>,
508 ts_event: UnixNanos,
509}
510
511pub(crate) struct WsDispatchContext<'a> {
513 pub token_instruments: &'a AtomicMap<Ustr, InstrumentAny>,
514 pub fill_tracker: &'a OrderFillTrackerMap,
515 pub pending_submits: &'a PendingSubmitTracker,
516 pub order_contexts: &'a OrderContextRegistry,
517 pub emitter: &'a ExecutionEventEmitter,
518 pub account_id: AccountId,
519 pub clock: &'static AtomicTime,
520 pub signer_type: PolymarketSignerType,
521 pub user_address: &'a str,
522 pub user_api_key: &'a str,
523}
524
525impl Debug for WsDispatchContext<'_> {
526 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
527 f.debug_struct(stringify!(WsDispatchContext))
528 .field("token_instruments", &self.token_instruments)
529 .field("fill_tracker", &self.fill_tracker)
530 .field("pending_submits", &self.pending_submits)
531 .field("order_contexts", &self.order_contexts)
532 .field("emitter", &self.emitter)
533 .field("account_id", &self.account_id)
534 .field("clock", &self.clock)
535 .field("user_address", &self.user_address)
536 .field("user_api_key", &REDACTED)
537 .finish()
538 }
539}
540
541pub(crate) fn dispatch_user_message(
543 message: &UserWsMessage,
544 ctx: &WsDispatchContext<'_>,
545 state: &mut WsDispatchState,
546) -> Option<AccountRefreshRequest> {
547 match message {
548 UserWsMessage::Order(order) => {
549 dispatch_order_update(order, ctx, state);
550 None
551 }
552 UserWsMessage::Trade(trade) => dispatch_trade_update(trade, ctx, state),
553 }
554}
555
556fn dispatch_order_update(
557 order: &PolymarketUserOrder,
558 ctx: &WsDispatchContext<'_>,
559 state: &mut WsDispatchState,
560) {
561 let Some(status) = order.status.as_ref() else {
562 log::warn!("Ignoring order update without status: {}", order.id);
563 return;
564 };
565
566 let Some(order_type) = order.order_type else {
567 log::warn!("Ignoring order update without order_type: {}", order.id);
568 return;
569 };
570
571 let instruments = ctx.token_instruments.load();
572 let instrument = match instruments.get(&order.asset_id) {
573 Some(i) => i,
574 None => {
575 log::warn!("Unknown asset_id in order update: {}", order.asset_id);
576 return;
577 }
578 };
579
580 let ts_event = parse_timestamp_ms(&order.timestamp).unwrap_or_else(|_| ctx.clock.get_time_ns());
581 let venue_order_id = VenueOrderId::from(order.id.as_str());
582
583 let ts_init = ctx.clock.get_time_ns();
584 let mut report = match build_ws_order_status_report(
585 order,
586 status,
587 order_type,
588 instrument,
589 ctx.account_id,
590 ts_event,
591 ts_init,
592 ) {
593 Ok(report) => report,
594 Err(e) => {
595 log::warn!("Ignoring invalid order update {}: {e}", order.id);
596 return;
597 }
598 };
599 let mut promoted_fills = Vec::new();
600 let mut promoted_reports = Vec::new();
601 let promoted_client_order_id = if state.pending_modify_promotion(venue_order_id).is_some() {
602 if report.order_status == OrderStatus::Rejected {
603 reject_modify_replacement(venue_order_id, &report, ts_event, ctx, state);
604 return;
605 }
606
607 promote_modify_replacement_from_ws(
608 venue_order_id,
609 ts_event,
610 ctx,
611 state,
612 &mut promoted_fills,
613 &mut promoted_reports,
614 )
615 } else {
616 None
617 };
618
619 let local_client_order_id =
620 promoted_client_order_id.or_else(|| ctx.pending_submits.client_order_id(&venue_order_id));
621 let mut is_accepted = ctx.fill_tracker.contains(&venue_order_id);
622 report.client_order_id = local_client_order_id;
623
624 let mut buffered_fills = if local_client_order_id.is_some()
626 && !is_accepted
627 && report.order_status != OrderStatus::Rejected
628 {
629 is_accepted = true;
630 ctx.fill_tracker.register_and_take_pending_fills(
631 venue_order_id,
632 local_client_order_id,
633 report.quantity,
634 report
635 .order_side
636 .expect("WebSocket order report side must be Buy or Sell"),
637 )
638 } else if is_accepted {
639 ctx.fill_tracker
640 .take_pending_fills(venue_order_id, local_client_order_id)
641 } else {
642 Vec::new()
643 };
644
645 buffered_fills.splice(0..0, promoted_fills);
646
647 if let Some(tracked_filled) = ctx.fill_tracker.get_cumulative_filled(&venue_order_id)
650 && report.filled_qty > tracked_filled
651 {
652 log::debug!(
653 "Capping filled_qty for {venue_order_id} from {} to {} (awaiting trade messages)",
654 report.filled_qty,
655 tracked_filled,
656 );
657 report.filled_qty = tracked_filled;
658 }
659
660 if report.order_status == OrderStatus::Canceled {
664 state
665 .terminal_cancel_reports
666 .insert(venue_order_id, report.clone());
667 }
668
669 let suppress_cancel = report.order_status == OrderStatus::Canceled
670 && state.suppress_modify_cancel(venue_order_id);
671
672 let context = ctx.order_contexts.get(&venue_order_id);
675
676 for fill in buffered_fills {
678 match context {
679 Some(context) => {
680 emit_buffered_order_filled(&context, &fill, ctx);
681 }
682 None => ctx.emitter.send_fill_report(fill.report),
683 }
684 }
685
686 for buffered in promoted_reports {
687 if buffered.order_status == OrderStatus::Canceled {
688 state
689 .terminal_cancel_reports
690 .insert(venue_order_id, buffered.clone());
691 }
692
693 if let Some(context) = context {
694 emit_tracked_order_status(&buffered, &context, buffered.ts_last, ctx);
695 }
696 }
697
698 if suppress_cancel {
699 log::debug!("Suppressing stale cancel for modified venue leg {venue_order_id}");
700 return;
701 }
702
703 if is_accepted || local_client_order_id.is_some() {
704 match context {
705 Some(context) => emit_tracked_order_status(&report, &context, ts_event, ctx),
706 None => ctx.emitter.send_order_status_report(report),
707 }
708 } else if let Some(report) = ctx
709 .fill_tracker
710 .accept_or_buffer_report(venue_order_id, report)
711 {
712 match ctx.order_contexts.get(&venue_order_id) {
714 Some(context) => emit_tracked_order_status(&report, &context, ts_event, ctx),
715 None => ctx.emitter.send_order_status_report(report),
716 }
717 }
718
719 if status.status == PolymarketOrderStatus::Matched
720 && let Some(trade_ids) = order.associate_trades.clone().filter(|ids| !ids.is_empty())
721 {
722 state.pending_terminal_orders.insert(
723 venue_order_id,
724 PendingTerminalOrder {
725 trade_ids,
726 ts_event,
727 },
728 );
729 emit_quantity_normalization_if_ready(venue_order_id, ctx, state);
730 }
731}
732
733fn promote_modify_replacement_from_ws(
734 venue_order_id: VenueOrderId,
735 ts_event: UnixNanos,
736 ctx: &WsDispatchContext<'_>,
737 state: &mut WsDispatchState,
738 buffered_fills: &mut Vec<BufferedFill>,
739 buffered_reports: &mut Vec<OrderStatusReport>,
740) -> Option<ClientOrderId> {
741 let promotion = state.claim_modify_replacement(venue_order_id)?;
742 let Some(mut context) = ctx.order_contexts.get(&promotion.old_venue_order_id) else {
743 log::error!(
744 "Cannot promote Polymarket replacement {venue_order_id}: old venue leg {} has no context",
745 promotion.old_venue_order_id,
746 );
747 return None;
748 };
749
750 context.quantity = promotion.quantity;
751 context.price = Some(promotion.price);
752 ctx.order_contexts.register_context(venue_order_id, context);
753 ctx.order_contexts.mark_accepted(venue_order_id);
754
755 let updated = OrderUpdated::new(
756 ctx.emitter.trader_id(),
757 context.identity.strategy_id,
758 context.identity.instrument_id,
759 promotion.client_order_id,
760 promotion.quantity,
761 UUID4::new(),
762 ts_event,
763 ctx.clock.get_time_ns(),
764 false,
765 Some(venue_order_id),
766 Some(ctx.account_id),
767 Some(promotion.price),
768 None,
769 None,
770 false,
771 );
772 ctx.emitter
773 .send_order_event(OrderEventAny::Updated(updated));
774
775 buffered_fills.extend(ctx.fill_tracker.register_and_take_pending_fills(
776 venue_order_id,
777 Some(promotion.client_order_id),
778 promotion.leg_quantity,
779 context.identity.order_side,
780 ));
781 buffered_reports.extend(ctx.fill_tracker.take_pending_reports(&venue_order_id));
782 Some(promotion.client_order_id)
783}
784
785fn reject_modify_replacement(
786 venue_order_id: VenueOrderId,
787 report: &OrderStatusReport,
788 ts_event: UnixNanos,
789 ctx: &WsDispatchContext<'_>,
790 state: &mut WsDispatchState,
791) {
792 let Some(promotion) = state.pending_modify_promotion(venue_order_id) else {
793 return;
794 };
795
796 let Some((old_venue_order_id, cancel_ts)) = state.finish_modify_without_replacement(
797 promotion.client_order_id,
798 promotion.old_venue_order_id,
799 true,
800 ts_event,
801 ) else {
802 return;
803 };
804
805 let Some(context) = ctx.order_contexts.get(&old_venue_order_id) else {
806 return;
807 };
808
809 let reason = report
810 .cancel_reason
811 .as_deref()
812 .unwrap_or("replacement order rejected");
813 ctx.emitter.emit_order_modify_rejected_event(
814 context.identity.strategy_id,
815 context.identity.instrument_id,
816 context.identity.client_order_id,
817 Some(old_venue_order_id),
818 &sanitize_error_text(reason),
819 ts_event,
820 );
821
822 if let Some(cancel_ts) = cancel_ts {
823 emit_order_canceled(&context, old_venue_order_id, cancel_ts, ctx);
824 }
825}
826
827fn emit_buffered_order_filled(
828 context: &OrderContext,
829 buffered: &BufferedFill,
830 ctx: &WsDispatchContext<'_>,
831) {
832 let fill = &buffered.report;
833 ensure_accepted(context, fill.venue_order_id, fill.ts_event, ctx);
834
835 let info = buffered
836 .correction
837 .as_ref()
838 .and_then(|correction| correction.info.clone());
839 let filled = build_order_filled(context, fill, info, ctx);
840 ctx.fill_tracker
841 .emit_buffered_fill(filled, buffered.correction.as_ref(), |filled, new_qty| {
842 if let Some(new_qty) = new_qty {
843 emit_buy_overfill_update(context, fill.venue_order_id, new_qty, fill.ts_event, ctx);
844 }
845 ctx.emitter.send_order_event(OrderEventAny::Filled(filled));
846 });
847}
848
849fn emit_quantity_normalization_if_ready(
850 venue_order_id: VenueOrderId,
851 ctx: &WsDispatchContext<'_>,
852 state: &mut WsDispatchState,
853) {
854 let is_ready = state
855 .pending_terminal_orders
856 .get(&venue_order_id)
857 .is_some_and(|pending| {
858 pending
859 .trade_ids
860 .iter()
861 .all(|trade_id| state.confirmed_trades.contains(trade_id))
862 });
863
864 if !is_ready {
865 return;
866 }
867
868 let Some(pending) = state.pending_terminal_orders.remove(&venue_order_id) else {
869 return;
870 };
871
872 let Some(context) = ctx.order_contexts.get(&venue_order_id) else {
873 log::warn!("Cannot normalize terminal order {venue_order_id} without a local context");
874 return;
875 };
876
877 if let Some(quantity) = ctx
878 .fill_tracker
879 .check_terminal_quantity_normalization(&venue_order_id)
880 {
881 emit_terminal_quantity_update(&context, venue_order_id, quantity, pending.ts_event, ctx);
882 }
883}
884
885fn emit_taker_terminal_status(
891 trade: &PolymarketUserTrade,
892 ctx: &WsDispatchContext<'_>,
893 ts_event: UnixNanos,
894) {
895 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
896
897 let Some(context) = ctx.order_contexts.get(&venue_order_id) else {
898 return;
899 };
900
901 if context.time_in_force == TimeInForce::Fok {
902 if let Some(quantity) = ctx
903 .fill_tracker
904 .check_terminal_quantity_normalization(&venue_order_id)
905 {
906 emit_terminal_quantity_update(&context, venue_order_id, quantity, ts_event, ctx);
907 }
908 return;
909 }
910
911 if context.time_in_force == TimeInForce::Ioc
912 && let Some(remainder) = ctx
913 .fill_tracker
914 .take_terminal_ioc_remainder(&venue_order_id)
915 {
916 log::debug!(
917 "Closing terminal IOC order {venue_order_id} as Canceled (unfilled remainder={remainder})"
918 );
919 emit_order_canceled(&context, venue_order_id, ts_event, ctx);
920 }
921}
922
923fn dispatch_trade_update(
924 trade: &PolymarketUserTrade,
925 ctx: &WsDispatchContext<'_>,
926 state: &mut WsDispatchState,
927) -> Option<AccountRefreshRequest> {
928 let dedup_key = format!("{}-{}", trade.id, trade.taker_order_id);
929 if trade.status == PolymarketTradeStatus::Failed {
930 void_failed_trade(trade, dedup_key, ctx, state);
931 return Some(AccountRefreshRequest);
932 }
933
934 if matches!(
935 trade.status,
936 PolymarketTradeStatus::Mined | PolymarketTradeStatus::Retrying
937 ) {
938 log::debug!("Waiting for terminal trade status: {}", trade.id);
939 return None;
940 }
941
942 if has_unknown_trade_instrument(trade, ctx) {
943 log::warn!(
944 "Deferring trade {} until its instrument is available",
945 trade.id
946 );
947 return None;
948 }
949
950 let is_confirmed = trade.status == PolymarketTradeStatus::Confirmed;
951 if !dispatch_trade_fills(trade, &dedup_key, is_confirmed, ctx, state) {
952 return None;
953 }
954
955 if !is_confirmed {
956 return None;
957 }
958
959 confirm_trade(trade, &dedup_key, ctx, state);
960 Some(AccountRefreshRequest)
961}
962
963fn void_failed_trade(
964 trade: &PolymarketUserTrade,
965 dedup_key: String,
966 ctx: &WsDispatchContext<'_>,
967 state: &mut WsDispatchState,
968) {
969 if state.voided_trades.contains(&dedup_key) {
970 return;
971 }
972
973 let direct_fills = state.matched_fills.remove(&dedup_key).unwrap_or_default();
974 for fill in &direct_fills {
975 ctx.fill_tracker
976 .reverse_fill(&fill.venue_order_id, fill.last_qty);
977 }
978
979 let mut fills = direct_fills;
980 fills.extend(ctx.fill_tracker.void_buffered_trade(&dedup_key));
981 for fill in fills {
982 emit_order_fill_voided(&fill, trade, Some(fill.event_id), ctx);
983 }
984
985 state.processed_fills.add(dedup_key.clone());
986 state.voided_trades.add(dedup_key);
987 state.confirmed_trades.remove(&trade.id);
988}
989
990fn has_unknown_trade_instrument(trade: &PolymarketUserTrade, ctx: &WsDispatchContext<'_>) -> bool {
991 let instruments = ctx.token_instruments.load();
992
993 if trade.trader_side == PolymarketLiquiditySide::Maker {
994 trade
995 .maker_orders
996 .iter()
997 .filter(|order| is_user_maker_order(order, ctx))
998 .any(|order| !instruments.contains_key(&order.asset_id))
999 } else {
1000 !instruments.contains_key(&trade.asset_id)
1001 }
1002}
1003
1004fn dispatch_trade_fills(
1005 trade: &PolymarketUserTrade,
1006 dedup_key: &String,
1007 is_confirmed: bool,
1008 ctx: &WsDispatchContext<'_>,
1009 state: &mut WsDispatchState,
1010) -> bool {
1011 if state.processed_fills.contains(dedup_key) {
1012 log::debug!("Duplicate fill skipped: {dedup_key}");
1013 return true;
1014 }
1015
1016 let fills = if trade.trader_side == PolymarketLiquiditySide::Maker {
1017 let reports = match build_ws_maker_fill_reports(trade, ctx) {
1018 Ok(reports) => reports,
1019 Err(e) => {
1020 log::error!("Cannot build maker fills for trade {}: {e}", trade.id);
1021 return false;
1022 }
1023 };
1024 dispatch_maker_fill_reports(reports, trade, dedup_key, is_confirmed, ctx, state)
1025 } else {
1026 let report = match build_ws_taker_fill_report_for_trade(trade, ctx) {
1027 Ok(report) => report,
1028 Err(e) => {
1029 log::error!("Cannot build taker fill for trade {}: {e}", trade.id);
1030 return false;
1031 }
1032 };
1033 dispatch_taker_fill_report(report, trade, dedup_key, is_confirmed, ctx, state)
1034 };
1035
1036 if !fills.is_empty() {
1037 state.matched_fills.insert(dedup_key.clone(), fills);
1038 }
1039 state.processed_fills.add(dedup_key.clone());
1040 true
1041}
1042
1043fn confirm_trade(
1044 trade: &PolymarketUserTrade,
1045 dedup_key: &str,
1046 ctx: &WsDispatchContext<'_>,
1047 state: &mut WsDispatchState,
1048) {
1049 let ts_event = parse_timestamp_ms(&trade.timestamp).unwrap_or_else(|_| ctx.clock.get_time_ns());
1050 ctx.fill_tracker.mark_trade_confirmed(dedup_key);
1051 state.confirmed_trades.add(trade.id.clone());
1052 if trade.trader_side == PolymarketLiquiditySide::Maker {
1053 for order in trade
1054 .maker_orders
1055 .iter()
1056 .filter(|order| is_user_maker_order(order, ctx))
1057 {
1058 emit_quantity_normalization_if_ready(
1059 VenueOrderId::from(order.order_id.as_str()),
1060 ctx,
1061 state,
1062 );
1063 }
1064 } else {
1065 emit_quantity_normalization_if_ready(
1066 VenueOrderId::from(trade.taker_order_id.as_str()),
1067 ctx,
1068 state,
1069 );
1070 emit_taker_terminal_status(trade, ctx, ts_event);
1071 }
1072}
1073
1074fn build_ws_maker_fill_reports(
1075 trade: &PolymarketUserTrade,
1076 ctx: &WsDispatchContext<'_>,
1077) -> anyhow::Result<Vec<FillReport>> {
1078 let user_orders: Vec<_> = trade
1079 .maker_orders
1080 .iter()
1081 .filter(|order| is_user_maker_order(order, ctx))
1082 .collect();
1083
1084 if user_orders.is_empty() {
1085 log::warn!("No matching maker orders for user in trade: {}", trade.id);
1086 return Ok(Vec::new());
1087 }
1088
1089 let instruments = ctx.token_instruments.load();
1090 let liquidity_side = parse_liquidity_side(trade.trader_side);
1091 let ts_event = parse_timestamp_ms(&trade.timestamp).unwrap_or_else(|_| ctx.clock.get_time_ns());
1092 let ts_init = ctx.clock.get_time_ns();
1093 let mut reports = Vec::with_capacity(user_orders.len());
1094
1095 for mo in user_orders {
1096 let asset_id = mo.asset_id;
1097 let instrument = instruments
1098 .get(&asset_id)
1099 .with_context(|| format!("unknown asset_id in maker order: {asset_id}"))?;
1100 let mut report = build_maker_fill_report(
1101 mo,
1102 &trade.id,
1103 trade.trader_side,
1104 trade.side,
1105 trade.asset_id.as_str(),
1106 ctx.account_id,
1107 instrument.id(),
1108 instrument.price_precision(),
1109 instrument.size_precision(),
1110 crate::execution::get_pusd_currency(),
1111 liquidity_side,
1112 ts_event,
1113 ts_init,
1114 )
1115 .with_context(|| format!("failed to build maker fill for asset {asset_id}"))?;
1116
1117 let maker_venue_order_id = report.venue_order_id;
1118 report.client_order_id = ctx.pending_submits.client_order_id(&maker_venue_order_id);
1119 report.last_qty = ctx
1120 .fill_tracker
1121 .snap_fill_qty(&maker_venue_order_id, report.last_qty);
1122 reports.push(report);
1123 }
1124
1125 Ok(reports)
1126}
1127
1128fn dispatch_maker_fill_reports(
1129 reports: Vec<FillReport>,
1130 trade: &PolymarketUserTrade,
1131 correction_key: &str,
1132 is_confirmed: bool,
1133 ctx: &WsDispatchContext<'_>,
1134 state: &mut WsDispatchState,
1135) -> Vec<OrderFilled> {
1136 let fill_info = trade_fill_info(trade);
1137 let mut fills = Vec::new();
1138
1139 for mut report in reports {
1140 let maker_venue_order_id = report.venue_order_id;
1141
1142 if state
1143 .reconciled_fills
1144 .contains(&(report.trade_id, maker_venue_order_id))
1145 {
1146 continue;
1147 }
1148
1149 let mut promoted_reports = Vec::new();
1150
1151 if state
1152 .pending_modify_promotion(maker_venue_order_id)
1153 .is_some()
1154 {
1155 let mut buffered_fills = Vec::new();
1156 report.client_order_id = promote_modify_replacement_from_ws(
1157 maker_venue_order_id,
1158 report.ts_event,
1159 ctx,
1160 state,
1161 &mut buffered_fills,
1162 &mut promoted_reports,
1163 );
1164 emit_promoted_ws_fills(maker_venue_order_id, buffered_fills, ctx);
1165 }
1166
1167 if let Some(report) = ctx.fill_tracker.accept_or_buffer_fill(
1168 maker_venue_order_id,
1169 report,
1170 FillCorrectionMetadata {
1171 correction_key: correction_key.to_string(),
1172 info: fill_info.clone(),
1173 is_confirmed,
1174 },
1175 ) {
1176 match ctx.order_contexts.get(&maker_venue_order_id) {
1177 Some(context) => {
1178 fills.push(emit_order_filled(&context, &report, fill_info.clone(), ctx));
1179 }
1180 None => ctx.emitter.send_fill_report(report),
1181 }
1182 reemit_terminal_cancel(maker_venue_order_id, state, ctx);
1183 }
1184
1185 emit_promoted_ws_reports(maker_venue_order_id, promoted_reports, ctx, state);
1186 }
1187 fills
1188}
1189
1190fn is_user_maker_order(order: &PolymarketMakerOrder, ctx: &WsDispatchContext<'_>) -> bool {
1191 order.is_owned_by(ctx.user_address, ctx.user_api_key, ctx.signer_type)
1192}
1193
1194fn build_ws_taker_fill_report_for_trade(
1195 trade: &PolymarketUserTrade,
1196 ctx: &WsDispatchContext<'_>,
1197) -> anyhow::Result<FillReport> {
1198 let instruments = ctx.token_instruments.load();
1199 let instrument = instruments
1200 .get(&trade.asset_id)
1201 .with_context(|| format!("unknown asset_id in trade: {}", trade.asset_id))?;
1202 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
1203 let liquidity_side = parse_liquidity_side(trade.trader_side);
1204 let ts_event = parse_timestamp_ms(&trade.timestamp).unwrap_or_else(|_| ctx.clock.get_time_ns());
1205 let ts_init = ctx.clock.get_time_ns();
1206
1207 let mut report = build_ws_taker_fill_report(
1208 trade,
1209 instrument,
1210 ctx.account_id,
1211 liquidity_side,
1212 ts_event,
1213 ts_init,
1214 )?;
1215 report.client_order_id = ctx.pending_submits.client_order_id(&venue_order_id);
1216 report.last_qty = ctx
1217 .fill_tracker
1218 .snap_fill_qty(&venue_order_id, report.last_qty);
1219 Ok(report)
1220}
1221
1222fn dispatch_taker_fill_report(
1223 mut report: FillReport,
1224 trade: &PolymarketUserTrade,
1225 correction_key: &str,
1226 is_confirmed: bool,
1227 ctx: &WsDispatchContext<'_>,
1228 state: &mut WsDispatchState,
1229) -> Vec<OrderFilled> {
1230 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
1231
1232 if state
1233 .reconciled_fills
1234 .contains(&(report.trade_id, venue_order_id))
1235 {
1236 return Vec::new();
1237 }
1238
1239 let mut promoted_reports = Vec::new();
1240 if state.pending_modify_promotion(venue_order_id).is_some() {
1241 let mut buffered_fills = Vec::new();
1242 report.client_order_id = promote_modify_replacement_from_ws(
1243 venue_order_id,
1244 report.ts_event,
1245 ctx,
1246 state,
1247 &mut buffered_fills,
1248 &mut promoted_reports,
1249 );
1250 emit_promoted_ws_fills(venue_order_id, buffered_fills, ctx);
1251 }
1252
1253 let mut fills = Vec::new();
1254
1255 if let Some(report) = ctx.fill_tracker.accept_or_buffer_fill(
1256 venue_order_id,
1257 report,
1258 FillCorrectionMetadata {
1259 correction_key: correction_key.to_string(),
1260 info: trade_fill_info(trade),
1261 is_confirmed,
1262 },
1263 ) {
1264 match ctx.order_contexts.get(&venue_order_id) {
1265 Some(context) => {
1266 let fill = emit_order_filled(&context, &report, trade_fill_info(trade), ctx);
1267 fills.push(fill);
1268 }
1269 None => ctx.emitter.send_fill_report(report),
1270 }
1271 reemit_terminal_cancel(venue_order_id, state, ctx);
1272 }
1273
1274 emit_promoted_ws_reports(venue_order_id, promoted_reports, ctx, state);
1275 fills
1276}
1277
1278fn emit_promoted_ws_fills(
1279 venue_order_id: VenueOrderId,
1280 buffered_fills: Vec<BufferedFill>,
1281 ctx: &WsDispatchContext<'_>,
1282) {
1283 let context = ctx.order_contexts.get(&venue_order_id);
1284 for fill in buffered_fills {
1285 match context {
1286 Some(context) => emit_buffered_order_filled(&context, &fill, ctx),
1287 None => ctx.emitter.send_fill_report(fill.report),
1288 }
1289 }
1290}
1291
1292fn emit_promoted_ws_reports(
1293 venue_order_id: VenueOrderId,
1294 buffered_reports: Vec<OrderStatusReport>,
1295 ctx: &WsDispatchContext<'_>,
1296 state: &mut WsDispatchState,
1297) {
1298 let context = ctx.order_contexts.get(&venue_order_id);
1299
1300 for report in buffered_reports {
1301 if report.order_status == OrderStatus::Canceled {
1302 state.record_terminal_cancel_report(report.clone());
1303 }
1304
1305 match context {
1306 Some(context) => emit_tracked_order_status(&report, &context, report.ts_last, ctx),
1307 None => ctx.emitter.send_order_status_report(report),
1308 }
1309 }
1310}
1311
1312fn reemit_terminal_cancel(
1322 venue_order_id: VenueOrderId,
1323 state: &WsDispatchState,
1324 ctx: &WsDispatchContext<'_>,
1325) {
1326 if ctx.fill_tracker.is_fully_filled(&venue_order_id) {
1327 return;
1328 }
1329
1330 if state.suppress_modify_cancel_reemit(venue_order_id) {
1331 return;
1332 }
1333
1334 let cancel_ts = state
1335 .closed_modify_venue_order_ids
1336 .get(&venue_order_id)
1337 .copied()
1338 .or_else(|| {
1339 state
1340 .terminal_cancel_reports
1341 .get(&venue_order_id)
1342 .map(|report| report.ts_last)
1343 });
1344
1345 if let Some(cancel_ts) = cancel_ts {
1346 log::debug!("Re-emitting cancel for {venue_order_id} after fill to restore terminal state");
1347 match ctx.order_contexts.get(&venue_order_id) {
1348 Some(context) => {
1349 emit_order_canceled(&context, venue_order_id, cancel_ts, ctx);
1350 }
1351 None => {
1352 if let Some(cancel_report) = state.terminal_cancel_reports.get(&venue_order_id) {
1353 ctx.emitter.send_order_status_report(cancel_report.clone());
1354 }
1355 }
1356 }
1357 }
1358}
1359
1360fn build_ws_order_status_report(
1361 order: &PolymarketUserOrder,
1362 status: &PolymarketUserOrderStatus,
1363 order_type: PolymarketOrderType,
1364 instrument: &InstrumentAny,
1365 account_id: AccountId,
1366 ts_event: UnixNanos,
1367 ts_init: UnixNanos,
1368) -> anyhow::Result<OrderStatusReport> {
1369 let venue_order_id = VenueOrderId::from(order.id.as_str());
1370 let order_status =
1371 crate::execution::parse::resolve_order_status(status.status, order.event_type);
1372 let order_side = OrderSide::from(order.side);
1373 let time_in_force = TimeInForce::from(order_type);
1374 let size_precision = instrument.size_precision();
1375 let price_precision = instrument.price_precision();
1376 let price_dec = parse_decimal_exact(&order.price)?;
1377 anyhow::ensure!(
1378 price_dec > Decimal::ZERO && price_dec < Decimal::ONE,
1379 "order price must be in (0, 1)"
1380 );
1381 let quantity_dec = parse_decimal_exact(&order.original_size)?;
1382 let filled_dec = if order.size_matched.is_empty() {
1384 Decimal::ZERO
1385 } else {
1386 parse_decimal_exact(&order.size_matched)?
1387 };
1388 anyhow::ensure!(
1389 quantity_dec > Decimal::ZERO && filled_dec >= Decimal::ZERO,
1390 "invalid order quantity"
1391 );
1392 let quantity = Quantity::from_decimal_dp(
1393 original_size_to_shares(quantity_dec, price_dec, order.side, order_type)?,
1394 size_precision,
1395 )?;
1396 let filled_qty = Quantity::from_decimal_dp(filled_dec, size_precision)?;
1397 let price = Price::from_decimal_dp(price_dec, price_precision)?;
1398
1399 let mut report = OrderStatusReport::new(
1400 account_id,
1401 instrument.id(),
1402 None,
1403 venue_order_id,
1404 order_side.into(),
1405 OrderType::Limit,
1406 time_in_force,
1407 order_status,
1408 quantity,
1409 filled_qty,
1410 ts_event,
1411 ts_event,
1412 ts_init,
1413 None,
1414 );
1415 report.price = Some(price);
1416
1417 if order_status == OrderStatus::Rejected {
1418 report.cancel_reason.clone_from(&status.reason);
1419 }
1420
1421 Ok(report)
1422}
1423
1424fn original_size_to_shares(
1435 original_size: Decimal,
1436 price: Decimal,
1437 side: PolymarketOrderSide,
1438 order_type: PolymarketOrderType,
1439) -> anyhow::Result<Decimal> {
1440 if side != PolymarketOrderSide::Buy
1441 || !matches!(
1442 order_type,
1443 PolymarketOrderType::FAK | PolymarketOrderType::FOK
1444 )
1445 {
1446 return Ok(original_size);
1447 }
1448
1449 original_size
1450 .checked_div(price)
1451 .context("order share quantity overflow")
1452}
1453
1454fn build_ws_taker_fill_report(
1455 trade: &PolymarketUserTrade,
1456 instrument: &InstrumentAny,
1457 account_id: AccountId,
1458 liquidity_side: LiquiditySide,
1459 ts_event: UnixNanos,
1460 ts_init: UnixNanos,
1461) -> anyhow::Result<FillReport> {
1462 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
1463 let trade_id = TradeId::from(trade.id.as_str());
1464 let order_side = determine_order_side(
1465 trade.trader_side,
1466 trade.side,
1467 trade.asset_id.as_str(),
1468 trade.asset_id.as_str(),
1469 );
1470
1471 let size_precision = instrument.size_precision();
1472 let price_precision = instrument.price_precision();
1473 let size_dec = parse_decimal_exact(&trade.size)?;
1474 let price_dec = parse_decimal_exact(&trade.price)?;
1475 anyhow::ensure!(size_dec > Decimal::ZERO, "trade quantity must be positive");
1476 anyhow::ensure!(
1477 price_dec > Decimal::ZERO && price_dec < Decimal::ONE,
1478 "trade price must be in (0, 1)"
1479 );
1480 let last_qty = Quantity::from_decimal_dp(size_dec, size_precision)?;
1481 let last_px = Price::from_decimal_dp(price_dec, price_precision)?;
1482
1483 let fee_rate = instrument_taker_fee(instrument);
1484 let commission_value = compute_commission(
1485 fee_rate,
1486 instrument_fee_exponent(instrument)?,
1487 size_dec,
1488 price_dec,
1489 liquidity_side,
1490 )?;
1491 let pusd = crate::execution::get_pusd_currency();
1492
1493 Ok(FillReport {
1494 account_id,
1495 instrument_id: instrument.id(),
1496 venue_order_id,
1497 trade_id,
1498 order_side,
1499 last_qty,
1500 last_px,
1501 commission: Money::from_decimal(commission_value, pusd)
1502 .context("commission is not representable as Money")?,
1503 liquidity_side,
1504 avg_px: None,
1505 report_id: UUID4::new(),
1506 ts_event,
1507 ts_init,
1508 client_order_id: None,
1509 venue_position_id: None,
1510 })
1511}
1512
1513fn emit_tracked_order_status(
1519 report: &OrderStatusReport,
1520 context: &OrderContext,
1521 ts_event: UnixNanos,
1522 ctx: &WsDispatchContext<'_>,
1523) {
1524 let venue_order_id = report.venue_order_id;
1525 match report.order_status {
1526 OrderStatus::Accepted => ensure_accepted(context, venue_order_id, ts_event, ctx),
1527 OrderStatus::PartiallyFilled | OrderStatus::Filled => {
1528 ensure_accepted(context, venue_order_id, ts_event, ctx);
1529 }
1530 OrderStatus::Canceled => {
1531 ensure_accepted(context, venue_order_id, ts_event, ctx);
1532 emit_order_canceled(context, venue_order_id, ts_event, ctx);
1533 }
1534 OrderStatus::Expired => {
1535 ensure_accepted(context, venue_order_id, ts_event, ctx);
1536 emit_order_expired(context, venue_order_id, ts_event, ctx);
1537 }
1538 OrderStatus::Rejected => {
1539 let reason = report
1540 .cancel_reason
1541 .clone()
1542 .unwrap_or_else(|| "REJECTED".to_string());
1543
1544 emit_order_rejected(context, &reason, ts_event, ctx);
1545 }
1546 other => log::debug!("No order event for status {other:?} on {venue_order_id}"),
1547 }
1548}
1549
1550fn ensure_accepted(
1556 context: &OrderContext,
1557 venue_order_id: VenueOrderId,
1558 ts_event: UnixNanos,
1559 ctx: &WsDispatchContext<'_>,
1560) {
1561 if !ctx.order_contexts.mark_accepted(venue_order_id) {
1562 return;
1563 }
1564
1565 let accepted = OrderAccepted::new(
1566 ctx.emitter.trader_id(),
1567 context.identity.strategy_id,
1568 context.identity.instrument_id,
1569 context.identity.client_order_id,
1570 venue_order_id,
1571 ctx.account_id,
1572 UUID4::new(),
1573 ts_event,
1574 ctx.clock.get_time_ns(),
1575 false,
1576 );
1577 ctx.emitter
1578 .send_order_event(OrderEventAny::Accepted(accepted));
1579}
1580
1581fn emit_order_filled(
1586 context: &OrderContext,
1587 fill: &FillReport,
1588 info: Option<IndexMap<Ustr, Ustr>>,
1589 ctx: &WsDispatchContext<'_>,
1590) -> OrderFilled {
1591 ensure_accepted(context, fill.venue_order_id, fill.ts_event, ctx);
1592
1593 if let Some(new_qty) = ctx.fill_tracker.buy_overfill_bump(&fill.venue_order_id) {
1594 emit_buy_overfill_update(context, fill.venue_order_id, new_qty, fill.ts_event, ctx);
1595 }
1596
1597 let filled = build_order_filled(context, fill, info, ctx);
1598 ctx.emitter
1599 .send_order_event(OrderEventAny::Filled(filled.clone()));
1600 filled
1601}
1602
1603fn build_order_filled(
1604 context: &OrderContext,
1605 fill: &FillReport,
1606 info: Option<IndexMap<Ustr, Ustr>>,
1607 ctx: &WsDispatchContext<'_>,
1608) -> OrderFilled {
1609 OrderFilled::new(
1610 ctx.emitter.trader_id(),
1611 context.identity.strategy_id,
1612 context.identity.instrument_id,
1613 context.identity.client_order_id,
1614 fill.venue_order_id,
1615 ctx.account_id,
1616 fill.trade_id,
1617 context.identity.order_side,
1618 context.identity.order_type,
1619 fill.last_qty,
1620 fill.last_px,
1621 get_pusd_currency(),
1622 fill.liquidity_side,
1623 UUID4::new(),
1624 fill.ts_event,
1625 fill.ts_init,
1626 false,
1627 fill.venue_position_id,
1628 Some(fill.commission),
1629 info,
1630 )
1631}
1632
1633fn emit_order_fill_voided(
1634 fill: &OrderFilled,
1635 trade: &PolymarketUserTrade,
1636 causation_id: Option<UUID4>,
1637 ctx: &WsDispatchContext<'_>,
1638) {
1639 let ts_event = parse_timestamp_ms(&trade.timestamp).unwrap_or_else(|_| ctx.clock.get_time_ns());
1640 let mut voided = OrderFillVoided::new(
1641 fill.trader_id,
1642 fill.strategy_id,
1643 fill.instrument_id,
1644 fill.client_order_id,
1645 fill.venue_order_id,
1646 fill.account_id,
1647 Ustr::from(&format!("{}-FAILED-{}", trade.id, fill.client_order_id)),
1648 fill.trade_id,
1649 fill.last_qty,
1650 fill.commission,
1651 fill.order_side,
1652 fill.order_type,
1653 fill.last_px,
1654 fill.currency,
1655 fill.liquidity_side,
1656 fill.position_id,
1657 Some(Ustr::from("FAILED")),
1658 trade_fill_info(trade),
1659 UUID4::new(),
1660 ts_event,
1661 ctx.clock.get_time_ns(),
1662 false,
1663 false,
1664 );
1665 voided.causation_id = causation_id;
1666 ctx.emitter
1667 .send_order_event(OrderEventAny::FillVoided(voided));
1668}
1669
1670fn trade_fill_info(trade: &PolymarketUserTrade) -> Option<IndexMap<Ustr, Ustr>> {
1675 let value = serde_json::to_value(trade).ok()?;
1676 let object = value.as_object()?;
1677 let mut info = IndexMap::with_capacity(object.len());
1678 for (key, val) in object {
1679 let val_str = match val {
1680 serde_json::Value::String(s) => s.clone(),
1681 other => other.to_string(),
1682 };
1683 info.insert(Ustr::from(key.as_str()), Ustr::from(val_str.as_str()));
1684 }
1685 Some(info)
1686}
1687
1688fn emit_buy_overfill_update(
1694 context: &OrderContext,
1695 venue_order_id: VenueOrderId,
1696 new_qty: Quantity,
1697 ts_event: UnixNanos,
1698 ctx: &WsDispatchContext<'_>,
1699) {
1700 let updated = OrderUpdated::new(
1701 ctx.emitter.trader_id(),
1702 context.identity.strategy_id,
1703 context.identity.instrument_id,
1704 context.identity.client_order_id,
1705 new_qty,
1706 UUID4::new(),
1707 ts_event,
1708 ctx.clock.get_time_ns(),
1709 false,
1710 Some(venue_order_id),
1711 Some(ctx.account_id),
1712 None,
1713 None,
1714 None,
1715 false,
1716 );
1717 ctx.emitter
1718 .send_order_event(OrderEventAny::Updated(updated));
1719}
1720
1721fn emit_terminal_quantity_update(
1723 context: &OrderContext,
1724 venue_order_id: VenueOrderId,
1725 quantity: Quantity,
1726 ts_event: UnixNanos,
1727 ctx: &WsDispatchContext<'_>,
1728) {
1729 let updated = OrderUpdated::new(
1730 ctx.emitter.trader_id(),
1731 context.identity.strategy_id,
1732 context.identity.instrument_id,
1733 context.identity.client_order_id,
1734 quantity,
1735 UUID4::new(),
1736 ts_event,
1737 ctx.clock.get_time_ns(),
1738 true,
1739 Some(venue_order_id),
1740 Some(ctx.account_id),
1741 None,
1742 None,
1743 None,
1744 false,
1745 );
1746 ctx.emitter
1747 .send_order_event(OrderEventAny::Updated(updated));
1748}
1749
1750fn emit_order_canceled(
1751 context: &OrderContext,
1752 venue_order_id: VenueOrderId,
1753 ts_event: UnixNanos,
1754 ctx: &WsDispatchContext<'_>,
1755) {
1756 let canceled = OrderCanceled::new(
1757 ctx.emitter.trader_id(),
1758 context.identity.strategy_id,
1759 context.identity.instrument_id,
1760 context.identity.client_order_id,
1761 UUID4::new(),
1762 ts_event,
1763 ctx.clock.get_time_ns(),
1764 false,
1765 Some(venue_order_id),
1766 Some(ctx.account_id),
1767 None,
1768 );
1769 ctx.emitter
1770 .send_order_event(OrderEventAny::Canceled(canceled));
1771}
1772
1773fn emit_order_expired(
1774 context: &OrderContext,
1775 venue_order_id: VenueOrderId,
1776 ts_event: UnixNanos,
1777 ctx: &WsDispatchContext<'_>,
1778) {
1779 let expired = OrderExpired::new(
1780 ctx.emitter.trader_id(),
1781 context.identity.strategy_id,
1782 context.identity.instrument_id,
1783 context.identity.client_order_id,
1784 UUID4::new(),
1785 ts_event,
1786 ctx.clock.get_time_ns(),
1787 false,
1788 Some(venue_order_id),
1789 Some(ctx.account_id),
1790 );
1791 ctx.emitter
1792 .send_order_event(OrderEventAny::Expired(expired));
1793}
1794
1795fn emit_order_rejected(
1796 context: &OrderContext,
1797 reason: &str,
1798 ts_event: UnixNanos,
1799 ctx: &WsDispatchContext<'_>,
1800) {
1801 let reason = sanitize_error_text(reason);
1802
1803 let rejected = OrderRejected::new(
1804 ctx.emitter.trader_id(),
1805 context.identity.strategy_id,
1806 context.identity.instrument_id,
1807 context.identity.client_order_id,
1808 ctx.account_id,
1809 Ustr::from(&reason),
1810 UUID4::new(),
1811 ts_event,
1812 ctx.clock.get_time_ns(),
1813 false,
1814 is_post_only_crossing(&reason),
1815 );
1816 ctx.emitter
1817 .send_order_event(OrderEventAny::Rejected(rejected));
1818}
1819
1820#[cfg(test)]
1821mod tests {
1822 use nautilus_common::messages::{ExecutionEvent, ExecutionReport};
1823 use nautilus_core::time::AtomicTime;
1824 use nautilus_live::execution::context::OrderIdentity;
1825 use nautilus_model::{
1826 enums::{AccountType, OrderSide, OrderStatus},
1827 events::OrderEventAny,
1828 identifiers::{ClientOrderId, InstrumentId, StrategyId, TraderId},
1829 orders::{Order, builder::OrderTestBuilder},
1830 types::Currency,
1831 };
1832 use rstest::rstest;
1833 use rust_decimal_macros::dec;
1834
1835 use super::*;
1836 use crate::http::{
1837 models::GammaMarket,
1838 parse::{create_instrument_from_def, parse_gamma_market},
1839 };
1840
1841 fn register_context(
1843 order_contexts: &OrderContextRegistry,
1844 venue_order_id: VenueOrderId,
1845 instrument_id: InstrumentId,
1846 client_order_id: &str,
1847 ) {
1848 order_contexts.register_context(
1849 venue_order_id,
1850 OrderContext {
1851 identity: OrderIdentity {
1852 client_order_id: ClientOrderId::from(client_order_id),
1853 strategy_id: StrategyId::from("S-001"),
1854 instrument_id,
1855 order_side: OrderSide::Buy,
1856 order_type: OrderType::Limit,
1857 },
1858 quantity: Quantity::from("10"),
1859 price: Some(Price::from("0.50")),
1860 trigger_price: None,
1861 trigger_type: None,
1862 time_in_force: TimeInForce::Gtc,
1863 is_post_only: false,
1864 is_reduce_only: false,
1865 is_quote_quantity: false,
1866 },
1867 );
1868 }
1869
1870 fn load<T: serde::de::DeserializeOwned>(filename: &str) -> T {
1871 let path = format!("test_data/{filename}");
1872 let content = std::fs::read_to_string(path).expect("Failed to read test data");
1873 serde_json::from_str(&content).expect("Failed to parse test data")
1874 }
1875
1876 fn test_instrument() -> InstrumentAny {
1877 let market: GammaMarket = load("gamma_market.json");
1878 let defs = parse_gamma_market(&market).unwrap();
1879 create_instrument_from_def(&defs[0], UnixNanos::from(1_000_000_000u64)).unwrap()
1880 }
1881
1882 fn test_emitter() -> ExecutionEventEmitter {
1883 ExecutionEventEmitter::new(
1884 nautilus_core::time::get_atomic_clock_realtime(),
1885 TraderId::from("TESTER-001"),
1886 AccountId::from("POLY-001"),
1887 AccountType::Cash,
1888 Some(Currency::pUSD()),
1889 )
1890 }
1891
1892 #[rstest]
1893 fn test_emit_order_rejected_uses_bounded_clean_reason() {
1894 let instrument = test_instrument();
1895 let token_instruments = AtomicMap::new();
1896 let fill_tracker = OrderFillTrackerMap::new();
1897 let pending_submits = PendingSubmitTracker::default();
1898 let order_contexts = OrderContextRegistry::default();
1899 let mut emitter = test_emitter();
1900 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
1901
1902 emitter.set_sender(sender);
1903
1904 let ctx = WsDispatchContext {
1905 signer_type: PolymarketSignerType::Owner,
1906 token_instruments: &token_instruments,
1907 fill_tracker: &fill_tracker,
1908 pending_submits: &pending_submits,
1909 order_contexts: &order_contexts,
1910 emitter: &emitter,
1911 account_id: AccountId::from("POLY-001"),
1912 clock: nautilus_core::time::get_atomic_clock_realtime(),
1913 user_address: "0xtest",
1914 user_api_key: "test-key",
1915 };
1916 let context = OrderContext {
1917 identity: OrderIdentity {
1918 client_order_id: ClientOrderId::from("O-WS-REJECT"),
1919 strategy_id: StrategyId::from("S-001"),
1920 instrument_id: instrument.id(),
1921 order_side: OrderSide::Buy,
1922 order_type: OrderType::Limit,
1923 },
1924 quantity: Quantity::from("10"),
1925 price: Some(Price::from("0.50")),
1926 trigger_price: None,
1927 trigger_type: None,
1928 time_in_force: TimeInForce::Gtc,
1929 is_post_only: false,
1930 is_reduce_only: false,
1931 is_quote_quantity: false,
1932 };
1933
1934 emit_order_rejected(
1935 &context,
1936 " invalid post-only order:\norder crosses book ",
1937 UnixNanos::from(1_000_000_000),
1938 &ctx,
1939 );
1940
1941 match receiver.try_recv().expect("expected rejected event") {
1942 ExecutionEvent::Order(OrderEventAny::Rejected(event)) => {
1943 assert_eq!(event.reason, "invalid post-only order: order crosses book");
1944 assert!(event.due_post_only);
1945 }
1946 other => panic!("expected rejected event, was {other:?}"),
1947 }
1948 }
1949
1950 #[rstest]
1951 #[case::empty_price("", "1.01", "0")]
1952 #[case::zero_price("0", "1.01", "0")]
1953 #[case::malformed("bad", "100", "0")]
1954 #[case::too_precise("0.50000000000000000000000000001", "100", "0")]
1955 #[case::quantity_overflow("0.5", "79228162514264337593543950335", "0")]
1956 #[case::filled_overflow("0.5", "100", "79228162514264337593543950335")]
1957 fn test_ws_order_report_rejects_invalid_values(
1958 #[case] price: &str,
1959 #[case] quantity: &str,
1960 #[case] filled: &str,
1961 ) {
1962 let mut order: PolymarketUserOrder = load("ws_user_order_placement.json");
1963 order.price = price.into();
1964 order.original_size = quantity.into();
1965 order.size_matched = filled.into();
1966 assert!(
1967 build_ws_order_status_report(
1968 &order,
1969 order.status.as_ref().unwrap(),
1970 order.order_type.unwrap(),
1971 &test_instrument(),
1972 AccountId::from("POLY-001"),
1973 UnixNanos::default(),
1974 UnixNanos::default()
1975 )
1976 .is_err()
1977 );
1978 }
1979
1980 #[rstest]
1981 fn test_ws_fok_order_report_rejects_share_overflow() {
1982 let mut order: PolymarketUserOrder = load("ws_user_order_fok_buy_pusd_size.json");
1983 order.original_size = Decimal::MAX.to_string();
1984 order.price = "0.5".into();
1985 let result = build_ws_order_status_report(
1986 &order,
1987 order.status.as_ref().unwrap(),
1988 PolymarketOrderType::FOK,
1989 &test_instrument(),
1990 AccountId::from("POLY-001"),
1991 UnixNanos::default(),
1992 UnixNanos::default(),
1993 );
1994
1995 assert_eq!(
1996 result.unwrap_err().to_string(),
1997 "order share quantity overflow"
1998 );
1999 }
2000
2001 #[rstest]
2002 fn test_build_ws_order_status_report() {
2003 let order: PolymarketUserOrder = load("ws_user_order_placement.json");
2004 let instrument = test_instrument();
2005 let ts_event = UnixNanos::from(1_000_000_000u64);
2006 let ts_init = UnixNanos::from(2_000_000_000u64);
2007
2008 let report = build_ws_order_status_report(
2009 &order,
2010 order.status.as_ref().unwrap(),
2011 order.order_type.unwrap(),
2012 &instrument,
2013 AccountId::from("POLY-001"),
2014 ts_event,
2015 ts_init,
2016 )
2017 .unwrap();
2018
2019 assert_eq!(report.order_side, Some(OrderSide::Buy));
2020 assert_eq!(report.order_type, OrderType::Limit);
2021 assert_eq!(report.quantity.as_decimal(), dec!(100));
2023 assert_eq!(
2024 report.price.map(|price| price.as_decimal()),
2025 Some(dec!(0.5))
2026 );
2027 assert_eq!(report.ts_accepted, ts_event);
2028 assert_eq!(report.ts_init, ts_init);
2029 }
2030
2031 #[rstest]
2032 fn test_build_ws_order_status_report_venue_cancel_maps_to_canceled() {
2033 let order: PolymarketUserOrder = load("ws_user_order_venue_cancel.json");
2034 let instrument = test_instrument();
2035 let ts_event = UnixNanos::from(1_000_000_000u64);
2036 let ts_init = UnixNanos::from(2_000_000_000u64);
2037
2038 let report = build_ws_order_status_report(
2039 &order,
2040 order.status.as_ref().unwrap(),
2041 order.order_type.unwrap(),
2042 &instrument,
2043 AccountId::from("POLY-001"),
2044 ts_event,
2045 ts_init,
2046 )
2047 .unwrap();
2048
2049 assert_eq!(report.order_status, OrderStatus::Canceled);
2050 }
2051
2052 #[rstest]
2055 #[case(
2056 PolymarketOrderSide::Buy,
2057 PolymarketOrderType::FOK,
2058 dec!(1.01),
2059 dec!(0.01),
2060 dec!(101)
2061 )]
2062 #[case(
2063 PolymarketOrderSide::Buy,
2064 PolymarketOrderType::FOK,
2065 dec!(12),
2066 dec!(0.6),
2067 dec!(20)
2068 )]
2069 #[case(
2070 PolymarketOrderSide::Buy,
2071 PolymarketOrderType::FAK,
2072 dec!(1),
2073 dec!(0.01),
2074 dec!(100)
2075 )]
2076 #[case(
2077 PolymarketOrderSide::Buy,
2078 PolymarketOrderType::GTC,
2079 dec!(20),
2080 dec!(0.18),
2081 dec!(20)
2082 )]
2083 #[case(
2084 PolymarketOrderSide::Buy,
2085 PolymarketOrderType::GTD,
2086 dec!(20),
2087 dec!(0.18),
2088 dec!(20)
2089 )]
2090 #[case(
2091 PolymarketOrderSide::Sell,
2092 PolymarketOrderType::FOK,
2093 dec!(20),
2094 dec!(0.6),
2095 dec!(20)
2096 )]
2097 fn test_original_size_to_shares(
2098 #[case] side: PolymarketOrderSide,
2099 #[case] order_type: PolymarketOrderType,
2100 #[case] original_size: Decimal,
2101 #[case] price: Decimal,
2102 #[case] expected: Decimal,
2103 ) {
2104 let shares = original_size_to_shares(original_size, price, side, order_type).unwrap();
2105
2106 assert_eq!(shares, expected);
2107 }
2108
2109 #[rstest]
2111 #[case("1", "0.03", "33.333333", "0.03")]
2112 fn test_build_ws_order_status_report_fok_buy_quantity(
2113 #[case] original_size: &str,
2114 #[case] price: &str,
2115 #[case] expected_quantity: &str,
2116 #[case] expected_price: &str,
2117 ) {
2118 let mut order: PolymarketUserOrder = load("ws_user_order_fok_buy_pusd_size.json");
2119 order.original_size = original_size.to_string();
2120 order.price = price.to_string();
2121 let instrument = test_instrument();
2122
2123 let report = build_ws_order_status_report(
2124 &order,
2125 order.status.as_ref().unwrap(),
2126 order.order_type.unwrap(),
2127 &instrument,
2128 AccountId::from("POLY-001"),
2129 UnixNanos::from(1_000_000_000u64),
2130 UnixNanos::from(2_000_000_000u64),
2131 )
2132 .unwrap();
2133
2134 assert_eq!(
2135 report.quantity.as_decimal(),
2136 Decimal::from_str_exact(expected_quantity).unwrap()
2137 );
2138 assert_eq!(
2139 report.price.map(|price| price.as_decimal()),
2140 Some(Decimal::from_str_exact(expected_price).unwrap())
2141 );
2142 }
2143
2144 #[rstest]
2145 fn test_dispatch_fok_buy_registers_share_quantity_for_in_flight_submit() {
2146 let order: PolymarketUserOrder = load("ws_user_order_fok_buy_pusd_size.json");
2147 let instrument = test_instrument();
2148
2149 let token_instruments = AtomicMap::new();
2150 token_instruments.insert(order.asset_id, instrument.clone());
2151
2152 let fill_tracker = OrderFillTrackerMap::new();
2154 let pending_submits = PendingSubmitTracker::default();
2155 let order_contexts = OrderContextRegistry::default();
2156 let emitter = test_emitter();
2157
2158 let venue_order_id = VenueOrderId::from(order.id.as_str());
2159 let client_order_id = ClientOrderId::from("O-FOK-IN-FLIGHT");
2160 pending_submits.insert(venue_order_id, client_order_id);
2161 register_context(
2162 &order_contexts,
2163 venue_order_id,
2164 instrument.id(),
2165 client_order_id.as_str(),
2166 );
2167
2168 let ctx = WsDispatchContext {
2169 signer_type: PolymarketSignerType::Owner,
2170 token_instruments: &token_instruments,
2171 fill_tracker: &fill_tracker,
2172 pending_submits: &pending_submits,
2173 order_contexts: &order_contexts,
2174 emitter: &emitter,
2175 account_id: AccountId::from("POLY-001"),
2176 clock: nautilus_core::time::get_atomic_clock_realtime(),
2177 user_address: "0xtest",
2178 user_api_key: "test-key",
2179 };
2180 let mut state = WsDispatchState::default();
2181
2182 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
2183
2184 assert_eq!(
2186 fill_tracker
2187 .submitted_qty(&venue_order_id)
2188 .map(|qty| qty.as_decimal()),
2189 Some(dec!(101)),
2190 );
2191 }
2192
2193 #[rstest]
2194 fn test_dispatch_fok_buy_report_quantity_is_shares_without_identity() {
2195 let order: PolymarketUserOrder = load("ws_user_order_fok_buy_pusd_size.json");
2196 let instrument = test_instrument();
2197
2198 let token_instruments = AtomicMap::new();
2199 token_instruments.insert(order.asset_id, instrument.clone());
2200
2201 let fill_tracker = OrderFillTrackerMap::new();
2202 let venue_order_id = VenueOrderId::from(order.id.as_str());
2203 fill_tracker.register(
2204 venue_order_id,
2205 Quantity::from("101"),
2206 OrderSide::Buy,
2207 instrument.id(),
2208 instrument.size_precision(),
2209 instrument.price_precision(),
2210 );
2211
2212 let pending_submits = PendingSubmitTracker::default();
2213 let order_contexts = OrderContextRegistry::default();
2215 let mut emitter = test_emitter();
2216 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2217 emitter.set_sender(sender);
2218
2219 let ctx = WsDispatchContext {
2220 signer_type: PolymarketSignerType::Owner,
2221 token_instruments: &token_instruments,
2222 fill_tracker: &fill_tracker,
2223 pending_submits: &pending_submits,
2224 order_contexts: &order_contexts,
2225 emitter: &emitter,
2226 account_id: AccountId::from("POLY-001"),
2227 clock: nautilus_core::time::get_atomic_clock_realtime(),
2228 user_address: "0xtest",
2229 user_api_key: "test-key",
2230 };
2231 let mut state = WsDispatchState::default();
2232
2233 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
2234
2235 let event = receiver.try_recv().expect("expected order report");
2236 let ExecutionEvent::Report(ExecutionReport::Order(report)) = event else {
2237 panic!("expected an order report, was {event:?}");
2238 };
2239
2240 assert_eq!(report.venue_order_id, venue_order_id);
2241 assert_eq!(report.order_side, Some(OrderSide::Buy));
2242 assert_eq!(report.time_in_force, TimeInForce::Fok);
2243 assert_eq!(report.order_status, OrderStatus::Canceled);
2244 assert_eq!(report.quantity.as_decimal(), dec!(101));
2245 assert_eq!(report.filled_qty.as_decimal(), dec!(0));
2246 assert_eq!(
2247 report.price.map(|price| price.as_decimal()),
2248 Some(dec!(0.01))
2249 );
2250 }
2251
2252 #[rstest]
2253 fn test_build_ws_taker_fill_report() {
2254 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2255 let instrument = test_instrument();
2256 let ts_event = UnixNanos::from(1_000_000_000u64);
2257 let ts_init = UnixNanos::from(2_000_000_000u64);
2258
2259 let report = build_ws_taker_fill_report(
2260 &trade,
2261 &instrument,
2262 AccountId::from("POLY-001"),
2263 LiquiditySide::Taker,
2264 ts_event,
2265 ts_init,
2266 )
2267 .expect("representable commission builds a fill report");
2268
2269 assert_eq!(report.order_side, OrderSide::Buy);
2270 assert_eq!(report.liquidity_side, LiquiditySide::Taker);
2271 assert_eq!(report.trade_id.as_str(), trade.id);
2272 assert_eq!(report.ts_event, ts_event);
2273 assert_eq!(report.ts_init, ts_init);
2274 }
2275
2276 #[rstest]
2277 fn test_trade_fill_info_flattens_raw_trade() {
2278 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2279
2280 let info = trade_fill_info(&trade).expect("info should be present");
2281
2282 assert_eq!(info.len(), 21);
2284 assert_eq!(info[&Ustr::from("id")], Ustr::from("trade-0xabcdef1234"));
2285 assert_eq!(info[&Ustr::from("fee_rate_bps")], Ustr::from("0"));
2286 assert_eq!(
2287 info[&Ustr::from("transaction_hash")],
2288 Ustr::from("0xabcdef1234567890abcdef1234567890abcdef1234567890abcdef1234567890ab")
2289 );
2290 assert_eq!(info[&Ustr::from("bucket_index")], Ustr::from("1"));
2292 assert_eq!(info[&Ustr::from("size")], Ustr::from("25.0"));
2293 assert_eq!(
2294 info[&Ustr::from("taker_order_id")],
2295 Ustr::from("0x1234567890abcdef1234567890abcdef1234567890abcdef1234567890abcdef12")
2296 );
2297 assert_eq!(info[&Ustr::from("type")], Ustr::from("TRADE"));
2299 let maker_orders = info[&Ustr::from("maker_orders")].as_str();
2301 assert!(maker_orders.starts_with('['));
2302 assert!(maker_orders.contains("order_id"));
2303
2304 let empty_hash_trade: PolymarketUserTrade = load("ws_user_trade_msg.json");
2305 let empty_hash_info =
2306 trade_fill_info(&empty_hash_trade).expect("empty hash info should be present");
2307 assert!(!empty_hash_info.contains_key(&Ustr::from("transaction_hash")));
2308 }
2309
2310 #[rstest]
2311 fn test_dispatch_order_message_buffers_when_not_accepted() {
2312 let order: PolymarketUserOrder = load("ws_user_order_placement.json");
2313 let instrument = test_instrument();
2314
2315 let token_instruments = AtomicMap::new();
2316 token_instruments.insert(order.asset_id, instrument);
2317
2318 let fill_tracker = OrderFillTrackerMap::new();
2319 let pending_submits = PendingSubmitTracker::default();
2320 let order_contexts = OrderContextRegistry::default();
2321 let emitter = test_emitter();
2322
2323 let ctx = WsDispatchContext {
2324 signer_type: PolymarketSignerType::Owner,
2325 token_instruments: &token_instruments,
2326 fill_tracker: &fill_tracker,
2327 pending_submits: &pending_submits,
2328 order_contexts: &order_contexts,
2329 emitter: &emitter,
2330 account_id: AccountId::from("POLY-001"),
2331 clock: nautilus_core::time::get_atomic_clock_realtime(),
2332 user_address: "0xtest",
2333 user_api_key: "test-key",
2334 };
2335 let mut state = WsDispatchState::default();
2336
2337 let result = dispatch_user_message(&UserWsMessage::Order(order.clone()), &ctx, &mut state);
2338 assert!(result.is_none());
2339
2340 let venue_order_id = VenueOrderId::from(order.id.as_str());
2342 assert!(fill_tracker.has_pending_report(&venue_order_id));
2343 }
2344
2345 #[rstest]
2346 fn test_dispatch_order_message_ignores_missing_lifecycle_fields() {
2347 let order: PolymarketUserOrder = load("ws_user_order_placement.json");
2348 let instrument = test_instrument();
2349 let token_instruments = AtomicMap::new();
2350 token_instruments.insert(order.asset_id, instrument);
2351 let fill_tracker = OrderFillTrackerMap::new();
2352 let pending_submits = PendingSubmitTracker::default();
2353 let order_contexts = OrderContextRegistry::default();
2354 let emitter = test_emitter();
2355 let ctx = WsDispatchContext {
2356 signer_type: PolymarketSignerType::Owner,
2357 token_instruments: &token_instruments,
2358 fill_tracker: &fill_tracker,
2359 pending_submits: &pending_submits,
2360 order_contexts: &order_contexts,
2361 emitter: &emitter,
2362 account_id: AccountId::from("POLY-001"),
2363 clock: nautilus_core::time::get_atomic_clock_realtime(),
2364 user_address: "0xtest",
2365 user_api_key: "test-key",
2366 };
2367 let venue_order_id = VenueOrderId::from(order.id.as_str());
2368
2369 let mut missing_status = order.clone();
2370 missing_status.status = None;
2371 dispatch_user_message(
2372 &UserWsMessage::Order(missing_status),
2373 &ctx,
2374 &mut WsDispatchState::default(),
2375 );
2376 assert!(!fill_tracker.has_pending_report(&venue_order_id));
2377
2378 let mut missing_order_type = order;
2379 missing_order_type.order_type = None;
2380 dispatch_user_message(
2381 &UserWsMessage::Order(missing_order_type),
2382 &ctx,
2383 &mut WsDispatchState::default(),
2384 );
2385 assert!(!fill_tracker.has_pending_report(&venue_order_id));
2386 }
2387
2388 #[rstest]
2389 fn test_dispatch_order_message_uses_pending_submit_client_order_id() {
2390 let order: PolymarketUserOrder = load("ws_user_order_placement.json");
2391 let instrument = test_instrument();
2392
2393 let token_instruments = AtomicMap::new();
2394 token_instruments.insert(order.asset_id, instrument);
2395
2396 let fill_tracker = OrderFillTrackerMap::new();
2397 let pending_submits = PendingSubmitTracker::default();
2398 let order_contexts = OrderContextRegistry::default();
2399 let mut emitter = test_emitter();
2400 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2401 emitter.set_sender(sender);
2402
2403 let venue_order_id = VenueOrderId::from(order.id.as_str());
2404 let client_order_id = ClientOrderId::from("O-UNKNOWN-SUBMIT");
2405 pending_submits.insert(venue_order_id, client_order_id);
2406 register_context(
2407 &order_contexts,
2408 venue_order_id,
2409 test_instrument().id(),
2410 "O-UNKNOWN-SUBMIT",
2411 );
2412
2413 let ctx = WsDispatchContext {
2414 signer_type: PolymarketSignerType::Owner,
2415 token_instruments: &token_instruments,
2416 fill_tracker: &fill_tracker,
2417 pending_submits: &pending_submits,
2418 order_contexts: &order_contexts,
2419 emitter: &emitter,
2420 account_id: AccountId::from("POLY-001"),
2421 clock: nautilus_core::time::get_atomic_clock_realtime(),
2422 user_address: "0xtest",
2423 user_api_key: "test-key",
2424 };
2425 let mut state = WsDispatchState::default();
2426
2427 let _ = dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
2428
2429 let event = receiver.try_recv().expect("expected accepted event");
2431 match event {
2432 ExecutionEvent::Order(OrderEventAny::Accepted(accepted)) => {
2433 assert_eq!(accepted.client_order_id, client_order_id);
2434 }
2435 other => panic!("Expected accepted event, was {other:?}"),
2436 }
2437
2438 assert!(!fill_tracker.has_pending_report(&venue_order_id));
2439 }
2440
2441 #[rstest]
2442 #[case(PolymarketSignerType::Owner, 1)]
2443 #[case(PolymarketSignerType::Session, 0)]
2444 fn test_dispatch_maker_fill_owned_by_case_variant_address(
2445 #[case] signer_type: PolymarketSignerType,
2446 #[case] expected_fills: usize,
2447 ) {
2448 let mut trade: PolymarketUserTrade = load("ws_user_trade.json");
2449 trade.trader_side = PolymarketLiquiditySide::Maker;
2450 let configured_address = trade.maker_orders[0].maker_address.clone();
2451 let case_variant_address = configured_address
2452 .to_ascii_uppercase()
2453 .replacen("0X", "0x", 1);
2454 assert_ne!(case_variant_address, configured_address);
2455 trade.maker_orders[0].maker_address = case_variant_address;
2456 let foreign_api_key = "ffffffff-ffff-ffff-ffff-ffffffffffff";
2457 assert_ne!(trade.maker_orders[0].owner, foreign_api_key);
2458
2459 let venue_order_id = VenueOrderId::from(trade.maker_orders[0].order_id.as_str());
2460 let token_instruments = AtomicMap::new();
2461 token_instruments.insert(trade.maker_orders[0].asset_id, test_instrument());
2462 let fill_tracker = OrderFillTrackerMap::new();
2463 let pending_submits = PendingSubmitTracker::default();
2464 let order_contexts = OrderContextRegistry::default();
2465 let emitter = test_emitter();
2466 let ctx = WsDispatchContext {
2467 signer_type,
2468 token_instruments: &token_instruments,
2469 fill_tracker: &fill_tracker,
2470 pending_submits: &pending_submits,
2471 order_contexts: &order_contexts,
2472 emitter: &emitter,
2473 account_id: AccountId::from("POLY-001"),
2474 clock: nautilus_core::time::get_atomic_clock_realtime(),
2475 user_address: &configured_address,
2476 user_api_key: foreign_api_key,
2477 };
2478 let mut state = WsDispatchState::default();
2479
2480 let _ = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2481
2482 let fills = fill_tracker.pending_fills_for(&venue_order_id);
2483 assert_eq!(fills.len(), expected_fills);
2484
2485 for fill in fills {
2486 assert_eq!(fill.venue_order_id, venue_order_id);
2487 }
2488 }
2489
2490 #[rstest]
2491 #[case(dec!(-1))]
2492 #[case(Decimal::MAX)]
2493 fn test_dispatch_maker_numeric_failure_preserves_replay(#[case] invalid_amount: Decimal) {
2494 let mut trade: PolymarketUserTrade = load("ws_user_trade.json");
2495 trade.trader_side = PolymarketLiquiditySide::Maker;
2496 let configured_address = trade.maker_orders[0].maker_address.clone();
2497 let foreign_api_key = "ffffffff-ffff-ffff-ffff-ffffffffffff";
2498 assert_ne!(trade.maker_orders[0].owner, foreign_api_key);
2499
2500 let venue_order_id = VenueOrderId::from(trade.maker_orders[0].order_id.as_str());
2501 let token_instruments = AtomicMap::new();
2502 token_instruments.insert(trade.maker_orders[0].asset_id, test_instrument());
2503 let fill_tracker = OrderFillTrackerMap::new();
2504 let pending_submits = PendingSubmitTracker::default();
2505 let order_contexts = OrderContextRegistry::default();
2506 let emitter = test_emitter();
2507 let ctx = WsDispatchContext {
2508 signer_type: PolymarketSignerType::Owner,
2509 token_instruments: &token_instruments,
2510 fill_tracker: &fill_tracker,
2511 pending_submits: &pending_submits,
2512 order_contexts: &order_contexts,
2513 emitter: &emitter,
2514 account_id: AccountId::from("POLY-001"),
2515 clock: nautilus_core::time::get_atomic_clock_realtime(),
2516 user_address: &configured_address,
2517 user_api_key: foreign_api_key,
2518 };
2519 let mut state = WsDispatchState::default();
2520
2521 let expected = trade.maker_orders[0].matched_amount;
2522 let mut invalid_trade = trade.clone();
2523 invalid_trade.maker_orders[0].matched_amount = invalid_amount;
2524 dispatch_user_message(&UserWsMessage::Trade(invalid_trade), &ctx, &mut state);
2525 assert_eq!(fill_tracker.pending_fills_for(&venue_order_id).len(), 0);
2526
2527 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2528
2529 let fills = fill_tracker.pending_fills_for(&venue_order_id);
2530 assert_eq!(fills.len(), 1);
2531 assert_eq!(fills[0].venue_order_id, venue_order_id);
2532 assert_eq!(fills[0].last_qty.as_decimal(), expected);
2533 }
2534
2535 #[rstest]
2536 fn test_dispatch_trade_dedup() {
2537 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2538 let instrument = test_instrument();
2539
2540 let token_instruments = AtomicMap::new();
2541 token_instruments.insert(trade.asset_id, instrument);
2542
2543 let fill_tracker = OrderFillTrackerMap::new();
2544 let pending_submits = PendingSubmitTracker::default();
2545 let order_contexts = OrderContextRegistry::default();
2546 let emitter = test_emitter();
2547
2548 let ctx = WsDispatchContext {
2549 signer_type: PolymarketSignerType::Owner,
2550 token_instruments: &token_instruments,
2551 fill_tracker: &fill_tracker,
2552 pending_submits: &pending_submits,
2553 order_contexts: &order_contexts,
2554 emitter: &emitter,
2555 account_id: AccountId::from("POLY-001"),
2556 clock: nautilus_core::time::get_atomic_clock_realtime(),
2557 user_address: "0xtest",
2558 user_api_key: "test-key",
2559 };
2560 let mut state = WsDispatchState::default();
2561
2562 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
2563
2564 let _ = dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
2566 assert_eq!(fill_tracker.pending_fills_for(&venue_order_id).len(), 1);
2567
2568 let _ = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2570 assert_eq!(fill_tracker.pending_fills_for(&venue_order_id).len(), 1);
2571 }
2572
2573 #[rstest]
2574 fn test_dispatch_taker_commission_failure_preserves_replay_state() {
2575 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2576 let valid_instrument = test_instrument();
2577 let mut invalid_instrument = valid_instrument.clone();
2578 let InstrumentAny::BinaryOption(binary_option) = &mut invalid_instrument else {
2579 panic!("expected binary option test instrument");
2580 };
2581 binary_option.taker_fee =
2582 Decimal::from_i128_with_scale(100_000_000_000_000_000_000_000_000i128, 0);
2583
2584 let token_instruments = AtomicMap::new();
2585 token_instruments.insert(trade.asset_id, invalid_instrument);
2586 let fill_tracker = OrderFillTrackerMap::new();
2587 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
2588 fill_tracker.register(
2589 venue_order_id,
2590 Quantity::from("100"),
2591 OrderSide::Buy,
2592 valid_instrument.id(),
2593 valid_instrument.size_precision(),
2594 valid_instrument.price_precision(),
2595 );
2596 let pending_submits = PendingSubmitTracker::default();
2597 let order_contexts = OrderContextRegistry::default();
2598 register_context(
2599 &order_contexts,
2600 venue_order_id,
2601 valid_instrument.id(),
2602 "O-COMMISSION-REPLAY",
2603 );
2604 order_contexts.mark_accepted(venue_order_id);
2605 let mut emitter = test_emitter();
2606 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2607 emitter.set_sender(sender);
2608 let ctx = WsDispatchContext {
2609 signer_type: PolymarketSignerType::Owner,
2610 token_instruments: &token_instruments,
2611 fill_tracker: &fill_tracker,
2612 pending_submits: &pending_submits,
2613 order_contexts: &order_contexts,
2614 emitter: &emitter,
2615 account_id: AccountId::from("POLY-001"),
2616 clock: nautilus_core::time::get_atomic_clock_realtime(),
2617 user_address: "0xtest",
2618 user_api_key: "test-key",
2619 };
2620 let mut state = WsDispatchState::default();
2621 let dedup_key = format!("{}-{}", trade.id, trade.taker_order_id);
2622
2623 let failed = dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
2624
2625 assert!(failed.is_none());
2626 assert!(!state.processed_fills.contains(&dedup_key));
2627 assert!(!state.confirmed_trades.contains(&trade.id));
2628 assert!(!fill_tracker.is_trade_confirmed(&dedup_key));
2629 assert_eq!(
2630 fill_tracker.get_cumulative_filled(&venue_order_id),
2631 Some(Quantity::zero(valid_instrument.size_precision()))
2632 );
2633 assert!(receiver.try_recv().is_err());
2634
2635 token_instruments.insert(trade.asset_id, valid_instrument);
2636 let replay = dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
2637 let emitted = receiver.try_recv().expect("valid replay emits one fill");
2638 let duplicate = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2639
2640 assert!(replay.is_some());
2641 assert!(duplicate.is_some());
2642 assert!(state.processed_fills.contains(&dedup_key));
2643 assert!(
2644 state
2645 .confirmed_trades
2646 .contains(&"trade-0xabcdef1234".to_string())
2647 );
2648 assert!(fill_tracker.is_trade_confirmed(&dedup_key));
2649 assert!(matches!(
2650 emitted,
2651 ExecutionEvent::Order(OrderEventAny::Filled(_))
2652 ));
2653 assert_eq!(
2654 fill_tracker.get_cumulative_filled(&venue_order_id),
2655 Some(Quantity::from("25.0"))
2656 );
2657 assert!(receiver.try_recv().is_err());
2658 }
2659
2660 #[rstest]
2661 fn test_dispatch_trade_replays_after_instrument_becomes_available() {
2662 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2663 let instrument = test_instrument();
2664 let token_instruments = AtomicMap::new();
2665 let fill_tracker = OrderFillTrackerMap::new();
2666 let pending_submits = PendingSubmitTracker::default();
2667 let order_contexts = OrderContextRegistry::default();
2668 let emitter = test_emitter();
2669 let ctx = WsDispatchContext {
2670 signer_type: PolymarketSignerType::Owner,
2671 token_instruments: &token_instruments,
2672 fill_tracker: &fill_tracker,
2673 pending_submits: &pending_submits,
2674 order_contexts: &order_contexts,
2675 emitter: &emitter,
2676 account_id: AccountId::from("POLY-001"),
2677 clock: nautilus_core::time::get_atomic_clock_realtime(),
2678 user_address: "0xtest",
2679 user_api_key: "test-key",
2680 };
2681 let mut state = WsDispatchState::default();
2682 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
2683
2684 let first_result =
2685 dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
2686 token_instruments.insert(trade.asset_id, instrument);
2687 let replay_result = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2688
2689 assert!(first_result.is_none());
2690 assert!(replay_result.is_some());
2691 assert_eq!(fill_tracker.pending_fills_for(&venue_order_id).len(), 1);
2692 }
2693
2694 #[rstest]
2695 #[case(crate::common::enums::PolymarketTradeStatus::Mined)]
2696 #[case(crate::common::enums::PolymarketTradeStatus::Retrying)]
2697 fn test_dispatch_trade_ignores_pending_settlement_status(
2698 #[case] status: crate::common::enums::PolymarketTradeStatus,
2699 ) {
2700 let mut trade: PolymarketUserTrade = load("ws_user_trade.json");
2701 trade.status = status;
2702 let instrument = test_instrument();
2703 let token_instruments = AtomicMap::new();
2704 token_instruments.insert(trade.asset_id, instrument);
2705 let fill_tracker = OrderFillTrackerMap::new();
2706 let pending_submits = PendingSubmitTracker::default();
2707 let order_contexts = OrderContextRegistry::default();
2708 let emitter = test_emitter();
2709 let ctx = WsDispatchContext {
2710 signer_type: PolymarketSignerType::Owner,
2711 token_instruments: &token_instruments,
2712 fill_tracker: &fill_tracker,
2713 pending_submits: &pending_submits,
2714 order_contexts: &order_contexts,
2715 emitter: &emitter,
2716 account_id: AccountId::from("POLY-001"),
2717 clock: nautilus_core::time::get_atomic_clock_realtime(),
2718 user_address: "0xtest",
2719 user_api_key: "test-key",
2720 };
2721 let mut state = WsDispatchState::default();
2722 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
2723
2724 let result = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2725
2726 assert!(result.is_none());
2727 assert!(fill_tracker.pending_fills_for(&venue_order_id).is_empty());
2728 }
2729
2730 #[rstest]
2731 fn test_dispatch_matched_trade_emits_fill_and_failed_trade_voids_it() {
2732 let mut trade: PolymarketUserTrade = load("ws_user_trade.json");
2733 trade.status = crate::common::enums::PolymarketTradeStatus::Matched;
2734 let instrument = test_instrument();
2735 let token_instruments = AtomicMap::new();
2736 token_instruments.insert(trade.asset_id, instrument.clone());
2737 let fill_tracker = OrderFillTrackerMap::new();
2738 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
2739 fill_tracker.register(
2740 venue_order_id,
2741 Quantity::from("100"),
2742 OrderSide::Buy,
2743 instrument.id(),
2744 instrument.size_precision(),
2745 instrument.price_precision(),
2746 );
2747 let pending_submits = PendingSubmitTracker::default();
2748 let order_contexts = OrderContextRegistry::default();
2749 register_context(
2750 &order_contexts,
2751 venue_order_id,
2752 instrument.id(),
2753 "O-MATCHED-FAILED",
2754 );
2755 order_contexts.mark_accepted(venue_order_id);
2756 let mut emitter = test_emitter();
2757 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2758 emitter.set_sender(sender);
2759 let ctx = WsDispatchContext {
2760 signer_type: PolymarketSignerType::Owner,
2761 token_instruments: &token_instruments,
2762 fill_tracker: &fill_tracker,
2763 pending_submits: &pending_submits,
2764 order_contexts: &order_contexts,
2765 emitter: &emitter,
2766 account_id: AccountId::from("POLY-001"),
2767 clock: nautilus_core::time::get_atomic_clock_realtime(),
2768 user_address: "0xtest",
2769 user_api_key: "test-key",
2770 };
2771 let mut state = WsDispatchState::default();
2772
2773 let matched = dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
2774 let filled = match receiver.try_recv().unwrap() {
2775 ExecutionEvent::Order(OrderEventAny::Filled(event)) => event,
2776 other => panic!("expected matched fill, was {other:?}"),
2777 };
2778 trade.status = crate::common::enums::PolymarketTradeStatus::Failed;
2779 let failed = dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
2780 let voided = match receiver.try_recv().unwrap() {
2781 ExecutionEvent::Order(OrderEventAny::FillVoided(event)) => event,
2782 other => panic!("expected failed fill correction, was {other:?}"),
2783 };
2784
2785 let mut failed_first_state = WsDispatchState::default();
2786 let failed_first = dispatch_user_message(
2787 &UserWsMessage::Trade(trade.clone()),
2788 &ctx,
2789 &mut failed_first_state,
2790 );
2791 let dedup_key = format!("{}-{}", trade.id, trade.taker_order_id);
2792 trade.status = crate::common::enums::PolymarketTradeStatus::Matched;
2793 let matched_after_failure =
2794 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut failed_first_state);
2795
2796 assert!(matched.is_none());
2797 assert!(failed.is_some());
2798 assert!(failed_first.is_some());
2799 assert!(matched_after_failure.is_none());
2800 assert_eq!(voided.trade_id, filled.trade_id);
2801 assert_eq!(voided.voided_qty, filled.last_qty);
2802 assert_eq!(voided.commission_voided, filled.commission);
2803 assert_eq!(voided.last_px, filled.last_px);
2804 assert!(!voided.is_reopened);
2805 assert_eq!(voided.causation_id, Some(filled.event_id));
2806 assert_eq!(
2807 fill_tracker.get_cumulative_filled(&venue_order_id),
2808 Some(Quantity::zero(instrument.size_precision()))
2809 );
2810 assert!(failed_first_state.processed_fills.contains(&dedup_key));
2811 assert!(failed_first_state.is_voided_trade(&dedup_key));
2812 assert!(receiver.try_recv().is_err());
2813 }
2814
2815 #[rstest]
2816 fn test_dispatch_trade_uses_pending_submit_client_order_id() {
2817 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2818 let instrument = test_instrument();
2819
2820 let token_instruments = AtomicMap::new();
2821 token_instruments.insert(trade.asset_id, instrument);
2822
2823 let fill_tracker = OrderFillTrackerMap::new();
2824 let pending_submits = PendingSubmitTracker::default();
2825 let order_contexts = OrderContextRegistry::default();
2826 let emitter = test_emitter();
2827
2828 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
2829 let client_order_id = ClientOrderId::from("O-UNKNOWN-FILL");
2830 pending_submits.insert(venue_order_id, client_order_id);
2831
2832 let ctx = WsDispatchContext {
2833 signer_type: PolymarketSignerType::Owner,
2834 token_instruments: &token_instruments,
2835 fill_tracker: &fill_tracker,
2836 pending_submits: &pending_submits,
2837 order_contexts: &order_contexts,
2838 emitter: &emitter,
2839 account_id: AccountId::from("POLY-001"),
2840 clock: nautilus_core::time::get_atomic_clock_realtime(),
2841 user_address: "0xtest",
2842 user_api_key: "test-key",
2843 };
2844 let mut state = WsDispatchState::default();
2845
2846 let _ = dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
2847
2848 let fills = fill_tracker.pending_fills_for(&venue_order_id);
2849 assert_eq!(fills[0].client_order_id, Some(client_order_id));
2850 }
2851
2852 #[rstest]
2853 fn test_dispatch_late_fill_stays_tracked_after_later_registrations() {
2854 let trade: PolymarketUserTrade = load("ws_user_trade.json");
2855 let market: GammaMarket = load("gamma_market_sports_market_money_line.json");
2856 let defs = parse_gamma_market(&market).unwrap();
2857 let instrument =
2858 create_instrument_from_def(&defs[0], UnixNanos::from(1_000_000_000u64)).unwrap();
2859 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
2860
2861 let token_instruments = AtomicMap::new();
2862 token_instruments.insert(trade.asset_id, instrument.clone());
2863
2864 let fill_tracker = OrderFillTrackerMap::new();
2865 fill_tracker.register(
2866 venue_order_id,
2867 Quantity::from("100"),
2868 OrderSide::Buy,
2869 instrument.id(),
2870 instrument.size_precision(),
2871 instrument.price_precision(),
2872 );
2873
2874 let pending_submits = PendingSubmitTracker::default();
2875 let order_contexts = OrderContextRegistry::default();
2876 register_context(
2877 &order_contexts,
2878 venue_order_id,
2879 instrument.id(),
2880 "O-LATE-FILL",
2881 );
2882 order_contexts.mark_accepted(venue_order_id);
2883 assert!(order_contexts.get(&venue_order_id).is_some());
2884
2885 for index in 0..10_000 {
2886 let later_venue_order_id = VenueOrderId::from(format!("V-LATER-{index}").as_str());
2887 let later_client_order_id = format!("O-LATER-{index}");
2888 register_context(
2889 &order_contexts,
2890 later_venue_order_id,
2891 instrument.id(),
2892 &later_client_order_id,
2893 );
2894 order_contexts.mark_accepted(later_venue_order_id);
2895 fill_tracker.register(
2896 later_venue_order_id,
2897 Quantity::from("1"),
2898 OrderSide::Sell,
2899 instrument.id(),
2900 instrument.size_precision(),
2901 instrument.price_precision(),
2902 );
2903 }
2904 assert!(order_contexts.get(&venue_order_id).is_some());
2905 assert!(fill_tracker.contains(&venue_order_id));
2906
2907 let mut emitter = test_emitter();
2908 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2909 emitter.set_sender(sender);
2910
2911 let ctx = WsDispatchContext {
2912 signer_type: PolymarketSignerType::Owner,
2913 token_instruments: &token_instruments,
2914 fill_tracker: &fill_tracker,
2915 pending_submits: &pending_submits,
2916 order_contexts: &order_contexts,
2917 emitter: &emitter,
2918 account_id: AccountId::from("POLY-001"),
2919 clock: nautilus_core::time::get_atomic_clock_realtime(),
2920 user_address: "0xtest",
2921 user_api_key: "test-key",
2922 };
2923 let mut state = WsDispatchState::default();
2924
2925 dispatch_user_message(&UserWsMessage::Trade(trade.clone()), &ctx, &mut state);
2926
2927 let event = receiver.try_recv().expect("expected tracked late fill");
2928 let ExecutionEvent::Order(OrderEventAny::Filled(filled)) = event else {
2929 panic!("expected tracked OrderFilled after later registrations, was {event:?}");
2930 };
2931
2932 assert_eq!(filled.client_order_id, ClientOrderId::from("O-LATE-FILL"));
2933 assert_eq!(filled.venue_order_id, venue_order_id);
2934 assert_eq!(filled.trade_id, TradeId::from(trade.id.as_str()));
2935 assert_eq!(filled.instrument_id, instrument.id());
2936 assert_eq!(
2937 filled.last_qty.as_decimal(),
2938 Decimal::from_str_exact(&trade.size).unwrap()
2939 );
2940 assert_eq!(
2941 filled.last_px.as_decimal(),
2942 Decimal::from_str_exact(&trade.price).unwrap()
2943 );
2944 assert_eq!(filled.order_side, OrderSide::Buy);
2945 assert_eq!(filled.liquidity_side, LiquiditySide::Taker);
2946 let commission = filled.commission.expect("tracked fill has commission");
2947 assert_eq!(commission.as_decimal(), dec!(0.1875));
2948 assert_eq!(commission.currency, Currency::pUSD());
2949 assert!(receiver.try_recv().is_err());
2950 }
2951
2952 #[rstest]
2953 fn test_dispatch_order_matched_caps_filled_qty_when_no_trades_tracked() {
2954 let order: PolymarketUserOrder = load("ws_user_order_matched.json");
2955 let instrument = test_instrument();
2956
2957 let token_instruments = AtomicMap::new();
2958 token_instruments.insert(order.asset_id, instrument.clone());
2959
2960 let fill_tracker = OrderFillTrackerMap::new();
2961 let venue_order_id = VenueOrderId::from(order.id.as_str());
2962
2963 fill_tracker.register(
2965 venue_order_id,
2966 Quantity::from("100"),
2967 OrderSide::Buy,
2968 instrument.id(),
2969 instrument.size_precision(),
2970 instrument.price_precision(),
2971 );
2972
2973 let pending_submits = PendingSubmitTracker::default();
2974 let order_contexts = OrderContextRegistry::default();
2977 let mut emitter = test_emitter();
2978 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
2979 emitter.set_sender(sender);
2980
2981 let ctx = WsDispatchContext {
2982 signer_type: PolymarketSignerType::Owner,
2983 token_instruments: &token_instruments,
2984 fill_tracker: &fill_tracker,
2985 pending_submits: &pending_submits,
2986 order_contexts: &order_contexts,
2987 emitter: &emitter,
2988 account_id: AccountId::from("POLY-001"),
2989 clock: nautilus_core::time::get_atomic_clock_realtime(),
2990 user_address: "0xtest",
2991 user_api_key: "test-key",
2992 };
2993 let mut state = WsDispatchState::default();
2994
2995 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
2996
2997 let event = receiver.try_recv().expect("Expected report");
2998 match event {
2999 ExecutionEvent::Report(report) => match report {
3000 ExecutionReport::Order(order_report) => {
3001 assert_eq!(order_report.filled_qty, Quantity::from("0"));
3002 }
3003 other => panic!("Expected order report, was {other:?}"),
3004 },
3005 other => panic!("Expected report event, was {other:?}"),
3006 }
3007 }
3008
3009 #[rstest]
3010 fn test_dispatch_order_matched_uses_tracked_fills_for_filled_qty() {
3011 let order: PolymarketUserOrder = load("ws_user_order_matched.json");
3012 let instrument = test_instrument();
3013
3014 let token_instruments = AtomicMap::new();
3015 token_instruments.insert(order.asset_id, instrument.clone());
3016
3017 let fill_tracker = OrderFillTrackerMap::new();
3018 let venue_order_id = VenueOrderId::from(order.id.as_str());
3019
3020 fill_tracker.register(
3022 venue_order_id,
3023 Quantity::from("100"),
3024 OrderSide::Buy,
3025 instrument.id(),
3026 instrument.size_precision(),
3027 instrument.price_precision(),
3028 );
3029 fill_tracker.record_fill(&venue_order_id, Quantity::new(50.0, 6));
3030
3031 let pending_submits = PendingSubmitTracker::default();
3032 let order_contexts = OrderContextRegistry::default();
3035 let mut emitter = test_emitter();
3036 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3037 emitter.set_sender(sender);
3038
3039 let ctx = WsDispatchContext {
3040 signer_type: PolymarketSignerType::Owner,
3041 token_instruments: &token_instruments,
3042 fill_tracker: &fill_tracker,
3043 pending_submits: &pending_submits,
3044 order_contexts: &order_contexts,
3045 emitter: &emitter,
3046 account_id: AccountId::from("POLY-001"),
3047 clock: nautilus_core::time::get_atomic_clock_realtime(),
3048 user_address: "0xtest",
3049 user_api_key: "test-key",
3050 };
3051 let mut state = WsDispatchState::default();
3052
3053 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
3054
3055 let event = receiver.try_recv().expect("Expected report");
3056 match event {
3057 ExecutionEvent::Report(report) => match report {
3058 ExecutionReport::Order(order_report) => {
3059 assert_eq!(order_report.filled_qty, Quantity::from("50"));
3060 }
3061 other => panic!("Expected order report, was {other:?}"),
3062 },
3063 other => panic!("Expected report event, was {other:?}"),
3064 }
3065 }
3066
3067 #[rstest]
3068 fn test_dispatch_order_matched_normalizes_quantity_without_fill() {
3069 let order: PolymarketUserOrder = load("ws_user_order_matched.json");
3070 let instrument = test_instrument();
3071
3072 let token_instruments = AtomicMap::new();
3073 token_instruments.insert(order.asset_id, instrument.clone());
3074
3075 let fill_tracker = OrderFillTrackerMap::new();
3076 let venue_order_id = VenueOrderId::from(order.id.as_str());
3077 fill_tracker.register(
3078 venue_order_id,
3079 Quantity::from("100"),
3080 OrderSide::Buy,
3081 instrument.id(),
3082 instrument.size_precision(),
3083 instrument.price_precision(),
3084 );
3085 fill_tracker.record_fill(&venue_order_id, Quantity::new(99.995, 6));
3086
3087 let pending_submits = PendingSubmitTracker::default();
3088 let order_contexts = OrderContextRegistry::default();
3089 register_context(
3090 &order_contexts,
3091 venue_order_id,
3092 instrument.id(),
3093 "O-MATCHED",
3094 );
3095 order_contexts.mark_accepted(venue_order_id);
3096 let mut emitter = test_emitter();
3097 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3098 emitter.set_sender(sender);
3099
3100 let clock = Box::leak(Box::new(AtomicTime::new(
3101 false,
3102 UnixNanos::from(2_000_000_000u64),
3103 )));
3104
3105 let ctx = WsDispatchContext {
3106 signer_type: PolymarketSignerType::Owner,
3107 token_instruments: &token_instruments,
3108 fill_tracker: &fill_tracker,
3109 pending_submits: &pending_submits,
3110 order_contexts: &order_contexts,
3111 emitter: &emitter,
3112 account_id: AccountId::from("POLY-001"),
3113 clock,
3114 user_address: "0xtest",
3115 user_api_key: "test-key",
3116 };
3117 let mut state = WsDispatchState::default();
3118 state.confirmed_trades.add("trade-0xfill1".to_string());
3119
3120 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
3121
3122 let event = receiver.try_recv().expect("expected quantity update");
3123 match event {
3124 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
3125 assert_eq!(
3126 updated.ts_event,
3127 UnixNanos::from(1_703_875_201_000_000_000u64)
3128 );
3129 assert_eq!(updated.ts_init, UnixNanos::from(2_000_000_000u64));
3130 assert_eq!(updated.quantity, Quantity::new(99.995, 6));
3131 assert!(updated.reconciliation);
3132 }
3133 other => panic!("expected updated event, was {other:?}"),
3134 }
3135 assert!(receiver.try_recv().is_err());
3136 }
3137
3138 #[rstest]
3139 fn test_confirmed_trade_normalizes_pending_matched_quantity() {
3140 let mut order: PolymarketUserOrder = load("ws_user_order_matched.json");
3141 let mut trade: PolymarketUserTrade = load("ws_user_trade.json");
3142 let instrument = test_instrument();
3143 order.associate_trades = Some(vec![trade.id.clone()]);
3144 trade.size = "99.995".to_string();
3145 trade.price = order.price.clone();
3146
3147 let token_instruments = AtomicMap::new();
3148 token_instruments.insert(order.asset_id, instrument.clone());
3149 let fill_tracker = OrderFillTrackerMap::new();
3150 let venue_order_id = VenueOrderId::from(order.id.as_str());
3151 fill_tracker.register(
3152 venue_order_id,
3153 Quantity::from("100"),
3154 OrderSide::Buy,
3155 instrument.id(),
3156 instrument.size_precision(),
3157 instrument.price_precision(),
3158 );
3159 let pending_submits = PendingSubmitTracker::default();
3160 let order_contexts = OrderContextRegistry::default();
3161 register_context(
3162 &order_contexts,
3163 venue_order_id,
3164 instrument.id(),
3165 "O-CONFIRMED-DUST",
3166 );
3167 order_contexts.mark_accepted(venue_order_id);
3168 let mut emitter = test_emitter();
3169 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3170 emitter.set_sender(sender);
3171 let ctx = WsDispatchContext {
3172 signer_type: PolymarketSignerType::Owner,
3173 token_instruments: &token_instruments,
3174 fill_tracker: &fill_tracker,
3175 pending_submits: &pending_submits,
3176 order_contexts: &order_contexts,
3177 emitter: &emitter,
3178 account_id: AccountId::from("POLY-001"),
3179 clock: nautilus_core::time::get_atomic_clock_realtime(),
3180 user_address: "0xtest",
3181 user_api_key: "test-key",
3182 };
3183 let mut state = WsDispatchState::default();
3184
3185 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
3186 assert!(receiver.try_recv().is_err());
3187
3188 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
3189
3190 let real_fill = receiver.try_recv().expect("expected confirmed venue fill");
3191 let normalized = receiver
3192 .try_recv()
3193 .expect("expected quantity normalization");
3194
3195 match (real_fill, normalized) {
3196 (
3197 ExecutionEvent::Order(OrderEventAny::Filled(real)),
3198 ExecutionEvent::Order(OrderEventAny::Updated(updated)),
3199 ) => {
3200 assert_eq!(real.last_qty, Quantity::from("99.995"));
3201 assert_eq!(updated.quantity, Quantity::from("99.995"));
3202 assert!(updated.reconciliation);
3203 }
3204 other => panic!("expected fill then quantity update, was {other:?}"),
3205 }
3206 assert!(receiver.try_recv().is_err());
3207 }
3208
3209 #[rstest]
3210 fn test_cancel_reemitted_after_fill_for_canceled_order() {
3211 let cancel_order: PolymarketUserOrder = load("ws_user_order_cancellation.json");
3212 let trade: PolymarketUserTrade = load("ws_user_trade.json");
3213 let instrument = test_instrument();
3214
3215 let token_instruments = AtomicMap::new();
3216 token_instruments.insert(cancel_order.asset_id, instrument.clone());
3217
3218 let fill_tracker = OrderFillTrackerMap::new();
3219 let venue_order_id = VenueOrderId::from(cancel_order.id.as_str());
3220
3221 fill_tracker.register(
3223 venue_order_id,
3224 Quantity::from("100"),
3225 OrderSide::Buy,
3226 instrument.id(),
3227 instrument.size_precision(),
3228 instrument.price_precision(),
3229 );
3230
3231 let pending_submits = PendingSubmitTracker::default();
3232 let order_contexts = OrderContextRegistry::default();
3233 register_context(&order_contexts, venue_order_id, instrument.id(), "O-CANCEL");
3234 order_contexts.mark_accepted(venue_order_id);
3235 let mut emitter = test_emitter();
3236 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3237 emitter.set_sender(sender);
3238
3239 let ctx = WsDispatchContext {
3240 signer_type: PolymarketSignerType::Owner,
3241 token_instruments: &token_instruments,
3242 fill_tracker: &fill_tracker,
3243 pending_submits: &pending_submits,
3244 order_contexts: &order_contexts,
3245 emitter: &emitter,
3246 account_id: AccountId::from("POLY-001"),
3247 clock: nautilus_core::time::get_atomic_clock_realtime(),
3248 user_address: "0xtest",
3249 user_api_key: "test-key",
3250 };
3251 let mut state = WsDispatchState::default();
3252
3253 dispatch_user_message(&UserWsMessage::Order(cancel_order), &ctx, &mut state);
3255 let cancel_event = receiver.try_recv().expect("Expected canceled event");
3256 match &cancel_event {
3257 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
3258 assert_eq!(c.venue_order_id, Some(venue_order_id));
3259 }
3260 other => panic!("Expected canceled event, was {other:?}"),
3261 }
3262
3263 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
3265
3266 let fill_event = receiver.try_recv().expect("Expected filled event");
3268 match &fill_event {
3269 ExecutionEvent::Order(OrderEventAny::Filled(f)) => {
3270 assert_eq!(f.venue_order_id, venue_order_id);
3271 }
3272 other => panic!("Expected filled event, was {other:?}"),
3273 }
3274
3275 let reemitted_cancel = receiver
3276 .try_recv()
3277 .expect("Expected re-emitted canceled event");
3278
3279 match &reemitted_cancel {
3280 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
3281 assert_eq!(c.venue_order_id, Some(venue_order_id));
3282 }
3283 other => panic!("Expected canceled event, was {other:?}"),
3284 }
3285 }
3286
3287 #[rstest]
3288 fn test_modified_old_leg_suppresses_cancel_but_still_emits_late_fill() {
3289 let cancel_order: PolymarketUserOrder = load("ws_user_order_cancellation.json");
3290 let trade: PolymarketUserTrade = load("ws_user_trade.json");
3291 let instrument = test_instrument();
3292 let token_instruments = AtomicMap::new();
3293 token_instruments.insert(cancel_order.asset_id, instrument.clone());
3294
3295 let fill_tracker = OrderFillTrackerMap::new();
3296 let old_venue_order_id = VenueOrderId::from(cancel_order.id.as_str());
3297 fill_tracker.register(
3298 old_venue_order_id,
3299 Quantity::from("100"),
3300 OrderSide::Buy,
3301 instrument.id(),
3302 instrument.size_precision(),
3303 instrument.price_precision(),
3304 );
3305
3306 let client_order_id = ClientOrderId::from("O-MODIFIED-OLD-LEG");
3307 let pending_submits = PendingSubmitTracker::default();
3308 let order_contexts = OrderContextRegistry::default();
3309 register_context(
3310 &order_contexts,
3311 old_venue_order_id,
3312 instrument.id(),
3313 client_order_id.as_str(),
3314 );
3315 order_contexts.mark_accepted(old_venue_order_id);
3316 let mut emitter = test_emitter();
3317 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3318 emitter.set_sender(sender);
3319 let ctx = WsDispatchContext {
3320 signer_type: PolymarketSignerType::Owner,
3321 token_instruments: &token_instruments,
3322 fill_tracker: &fill_tracker,
3323 pending_submits: &pending_submits,
3324 order_contexts: &order_contexts,
3325 emitter: &emitter,
3326 account_id: AccountId::from("POLY-001"),
3327 clock: nautilus_core::time::get_atomic_clock_realtime(),
3328 user_address: "0xtest",
3329 user_api_key: "test-key",
3330 };
3331
3332 let mut state = WsDispatchState::default();
3333 let replacement_venue_order_id = VenueOrderId::from("0xreplacement");
3334 assert!(state.begin_modify(client_order_id, old_venue_order_id, instrument.id()));
3335 assert!(state.set_modify_replacement(
3336 client_order_id,
3337 replacement_venue_order_id,
3338 Quantity::from("100"),
3339 Quantity::from("100"),
3340 Price::from("0.5"),
3341 ));
3342 assert!(
3343 state
3344 .claim_modify_replacement(replacement_venue_order_id)
3345 .is_some()
3346 );
3347
3348 dispatch_user_message(&UserWsMessage::Order(cancel_order), &ctx, &mut state);
3349 assert!(receiver.try_recv().is_err());
3350
3351 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
3352
3353 match receiver.try_recv().expect("expected late old-leg fill") {
3354 ExecutionEvent::Order(OrderEventAny::Filled(fill)) => {
3355 assert_eq!(fill.client_order_id, client_order_id);
3356 assert_eq!(fill.venue_order_id, old_venue_order_id);
3357 }
3358 other => panic!("expected late old-leg fill, was {other:?}"),
3359 }
3360
3361 assert!(receiver.try_recv().is_err());
3362 }
3363
3364 #[rstest]
3365 fn test_pending_modify_replacement_ws_activity_emits_updated_without_accepted() {
3366 let mut replacement: PolymarketUserOrder = load("ws_user_order_placement.json");
3367 let instrument = test_instrument();
3368 let old_venue_order_id = VenueOrderId::from("0xold-modify-leg");
3369 let replacement_venue_order_id = VenueOrderId::from("0xreplacement-modify-leg");
3370 replacement.id = replacement_venue_order_id.to_string();
3371
3372 let token_instruments = AtomicMap::new();
3373 token_instruments.insert(replacement.asset_id, instrument.clone());
3374 let fill_tracker = OrderFillTrackerMap::new();
3375 fill_tracker.register(
3376 old_venue_order_id,
3377 Quantity::from("100"),
3378 OrderSide::Buy,
3379 instrument.id(),
3380 instrument.size_precision(),
3381 instrument.price_precision(),
3382 );
3383 let client_order_id = ClientOrderId::from("O-PENDING-MODIFY");
3384 let pending_submits = PendingSubmitTracker::default();
3385 let order_contexts = OrderContextRegistry::default();
3386 register_context(
3387 &order_contexts,
3388 old_venue_order_id,
3389 instrument.id(),
3390 client_order_id.as_str(),
3391 );
3392 order_contexts.mark_accepted(old_venue_order_id);
3393 let mut emitter = test_emitter();
3394 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3395 emitter.set_sender(sender);
3396 let ctx = WsDispatchContext {
3397 signer_type: PolymarketSignerType::Owner,
3398 token_instruments: &token_instruments,
3399 fill_tracker: &fill_tracker,
3400 pending_submits: &pending_submits,
3401 order_contexts: &order_contexts,
3402 emitter: &emitter,
3403 account_id: AccountId::from("POLY-001"),
3404 clock: nautilus_core::time::get_atomic_clock_realtime(),
3405 user_address: "0xtest",
3406 user_api_key: "test-key",
3407 };
3408
3409 let mut state = WsDispatchState::default();
3410 assert!(state.begin_modify(client_order_id, old_venue_order_id, instrument.id()));
3411 assert!(state.set_modify_replacement(
3412 client_order_id,
3413 replacement_venue_order_id,
3414 Quantity::from("120"),
3415 Quantity::from("100"),
3416 Price::from("0.5"),
3417 ));
3418
3419 let original_context = order_contexts.get(&old_venue_order_id).unwrap();
3420
3421 dispatch_user_message(&UserWsMessage::Order(replacement), &ctx, &mut state);
3422
3423 match receiver.try_recv().expect("expected replacement update") {
3424 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
3425 assert_eq!(updated.client_order_id, client_order_id);
3426 assert_eq!(updated.venue_order_id, Some(replacement_venue_order_id));
3427 assert_eq!(updated.quantity, Quantity::from("120"));
3428 assert_eq!(updated.price, Some(Price::from("0.5")));
3429 }
3430 other => panic!("expected replacement update, was {other:?}"),
3431 }
3432
3433 assert_eq!(
3434 order_contexts.venue_order_id(&client_order_id),
3435 Some(replacement_venue_order_id)
3436 );
3437 assert_eq!(
3438 order_contexts.get(&old_venue_order_id),
3439 Some(original_context)
3440 );
3441 assert_eq!(
3442 order_contexts.get(&replacement_venue_order_id),
3443 Some(OrderContext {
3444 quantity: Quantity::from("120"),
3445 price: Some(Price::from("0.5")),
3446 ..original_context
3447 })
3448 );
3449 assert!(receiver.try_recv().is_err());
3450 }
3451
3452 #[rstest]
3453 fn test_pending_modify_replacement_ws_rejection_emits_once_and_closes_old_leg() {
3454 let mut replacement: PolymarketUserOrder = load("ws_user_order_placement.json");
3455 let instrument = test_instrument();
3456 let old_venue_order_id = VenueOrderId::from("0xold-rejected-modify-leg");
3457 let replacement_venue_order_id = VenueOrderId::from("0xrejected-replacement-modify-leg");
3458 replacement.id = replacement_venue_order_id.to_string();
3459 replacement.status = Some(PolymarketUserOrderStatus::new(
3460 PolymarketOrderStatus::Unmatched,
3461 Some("replacement rejected"),
3462 ));
3463
3464 let token_instruments = AtomicMap::new();
3465 token_instruments.insert(replacement.asset_id, instrument.clone());
3466 let fill_tracker = OrderFillTrackerMap::new();
3467 fill_tracker.register(
3468 old_venue_order_id,
3469 Quantity::from("100"),
3470 OrderSide::Buy,
3471 instrument.id(),
3472 instrument.size_precision(),
3473 instrument.price_precision(),
3474 );
3475 let client_order_id = ClientOrderId::from("O-REJECTED-MODIFY");
3476 let pending_submits = PendingSubmitTracker::default();
3477 let order_contexts = OrderContextRegistry::default();
3478 register_context(
3479 &order_contexts,
3480 old_venue_order_id,
3481 instrument.id(),
3482 client_order_id.as_str(),
3483 );
3484 order_contexts.mark_accepted(old_venue_order_id);
3485 let mut emitter = test_emitter();
3486 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3487 emitter.set_sender(sender);
3488 let ctx = WsDispatchContext {
3489 signer_type: PolymarketSignerType::Owner,
3490 token_instruments: &token_instruments,
3491 fill_tracker: &fill_tracker,
3492 pending_submits: &pending_submits,
3493 order_contexts: &order_contexts,
3494 emitter: &emitter,
3495 account_id: AccountId::from("POLY-001"),
3496 clock: nautilus_core::time::get_atomic_clock_realtime(),
3497 user_address: "0xtest",
3498 user_api_key: "test-key",
3499 };
3500
3501 let mut state = WsDispatchState::default();
3502 assert!(state.begin_modify(client_order_id, old_venue_order_id, instrument.id()));
3503 let cancel_ts = UnixNanos::from(123);
3504 assert!(state.confirm_modify_cancel(client_order_id, old_venue_order_id, cancel_ts,));
3505 assert!(state.set_modify_replacement(
3506 client_order_id,
3507 replacement_venue_order_id,
3508 Quantity::from("120"),
3509 Quantity::from("100"),
3510 Price::from("0.5"),
3511 ));
3512
3513 dispatch_user_message(&UserWsMessage::Order(replacement), &ctx, &mut state);
3514
3515 match receiver.try_recv().expect("expected modify rejection") {
3516 ExecutionEvent::Order(OrderEventAny::ModifyRejected(rejected)) => {
3517 assert_eq!(rejected.client_order_id, client_order_id);
3518 assert_eq!(rejected.venue_order_id, Some(old_venue_order_id));
3519 assert_eq!(rejected.reason, "replacement rejected");
3520 }
3521 other => panic!("expected modify rejection, was {other:?}"),
3522 }
3523
3524 match receiver.try_recv().expect("expected old-leg cancellation") {
3525 ExecutionEvent::Order(OrderEventAny::Canceled(canceled)) => {
3526 assert_eq!(canceled.client_order_id, client_order_id);
3527 assert_eq!(canceled.venue_order_id, Some(old_venue_order_id));
3528 assert_eq!(canceled.ts_event, cancel_ts);
3529 }
3530 other => panic!("expected old-leg cancellation, was {other:?}"),
3531 }
3532
3533 assert!(
3534 state
3535 .pending_modify_promotion(replacement_venue_order_id)
3536 .is_none()
3537 );
3538 assert!(receiver.try_recv().is_err());
3539 }
3540
3541 #[rstest]
3542 fn test_pending_modify_replacement_fill_promotes_before_fill() {
3543 let trade: PolymarketUserTrade = load("ws_user_trade.json");
3544 let instrument = test_instrument();
3545 let old_venue_order_id = VenueOrderId::from("0xold-fill-leg");
3546 let replacement_venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
3547 let token_instruments = AtomicMap::new();
3548 token_instruments.insert(trade.asset_id, instrument.clone());
3549 let fill_tracker = OrderFillTrackerMap::new();
3550 fill_tracker.register(
3551 old_venue_order_id,
3552 Quantity::from("100"),
3553 OrderSide::Buy,
3554 instrument.id(),
3555 instrument.size_precision(),
3556 instrument.price_precision(),
3557 );
3558 let client_order_id = ClientOrderId::from("O-PENDING-MODIFY-FILL");
3559 let pending_submits = PendingSubmitTracker::default();
3560 let order_contexts = OrderContextRegistry::default();
3561 register_context(
3562 &order_contexts,
3563 old_venue_order_id,
3564 instrument.id(),
3565 client_order_id.as_str(),
3566 );
3567 order_contexts.mark_accepted(old_venue_order_id);
3568 let mut emitter = test_emitter();
3569 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3570 emitter.set_sender(sender);
3571 let ctx = WsDispatchContext {
3572 signer_type: PolymarketSignerType::Owner,
3573 token_instruments: &token_instruments,
3574 fill_tracker: &fill_tracker,
3575 pending_submits: &pending_submits,
3576 order_contexts: &order_contexts,
3577 emitter: &emitter,
3578 account_id: AccountId::from("POLY-001"),
3579 clock: nautilus_core::time::get_atomic_clock_realtime(),
3580 user_address: "0xtest",
3581 user_api_key: "test-key",
3582 };
3583
3584 let mut state = WsDispatchState::default();
3585 assert!(state.begin_modify(client_order_id, old_venue_order_id, instrument.id()));
3586 assert!(state.set_modify_replacement(
3587 client_order_id,
3588 replacement_venue_order_id,
3589 Quantity::from("120"),
3590 Quantity::from("100"),
3591 Price::from("0.5"),
3592 ));
3593 let mut cancellation: PolymarketUserOrder = load("ws_user_order_cancellation.json");
3594 cancellation.id = replacement_venue_order_id.to_string();
3595 let cancellation_report = build_ws_order_status_report(
3596 &cancellation,
3597 cancellation.status.as_ref().unwrap(),
3598 cancellation.order_type.unwrap(),
3599 &instrument,
3600 ctx.account_id,
3601 UnixNanos::from(2_000_000_000),
3602 UnixNanos::from(3_000_000_000),
3603 )
3604 .unwrap();
3605 assert!(
3606 fill_tracker
3607 .accept_or_buffer_report(replacement_venue_order_id, cancellation_report)
3608 .is_none()
3609 );
3610
3611 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
3612
3613 match receiver.try_recv().expect("expected replacement update") {
3614 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
3615 assert_eq!(updated.client_order_id, client_order_id);
3616 assert_eq!(updated.venue_order_id, Some(replacement_venue_order_id));
3617 assert_eq!(updated.quantity, Quantity::from("120"));
3618 }
3619 other => panic!("expected replacement update, was {other:?}"),
3620 }
3621
3622 match receiver.try_recv().expect("expected replacement fill") {
3623 ExecutionEvent::Order(OrderEventAny::Filled(fill)) => {
3624 assert_eq!(fill.client_order_id, client_order_id);
3625 assert_eq!(fill.venue_order_id, replacement_venue_order_id);
3626 }
3627 other => panic!("expected replacement fill, was {other:?}"),
3628 }
3629
3630 match receiver.try_recv().expect("expected replacement cancel") {
3631 ExecutionEvent::Order(OrderEventAny::Canceled(cancel)) => {
3632 assert_eq!(cancel.client_order_id, client_order_id);
3633 assert_eq!(cancel.venue_order_id, Some(replacement_venue_order_id));
3634 }
3635 other => panic!("expected replacement cancel, was {other:?}"),
3636 }
3637
3638 assert!(receiver.try_recv().is_err());
3639 }
3640
3641 #[rstest]
3642 fn test_late_modify_completion_does_not_finish_newer_modify() {
3643 let instrument_id = InstrumentId::from("TEST.POLYMARKET");
3644 let client_order_id = ClientOrderId::from("O-MODIFY-GENERATION");
3645 let old_venue_order_id = VenueOrderId::from("0xmodify-generation-old");
3646 let first_replacement_venue_order_id = VenueOrderId::from("0xmodify-generation-first");
3647 let second_replacement_venue_order_id = VenueOrderId::from("0xmodify-generation-second");
3648 let mut state = WsDispatchState::default();
3649
3650 assert!(state.begin_modify(client_order_id, old_venue_order_id, instrument_id));
3651 assert!(state.set_modify_replacement(
3652 client_order_id,
3653 first_replacement_venue_order_id,
3654 Quantity::from("12"),
3655 Quantity::from("12"),
3656 Price::from("0.5"),
3657 ));
3658 assert!(
3659 state
3660 .claim_modify_replacement(first_replacement_venue_order_id)
3661 .is_some()
3662 );
3663 assert!(state.begin_modify(
3664 client_order_id,
3665 first_replacement_venue_order_id,
3666 instrument_id,
3667 ));
3668 assert!(state.set_modify_replacement(
3669 client_order_id,
3670 second_replacement_venue_order_id,
3671 Quantity::from("15"),
3672 Quantity::from("15"),
3673 Price::from("0.6"),
3674 ));
3675
3676 assert!(
3677 state
3678 .finish_modify_without_replacement(
3679 client_order_id,
3680 old_venue_order_id,
3681 true,
3682 UnixNanos::from(123),
3683 )
3684 .is_none()
3685 );
3686 let promotion = state
3687 .pending_modify_promotion(second_replacement_venue_order_id)
3688 .expect("newer modification must remain pending");
3689 assert_eq!(promotion.client_order_id, client_order_id);
3690 assert_eq!(
3691 promotion.old_venue_order_id,
3692 first_replacement_venue_order_id
3693 );
3694 assert_eq!(promotion.venue_order_id, second_replacement_venue_order_id);
3695 assert_eq!(promotion.quantity, Quantity::from("15"));
3696 assert_eq!(promotion.leg_quantity, Quantity::from("15"));
3697 assert_eq!(promotion.price, Price::from("0.6"));
3698 }
3699
3700 #[rstest]
3701 fn test_pending_modify_lookup_selects_matching_replacement() {
3702 let instrument_id = InstrumentId::from("TEST.POLYMARKET");
3703 let first_client_order_id = ClientOrderId::from("O-MODIFY-LOOKUP-1");
3704 let second_client_order_id = ClientOrderId::from("O-MODIFY-LOOKUP-2");
3705 let first_old_venue_order_id = VenueOrderId::from("0xmodify-lookup-old-1");
3706 let second_old_venue_order_id = VenueOrderId::from("0xmodify-lookup-old-2");
3707 let first_replacement_venue_order_id = VenueOrderId::from("0xmodify-lookup-new-1");
3708 let second_replacement_venue_order_id = VenueOrderId::from("0xmodify-lookup-new-2");
3709 let mut state = WsDispatchState::default();
3710
3711 assert!(state.begin_modify(
3712 first_client_order_id,
3713 first_old_venue_order_id,
3714 instrument_id,
3715 ));
3716 assert!(state.set_modify_replacement(
3717 first_client_order_id,
3718 first_replacement_venue_order_id,
3719 Quantity::from("11"),
3720 Quantity::from("10"),
3721 Price::from("0.4"),
3722 ));
3723 assert!(state.begin_modify(
3724 second_client_order_id,
3725 second_old_venue_order_id,
3726 instrument_id,
3727 ));
3728 assert!(state.set_modify_replacement(
3729 second_client_order_id,
3730 second_replacement_venue_order_id,
3731 Quantity::from("22"),
3732 Quantity::from("20"),
3733 Price::from("0.6"),
3734 ));
3735
3736 let second = state
3737 .claim_modify_replacement(second_replacement_venue_order_id)
3738 .expect("second replacement must be selected");
3739 assert_eq!(second.client_order_id, second_client_order_id);
3740 assert_eq!(second.old_venue_order_id, second_old_venue_order_id);
3741 assert_eq!(second.venue_order_id, second_replacement_venue_order_id);
3742 assert_eq!(second.quantity, Quantity::from("22"));
3743 assert_eq!(second.leg_quantity, Quantity::from("20"));
3744 assert_eq!(second.price, Price::from("0.6"));
3745
3746 let first = state
3747 .pending_modify_promotion(first_replacement_venue_order_id)
3748 .expect("first replacement must remain pending");
3749 assert_eq!(first.client_order_id, first_client_order_id);
3750 assert_eq!(first.old_venue_order_id, first_old_venue_order_id);
3751 assert_eq!(first.venue_order_id, first_replacement_venue_order_id);
3752 assert_eq!(first.quantity, Quantity::from("11"));
3753 assert_eq!(first.leg_quantity, Quantity::from("10"));
3754 assert_eq!(first.price, Price::from("0.4"));
3755 }
3756
3757 #[rstest]
3758 fn test_begin_cancels_is_atomic_when_one_order_conflicts() {
3759 let instrument_id = InstrumentId::from("TEST.POLYMARKET");
3760 let available_client_order_id = ClientOrderId::from("O-CANCEL-AVAILABLE");
3761 let conflicting_client_order_id = ClientOrderId::from("O-CANCEL-CONFLICT");
3762 let mut state = WsDispatchState::default();
3763
3764 assert!(state.begin_modify(
3765 conflicting_client_order_id,
3766 VenueOrderId::from("0xmodify-conflict"),
3767 instrument_id,
3768 ));
3769 assert!(!state.begin_cancels(&[
3770 (available_client_order_id, instrument_id),
3771 (conflicting_client_order_id, instrument_id),
3772 ]));
3773 assert!(state.begin_modify(
3774 available_client_order_id,
3775 VenueOrderId::from("0xmodify-available"),
3776 instrument_id,
3777 ));
3778 }
3779
3780 #[rstest]
3781 fn test_begin_available_cancels_skips_existing_cancel() {
3782 let instrument_id = InstrumentId::from("TEST.POLYMARKET");
3783 let pending_client_order_id = ClientOrderId::from("O-CANCEL-PENDING");
3784 let available_client_order_id = ClientOrderId::from("O-CANCEL-AVAILABLE");
3785 let modifying_client_order_id = ClientOrderId::from("O-MODIFY-PENDING");
3786 let unreserved_client_order_id = ClientOrderId::from("O-CANCEL-UNRESERVED");
3787 let mut state = WsDispatchState::default();
3788
3789 assert!(state.begin_cancels(&[(pending_client_order_id, instrument_id)]));
3790 assert_eq!(
3791 state
3792 .begin_available_cancels(&[
3793 (pending_client_order_id, instrument_id),
3794 (available_client_order_id, instrument_id),
3795 ])
3796 .unwrap(),
3797 vec![available_client_order_id],
3798 );
3799 assert!(!state.begin_modify(
3800 pending_client_order_id,
3801 VenueOrderId::from("0xcancel-pending"),
3802 instrument_id,
3803 ));
3804 assert!(!state.begin_modify(
3805 available_client_order_id,
3806 VenueOrderId::from("0xcancel-available"),
3807 instrument_id,
3808 ));
3809
3810 state.finish_cancels(&[pending_client_order_id, available_client_order_id]);
3811 assert!(state.begin_modify(
3812 modifying_client_order_id,
3813 VenueOrderId::from("0xmodify-pending"),
3814 instrument_id,
3815 ));
3816 assert!(
3817 state
3818 .begin_available_cancels(&[
3819 (unreserved_client_order_id, instrument_id),
3820 (modifying_client_order_id, instrument_id),
3821 ])
3822 .is_none()
3823 );
3824 assert!(state.begin_modify(
3825 unreserved_client_order_id,
3826 VenueOrderId::from("0xcancel-unreserved"),
3827 instrument_id,
3828 ));
3829 }
3830
3831 #[rstest]
3832 fn test_cancel_and_modify_are_mutually_exclusive() {
3833 let instrument_id = InstrumentId::from("TEST.POLYMARKET");
3834 let client_order_id = ClientOrderId::from("O-MODIFY-CANCEL-ALL");
3835 let other_client_order_id = ClientOrderId::from("O-MODIFY-MARKET-CANCEL");
3836 let venue_order_id = VenueOrderId::from("0xmodify-cancel-all");
3837 let mut state = WsDispatchState::default();
3838
3839 assert!(state.begin_cancels(&[(client_order_id, instrument_id)]));
3840 assert!(!state.begin_modify(client_order_id, venue_order_id, instrument_id));
3841 assert!(state.begin_market_cancel(instrument_id));
3842 assert!(!state.begin_modify(
3843 other_client_order_id,
3844 VenueOrderId::from("0xmodify-market-cancel"),
3845 instrument_id,
3846 ));
3847 state.finish_market_cancel(instrument_id);
3848 state.finish_cancels(&[client_order_id]);
3849 assert!(state.begin_modify(client_order_id, venue_order_id, instrument_id));
3850 assert!(!state.begin_cancels(&[(client_order_id, instrument_id)]));
3851 assert!(!state.begin_market_cancel(instrument_id));
3852 assert!(state.set_modify_replacement(
3853 client_order_id,
3854 VenueOrderId::from("0xmodify-cancel-all-new"),
3855 Quantity::from("12"),
3856 Quantity::from("12"),
3857 Price::from("0.5"),
3858 ));
3859 assert!(!state.begin_cancels(&[(client_order_id, instrument_id)]));
3860 assert!(!state.begin_market_cancel(instrument_id));
3861 assert!(
3862 state
3863 .finish_modify_without_replacement(
3864 client_order_id,
3865 venue_order_id,
3866 false,
3867 UnixNanos::default(),
3868 )
3869 .is_some()
3870 );
3871 assert!(state.begin_market_cancel(instrument_id));
3872 assert!(!state.begin_modify(client_order_id, venue_order_id, instrument_id));
3873 state.finish_market_cancel(instrument_id);
3874 }
3875
3876 #[rstest]
3877 fn test_reset_preserves_modify_recovery_and_stale_leg_safety() {
3878 let instrument_id = InstrumentId::from("TEST.POLYMARKET");
3879 let pending_client_order_id = ClientOrderId::from("O-PENDING-RESET");
3880 let pending_old_venue_order_id = VenueOrderId::from("0xpending-old-reset");
3881 let pending_new_venue_order_id = VenueOrderId::from("0xpending-new-reset");
3882 let replaced_client_order_id = ClientOrderId::from("O-REPLACED-RESET");
3883 let replaced_old_venue_order_id = VenueOrderId::from("0xreplaced-old-reset");
3884 let replaced_new_venue_order_id = VenueOrderId::from("0xreplaced-new-reset");
3885 let closed_client_order_id = ClientOrderId::from("O-CLOSED-RESET");
3886 let closed_venue_order_id = VenueOrderId::from("0xclosed-reset");
3887 let cancel_client_order_id = ClientOrderId::from("O-CANCEL-RESET");
3888 let trade_id = TradeId::from("T-RESET");
3889 let mut state = WsDispatchState::default();
3890
3891 assert!(state.begin_modify(
3892 pending_client_order_id,
3893 pending_old_venue_order_id,
3894 instrument_id,
3895 ));
3896 assert!(state.set_modify_replacement(
3897 pending_client_order_id,
3898 pending_new_venue_order_id,
3899 Quantity::from("12"),
3900 Quantity::from("10"),
3901 Price::from("0.5"),
3902 ));
3903 assert!(state.begin_modify(
3904 replaced_client_order_id,
3905 replaced_old_venue_order_id,
3906 instrument_id,
3907 ));
3908 assert!(state.set_modify_replacement(
3909 replaced_client_order_id,
3910 replaced_new_venue_order_id,
3911 Quantity::from("12"),
3912 Quantity::from("10"),
3913 Price::from("0.5"),
3914 ));
3915 assert!(
3916 state
3917 .claim_modify_replacement(replaced_new_venue_order_id)
3918 .is_some()
3919 );
3920 assert!(state.begin_modify(closed_client_order_id, closed_venue_order_id, instrument_id,));
3921 assert!(
3922 state
3923 .finish_modify_without_replacement(
3924 closed_client_order_id,
3925 closed_venue_order_id,
3926 true,
3927 UnixNanos::from(1),
3928 )
3929 .is_some()
3930 );
3931 state.record_reconciled_fill(trade_id, pending_old_venue_order_id);
3932 assert!(state.begin_cancels(&[(cancel_client_order_id, instrument_id)]));
3933
3934 state.reset_session();
3935
3936 assert_eq!(
3937 state
3938 .pending_modify_promotion(pending_new_venue_order_id)
3939 .unwrap()
3940 .client_order_id,
3941 pending_client_order_id,
3942 );
3943 assert!(state.suppress_modify_cancel(replaced_old_venue_order_id));
3944 assert!(state.suppress_modify_cancel(closed_venue_order_id));
3945 assert!(state.suppress_modify_cancel_reemit(pending_old_venue_order_id));
3946 assert!(state.suppress_modify_cancel_reemit(replaced_old_venue_order_id));
3947 assert!(!state.suppress_modify_cancel_reemit(closed_venue_order_id));
3948 assert!(
3949 state
3950 .reconciled_fills
3951 .contains(&(trade_id, pending_old_venue_order_id))
3952 );
3953 assert!(state.begin_cancels(&[(cancel_client_order_id, instrument_id)]));
3954 }
3955
3956 #[rstest]
3957 fn test_reconciled_fill_is_not_reapplied_from_websocket() {
3958 let trade: PolymarketUserTrade = load("ws_user_trade.json");
3959 let instrument = test_instrument();
3960 let venue_order_id = VenueOrderId::from(trade.taker_order_id.as_str());
3961 let token_instruments = AtomicMap::new();
3962 token_instruments.insert(trade.asset_id, instrument.clone());
3963 let fill_tracker = OrderFillTrackerMap::new();
3964 fill_tracker.restore_order(
3965 venue_order_id,
3966 Quantity::from("100"),
3967 Quantity::from("25"),
3968 OrderSide::Buy,
3969 );
3970 let pending_submits = PendingSubmitTracker::default();
3971 let order_contexts = OrderContextRegistry::default();
3972 register_context(
3973 &order_contexts,
3974 venue_order_id,
3975 instrument.id(),
3976 "O-RECONCILED-FILL",
3977 );
3978 let mut emitter = test_emitter();
3979 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
3980 emitter.set_sender(sender);
3981 let ctx = WsDispatchContext {
3982 signer_type: PolymarketSignerType::Owner,
3983 token_instruments: &token_instruments,
3984 fill_tracker: &fill_tracker,
3985 pending_submits: &pending_submits,
3986 order_contexts: &order_contexts,
3987 emitter: &emitter,
3988 account_id: AccountId::from("POLY-001"),
3989 clock: nautilus_core::time::get_atomic_clock_realtime(),
3990 user_address: "0xtest",
3991 user_api_key: "test-key",
3992 };
3993
3994 let mut state = WsDispatchState::default();
3995 state.record_reconciled_fill(TradeId::from(trade.id.as_str()), venue_order_id);
3996
3997 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
3998
3999 assert_eq!(
4000 fill_tracker.get_cumulative_filled(&venue_order_id),
4001 Some(Quantity::from("25")),
4002 );
4003 assert!(receiver.try_recv().is_err());
4004 }
4005
4006 #[rstest]
4007 fn test_cancel_not_reemitted_when_fill_completes_order() {
4008 let cancel_order: PolymarketUserOrder = load("ws_user_order_cancellation.json");
4009 let trade: PolymarketUserTrade = load("ws_user_trade.json");
4010 let instrument = test_instrument();
4011
4012 let token_instruments = AtomicMap::new();
4013 token_instruments.insert(cancel_order.asset_id, instrument.clone());
4014
4015 let fill_tracker = OrderFillTrackerMap::new();
4016 let venue_order_id = VenueOrderId::from(cancel_order.id.as_str());
4017
4018 fill_tracker.register(
4020 venue_order_id,
4021 Quantity::from("25"),
4022 OrderSide::Buy,
4023 instrument.id(),
4024 instrument.size_precision(),
4025 instrument.price_precision(),
4026 );
4027
4028 let pending_submits = PendingSubmitTracker::default();
4029 let order_contexts = OrderContextRegistry::default();
4030 register_context(
4031 &order_contexts,
4032 venue_order_id,
4033 instrument.id(),
4034 "O-CANCEL-FULL",
4035 );
4036 order_contexts.mark_accepted(venue_order_id);
4037 let mut emitter = test_emitter();
4038 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
4039 emitter.set_sender(sender);
4040
4041 let ctx = WsDispatchContext {
4042 signer_type: PolymarketSignerType::Owner,
4043 token_instruments: &token_instruments,
4044 fill_tracker: &fill_tracker,
4045 pending_submits: &pending_submits,
4046 order_contexts: &order_contexts,
4047 emitter: &emitter,
4048 account_id: AccountId::from("POLY-001"),
4049 clock: nautilus_core::time::get_atomic_clock_realtime(),
4050 user_address: "0xtest",
4051 user_api_key: "test-key",
4052 };
4053 let mut state = WsDispatchState::default();
4054
4055 dispatch_user_message(&UserWsMessage::Order(cancel_order), &ctx, &mut state);
4057 let _cancel = receiver.try_recv().expect("Expected canceled event");
4058
4059 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
4060 let _fill = receiver.try_recv().expect("Expected filled event");
4061
4062 assert!(
4064 receiver.try_recv().is_err(),
4065 "Should not re-emit cancel when fill completes the order"
4066 );
4067 }
4068
4069 #[rstest]
4070 fn test_cancel_saved_before_acceptance_and_retained_for_modify_reset() {
4071 let cancel_order: PolymarketUserOrder = load("ws_user_order_cancellation.json");
4072 let instrument = test_instrument();
4073 let instrument_id = instrument.id();
4074
4075 let token_instruments = AtomicMap::new();
4076 token_instruments.insert(cancel_order.asset_id, instrument);
4077
4078 let fill_tracker = OrderFillTrackerMap::new();
4080 let venue_order_id = VenueOrderId::from(cancel_order.id.as_str());
4081
4082 let pending_submits = PendingSubmitTracker::default();
4083 let order_contexts = OrderContextRegistry::default();
4084 let emitter = test_emitter();
4085
4086 let ctx = WsDispatchContext {
4087 signer_type: PolymarketSignerType::Owner,
4088 token_instruments: &token_instruments,
4089 fill_tracker: &fill_tracker,
4090 pending_submits: &pending_submits,
4091 order_contexts: &order_contexts,
4092 emitter: &emitter,
4093 account_id: AccountId::from("POLY-001"),
4094 clock: nautilus_core::time::get_atomic_clock_realtime(),
4095 user_address: "0xtest",
4096 user_api_key: "test-key",
4097 };
4098 let mut state = WsDispatchState::default();
4099
4100 dispatch_user_message(&UserWsMessage::Order(cancel_order), &ctx, &mut state);
4102
4103 assert!(fill_tracker.has_pending_report(&venue_order_id));
4105 assert!(state.terminal_cancel_reports.get(&venue_order_id).is_some());
4106
4107 let client_order_id = ClientOrderId::from("O-CANCEL-RESET");
4108 assert!(state.begin_modify(client_order_id, venue_order_id, instrument_id));
4109 state.reset_session();
4110 assert!(
4111 state
4112 .finish_modify_without_replacement(
4113 client_order_id,
4114 venue_order_id,
4115 false,
4116 UnixNanos::default(),
4117 )
4118 .unwrap()
4119 .1
4120 .is_some()
4121 );
4122 }
4123
4124 #[rstest]
4127 #[case(PolymarketOrderStatus::Canceled, "Canceled", OrderStatus::Canceled)]
4128 #[case(
4129 PolymarketOrderStatus::CanceledMarketResolved,
4130 "Expired",
4131 OrderStatus::Expired
4132 )]
4133 fn test_buffered_fill_emitted_before_terminal_status(
4134 #[case] status: PolymarketOrderStatus,
4135 #[case] expected_terminal: &str,
4136 #[case] expected_order_status: OrderStatus,
4137 ) {
4138 let mut terminal_order: PolymarketUserOrder = load("ws_user_order_cancellation.json");
4139 terminal_order.status = Some(status.into());
4140 let trade: PolymarketUserTrade = load("ws_user_trade.json");
4141 let instrument = test_instrument();
4142
4143 let token_instruments = AtomicMap::new();
4144 token_instruments.insert(terminal_order.asset_id, instrument.clone());
4145
4146 let fill_tracker = OrderFillTrackerMap::new();
4148 let venue_order_id = VenueOrderId::from(terminal_order.id.as_str());
4149 let client_order_id = ClientOrderId::from("O-BUFFERED");
4150
4151 let pending_submits = PendingSubmitTracker::default();
4152 pending_submits.insert(venue_order_id, client_order_id);
4153 let order_contexts = OrderContextRegistry::default();
4154 register_context(
4155 &order_contexts,
4156 venue_order_id,
4157 instrument.id(),
4158 client_order_id.as_str(),
4159 );
4160 let mut emitter = test_emitter();
4161 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
4162 emitter.set_sender(sender);
4163
4164 let ctx = WsDispatchContext {
4165 signer_type: PolymarketSignerType::Owner,
4166 token_instruments: &token_instruments,
4167 fill_tracker: &fill_tracker,
4168 pending_submits: &pending_submits,
4169 order_contexts: &order_contexts,
4170 emitter: &emitter,
4171 account_id: AccountId::from("POLY-001"),
4172 clock: nautilus_core::time::get_atomic_clock_realtime(),
4173 user_address: "0xtest",
4174 user_api_key: "test-key",
4175 };
4176 let mut state = WsDispatchState::default();
4177
4178 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
4180 assert!(
4181 receiver.try_recv().is_err(),
4182 "a buffered fill must emit no event before the order is registered",
4183 );
4184
4185 dispatch_user_message(&UserWsMessage::Order(terminal_order), &ctx, &mut state);
4187
4188 let mut emitted = Vec::new();
4189
4190 while let Ok(event) = receiver.try_recv() {
4191 match event {
4192 ExecutionEvent::Order(order_event) => emitted.push(order_event),
4193 other => panic!("expected only order events, was {other:?}"),
4194 }
4195 }
4196
4197 assert_eq!(emitted.len(), 3, "emitted sequence was {emitted:?}");
4198 match &emitted[0] {
4199 OrderEventAny::Accepted(accepted) => {
4200 assert_eq!(accepted.client_order_id, client_order_id);
4201 assert_eq!(accepted.venue_order_id, venue_order_id);
4202 }
4203 other => panic!("expected accepted event first, was {other:?}"),
4204 }
4205
4206 match &emitted[1] {
4207 OrderEventAny::Filled(filled) => {
4208 assert_eq!(filled.client_order_id, client_order_id);
4209 assert_eq!(filled.venue_order_id, venue_order_id);
4210 assert_eq!(filled.last_qty.as_decimal(), dec!(25));
4211 }
4212 other => panic!("expected filled event before the terminal status, was {other:?}"),
4213 }
4214
4215 let terminal = match &emitted[2] {
4216 OrderEventAny::Canceled(canceled) => {
4217 assert_eq!(canceled.client_order_id, client_order_id);
4218 assert_eq!(canceled.venue_order_id, Some(venue_order_id));
4219 "Canceled"
4220 }
4221 OrderEventAny::Expired(expired) => {
4222 assert_eq!(expired.client_order_id, client_order_id);
4223 assert_eq!(expired.venue_order_id, Some(venue_order_id));
4224 "Expired"
4225 }
4226 other => panic!("expected a terminal order event last, was {other:?}"),
4227 };
4228 assert_eq!(terminal, expected_terminal);
4229
4230 let mut order = OrderTestBuilder::new(OrderType::Limit)
4232 .instrument_id(instrument.id())
4233 .client_order_id(client_order_id)
4234 .strategy_id(StrategyId::from("S-001"))
4235 .side(OrderSide::Buy)
4236 .price(Price::from("0.5"))
4237 .quantity(Quantity::from("100"))
4238 .build();
4239
4240 for event in emitted {
4241 order.apply(event).expect("emitted sequence must be valid");
4242 }
4243
4244 assert_eq!(order.status(), expected_order_status);
4245 assert_eq!(order.filled_qty().as_decimal(), dec!(25));
4246 }
4247
4248 #[rstest]
4260 fn test_issue_3797_interleaved_cancel_fill_sequence() {
4261 use crate::common::{
4262 enums::{
4263 PolymarketEventType, PolymarketLiquiditySide, PolymarketOrderSide,
4264 PolymarketOrderStatus, PolymarketOrderType, PolymarketOutcome,
4265 PolymarketTradeStatus,
4266 },
4267 models::PolymarketMakerOrder,
4268 };
4269
4270 let instrument = test_instrument();
4271 let asset_id = instrument.id().symbol.inner();
4272
4273 let order_id =
4274 "0xe743f6c823ecdfa9ddaaf08673b2441d15a38d89e14dcb25b3b70c284be4f6ad".to_string();
4275 let venue_order_id = VenueOrderId::from(order_id.as_str());
4276
4277 let token_instruments = AtomicMap::new();
4278 token_instruments.insert(asset_id, instrument.clone());
4279
4280 let fill_tracker = OrderFillTrackerMap::new();
4281 fill_tracker.register(
4282 venue_order_id,
4283 Quantity::from("20"),
4284 OrderSide::Buy,
4285 instrument.id(),
4286 instrument.size_precision(),
4287 instrument.price_precision(),
4288 );
4289
4290 let pending_submits = PendingSubmitTracker::default();
4291 let order_contexts = OrderContextRegistry::default();
4292 register_context(&order_contexts, venue_order_id, instrument.id(), "O-3797");
4293 order_contexts.mark_accepted(venue_order_id);
4294 let mut emitter = test_emitter();
4295 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
4296 emitter.set_sender(sender);
4297
4298 let ctx = WsDispatchContext {
4299 signer_type: PolymarketSignerType::Owner,
4300 token_instruments: &token_instruments,
4301 fill_tracker: &fill_tracker,
4302 pending_submits: &pending_submits,
4303 order_contexts: &order_contexts,
4304 emitter: &emitter,
4305 account_id: AccountId::from("POLY-001"),
4306 clock: nautilus_core::time::get_atomic_clock_realtime(),
4307 user_address: "0xabc",
4308 user_api_key: "xxx",
4309 };
4310 let mut state = WsDispatchState::default();
4311
4312 let make_order =
4313 |size_matched: &str, ts: &str, event_type: PolymarketEventType| PolymarketUserOrder {
4314 asset_id,
4315 associate_trades: None,
4316 created_at: Some("1775074735".to_string()),
4317 expiration: Some("0".to_string()),
4318 id: order_id.clone(),
4319 maker_address: Some(Ustr::from("0xabc")),
4320 market: Ustr::from("0x4134"),
4321 order_owner: Some(Ustr::from("xxx")),
4322 order_type: Some(PolymarketOrderType::GTC),
4323 original_size: "20".to_string(),
4324 outcome: Some(PolymarketOutcome::yes()),
4325 owner: Ustr::from("xxx"),
4326 price: "0.18".to_string(),
4327 side: PolymarketOrderSide::Buy,
4328 size_matched: size_matched.to_string(),
4329 status: Some(PolymarketOrderStatus::Canceled.into()),
4330 timestamp: ts.to_string(),
4331 event_type,
4332 };
4333
4334 let make_trade = |trade_id: &str, matched_amount: f64, ts: &str| PolymarketUserTrade {
4335 asset_id,
4336 bucket_index: 0,
4337 fee_rate_bps: "1000".to_string(),
4338 id: trade_id.to_string(),
4339 last_update: "1775074738".to_string(),
4340 maker_address: Ustr::from("0xother"),
4341 maker_orders: vec![PolymarketMakerOrder {
4342 asset_id,
4343 maker_address: "0xabc".to_string(),
4344 matched_amount: Decimal::from_f64_retain(matched_amount).unwrap_or(Decimal::ZERO),
4345 order_id: order_id.clone(),
4346 outcome: PolymarketOutcome::yes(),
4347 owner: "xxx".to_string(),
4348 price: Decimal::from_f64_retain(0.18).unwrap_or(Decimal::ZERO),
4349 side: None,
4350 }],
4351 market: Ustr::from("0x4134"),
4352 match_time: "1775074735".to_string(),
4353 outcome: PolymarketOutcome::yes(),
4354 owner: Ustr::from("other-owner"),
4355 price: "0.82".to_string(),
4356 side: PolymarketOrderSide::Buy,
4357 size: "1.219511".to_string(),
4358 status: PolymarketTradeStatus::Confirmed,
4359 taker_order_id: "0xtaker01".to_string(),
4360 timestamp: ts.to_string(),
4361 trade_owner: Ustr::from("other-owner"),
4362 transaction_hash: None,
4363 trader_side: PolymarketLiquiditySide::Maker,
4364 event_type: PolymarketEventType::Trade,
4365 };
4366
4367 let msg_a = make_order("0", "1775074738031", PolymarketEventType::Cancellation);
4369 dispatch_user_message(&UserWsMessage::Order(msg_a), &ctx, &mut state);
4370
4371 let evt = receiver.try_recv().expect("(A) canceled event");
4372 match &evt {
4373 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
4374 assert_eq!(c.venue_order_id, Some(venue_order_id));
4375 }
4376 other => panic!("(A) expected canceled event, was {other:?}"),
4377 }
4378
4379 let msg_b = make_trade("trade-b", 1.219511, "1775074738032");
4381 dispatch_user_message(&UserWsMessage::Trade(msg_b), &ctx, &mut state);
4382
4383 let evt = receiver.try_recv().expect("(B) filled event");
4384 match &evt {
4385 ExecutionEvent::Order(OrderEventAny::Filled(f)) => {
4386 assert_eq!(f.venue_order_id, venue_order_id);
4387 }
4388 other => panic!("(B) expected filled event, was {other:?}"),
4389 }
4390 let evt = receiver.try_recv().expect("(B) re-emitted cancel");
4392 match &evt {
4393 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
4394 assert_eq!(c.venue_order_id, Some(venue_order_id));
4395 }
4396 other => panic!("(B) expected re-emitted cancel, was {other:?}"),
4397 }
4398
4399 let msg_c = make_order("1.219511", "1775074738034", PolymarketEventType::Update);
4401 dispatch_user_message(&UserWsMessage::Order(msg_c), &ctx, &mut state);
4402
4403 let evt = receiver.try_recv().expect("(C) canceled event");
4404 match &evt {
4405 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
4406 assert_eq!(c.venue_order_id, Some(venue_order_id));
4407 }
4408 other => panic!("(C) expected canceled event, was {other:?}"),
4409 }
4410
4411 let msg_d = make_order("2.560972", "1775074738038", PolymarketEventType::Update);
4413 dispatch_user_message(&UserWsMessage::Order(msg_d), &ctx, &mut state);
4414
4415 let evt = receiver.try_recv().expect("(D) canceled event");
4416 match &evt {
4417 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
4418 assert_eq!(c.venue_order_id, Some(venue_order_id));
4419 }
4420 other => panic!("(D) expected canceled event, was {other:?}"),
4421 }
4422
4423 let msg_e = make_trade("trade-e", 1.341461, "1775074738036");
4425 dispatch_user_message(&UserWsMessage::Trade(msg_e), &ctx, &mut state);
4426
4427 let evt = receiver.try_recv().expect("(E) filled event");
4428 match &evt {
4429 ExecutionEvent::Order(OrderEventAny::Filled(f)) => {
4430 assert_eq!(f.venue_order_id, venue_order_id);
4431 }
4432 other => panic!("(E) expected filled event, was {other:?}"),
4433 }
4434
4435 let evt = receiver.try_recv().expect("(E) re-emitted cancel");
4437 match &evt {
4438 ExecutionEvent::Order(OrderEventAny::Canceled(c)) => {
4439 assert_eq!(c.venue_order_id, Some(venue_order_id));
4440 }
4441 other => panic!("(E) expected re-emitted cancel, was {other:?}"),
4442 }
4443
4444 assert!(
4446 receiver.try_recv().is_err(),
4447 "No further events expected after the sequence"
4448 );
4449 }
4450
4451 #[rstest]
4452 fn test_dispatch_taker_fill_snaps_overfill_to_submitted_qty() {
4453 use crate::common::enums::{
4458 PolymarketEventType, PolymarketOrderSide, PolymarketOutcome, PolymarketTradeStatus,
4459 };
4460
4461 let instrument = test_instrument();
4462 let asset_id = instrument.id().symbol.inner();
4463 let token_instruments = AtomicMap::new();
4464 token_instruments.insert(asset_id, instrument.clone());
4465
4466 let fill_tracker = OrderFillTrackerMap::new();
4467 let venue_order_id = VenueOrderId::from("0xtaker-overfill");
4468 let submitted = Quantity::new(714.285710, instrument.size_precision());
4470 fill_tracker.register(
4471 venue_order_id,
4472 submitted,
4473 OrderSide::Buy,
4474 instrument.id(),
4475 instrument.size_precision(),
4476 instrument.price_precision(),
4477 );
4478
4479 let pending_submits = PendingSubmitTracker::default();
4480 let order_contexts = OrderContextRegistry::default();
4481 register_context(
4482 &order_contexts,
4483 venue_order_id,
4484 instrument.id(),
4485 "O-OVERFILL",
4486 );
4487 order_contexts.mark_accepted(venue_order_id);
4488 let mut emitter = test_emitter();
4489 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
4490 emitter.set_sender(sender);
4491
4492 let ctx = WsDispatchContext {
4493 signer_type: PolymarketSignerType::Owner,
4494 token_instruments: &token_instruments,
4495 fill_tracker: &fill_tracker,
4496 pending_submits: &pending_submits,
4497 order_contexts: &order_contexts,
4498 emitter: &emitter,
4499 account_id: AccountId::from("POLY-001"),
4500 clock: nautilus_core::time::get_atomic_clock_realtime(),
4501 user_address: "0xtest",
4502 user_api_key: "test-key",
4503 };
4504 let mut state = WsDispatchState::default();
4505
4506 let trade = PolymarketUserTrade {
4507 asset_id,
4508 bucket_index: 0,
4509 fee_rate_bps: "0".to_string(),
4510 id: "trade-overfill".to_string(),
4511 last_update: "1700000001".to_string(),
4512 maker_address: Ustr::from("0xmaker"),
4513 maker_orders: vec![],
4514 market: Ustr::from("0xmarket"),
4515 match_time: "1700000000".to_string(),
4516 outcome: PolymarketOutcome::yes(),
4517 owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
4518 price: "0.014".to_string(),
4519 side: PolymarketOrderSide::Buy,
4520 size: "714.285714".to_string(),
4523 status: PolymarketTradeStatus::Confirmed,
4524 taker_order_id: venue_order_id.as_str().to_string(),
4525 timestamp: "1700000000000".to_string(),
4526 trade_owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
4527 transaction_hash: None,
4528 trader_side: PolymarketLiquiditySide::Taker,
4529 event_type: PolymarketEventType::Trade,
4530 };
4531
4532 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
4533
4534 let cumulative = fill_tracker
4538 .get_cumulative_filled(&venue_order_id)
4539 .expect("order must be registered");
4540 assert_eq!(cumulative, submitted);
4541
4542 let event = receiver.try_recv().expect("expected a filled event");
4545 match event {
4546 ExecutionEvent::Order(OrderEventAny::Filled(filled)) => {
4547 assert_eq!(
4548 filled.last_qty, submitted,
4549 "filled qty must be snapped to submitted",
4550 );
4551 assert_eq!(filled.venue_order_id, venue_order_id);
4552 }
4553 other => panic!("expected filled event, was {other:?}"),
4554 }
4555 }
4556
4557 #[rstest]
4558 #[case(
4559 TimeInForce::Ioc,
4560 OrderType::Market,
4561 OrderSide::Buy,
4562 "5.202910",
4563 "5.202897",
4564 false,
4565 true
4566 )]
4567 #[case(
4568 TimeInForce::Fok,
4569 OrderType::Limit,
4570 OrderSide::Buy,
4571 "5.202910",
4572 "5.202897",
4573 true,
4574 false
4575 )]
4576 #[case(
4577 TimeInForce::Ioc,
4578 OrderType::Limit,
4579 OrderSide::Buy,
4580 "30",
4581 "20",
4582 false,
4583 true
4584 )]
4585 #[case(
4586 TimeInForce::Ioc,
4587 OrderType::Market,
4588 OrderSide::Sell,
4589 "5.202910",
4590 "5.202897",
4591 false,
4592 true
4593 )]
4594 #[case(
4595 TimeInForce::Gtc,
4596 OrderType::Limit,
4597 OrderSide::Buy,
4598 "5.202910",
4599 "5.202897",
4600 false,
4601 false
4602 )]
4603 fn test_taker_terminal_status_on_trade_confirm(
4604 #[case] time_in_force: TimeInForce,
4605 #[case] order_type: OrderType,
4606 #[case] order_side: OrderSide,
4607 #[case] submitted_qty: &str,
4608 #[case] fill_qty: &str,
4609 #[case] expect_normalization: bool,
4610 #[case] expect_cancel: bool,
4611 ) {
4612 use crate::common::enums::{
4616 PolymarketEventType, PolymarketOrderSide, PolymarketOutcome, PolymarketTradeStatus,
4617 };
4618
4619 let instrument = test_instrument();
4620 let asset_id = instrument.id().symbol.inner();
4621 let token_instruments = AtomicMap::new();
4622 token_instruments.insert(asset_id, instrument.clone());
4623
4624 let fill_tracker = OrderFillTrackerMap::new();
4625 let venue_order_id = VenueOrderId::from("0xtaker-one-shot-dust");
4626 let submitted = Quantity::from_decimal_dp(
4627 Decimal::from_str_exact(submitted_qty).unwrap(),
4628 instrument.size_precision(),
4629 )
4630 .unwrap();
4631 fill_tracker.register(
4632 venue_order_id,
4633 submitted,
4634 order_side,
4635 instrument.id(),
4636 instrument.size_precision(),
4637 instrument.price_precision(),
4638 );
4639
4640 let pending_submits = PendingSubmitTracker::default();
4641 let order_contexts = OrderContextRegistry::default();
4642 order_contexts.register_context(
4643 venue_order_id,
4644 OrderContext {
4645 identity: OrderIdentity {
4646 client_order_id: ClientOrderId::from("O-ONE-SHOT"),
4647 strategy_id: StrategyId::from("S-001"),
4648 instrument_id: instrument.id(),
4649 order_side,
4650 order_type,
4651 },
4652 quantity: submitted,
4653 price: Some(Price::from("0.50")),
4654 trigger_price: None,
4655 trigger_type: None,
4656 time_in_force,
4657 is_post_only: false,
4658 is_reduce_only: false,
4659 is_quote_quantity: false,
4660 },
4661 );
4662 order_contexts.mark_accepted(venue_order_id);
4663 let mut emitter = test_emitter();
4664 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
4665 emitter.set_sender(sender);
4666
4667 let ctx = WsDispatchContext {
4668 signer_type: PolymarketSignerType::Owner,
4669 token_instruments: &token_instruments,
4670 fill_tracker: &fill_tracker,
4671 pending_submits: &pending_submits,
4672 order_contexts: &order_contexts,
4673 emitter: &emitter,
4674 account_id: AccountId::from("POLY-001"),
4675 clock: nautilus_core::time::get_atomic_clock_realtime(),
4676 user_address: "0xtest",
4677 user_api_key: "test-key",
4678 };
4679 let mut state = WsDispatchState::default();
4680
4681 let trade = PolymarketUserTrade {
4682 asset_id,
4683 bucket_index: 0,
4684 fee_rate_bps: "0".to_string(),
4685 id: "trade-one-shot-dust".to_string(),
4686 last_update: "1700000001".to_string(),
4687 maker_address: Ustr::from("0xmaker"),
4688 maker_orders: vec![],
4689 market: Ustr::from("0xmarket"),
4690 match_time: "1700000000".to_string(),
4691 outcome: PolymarketOutcome::yes(),
4692 owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
4693 price: "0.963".to_string(),
4694 side: if order_side == OrderSide::Buy {
4695 PolymarketOrderSide::Buy
4696 } else {
4697 PolymarketOrderSide::Sell
4698 },
4699 size: fill_qty.to_string(),
4700 status: PolymarketTradeStatus::Confirmed,
4701 taker_order_id: venue_order_id.as_str().to_string(),
4702 timestamp: "1700000000000".to_string(),
4703 trade_owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
4704 transaction_hash: None,
4705 trader_side: PolymarketLiquiditySide::Taker,
4706 event_type: PolymarketEventType::Trade,
4707 };
4708
4709 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
4710
4711 let event = receiver.try_recv().expect("expected the venue fill event");
4712 match event {
4713 ExecutionEvent::Order(OrderEventAny::Filled(filled)) => {
4714 assert_eq!(
4715 filled.last_qty,
4716 Quantity::from_decimal_dp(
4717 Decimal::from_str_exact(fill_qty).unwrap(),
4718 instrument.size_precision(),
4719 )
4720 .unwrap(),
4721 );
4722 }
4723 other => panic!("expected filled event, was {other:?}"),
4724 }
4725
4726 if expect_normalization {
4727 let event = receiver.try_recv().expect("expected quantity update");
4728 match event {
4729 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
4730 assert_eq!(
4731 updated.quantity,
4732 Quantity::new(5.202897, instrument.size_precision()),
4733 );
4734 assert_eq!(updated.venue_order_id, Some(venue_order_id));
4735 assert!(updated.reconciliation);
4736 }
4737 other => panic!("expected updated event, was {other:?}"),
4738 }
4739 assert!(
4740 fill_tracker
4741 .get_cumulative_filled(&venue_order_id)
4742 .is_none(),
4743 "order must be settled and removed from the tracker",
4744 );
4745 } else if expect_cancel {
4746 let event = receiver.try_recv().expect("expected IOC cancellation");
4747 match event {
4748 ExecutionEvent::Order(OrderEventAny::Canceled(canceled)) => {
4749 assert_eq!(canceled.venue_order_id, Some(venue_order_id));
4750 }
4751 other => panic!("expected canceled event, was {other:?}"),
4752 }
4753 assert!(
4754 fill_tracker
4755 .get_cumulative_filled(&venue_order_id)
4756 .is_none(),
4757 "canceled IOC must be settled and removed from the tracker",
4758 );
4759 } else {
4760 assert!(
4761 receiver.try_recv().is_err(),
4762 "resting order must not receive a terminal event",
4763 );
4764 assert!(
4765 fill_tracker
4766 .get_cumulative_filled(&venue_order_id)
4767 .is_some(),
4768 "ineligible order must stay tracked with open leaves",
4769 );
4770 }
4771 }
4772
4773 #[rstest]
4774 fn test_dispatch_taker_fill_gross_overfill_raises_qty_then_fills() {
4775 use crate::common::enums::{
4779 PolymarketEventType, PolymarketOrderSide, PolymarketOutcome, PolymarketTradeStatus,
4780 };
4781
4782 let instrument = test_instrument();
4783 let asset_id = instrument.id().symbol.inner();
4784 let size_precision = instrument.size_precision();
4785 let token_instruments = AtomicMap::new();
4786 token_instruments.insert(asset_id, instrument.clone());
4787
4788 let fill_tracker = OrderFillTrackerMap::new();
4789 let venue_order_id = VenueOrderId::from("0xtaker-gross-overfill");
4790 let submitted = Quantity::new(30.0, size_precision);
4791 fill_tracker.register(
4792 venue_order_id,
4793 submitted,
4794 OrderSide::Buy,
4795 instrument.id(),
4796 size_precision,
4797 instrument.price_precision(),
4798 );
4799
4800 let pending_submits = PendingSubmitTracker::default();
4801 let order_contexts = OrderContextRegistry::default();
4802 register_context(
4803 &order_contexts,
4804 venue_order_id,
4805 instrument.id(),
4806 "O-GROSS-OVERFILL",
4807 );
4808 order_contexts.mark_accepted(venue_order_id);
4809 let mut emitter = test_emitter();
4810 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
4811 emitter.set_sender(sender);
4812
4813 let ctx = WsDispatchContext {
4814 signer_type: PolymarketSignerType::Owner,
4815 token_instruments: &token_instruments,
4816 fill_tracker: &fill_tracker,
4817 pending_submits: &pending_submits,
4818 order_contexts: &order_contexts,
4819 emitter: &emitter,
4820 account_id: AccountId::from("POLY-001"),
4821 clock: nautilus_core::time::get_atomic_clock_realtime(),
4822 user_address: "0xtest",
4823 user_api_key: "test-key",
4824 };
4825 let mut state = WsDispatchState::default();
4826
4827 let trade = PolymarketUserTrade {
4829 asset_id,
4830 bucket_index: 0,
4831 fee_rate_bps: "0".to_string(),
4832 id: "trade-gross-overfill".to_string(),
4833 last_update: "1700000001".to_string(),
4834 maker_address: Ustr::from("0xmaker"),
4835 maker_orders: vec![],
4836 market: Ustr::from("0xmarket"),
4837 match_time: "1700000000".to_string(),
4838 outcome: PolymarketOutcome::yes(),
4839 owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
4840 price: "0.014".to_string(),
4841 side: PolymarketOrderSide::Buy,
4842 size: "33.846152".to_string(),
4843 status: PolymarketTradeStatus::Confirmed,
4844 taker_order_id: venue_order_id.as_str().to_string(),
4845 timestamp: "1700000000000".to_string(),
4846 trade_owner: Ustr::from("00000000-0000-0000-0000-000000000001"),
4847 transaction_hash: None,
4848 trader_side: PolymarketLiquiditySide::Taker,
4849 event_type: PolymarketEventType::Trade,
4850 };
4851
4852 dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
4853
4854 let expected_qty = Quantity::new(33.846152, size_precision);
4855
4856 match receiver.try_recv().expect("expected an updated event") {
4858 ExecutionEvent::Order(OrderEventAny::Updated(updated)) => {
4859 assert_eq!(updated.quantity, expected_qty);
4860 assert_eq!(updated.venue_order_id, Some(venue_order_id));
4861 }
4862 other => panic!("expected updated event raising qty to the fill, was {other:?}"),
4863 }
4864
4865 match receiver.try_recv().expect("expected a filled event") {
4866 ExecutionEvent::Order(OrderEventAny::Filled(filled)) => {
4867 assert_eq!(filled.last_qty, expected_qty);
4868 assert_eq!(filled.venue_order_id, venue_order_id);
4869 }
4870 other => panic!("expected filled event, was {other:?}"),
4871 }
4872 }
4873
4874 #[rstest]
4877 #[case(
4878 crate::common::enums::PolymarketOrderStatus::Unmatched,
4879 Some("invalid post-only order: order crosses book"),
4880 "Rejected"
4881 )]
4882 #[case(
4883 crate::common::enums::PolymarketOrderStatus::CanceledMarketResolved,
4884 None,
4885 "Expired"
4886 )]
4887 fn test_dispatch_order_terminal_status_emits_event(
4888 #[case] status: crate::common::enums::PolymarketOrderStatus,
4889 #[case] reason: Option<&str>,
4890 #[case] expected: &str,
4891 ) {
4892 use crate::common::enums::{
4893 PolymarketEventType, PolymarketOrderSide, PolymarketOrderType, PolymarketOutcome,
4894 };
4895
4896 let instrument = test_instrument();
4897 let asset_id = instrument.id().symbol.inner();
4898 let order_id = "0xterminal-order".to_string();
4899 let venue_order_id = VenueOrderId::from(order_id.as_str());
4900
4901 let token_instruments = AtomicMap::new();
4902 token_instruments.insert(asset_id, instrument.clone());
4903
4904 let fill_tracker = OrderFillTrackerMap::new();
4905 fill_tracker.register(
4906 venue_order_id,
4907 Quantity::from("10"),
4908 OrderSide::Buy,
4909 instrument.id(),
4910 instrument.size_precision(),
4911 instrument.price_precision(),
4912 );
4913
4914 let pending_submits = PendingSubmitTracker::default();
4915 let order_contexts = OrderContextRegistry::default();
4916 register_context(
4917 &order_contexts,
4918 venue_order_id,
4919 instrument.id(),
4920 "O-TERMINAL",
4921 );
4922 order_contexts.mark_accepted(venue_order_id);
4923 let mut emitter = test_emitter();
4924 let (sender, mut receiver) = tokio::sync::mpsc::unbounded_channel();
4925 emitter.set_sender(sender);
4926
4927 let ctx = WsDispatchContext {
4928 signer_type: PolymarketSignerType::Owner,
4929 token_instruments: &token_instruments,
4930 fill_tracker: &fill_tracker,
4931 pending_submits: &pending_submits,
4932 order_contexts: &order_contexts,
4933 emitter: &emitter,
4934 account_id: AccountId::from("POLY-001"),
4935 clock: nautilus_core::time::get_atomic_clock_realtime(),
4936 user_address: "0xabc",
4937 user_api_key: "xxx",
4938 };
4939 let mut state = WsDispatchState::default();
4940
4941 let order = PolymarketUserOrder {
4942 asset_id,
4943 associate_trades: None,
4944 created_at: Some("1775074735".to_string()),
4945 expiration: Some("0".to_string()),
4946 id: order_id,
4947 maker_address: Some(Ustr::from("0xabc")),
4948 market: Ustr::from("0x4134"),
4949 order_owner: Some(Ustr::from("xxx")),
4950 order_type: Some(PolymarketOrderType::FOK),
4951 original_size: "10".to_string(),
4952 outcome: Some(PolymarketOutcome::yes()),
4953 owner: Ustr::from("xxx"),
4954 price: "0.50".to_string(),
4955 side: PolymarketOrderSide::Buy,
4956 size_matched: "0".to_string(),
4957 status: Some(PolymarketUserOrderStatus::new(status, reason)),
4958 timestamp: "1775074738031".to_string(),
4959 event_type: PolymarketEventType::Placement,
4960 };
4961
4962 dispatch_user_message(&UserWsMessage::Order(order), &ctx, &mut state);
4963
4964 let event = receiver.try_recv().expect("expected terminal order event");
4965 match event {
4966 ExecutionEvent::Order(order_event) => {
4967 assert!(
4968 format!("{order_event:?}").starts_with(expected),
4969 "expected {expected}, was {order_event:?}"
4970 );
4971 assert_eq!(
4972 order_event.client_order_id(),
4973 ClientOrderId::from("O-TERMINAL")
4974 );
4975
4976 if let OrderEventAny::Rejected(rejected) = order_event {
4977 assert_eq!(
4978 rejected.reason,
4979 "invalid post-only order: order crosses book"
4980 );
4981 assert!(rejected.due_post_only);
4982 }
4983 }
4984 other => panic!("expected order event, was {other:?}"),
4985 }
4986 }
4987}