use std::{collections::HashMap, result::Result as StdResult, str::from_utf8, sync::Arc};
use nautilus_core::{
consts::NAUTILUS_USER_AGENT,
time::{AtomicTime, get_atomic_clock_realtime},
};
use nautilus_model::{
data::BookOrder,
enums::{BookType, OrderSide},
identifiers::InstrumentId,
orderbook::OrderBook,
};
use nautilus_network::{
http::{HttpClient, HttpClientError, Method, USER_AGENT},
websocket::proxy::ProxyUrl,
};
use serde::{Deserialize, Serialize, de::DeserializeOwned};
use crate::{
common::{credential::Credential, enums::PolymarketOrderType, urls::clob_http_url},
http::{
error::{Error, Result},
models::{
ClobBookResponse, ClobMarketResponse, FeeRateResponse, PolymarketOpenOrder,
PolymarketOrder, PolymarketTradeReport, TickSizeResponse,
},
query::{
BalanceAllowance, BatchCancelResponse, CancelMarketOrdersParams, CancelResponse,
GetBalanceAllowanceParams, GetOrdersParams, GetTradesParams, OrderResponse,
PaginatedResponse,
},
rate_limits::{PolymarketRateLimiter, RateLimitHeaders, TradingBucket},
},
websocket::parse::{parse_price, parse_quantity},
};
const CURSOR_START: &str = "MA==";
const CURSOR_END: &str = "LTE=";
const PATH_ORDERS: &str = "/data/orders";
const PATH_TRADES: &str = "/data/trades";
const PATH_BALANCE_ALLOWANCE: &str = "/balance-allowance";
const PATH_BALANCE_ALLOWANCE_UPDATE: &str = "/balance-allowance/update";
const PATH_POST_ORDER: &str = "/order";
const PATH_POST_ORDERS: &str = "/orders";
const PATH_CANCEL_ALL: &str = "/cancel-all";
const PATH_CANCEL_MARKET_ORDERS: &str = "/cancel-market-orders";
const PATH_HEARTBEATS: &str = "/heartbeats";
const CLOB_CANCEL_BATCH_LIMIT: usize = 1_000;
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct PostOrderBody<'a> {
order: &'a PolymarketOrder,
owner: &'a str,
order_type: PolymarketOrderType,
#[serde(skip_serializing_if = "std::ops::Not::not")]
post_only: bool,
}
#[derive(Serialize)]
struct CancelOrderBody<'a> {
#[serde(rename = "orderID")]
order_id: &'a str,
}
#[derive(Serialize)]
struct HeartbeatRequest<'a> {
heartbeat_id: &'a str,
}
#[derive(Deserialize)]
struct HeartbeatWireResponse {
heartbeat_id: Option<String>,
status: Option<String>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum HeartbeatResponse {
Acknowledged(Option<String>),
Resynchronize(String),
}
#[derive(Debug, Clone)]
pub struct PolymarketClobHttpClient {
client: HttpClient,
rate_limiter: Arc<PolymarketRateLimiter>,
base_url: String,
credential: Credential,
address: String,
clock: &'static AtomicTime,
}
impl PolymarketClobHttpClient {
pub fn new(
credential: Credential,
address: String,
base_url: Option<String>,
timeout_secs: u64,
) -> StdResult<Self, HttpClientError> {
Self::new_with_proxy(credential, address, base_url, timeout_secs, None)
}
pub fn new_with_proxy(
credential: Credential,
address: String,
base_url: Option<String>,
timeout_secs: u64,
proxy_url: Option<ProxyUrl>,
) -> StdResult<Self, HttpClientError> {
let rate_limiter = PolymarketRateLimiter::for_signer(&address);
Ok(Self {
client: HttpClient::new(
Self::default_headers(),
RateLimitHeaders::names(),
vec![],
None,
Some(timeout_secs),
proxy_url.map(|url| url.expose().to_string()),
)?,
rate_limiter,
base_url: base_url
.unwrap_or_else(|| clob_http_url().to_string())
.trim_end_matches('/')
.to_string(),
credential,
address,
clock: get_atomic_clock_realtime(),
})
}
fn default_headers() -> HashMap<String, String> {
HashMap::from([
(USER_AGENT.to_string(), NAUTILUS_USER_AGENT.to_string()),
("Content-Type".to_string(), "application/json".to_string()),
])
}
fn url(&self, path: &str) -> String {
format!("{}{path}", self.base_url)
}
fn timestamp(&self) -> String {
(self.clock.get_time_ns().as_u64() / 1_000_000_000).to_string()
}
fn auth_headers(&self, method: &str, path: &str, body: &str) -> HashMap<String, String> {
let timestamp = self.timestamp();
let signature = self.credential.sign(×tamp, method, path, body);
HashMap::from([
("POLY_ADDRESS".to_string(), self.address.clone()),
("POLY_SIGNATURE".to_string(), signature),
("POLY_TIMESTAMP".to_string(), timestamp),
(
"POLY_API_KEY".to_string(),
self.credential.api_key().to_string(),
),
(
"POLY_PASSPHRASE".to_string(),
self.credential.passphrase().to_string(),
),
])
}
async fn send_get<P: Serialize, T: DeserializeOwned>(
&self,
path: &str,
params: Option<&P>,
auth: bool,
) -> Result<T> {
let headers = if auth {
Some(self.auth_headers("GET", path, ""))
} else {
None
};
let url = self.url(path);
let response = self
.client
.request_with_params(Method::GET, url, params, headers, None, None, None)
.await
.map_err(Error::from_http_client)?;
if response.status.is_success() {
serde_json::from_slice(&response.body).map_err(Error::Serde)
} else {
Err(Error::from_status_code(
response.status.as_u16(),
&response.body,
))
}
}
async fn send_get_optional<P: Serialize, T: DeserializeOwned>(
&self,
path: &str,
params: Option<&P>,
auth: bool,
) -> Result<Option<T>> {
let headers = if auth {
Some(self.auth_headers("GET", path, ""))
} else {
None
};
let url = self.url(path);
let response = self
.client
.request_with_params(Method::GET, url, params, headers, None, None, None)
.await
.map_err(Error::from_http_client)?;
if response.status.is_success() {
if response.body.is_empty() || response.body.as_ref() == b"null" {
Ok(None)
} else {
serde_json::from_slice(&response.body)
.map(Some)
.map_err(Error::Serde)
}
} else {
Err(Error::from_status_code(
response.status.as_u16(),
&response.body,
))
}
}
async fn send_post<T: DeserializeOwned>(
&self,
path: &'static str,
body_bytes: Vec<u8>,
cost: u32,
) -> Result<T> {
self.send_trading(
Method::POST,
path,
Some(body_bytes),
TradingBucket::Order,
cost,
|_| 0,
)
.await
}
async fn send_delete<T: DeserializeOwned>(
&self,
path: &'static str,
body_bytes: Option<Vec<u8>>,
cost: u32,
) -> Result<T> {
self.send_trading(
Method::DELETE,
path,
body_bytes,
TradingBucket::Cancel,
cost,
|_| 0,
)
.await
}
async fn send_delete_with_cancel_debit(
&self,
path: &'static str,
body_bytes: Option<Vec<u8>>,
) -> Result<BatchCancelResponse> {
self.send_trading(
Method::DELETE,
path,
body_bytes,
TradingBucket::Cancel,
1,
canceled_count,
)
.await
}
async fn send_trading<T: DeserializeOwned>(
&self,
method: Method,
path: &'static str,
body_bytes: Option<Vec<u8>>,
bucket: TradingBucket,
cost: u32,
post_response_cost: impl FnOnce(&T) -> u32,
) -> Result<T> {
self.rate_limiter.acquire(path, bucket, cost).await?;
let body_str = body_bytes
.as_deref()
.map(|b| from_utf8(b).map_err(|e| Error::decode(format!("UTF-8 error: {e}"))))
.transpose()?
.unwrap_or("");
let headers = Some(self.auth_headers(method.as_str(), path, body_str));
let url = self.url(path);
let response = self
.client
.request(method, url, None, headers, body_bytes, None, None)
.await
.map_err(Error::from_http_client)?;
let rate_limit_headers = RateLimitHeaders::parse(&response.headers);
if response.status.is_success() {
let decoded = serde_json::from_slice(&response.body);
let post_response_cost = decoded.as_ref().map_or(0, post_response_cost);
self.rate_limiter
.observe_response(
path,
bucket,
cost,
post_response_cost,
&rate_limit_headers,
false,
)
.await;
decoded.map_err(Error::Serde)
} else {
let rate_limited = response.status.as_u16() == 429;
self.rate_limiter
.observe_response(path, bucket, cost, 0, &rate_limit_headers, rate_limited)
.await;
if rate_limited {
Err(Error::rate_limit(
path,
cost,
rate_limit_headers.retry_after_ms(),
))
} else {
Err(Error::from_status_code(
response.status.as_u16(),
&response.body,
))
}
}
}
pub async fn post_heartbeat(&self, heartbeat_id: &str) -> Result<HeartbeatResponse> {
let body = HeartbeatRequest { heartbeat_id };
let body_bytes = serde_json::to_vec(&body).map_err(Error::Serde)?;
let body_str =
from_utf8(&body_bytes).map_err(|e| Error::decode(format!("UTF-8 error: {e}")))?;
let headers = Some(self.auth_headers("POST", PATH_HEARTBEATS, body_str));
let response = self
.client
.request(
Method::POST,
self.url(PATH_HEARTBEATS),
None,
headers,
Some(body_bytes),
None,
None,
)
.await
.map_err(Error::from_http_client)?;
let wire = serde_json::from_slice::<HeartbeatWireResponse>(&response.body);
if response.status.is_success() {
let wire = wire.map_err(Error::Serde)?;
if let Some(heartbeat_id) = wire.heartbeat_id.filter(|id| !id.is_empty()) {
return Ok(HeartbeatResponse::Acknowledged(Some(heartbeat_id)));
}
if wire.status.as_deref() == Some("ok") {
return Ok(HeartbeatResponse::Acknowledged(None));
}
return Err(Error::exchange("Heartbeat acknowledgment was invalid"));
}
if response.status.as_u16() == 400
&& let Ok(wire) = wire
&& let Some(heartbeat_id) = wire.heartbeat_id.filter(|id| !id.is_empty())
{
return Ok(HeartbeatResponse::Resynchronize(heartbeat_id));
}
Err(Error::from_status_code(
response.status.as_u16(),
&response.body,
))
}
pub async fn get_orders(
&self,
mut params: GetOrdersParams,
) -> Result<Vec<PolymarketOpenOrder>> {
if params.next_cursor.is_none() {
params.next_cursor = Some(CURSOR_START.to_string());
}
let mut all = Vec::new();
loop {
let page: PaginatedResponse<PolymarketOpenOrder> =
self.send_get(PATH_ORDERS, Some(¶ms), true).await?;
all.extend(page.data);
if page.next_cursor == CURSOR_END {
break;
}
params.next_cursor = Some(page.next_cursor);
}
Ok(all)
}
pub async fn get_order_optional(&self, order_id: &str) -> Result<Option<PolymarketOpenOrder>> {
let path = format!("/data/order/{order_id}");
self.send_get_optional::<(), _>(&path, None::<&()>, true)
.await
}
pub async fn get_order(&self, order_id: &str) -> Result<PolymarketOpenOrder> {
self.get_order_optional(order_id)
.await?
.ok_or_else(|| Error::decode(format!("Order {order_id} not found (empty response)")))
}
pub async fn get_trades(
&self,
mut params: GetTradesParams,
) -> Result<Vec<PolymarketTradeReport>> {
if params.next_cursor.is_none() {
params.next_cursor = Some(CURSOR_START.to_string());
}
let mut all = Vec::new();
loop {
let page: PaginatedResponse<PolymarketTradeReport> =
self.send_get(PATH_TRADES, Some(¶ms), true).await?;
all.extend(page.data);
if page.next_cursor == CURSOR_END {
break;
}
params.next_cursor = Some(page.next_cursor);
}
Ok(all)
}
pub async fn get_balance_allowance(
&self,
params: GetBalanceAllowanceParams,
) -> Result<BalanceAllowance> {
let headers = Some(self.auth_headers("GET", PATH_BALANCE_ALLOWANCE, ""));
let url = self.url(PATH_BALANCE_ALLOWANCE);
let response = self
.client
.request_with_params(Method::GET, url, Some(¶ms), headers, None, None, None)
.await
.map_err(Error::from_http_client)?;
if response.status.is_success() {
serde_json::from_slice(&response.body).map_err(Error::Serde)
} else {
Err(Error::from_status_code(
response.status.as_u16(),
&response.body,
))
}
}
pub async fn update_balance_allowance(&self, params: GetBalanceAllowanceParams) -> Result<()> {
self.send_get_optional::<_, serde_json::Value>(
PATH_BALANCE_ALLOWANCE_UPDATE,
Some(¶ms),
true,
)
.await?;
Ok(())
}
pub async fn post_order(
&self,
order: &PolymarketOrder,
order_type: PolymarketOrderType,
post_only: bool,
) -> Result<OrderResponse> {
let owner = self.credential.api_key().to_string();
let body = PostOrderBody {
order,
owner: &owner,
order_type,
post_only,
};
let body_bytes = serde_json::to_vec(&body).map_err(Error::Serde)?;
self.send_post(PATH_POST_ORDER, body_bytes, 1).await
}
pub async fn post_orders(
&self,
orders: &[(&PolymarketOrder, PolymarketOrderType, bool)],
) -> Result<Vec<OrderResponse>> {
let owner = self.credential.api_key().to_string();
let entries: Vec<PostOrderBody<'_>> = orders
.iter()
.map(|(order, order_type, post_only)| PostOrderBody {
order,
owner: &owner,
order_type: *order_type,
post_only: *post_only,
})
.collect();
let body_bytes = serde_json::to_vec(&entries).map_err(Error::Serde)?;
let cost = batch_cost(PATH_POST_ORDERS, entries.len())?;
self.send_post(PATH_POST_ORDERS, body_bytes, cost).await
}
pub async fn cancel_order(&self, order_id: &str) -> Result<CancelResponse> {
let body = CancelOrderBody { order_id };
let body_bytes = serde_json::to_vec(&body).map_err(Error::Serde)?;
self.send_delete(PATH_POST_ORDER, Some(body_bytes), 1).await
}
pub async fn cancel_orders(&self, order_ids: &[&str]) -> Result<BatchCancelResponse> {
let body_bytes = serde_json::to_vec(order_ids).map_err(Error::Serde)?;
let cost = batch_cost(PATH_POST_ORDERS, order_ids.len())?;
self.send_delete(PATH_POST_ORDERS, Some(body_bytes), cost)
.await
}
pub(crate) async fn cancel_batch_limit(&self) -> usize {
cancel_batch_limit(self.rate_limiter.burst(TradingBucket::Cancel).await)
}
pub async fn cancel_all(&self) -> Result<BatchCancelResponse> {
self.send_delete_with_cancel_debit(PATH_CANCEL_ALL, None)
.await
}
pub async fn cancel_market_orders(
&self,
params: CancelMarketOrdersParams,
) -> Result<BatchCancelResponse> {
let body_bytes = serde_json::to_vec(¶ms).map_err(Error::Serde)?;
self.send_delete_with_cancel_debit(PATH_CANCEL_MARKET_ORDERS, Some(body_bytes))
.await
}
pub async fn get_tick_size(&self, token_id: &str) -> Result<TickSizeResponse> {
let params = [("token_id", token_id)];
self.send_get("/tick-size", Some(¶ms), false).await
}
pub async fn get_fee_rate(&self, token_id: &str) -> Result<FeeRateResponse> {
let params = [("token_id", token_id)];
self.send_get("/fee-rate", Some(¶ms), false).await
}
pub async fn get_book(&self, token_id: &str) -> Result<ClobBookResponse> {
let params = [("token_id", token_id)];
self.send_get("/book", Some(¶ms), false).await
}
}
#[derive(Debug, Clone)]
pub struct PolymarketClobPublicClient {
client: HttpClient,
base_url: String,
}
impl PolymarketClobPublicClient {
pub fn new(base_url: Option<String>, timeout_secs: u64) -> StdResult<Self, HttpClientError> {
Self::new_with_proxy(base_url, timeout_secs, None)
}
pub fn new_with_proxy(
base_url: Option<String>,
timeout_secs: u64,
proxy_url: Option<ProxyUrl>,
) -> StdResult<Self, HttpClientError> {
Ok(Self {
client: HttpClient::new(
HashMap::from([
(USER_AGENT.to_string(), NAUTILUS_USER_AGENT.to_string()),
("Content-Type".to_string(), "application/json".to_string()),
]),
vec![],
vec![],
None,
Some(timeout_secs),
proxy_url.map(|url| url.expose().to_string()),
)?,
base_url: base_url
.unwrap_or_else(|| clob_http_url().to_string())
.trim_end_matches('/')
.to_string(),
})
}
pub async fn get_book(&self, token_id: &str) -> Result<ClobBookResponse> {
let params = [("token_id", token_id)];
let url = format!("{}/book", self.base_url);
let response = self
.client
.request_with_params(Method::GET, url, Some(¶ms), None, None, None, None)
.await
.map_err(Error::from_http_client)?;
if response.status.is_success() {
serde_json::from_slice(&response.body).map_err(Error::Serde)
} else {
Err(Error::from_status_code(
response.status.as_u16(),
&response.body,
))
}
}
pub async fn get_market(&self, condition_id: &str) -> Result<ClobMarketResponse> {
let url = format!("{}/markets/{condition_id}", self.base_url);
let response = self
.client
.request_with_params(Method::GET, url, None::<&()>, None, None, None, None)
.await
.map_err(Error::from_http_client)?;
if response.status.is_success() {
serde_json::from_slice(&response.body).map_err(Error::Serde)
} else {
Err(Error::from_status_code(
response.status.as_u16(),
&response.body,
))
}
}
pub async fn request_book_snapshot(
&self,
instrument_id: InstrumentId,
token_id: &str,
price_precision: u8,
size_precision: u8,
) -> anyhow::Result<OrderBook> {
let resp = self
.get_book(token_id)
.await
.map_err(|e| anyhow::anyhow!(e))?;
let mut book = OrderBook::new(instrument_id, BookType::L2_MBP);
for (i, level) in resp.bids.iter().enumerate() {
let price = parse_price(&level.price, price_precision)?;
let size = parse_quantity(&level.size, size_precision)?;
let order = BookOrder::new(OrderSide::Buy, price, size, i as u64);
book.add(order, 0, i as u64, Default::default());
}
let bids_len = resp.bids.len();
for (i, level) in resp.asks.iter().enumerate() {
let price = parse_price(&level.price, price_precision)?;
let size = parse_quantity(&level.size, size_precision)?;
let order = BookOrder::new(OrderSide::Sell, price, size, (bids_len + i) as u64);
book.add(order, 0, (bids_len + i) as u64, Default::default());
}
log::debug!(
"Fetched order book for {} with {} bids and {} asks",
instrument_id,
resp.bids.len(),
resp.asks.len(),
);
Ok(book)
}
}
fn batch_cost(endpoint: &'static str, len: usize) -> Result<u32> {
let cost = u32::try_from(len)
.map_err(|_| Error::bad_request(format!("{endpoint} batch length exceeds u32")))?;
if cost == 0 {
return Err(Error::bad_request(format!(
"{endpoint} batch must not be empty"
)));
}
Ok(cost)
}
fn cancel_batch_limit(burst: u32) -> usize {
usize::try_from(burst)
.unwrap_or(usize::MAX)
.min(CLOB_CANCEL_BATCH_LIMIT)
}
fn canceled_count(response: &BatchCancelResponse) -> u32 {
u32::try_from(response.canceled.len()).unwrap_or(u32::MAX)
}
#[cfg(test)]
mod tests {
use nautilus_model::{
enums::{BookType, OrderSide},
identifiers::InstrumentId,
types::{Price, Quantity},
};
use rstest::rstest;
use super::*;
use crate::http::models::{ClobBookLevel, ClobBookResponse};
fn build_book_from_response(resp: &ClobBookResponse) -> OrderBook {
let instrument_id = InstrumentId::from("TEST.POLYMARKET");
let price_precision = 2u8;
let size_precision = 2u8;
let mut book = OrderBook::new(instrument_id, BookType::L2_MBP);
for (i, level) in resp.bids.iter().enumerate() {
let price = parse_price(&level.price, price_precision).unwrap();
let size = parse_quantity(&level.size, size_precision).unwrap();
let order = BookOrder::new(OrderSide::Buy, price, size, i as u64);
book.add(order, 0, i as u64, Default::default());
}
let bids_len = resp.bids.len();
for (i, level) in resp.asks.iter().enumerate() {
let price = parse_price(&level.price, price_precision).unwrap();
let size = parse_quantity(&level.size, size_precision).unwrap();
let order = BookOrder::new(OrderSide::Sell, price, size, (bids_len + i) as u64);
book.add(order, 0, (bids_len + i) as u64, Default::default());
}
book
}
#[rstest]
fn test_build_order_book_from_clob_response() {
let resp = ClobBookResponse {
bids: vec![
ClobBookLevel {
price: "0.48".to_string(),
size: "100.00".to_string(),
},
ClobBookLevel {
price: "0.49".to_string(),
size: "200.00".to_string(),
},
ClobBookLevel {
price: "0.50".to_string(),
size: "150.00".to_string(),
},
],
asks: vec![
ClobBookLevel {
price: "0.51".to_string(),
size: "120.00".to_string(),
},
ClobBookLevel {
price: "0.52".to_string(),
size: "180.00".to_string(),
},
],
};
let book = build_book_from_response(&resp);
assert_eq!(book.instrument_id, InstrumentId::from("TEST.POLYMARKET"));
assert_eq!(book.book_type, BookType::L2_MBP);
assert_eq!(book.best_bid_price(), Some(Price::from("0.50")));
assert_eq!(book.best_ask_price(), Some(Price::from("0.51")));
assert_eq!(book.best_bid_size(), Some(Quantity::from("150.00")));
assert_eq!(book.best_ask_size(), Some(Quantity::from("120.00")));
assert_eq!(book.bids(None).count(), 3);
assert_eq!(book.asks(None).count(), 2);
}
#[rstest]
fn test_build_order_book_empty_response() {
let resp = ClobBookResponse {
bids: vec![],
asks: vec![],
};
let book = build_book_from_response(&resp);
assert!(book.best_bid_price().is_none());
assert!(book.best_ask_price().is_none());
}
#[rstest]
fn test_batch_cost_uses_entry_count_and_rejects_empty_batch() {
assert_eq!(batch_cost(PATH_POST_ORDERS, 15).unwrap(), 15);
assert_eq!(
batch_cost(PATH_POST_ORDERS, 0).unwrap_err().to_string(),
"bad request: /orders batch must not be empty"
);
}
#[rstest]
fn test_canceled_count_uses_only_successful_cancellations() {
let response = BatchCancelResponse {
canceled: vec!["order-1".to_string(), "order-2".to_string()],
not_canceled: ahash::AHashMap::from_iter([(
"order-3".to_string(),
Some("already canceled".to_string()),
)]),
};
assert_eq!(canceled_count(&response), 2);
}
#[rstest]
#[case::standard(120, 120)]
#[case::silver(600, 600)]
#[case::gold(1_200, 1_000)]
#[case::elite(1_800, 1_000)]
fn test_cancel_batch_limit_uses_tier_burst_and_venue_ceiling(
#[case] burst: u32,
#[case] expected: usize,
) {
assert_eq!(cancel_batch_limit(burst), expected);
}
}