use ahash::AHashMap;
use anyhow::Context;
use nautilus_core::{UnixNanos, collections::AtomicMap, time::AtomicTime};
use nautilus_model::{
enums::{LiquiditySide, OrderStatus, PositionSideSpecified},
identifiers::{AccountId, ClientId, InstrumentId, Venue, VenueOrderId},
instruments::{Instrument, InstrumentAny},
reports::{ExecutionMassStatus, FillReport, OrderStatusReport, PositionStatusReport},
types::{Currency, Quantity},
};
use rust_decimal::Decimal;
use ustr::Ustr;
use super::{
order_fill_tracker::OrderFillTrackerMap,
parse::{
build_maker_fill_report, instrument_fee_exponent, instrument_taker_fee, parse_fill_report,
parse_order_status_report, parse_timestamp,
},
};
use crate::{
common::{
consts::{DUST_POSITION_THRESHOLD, DUST_SNAP_THRESHOLD_DEC, USDC_DECIMALS},
enums::{PolymarketLiquiditySide, PolymarketTradeStatus},
},
http::{
clob::PolymarketClobHttpClient,
data_api::PolymarketDataApiHttpClient,
models::{DataApiPosition, PolymarketOpenOrder, PolymarketTradeReport},
query::{GetOrdersParams, GetTradesParams},
},
};
pub(crate) struct FillContext<'a> {
pub account_id: AccountId,
pub user_address: &'a str,
pub api_key: &'a str,
pub pusd: Currency,
pub clock: &'static AtomicTime,
}
pub(crate) fn build_fill_reports_from_trades(
trades: &[PolymarketTradeReport],
ctx: &FillContext<'_>,
instruments: &AtomicMap<Ustr, InstrumentAny>,
instrument_filter: Option<InstrumentId>,
ts_init: UnixNanos,
) -> (Vec<FillReport>, usize) {
let mut reports = Vec::new();
let mut filtered = 0usize;
for trade in trades {
if trade.status != PolymarketTradeStatus::Confirmed {
continue;
}
let is_maker = trade.trader_side == PolymarketLiquiditySide::Maker;
if is_maker {
for mo in &trade.maker_orders {
if mo.maker_address != ctx.user_address && mo.owner != ctx.api_key {
continue;
}
let token_id = Ustr::from(mo.asset_id.as_str());
let instrument = instruments.get_cloned(&token_id);
let (instrument_id, price_prec, size_prec) = match instrument {
Some(i) => (i.id(), i.price_precision(), i.size_precision()),
None => {
filtered += 1;
continue;
}
};
if let Some(filter_id) = instrument_filter
&& instrument_id != filter_id
{
continue;
}
let ts_event =
parse_timestamp(&trade.match_time).unwrap_or(ctx.clock.get_time_ns());
let report = build_maker_fill_report(
mo,
&trade.id,
trade.trader_side,
trade.side,
trade.asset_id.as_str(),
ctx.account_id,
instrument_id,
price_prec,
size_prec,
ctx.pusd,
LiquiditySide::Maker,
ts_event,
ts_init,
);
reports.push(report);
}
} else {
let token_id = Ustr::from(trade.asset_id.as_str());
let instrument = instruments.get_cloned(&token_id);
let (instrument_id, price_prec, size_prec, taker_fee_rate, fee_exponent) =
match instrument {
Some(i) => (
i.id(),
i.price_precision(),
i.size_precision(),
instrument_taker_fee(&i),
instrument_fee_exponent(&i),
),
None => {
filtered += 1;
continue;
}
};
if let Some(filter_id) = instrument_filter
&& instrument_id != filter_id
{
continue;
}
let report = parse_fill_report(
trade,
instrument_id,
ctx.account_id,
None,
price_prec,
size_prec,
ctx.pusd,
taker_fee_rate,
fee_exponent,
ts_init,
);
reports.push(report);
}
}
(reports, filtered)
}
pub(crate) fn build_order_reports_from_orders(
orders: &[PolymarketOpenOrder],
instruments: &AtomicMap<Ustr, InstrumentAny>,
account_id: AccountId,
instrument_filter: Option<InstrumentId>,
ts_init: UnixNanos,
) -> (Vec<OrderStatusReport>, usize) {
let mut reports = Vec::new();
let mut filtered = 0usize;
for order in orders {
let token_id = Ustr::from(order.asset_id.as_str());
let instrument = instruments.get_cloned(&token_id);
let (instrument_id, price_prec, size_prec) = match instrument {
Some(i) => (i.id(), i.price_precision(), i.size_precision()),
None => {
filtered += 1;
continue;
}
};
if let Some(filter_id) = instrument_filter
&& instrument_id != filter_id
{
continue;
}
let report = parse_order_status_report(
order,
instrument_id,
account_id,
None,
price_prec,
size_prec,
ts_init,
);
reports.push(report);
}
(reports, filtered)
}
pub(crate) fn apply_fill_filters(
mut reports: Vec<FillReport>,
venue_order_id: Option<VenueOrderId>,
start: Option<UnixNanos>,
end: Option<UnixNanos>,
) -> Vec<FillReport> {
if let Some(vid) = venue_order_id {
reports.retain(|r| r.venue_order_id == vid);
}
match (start, end) {
(Some(s), Some(e)) => reports.retain(|r| r.ts_event >= s && r.ts_event <= e),
(Some(s), None) => reports.retain(|r| r.ts_event >= s),
(None, Some(e)) => reports.retain(|r| r.ts_event <= e),
(None, None) => {}
}
reports
}
pub(crate) fn build_position_reports(
positions: &[DataApiPosition],
account_id: AccountId,
ts: UnixNanos,
) -> Vec<PositionStatusReport> {
positions
.iter()
.filter(|p| {
if p.size > Decimal::ZERO && p.size < DUST_POSITION_THRESHOLD {
log::debug!(
"Filtering dust position: {}-{}, size={}",
p.condition_id,
p.asset,
p.size
);
}
p.size >= DUST_POSITION_THRESHOLD
})
.filter_map(|p| {
let instrument_id =
InstrumentId::from(format!("{}-{}.POLYMARKET", p.condition_id, p.asset).as_str());
let quantity = match Quantity::from_decimal_dp(p.size, USDC_DECIMALS as u8) {
Ok(quantity) => quantity,
Err(e) => {
log::warn!(
"Skipping invalid Data API position {}-{} size {}: {e}",
p.condition_id,
p.asset,
p.size,
);
return None;
}
};
Some(PositionStatusReport::new(
account_id,
instrument_id,
PositionSideSpecified::Long,
quantity,
ts,
ts,
None,
None,
p.avg_price,
))
})
.collect()
}
#[expect(clippy::too_many_arguments)]
pub(crate) async fn generate_mass_status(
http_client: &PolymarketClobHttpClient,
data_api_client: &PolymarketDataApiHttpClient,
instruments: &AtomicMap<Ustr, InstrumentAny>,
fill_tracker: &OrderFillTrackerMap,
ctx: &FillContext<'_>,
client_id: ClientId,
venue: Venue,
lookback_mins: Option<u64>,
) -> anyhow::Result<Option<ExecutionMassStatus>> {
let ts_init = ctx.clock.get_time_ns();
let orders = http_client
.get_orders(GetOrdersParams::default())
.await
.context("failed to fetch orders for mass status")?;
let (mut order_reports, orders_filtered) =
build_order_reports_from_orders(&orders, instruments, ctx.account_id, None, ts_init);
let trades = http_client
.get_trades(GetTradesParams::default())
.await
.context("failed to fetch trades for mass status")?;
let (mut fill_reports, fills_filtered) =
build_fill_reports_from_trades(&trades, ctx, instruments, None, ts_init);
fill_tracker.snap_fill_reports(&mut fill_reports);
let positions = data_api_client
.get_positions(ctx.user_address)
.await
.context("failed to fetch positions for mass status")?;
let position_reports = build_position_reports(&positions, ctx.account_id, ts_init);
if let Some(mins) = lookback_mins {
let now_ns = ctx.clock.get_time_ns();
let cutoff_ns = now_ns.as_u64().saturating_sub(mins * 60 * 1_000_000_000);
let cutoff = UnixNanos::from(cutoff_ns);
let orders_before = order_reports.len();
order_reports.retain(|r| r.ts_last >= cutoff);
let orders_removed = orders_before - order_reports.len();
let fills_before = fill_reports.len();
fill_reports.retain(|r| r.ts_event >= cutoff);
let fills_removed = fills_before - fill_reports.len();
log::debug!(
"Lookback filter ({}min): orders {}->{} (removed {}), fills {}->{} (removed {})",
mins,
orders_before,
order_reports.len(),
orders_removed,
fills_before,
fill_reports.len(),
fills_removed,
);
} else {
log::debug!(
"Generated mass status: {} orders ({} filtered), {} fills ({} filtered), {} positions",
order_reports.len(),
orders_filtered,
fill_reports.len(),
fills_filtered,
position_reports.len(),
);
}
cap_order_reports_to_confirmed_fills(&mut order_reports, &fill_reports);
let mut mass_status = ExecutionMassStatus::new(client_id, ctx.account_id, venue, ts_init, None);
mass_status.add_order_reports(order_reports);
mass_status.add_position_reports(position_reports);
mass_status.add_fill_reports(fill_reports);
Ok(Some(mass_status))
}
fn cap_order_reports_to_confirmed_fills(
order_reports: &mut [OrderStatusReport],
fill_reports: &[FillReport],
) {
let confirmed_by_order = confirmed_filled_quantities(fill_reports);
for report in order_reports {
let local_filled = Quantity::zero(report.quantity.precision);
cap_order_report_filled_qty(
report,
local_filled,
confirmed_by_order.get(&report.venue_order_id).copied(),
);
}
}
pub(crate) fn confirmed_filled_quantities(
fill_reports: &[FillReport],
) -> AHashMap<VenueOrderId, Decimal> {
let mut confirmed_by_order = AHashMap::new();
for fill in fill_reports {
*confirmed_by_order.entry(fill.venue_order_id).or_default() += fill.last_qty.as_decimal();
}
confirmed_by_order
}
pub(crate) fn cap_order_report_filled_qty(
report: &mut OrderStatusReport,
local_filled: Quantity,
confirmed_filled: Option<Decimal>,
) {
let confirmed_filled = confirmed_filled
.and_then(|qty| Quantity::from_decimal_dp(qty, report.quantity.precision).ok())
.unwrap_or_else(|| Quantity::zero(report.quantity.precision));
let capped = report.filled_qty.min(local_filled.max(confirmed_filled));
report.filled_qty = capped;
normalize_terminal_order_report_quantity(report);
}
pub(crate) fn normalize_terminal_order_report_quantity(report: &mut OrderStatusReport) {
if report.order_status != OrderStatus::Filled
|| report.filled_qty.is_zero()
|| report.filled_qty >= report.quantity
{
return;
}
let leaves = report.quantity.as_decimal() - report.filled_qty.as_decimal();
if leaves < DUST_SNAP_THRESHOLD_DEC {
log::debug!(
"Normalizing terminal order report {} quantity from {} to confirmed fills {}",
report.venue_order_id,
report.quantity,
report.filled_qty,
);
report.quantity = report.filled_qty;
}
}
#[cfg(test)]
mod tests {
use nautilus_model::{
enums::{LiquiditySide, OrderSide, OrderStatus, OrderType, TimeInForce},
identifiers::TradeId,
types::{Money, Price},
};
use rstest::rstest;
use super::*;
#[rstest]
fn caps_order_report_to_confirmed_companion_fills() {
let account_id = AccountId::from("POLY-001");
let instrument_id = InstrumentId::from("TEST.POLYMARKET");
let venue_order_id = VenueOrderId::from("V-1");
let mut reports = vec![OrderStatusReport::new(
account_id,
instrument_id,
None,
venue_order_id,
OrderSide::Buy,
OrderType::Limit,
TimeInForce::Gtc,
OrderStatus::PartiallyFilled,
Quantity::from("10.0000"),
Quantity::from("10.0000"),
UnixNanos::from(1),
UnixNanos::from(1),
UnixNanos::from(1),
None,
)];
let fills = vec![FillReport::new(
account_id,
instrument_id,
venue_order_id,
TradeId::from("T-1"),
OrderSide::Buy,
Quantity::from("4.0000"),
Price::from("0.5000"),
Money::new(0.0, Currency::pUSD()),
LiquiditySide::Taker,
None,
None,
UnixNanos::from(1),
UnixNanos::from(1),
None,
)];
cap_order_reports_to_confirmed_fills(&mut reports, &fills);
assert_eq!(reports[0].filled_qty, Quantity::from("4.0000"));
}
#[rstest]
#[case::below_threshold("99.995", "99.995")]
#[case::at_threshold("99.990", "100.000")]
fn normalizes_confirmed_dust_residual_to_order_quantity(
#[case] confirmed: &str,
#[case] expected_quantity: &str,
) {
let account_id = AccountId::from("POLY-001");
let instrument_id = InstrumentId::from("TEST.POLYMARKET");
let venue_order_id = VenueOrderId::from("V-DUST");
let mut reports = vec![OrderStatusReport::new(
account_id,
instrument_id,
None,
venue_order_id,
OrderSide::Buy,
OrderType::Limit,
TimeInForce::Gtc,
OrderStatus::Filled,
Quantity::from("100.000"),
Quantity::from("100.000"),
UnixNanos::from(1),
UnixNanos::from(1),
UnixNanos::from(1),
None,
)];
let fills = vec![FillReport::new(
account_id,
instrument_id,
venue_order_id,
TradeId::from("T-DUST"),
OrderSide::Buy,
Quantity::from(confirmed),
Price::from("0.5000"),
Money::zero(Currency::pUSD()),
LiquiditySide::Taker,
None,
None,
UnixNanos::from(1),
UnixNanos::from(1),
None,
)];
cap_order_reports_to_confirmed_fills(&mut reports, &fills);
assert_eq!(reports[0].quantity, Quantity::from(expected_quantity));
assert_eq!(reports[0].filled_qty, Quantity::from(confirmed));
}
}