use std::{
error::Error as StdError,
fmt::Display,
sync::{
Arc,
atomic::{AtomicBool, Ordering},
},
};
use nautilus_core::UnixNanos;
use nautilus_model::{
enums::{OrderSide, TimeInForce},
identifiers::VenueOrderId,
types::Quantity,
};
use nautilus_network::retry::{RetryConfig, RetryManager};
use rust_decimal::Decimal;
use super::{
order_builder::PolymarketOrderBuilder,
parse::{adjust_market_buy_amount, calculate_market_price},
types::{LimitOrderSubmitRequest, SignedLimitOrderSubmission},
};
use crate::{
common::enums::{PolymarketOrderSide, PolymarketOrderType},
http::{
clob::PolymarketClobHttpClient,
error::{Error, Result as HttpResult},
models::{PolymarketOpenOrder, PolymarketOrder},
query::{CancelResponse, OrderResponse},
},
};
#[derive(Debug, Clone)]
pub(crate) struct MarketBuyFeeContext {
pub user_pusd_balance: Decimal,
pub fee_rate: Decimal,
pub fee_exponent: f64,
pub builder_taker_fee_rate: Decimal,
}
#[derive(Debug, Clone)]
pub(crate) struct MarketOrderSubmitRequest {
pub(crate) token_id: String,
pub(crate) side: OrderSide,
pub(crate) amount: Quantity,
pub(crate) time_in_force: TimeInForce,
pub(crate) neg_risk: bool,
pub(crate) tick_decimals: u32,
pub(crate) fee_context: Option<MarketBuyFeeContext>,
}
#[derive(Debug, Clone)]
pub(crate) struct MarketOrderSubmitResult {
pub response: OrderResponse,
pub expected_base_qty: Decimal,
pub expected_venue_order_id: VenueOrderId,
}
#[derive(Debug, Clone)]
pub(crate) struct UnknownSubmitError {
pub reason: String,
pub expected_venue_order_id: VenueOrderId,
pub expected_base_qty: Option<Decimal>,
}
impl Display for UnknownSubmitError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"submit outcome unknown for {}: {}",
self.expected_venue_order_id, self.reason
)
}
}
impl StdError for UnknownSubmitError {}
#[derive(Debug, Clone)]
pub(crate) struct OrderSubmitter {
http_client: PolymarketClobHttpClient,
order_builder: Arc<PolymarketOrderBuilder>,
retry_manager: Arc<RetryManager<Error>>,
}
impl OrderSubmitter {
pub(crate) fn new(
http_client: PolymarketClobHttpClient,
order_builder: Arc<PolymarketOrderBuilder>,
retry_config: RetryConfig,
) -> Self {
Self {
http_client,
order_builder,
retry_manager: Arc::new(RetryManager::new(retry_config)),
}
}
pub(crate) async fn submit_market_order(
&self,
request: MarketOrderSubmitRequest,
) -> anyhow::Result<MarketOrderSubmitResult> {
let MarketOrderSubmitRequest {
token_id,
side,
amount,
time_in_force,
neg_risk,
tick_decimals,
fee_context,
} = request;
let poly_side = PolymarketOrderSide::try_from(side)
.map_err(|e| anyhow::anyhow!("Invalid order side: {e}"))?;
let order_type = PolymarketOrderType::from_market_time_in_force(time_in_force)
.map_err(|e| anyhow::anyhow!("{e}"))?;
let amount_dec = amount.as_decimal();
let book = self
.http_client
.get_book(&token_id)
.await
.map_err(|e| anyhow::anyhow!("Failed to fetch order book: {e}"))?;
let levels = match poly_side {
PolymarketOrderSide::Buy => &book.asks,
PolymarketOrderSide::Sell => &book.bids,
};
let result = calculate_market_price(levels, amount_dec, poly_side)
.map_err(|e| anyhow::anyhow!("Market price calculation failed: {e}"))?;
let signed_amount = match (poly_side, fee_context) {
(PolymarketOrderSide::Buy, Some(ctx)) => adjust_market_buy_amount(
amount_dec,
ctx.user_pusd_balance,
result.crossing_price,
ctx.fee_rate,
ctx.fee_exponent,
ctx.builder_taker_fee_rate,
)?,
_ => amount_dec,
};
let poly_order = self
.order_builder
.build_market_order(
&token_id,
poly_side,
result.crossing_price,
signed_amount,
neg_risk,
tick_decimals,
)
.map_err(|e| anyhow::anyhow!("Failed to build market order: {e}"))?;
let signed_base_qty =
signed_base_quantity(poly_order.maker_amount, poly_order.taker_amount, poly_side);
let expected_venue_order_id = self
.order_builder
.expected_order_id(&poly_order, neg_risk)?;
let http_client = self.http_client.clone();
let saw_unknown_outcome = Arc::new(AtomicBool::new(false));
let response = match self
.retry_manager
.execute_with_retry_with_delay(
"submit_market_order",
|| {
let http_client = http_client.clone();
let poly_order = poly_order.clone();
let saw_unknown_outcome = saw_unknown_outcome.clone();
async move {
let result = http_client.post_order(&poly_order, order_type, false).await;
if result.as_ref().is_err_and(Error::is_submit_outcome_unknown) {
saw_unknown_outcome.store(true, Ordering::Release);
}
result
}
},
|e| e.is_retryable(),
Error::retry_after,
Error::transport,
)
.await
{
Ok(response) => response,
Err(e)
if submit_outcome_is_unknown(&e, saw_unknown_outcome.load(Ordering::Acquire)) =>
{
return Err(UnknownSubmitError {
reason: e.to_string(),
expected_venue_order_id,
expected_base_qty: Some(signed_base_qty),
}
.into());
}
Err(e) => anyhow::bail!("{e}"),
};
Ok(MarketOrderSubmitResult {
response,
expected_base_qty: signed_base_qty,
expected_venue_order_id,
})
}
pub(crate) async fn cancel_order(&self, venue_order_id: &str) -> HttpResult<CancelResponse> {
let http_client = self.http_client.clone();
let order_id = venue_order_id.to_string();
self.retry_manager
.execute_with_retry_with_delay(
"cancel_order",
|| {
let http_client = http_client.clone();
let order_id = order_id.clone();
async move { http_client.cancel_order(&order_id).await }
},
|e| e.is_retryable(),
Error::retry_after,
Error::transport,
)
.await
}
pub(crate) async fn cancel_orders(
&self,
venue_order_ids: &[&str],
) -> HttpResult<CancelResponse> {
let order_ids: Vec<String> = venue_order_ids.iter().map(|s| s.to_string()).collect();
if order_ids.is_empty() {
return self.cancel_orders_chunk(&order_ids).await;
}
let mut response = CancelResponse::default();
let mut offset = 0;
while offset < order_ids.len() {
let limit = self.http_client.cancel_batch_limit().await;
let end = offset.saturating_add(limit).min(order_ids.len());
match self.cancel_orders_chunk(&order_ids[offset..end]).await {
Ok(chunk_response) => {
response.merge(chunk_response);
offset = end;
}
Err(e @ Error::BurstExceeded { .. }) => {
log::warn!("Cancellation batch limit changed, reselecting chunk: {e}");
}
Err(e) => return Err(e),
}
}
Ok(response)
}
async fn cancel_orders_chunk(&self, order_ids: &[String]) -> HttpResult<CancelResponse> {
let http_client = self.http_client.clone();
let order_ids = order_ids.to_vec();
self.retry_manager
.execute_with_retry_with_delay(
"cancel_orders",
|| {
let http_client = http_client.clone();
let order_ids = order_ids.clone();
async move {
let refs: Vec<&str> = order_ids.iter().map(String::as_str).collect();
http_client.cancel_orders(&refs).await
}
},
|e| e.is_retryable(),
Error::retry_after,
Error::transport,
)
.await
}
pub(crate) async fn get_order(
&self,
order_id: &str,
) -> anyhow::Result<Option<PolymarketOpenOrder>> {
let http_client = self.http_client.clone();
let oid = order_id.to_string();
self.retry_manager
.execute_with_retry_with_delay(
"get_order",
|| {
let http_client = http_client.clone();
let oid = oid.clone();
async move { http_client.get_order_optional(&oid).await }
},
|e| e.is_retryable(),
Error::retry_after,
Error::transport,
)
.await
.map_err(|e| anyhow::anyhow!("Failed to fetch order status: {e}"))
}
pub(crate) async fn prepare_limit_order_submissions(
&self,
requests: &[LimitOrderSubmitRequest],
) -> Vec<anyhow::Result<SignedLimitOrderSubmission>> {
let futures = requests
.iter()
.map(|request| self.prepare_limit_order_submission(request));
futures_util::future::join_all(futures).await
}
pub(crate) async fn prepare_limit_order_submission(
&self,
request: &LimitOrderSubmitRequest,
) -> anyhow::Result<SignedLimitOrderSubmission> {
let order_type = PolymarketOrderType::try_from(request.time_in_force)
.map_err(|e| anyhow::anyhow!("Unsupported time in force: {e}"))?;
let side = PolymarketOrderSide::try_from(request.side)
.map_err(|e| anyhow::anyhow!("Invalid order side: {e}"))?;
let expiration = limit_order_expiration(request.expire_time);
let order = self
.order_builder
.build_limit_order(
&request.token_id,
side,
request.price.as_decimal(),
request.quantity.as_decimal(),
order_type,
&expiration,
request.neg_risk,
request.tick_decimals,
)
.map_err(|e| anyhow::anyhow!("{e}"))?;
let expected_venue_order_id = self
.order_builder
.expected_order_id(&order, request.neg_risk)?;
Ok(SignedLimitOrderSubmission {
order,
order_type,
post_only: request.post_only,
expected_venue_order_id,
})
}
pub(crate) async fn post_limit_order_submission(
&self,
submission: SignedLimitOrderSubmission,
) -> crate::http::error::Result<OrderResponse> {
let http_client = self.http_client.clone();
let saw_unknown_outcome = Arc::new(AtomicBool::new(false));
let result = self
.retry_manager
.execute_with_retry_with_delay(
"submit_limit_order",
|| {
let http_client = http_client.clone();
let submission = submission.clone();
let saw_unknown_outcome = saw_unknown_outcome.clone();
async move {
let result = http_client
.post_order(
&submission.order,
submission.order_type,
submission.post_only,
)
.await;
if result.as_ref().is_err_and(Error::is_submit_outcome_unknown) {
saw_unknown_outcome.store(true, Ordering::Release);
}
result
}
},
|e| e.is_retryable(),
Error::retry_after,
Error::transport,
)
.await;
match result {
Err(e) if saw_unknown_outcome.load(Ordering::Acquire) => Err(Error::transport(
format!("submit outcome unknown after an earlier attempt: {e}"),
)),
result => result,
}
}
pub(crate) async fn post_limit_order_submissions(
&self,
submissions: Vec<SignedLimitOrderSubmission>,
) -> crate::http::error::Result<Vec<OrderResponse>> {
let order_refs: Vec<(&PolymarketOrder, PolymarketOrderType, bool)> = submissions
.iter()
.map(|submission| {
(
&submission.order,
submission.order_type,
submission.post_only,
)
})
.collect();
self.http_client.post_orders(&order_refs).await
}
}
fn submit_outcome_is_unknown(error: &Error, earlier_attempt_unknown: bool) -> bool {
error.is_submit_outcome_unknown() || earlier_attempt_unknown
}
fn signed_base_quantity(
maker_amount: Decimal,
taker_amount: Decimal,
side: PolymarketOrderSide,
) -> Decimal {
let usdc_scale = Decimal::from(1_000_000u32);
match side {
PolymarketOrderSide::Buy => taker_amount / usdc_scale,
PolymarketOrderSide::Sell => maker_amount / usdc_scale,
}
}
fn limit_order_expiration(expire_time: Option<UnixNanos>) -> String {
match expire_time {
Some(ns) if ns.as_u64() > 0 => (ns.as_u64() / 1_000_000_000).to_string(),
_ => "0".to_string(),
}
}
#[cfg(test)]
mod tests {
use rstest::rstest;
use rust_decimal_macros::dec;
use super::*;
#[rstest]
#[case::none(None, "0")]
#[case::zero(Some(UnixNanos::from(0u64)), "0")]
#[case::one_second(Some(UnixNanos::from(1_000_000_000u64)), "1")]
#[case::sub_second_truncates(Some(UnixNanos::from(1_500_000_000u64)), "1")]
#[case::typical(Some(UnixNanos::from(1_735_689_600_000_000_000u64)), "1735689600")]
fn test_limit_order_expiration(#[case] expire_time: Option<UnixNanos>, #[case] expected: &str) {
assert_eq!(limit_order_expiration(expire_time), expected);
}
#[rstest]
#[case::buy(dec!(4_800_000), dec!(5_202_897), PolymarketOrderSide::Buy, dec!(5.202897))]
#[case::sell(dec!(5_200_000), dec!(4_992_000), PolymarketOrderSide::Sell, dec!(5.2))]
fn test_signed_base_quantity_uses_share_denominated_wire_leg(
#[case] maker_amount: Decimal,
#[case] taker_amount: Decimal,
#[case] side: PolymarketOrderSide,
#[case] expected: Decimal,
) {
assert_eq!(
signed_base_quantity(maker_amount, taker_amount, side),
expected
);
}
#[rstest]
fn test_rate_limit_is_definitive_unless_an_earlier_attempt_was_unknown() {
let rate_limit = Error::rate_limit("/order", 1, Some(2_000));
assert!(!submit_outcome_is_unknown(&rate_limit, false));
assert!(submit_outcome_is_unknown(&rate_limit, true));
assert!(submit_outcome_is_unknown(
&Error::transport("connection reset"),
false
));
assert!(submit_outcome_is_unknown(&Error::Timeout, false));
}
}