Skip to main content

nautilus_polymarket/websocket/
dispatch.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! WebSocket message dispatch for the Polymarket execution client.
17//!
18//! Routes user-channel WS messages (order updates and trades) for orders submitted through this
19//! client into Nautilus order events (`OrderAccepted` / `OrderFilled` / `OrderFillVoided` /
20//! `OrderCanceled` / `OrderRejected` / `OrderExpired`), building them from the context captured at
21//! submit (`OrderContextRegistry`). Order-channel messages drive lifecycle events; trade-channel
22//! messages drive fills, and acceptance is synthesized before a fill or cancel that races ahead.
23//! Messages are emitted once the order is known (accepted, or with a submit in flight), otherwise
24//! buffered until acceptance. Reports are reserved for the `generate_*` query and reconciliation
25//! methods. Trade fills are emitted at `MATCHED`, retained until terminal settlement, and reversed
26//! with `OrderFillVoided` if the trade reaches `FAILED`.
27
28use 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/// Signal returned when a finalized trade requires an async account refresh.
81#[derive(Debug)]
82pub(crate) struct AccountRefreshRequest;
83
84/// Mutable state retained across user WebSocket stream generations.
85///
86/// Terminal cancel reports are re-emitted after fills to restore terminal state when fills race
87/// ahead of or arrive after cancel messages.
88#[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
511/// Immutable context borrowed from the async block's owned values.
512pub(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
541/// Top-level router: synchronous, returns signal for async account refresh.
542pub(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    // A known own order (submit in flight) self-registers on its first WS update
625    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    // Order updates can race ahead of trade messages, so cap filled_qty
648    // to what the fill tracker has recorded to prevent duplicate inferred fills
649    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    // Track cancel reports so we can re-emit them after late-arriving fills.
661    // Saved regardless of acceptance state so that cancels arriving during
662    // the HTTP round-trip are available once the order is later accepted.
663    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    // Tracked own orders route through order events; externally-managed orders
673    // (no captured context) buffer until accepted or fall back to reports.
674    let context = ctx.order_contexts.get(&venue_order_id);
675
676    // Emit fills first: a terminal status would otherwise close the order ahead of them
677    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        // Registered between the early accepted-check and here: emit rather than buffer
713        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
885/// Emits the terminal order event for a taker order once its trade confirms.
886///
887/// Taker fills receive no order-channel `MATCHED` update. FOK is atomic, so a sub-cent quantity
888/// difference can be normalized. IOC maps to FAK, so every positive remainder was killed by the
889/// venue and must close as `Canceled` without changing the venue-reported fill quantity.
890fn 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
1312/// Re-emits a saved cancel report after a fill to restore terminal state.
1313///
1314/// When fills race ahead of (or arrive after) cancel messages, the order can
1315/// get stuck in `PartiallyFilled`. This re-emission ensures the execution
1316/// engine transitions the order back to `Canceled`.
1317///
1318/// Skips re-emission when the fill tracker shows the order is fully filled,
1319/// because `Filled` is already terminal and a spurious cancel would fail
1320/// the `Filled -> Canceled` state transition.
1321fn 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    // Unfilled FOK cancellations carry an empty size_matched in captured venue messages
1383    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
1424/// Converts a venue-reported `original_size` on a user-channel order message into shares.
1425///
1426/// The venue echoes the signed `makerAmount`, which for a BUY is the pUSD budget rather than a
1427/// share count (see `compute_maker_taker_amounts`). Dividing by the order price recovers the
1428/// signed `takerAmount`, which is the share quantity the client submitted.
1429///
1430/// This is confirmed for the market order types (`FAK` and `FOK`), where a BUY at 0.01 for 100
1431/// shares reports `1`. A SELL signs shares as its maker amount and needs no conversion. Resting
1432/// types pass through unchanged: their denomination is unconfirmed, and converting a
1433/// share-denominated size would misreport every externally-managed resting order.
1434fn 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
1513/// Emits order events for a tracked own-order status update.
1514///
1515/// Order-channel messages drive lifecycle events only; fills arrive separately on the trade
1516/// channel as `OrderFilled`. `PartiallyFilled` / `Filled` statuses therefore emit no fill here,
1517/// they only ensure acceptance has been emitted so the order lifecycle stays well-formed.
1518fn 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
1550/// Emits `OrderAccepted` for a tracked order if acceptance has not yet been emitted.
1551///
1552/// Acceptance is also emitted on the submit happy path; the registry's dedup set ensures it
1553/// fires exactly once across the submit confirmation and the WS stream, including when a fill or
1554/// cancel races ahead of the acceptance message.
1555fn 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
1581/// Builds and emits an `OrderFilled` event for a tracked order, synthesizing acceptance first.
1582///
1583/// `info` carries the venue fill metadata (the raw trade fields) for trade-sourced fills, and is
1584/// `None` for order-path fills that have no originating trade payload.
1585fn 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
1670/// Flattens a user trade into a string map of venue fill metadata for `OrderFilled.info`.
1671///
1672/// Mirrors the v1 adapter, which attaches the full raw trade to each fill it generates. Scalar
1673/// fields map to their string form; nested fields (such as `maker_orders`) become their JSON text.
1674fn 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
1688/// Emits an `OrderUpdated` raising the order quantity to the actual BUY fill, before the fill.
1689///
1690/// A Polymarket BUY is bounded by the USDC it spends, so a marketable fill below the limit price
1691/// returns more shares than the nominal quantity. The engine rejects a fill past the order
1692/// quantity, so the quantity is raised first. The price is left unchanged (`None`).
1693fn 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
1721/// Emits an order-only reconciliation update which cannot change strategy position.
1722fn 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    /// Registers a tracked-order context so the dispatch routes the order through events.
1842    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        // A resting BUY already reports shares, so its size passes through unconverted
2022        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    // A market-order-type BUY reports the signed pUSD maker amount, so shares come from
2053    // dividing by the price. A SELL and the resting types already report shares.
2054    #[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    // A non-terminating division still rounds to the instrument's size precision
2110    #[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        // No registration: the submit response has not landed, so the order update registers it
2153        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        // The venue reported 1.01 pUSD for the 101 shares submitted at 0.01
2185        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        // No context registered, so the order surfaces as a report for reconciliation
2214        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        // Every raw trade field is captured (mirrors v1 info=msg.to_dict()).
2283        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        // Numeric fields flatten to their string form.
2291        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        // The `type` serde-rename key is preserved.
2298        assert_eq!(info[&Ustr::from("type")], Ustr::from("TRADE"));
2299        // Nested fields become their JSON text.
2300        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        // Order not registered in fill_tracker, so should be buffered
2341        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        // The tracked own order emits an OrderAccepted event carrying the client order ID.
2430        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        // First dispatch processes the trade
2565        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        // Second dispatch should be deduped, no additional fill
2569        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        // Register order so it is "accepted" but with no fills tracked
2964        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        // No context registered, so the order surfaces as a report (the external/reconciliation
2975        // fallback), where filled_qty is capped to tracked fills.
2976        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        // Register and record a partial fill (50 of 100)
3021        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        // No context registered, so the order surfaces as a report (the external/reconciliation
3033        // fallback), where filled_qty is capped to tracked fills.
3034        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        // Register order as accepted with original qty=100
3222        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        // Step 1: Dispatch cancel (simulates message A from the bug)
3254        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        // Step 2: Dispatch trade fill (simulates trade arriving after cancel)
3264        dispatch_user_message(&UserWsMessage::Trade(trade), &ctx, &mut state);
3265
3266        // Should get: filled event, then re-emitted canceled event
3267        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        // Register with qty=25 matching the trade size so the fill completes the order
4019        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        // Cancel then fill that completes the order
4056        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        // Channel should be empty: no re-emitted cancel for a fully-filled order
4063        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        // Fill tracker has NO registration (simulates HTTP still in-flight)
4079        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 cancel while order is not yet accepted
4101        dispatch_user_message(&UserWsMessage::Order(cancel_order), &ctx, &mut state);
4102
4103        // Cancel should be buffered (not emitted) AND saved to terminal_cancel_reports
4104        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    // A trade landing before the submit response buffers its fill, so the order update that
4125    // registers the order must emit that fill before its own terminal status
4126    #[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        // No registration: the submit response has not landed
4147        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        // The trade arrives first and buffers its fill
4179        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        // The order update registers the order and drains the buffered fill
4186        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        // The engine's state machine is what proves the order actually closes
4231        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    /// Replays the exact 5-message WS sequence from issue #3797.
4249    ///
4250    /// Messages in arrival order:
4251    ///   (A) Order Canceled, size_matched=0
4252    ///   (B) Trade fill 1.219511 (maker side)
4253    ///   (C) Order Canceled, size_matched=1.219511
4254    ///   (D) Order Canceled, size_matched=2.560972 (capped to tracked)
4255    ///   (E) Trade fill 1.341461 (maker side)
4256    ///
4257    /// Without the fix, the order ends in PartiallyFilled after (E).
4258    /// With the fix, a re-emitted cancel after (E) restores Canceled.
4259    #[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        // (A) Cancel with size_matched=0
4368        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        // (B) Trade fill 1.219511
4380        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        // Re-emitted cancel after fill (B)
4391        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        // (C) Cancel with size_matched=1.219511
4400        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        // (D) Cancel with size_matched=2.560972 (capped to tracked 1.219511)
4412        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        // (E) Trade fill 1.341461
4424        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        // The fix: re-emitted cancel after (E) restores terminal state
4436        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        // No more events
4445        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        // Reproduces the V2 market-BUY scenario that motivated the dust-snap
4454        // fix: SDK truncates the registered qty to USDC scale, but the
4455        // on-chain fill comes back at full precision and exceeds submitted
4456        // by microshares. Without the snap the engine rejects as overfill.
4457        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        // Submitted qty truncated to USDC scale.
4469        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            // Fill exceeds submitted_qty by 4 ulps at size_precision=6,
4521            // matching the production drift observed during smoke tests.
4522            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        // The dispatcher must record the snapped quantity in the tracker so
4535        // any subsequent ORDER MATCHED with size_matched > submitted_qty is
4536        // capped to it. record_fill happens before the FillReport is sent.
4537        let cumulative = fill_tracker
4538            .get_cumulative_filled(&venue_order_id)
4539            .expect("order must be registered");
4540        assert_eq!(cumulative, submitted);
4541
4542        // The emitted OrderFilled must carry the snapped qty so the engine
4543        // does not reject it as an overfill.
4544        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        // Takers receive no MATCHED order update. FOK is atomic, so a dust
4613        // difference normalizes the registered quantity. IOC maps to FAK, so
4614        // a positive remainder closes as Canceled without changing the fill.
4615        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        // A marketable BUY filled below its limit returns more shares than the nominal qty (a
4776        // gross overfill, beyond the dust band). The dispatcher must raise the order qty via
4777        // OrderUpdated before the OrderFilled, or the engine drops the fill as an overfill.
4778        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        // 33.846152 shares against a nominal 30: a marketable fill below the limit price.
4828        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        // The raise must precede the fill so the engine accepts the larger quantity.
4857        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    // Unmatched -> Rejected (placement never became live); CanceledMarketResolved -> Expired
4875    // (market settled). Both are tracked own-order terminal states emitted as order events.
4876    #[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}