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