Skip to main content

nautilus_polymarket/http/
rate_limits.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//! Trading rate limits for the Polymarket CLOB API.
17
18use std::{
19    collections::HashMap,
20    fmt::Display,
21    num::NonZeroU32,
22    sync::{Arc, LazyLock},
23    time::{Duration, SystemTime, UNIX_EPOCH},
24};
25
26use nautilus_network::ratelimiter::quota::Quota;
27use parking_lot::Mutex as BlockingMutex;
28use tokio::{
29    sync::Mutex,
30    time::{Instant, sleep},
31};
32
33use super::error::{Error, Result};
34use crate::common::consts::HTTP_RATE_LIMIT;
35
36pub(crate) const HEADER_RATE_LIMIT_REMAINING: &str = "Poly-RateLimit-Remaining";
37pub(crate) const HEADER_RATE_LIMIT_RESET: &str = "Poly-RateLimit-Reset";
38pub(crate) const HEADER_RATE_LIMIT_TIER: &str = "Poly-RateLimit-Tier";
39pub(crate) const HEADER_RATE_LIMIT_WARNING: &str = "Poly-RateLimit-Warning";
40pub(crate) const HEADER_RETRY_AFTER: &str = "Retry-After";
41
42const RATE_LIMIT_HEADERS: [&str; 5] = [
43    HEADER_RATE_LIMIT_REMAINING,
44    HEADER_RATE_LIMIT_RESET,
45    HEADER_RATE_LIMIT_TIER,
46    HEADER_RATE_LIMIT_WARNING,
47    HEADER_RETRY_AFTER,
48];
49
50static SIGNER_LIMITERS: LazyLock<BlockingMutex<HashMap<String, Arc<PolymarketRateLimiter>>>> =
51    LazyLock::new(|| BlockingMutex::new(HashMap::new()));
52
53/// Global REST quota for Polymarket Gamma API requests.
54pub static POLYMARKET_GAMMA_REST_QUOTA: LazyLock<Quota> =
55    LazyLock::new(|| Quota::per_minute(NonZeroU32::new(HTTP_RATE_LIMIT).unwrap()));
56
57#[derive(Clone, Copy, Debug, PartialEq, Eq)]
58pub(crate) enum TradingBucket {
59    Order,
60    Cancel,
61}
62
63impl Display for TradingBucket {
64    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
65        match self {
66            Self::Order => f.write_str("order"),
67            Self::Cancel => f.write_str("cancel"),
68        }
69    }
70}
71
72#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
73enum RateLimitTier {
74    #[default]
75    Standard,
76    Copper,
77    Bronze,
78    Silver,
79    Gold,
80    Platinum,
81    Diamond,
82    Elite,
83}
84
85impl RateLimitTier {
86    fn parse(value: &str) -> Option<Self> {
87        match value.trim().to_ascii_lowercase().as_str() {
88            "standard" => Some(Self::Standard),
89            "copper" => Some(Self::Copper),
90            "bronze" => Some(Self::Bronze),
91            "silver" => Some(Self::Silver),
92            "gold" => Some(Self::Gold),
93            "platinum" => Some(Self::Platinum),
94            "diamond" => Some(Self::Diamond),
95            "elite" => Some(Self::Elite),
96            _ => None,
97        }
98    }
99
100    const fn limits(self) -> TierLimits {
101        match self {
102            Self::Standard => TierLimits::new(40, 60, 80, 120, true),
103            Self::Copper => TierLimits::new(60, 90, 120, 180, true),
104            Self::Bronze => TierLimits::new(80, 120, 160, 240, true),
105            Self::Silver => TierLimits::new(200, 300, 400, 600, true),
106            Self::Gold => TierLimits::new(400, 600, 800, 1_200, true),
107            Self::Platinum => TierLimits::new(450, 675, 900, 1_350, false),
108            Self::Diamond => TierLimits::new(525, 787, 1_050, 1_575, false),
109            Self::Elite => TierLimits::new(600, 900, 1_200, 1_800, false),
110        }
111    }
112}
113
114impl Display for RateLimitTier {
115    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
116        match self {
117            Self::Standard => f.write_str("Standard"),
118            Self::Copper => f.write_str("Copper"),
119            Self::Bronze => f.write_str("Bronze"),
120            Self::Silver => f.write_str("Silver"),
121            Self::Gold => f.write_str("Gold"),
122            Self::Platinum => f.write_str("Platinum"),
123            Self::Diamond => f.write_str("Diamond"),
124            Self::Elite => f.write_str("Elite"),
125        }
126    }
127}
128
129#[derive(Clone, Copy, Debug)]
130struct TierLimits {
131    order: BucketLimits,
132    cancel: BucketLimits,
133    negative_cancel_balance: bool,
134}
135
136impl TierLimits {
137    const fn new(
138        order_rate: u32,
139        order_burst: u32,
140        cancel_rate: u32,
141        cancel_burst: u32,
142        negative_cancel_balance: bool,
143    ) -> Self {
144        Self {
145            order: BucketLimits::new(order_rate, order_burst),
146            cancel: BucketLimits::new(cancel_rate, cancel_burst),
147            negative_cancel_balance,
148        }
149    }
150
151    const fn bucket(self, bucket: TradingBucket) -> BucketLimits {
152        match bucket {
153            TradingBucket::Order => self.order,
154            TradingBucket::Cancel => self.cancel,
155        }
156    }
157}
158
159#[derive(Clone, Copy, Debug)]
160struct BucketLimits {
161    rate: f64,
162    burst: u32,
163}
164
165impl BucketLimits {
166    const fn new(rate: u32, burst: u32) -> Self {
167        Self {
168            rate: rate as f64,
169            burst,
170        }
171    }
172}
173
174#[derive(Clone, Debug, Default, PartialEq)]
175pub(crate) struct RateLimitHeaders {
176    pub(crate) remaining: Option<f64>,
177    pub(crate) reset: Option<f64>,
178    tier: Option<RateLimitTier>,
179    pub(crate) warning: bool,
180    pub(crate) retry_after: Option<Duration>,
181}
182
183impl RateLimitHeaders {
184    pub(crate) fn names() -> Vec<String> {
185        RATE_LIMIT_HEADERS.into_iter().map(str::to_string).collect()
186    }
187
188    pub(crate) fn parse(headers: &HashMap<String, String>) -> Self {
189        Self {
190            remaining: parse_number(headers, HEADER_RATE_LIMIT_REMAINING, true),
191            reset: parse_number(headers, HEADER_RATE_LIMIT_RESET, false),
192            tier: parse_tier(headers),
193            warning: parse_warning(headers),
194            retry_after: parse_retry_after(headers),
195        }
196    }
197
198    pub(crate) fn retry_after_ms(&self) -> Option<u64> {
199        self.retry_after.map(|duration| {
200            let partial_millisecond = u128::from(duration.subsec_nanos() % 1_000_000 != 0);
201            u64::try_from(duration.as_millis().saturating_add(partial_millisecond))
202                .unwrap_or(u64::MAX)
203        })
204    }
205
206    pub(crate) fn has_signer_headers(&self) -> bool {
207        self.remaining.is_some() || self.reset.is_some() || self.tier.is_some()
208    }
209}
210
211#[derive(Debug)]
212pub(crate) struct PolymarketRateLimiter {
213    state: Mutex<RateLimitState>,
214}
215
216impl PolymarketRateLimiter {
217    pub(crate) fn for_signer(signer: &str) -> Arc<Self> {
218        let signer = signer.to_ascii_lowercase();
219        let mut limiters = SIGNER_LIMITERS.lock();
220        limiters
221            .entry(signer)
222            .or_insert_with(|| {
223                Arc::new(Self {
224                    state: Mutex::new(RateLimitState::new()),
225                })
226            })
227            .clone()
228    }
229
230    pub(crate) async fn acquire(
231        &self,
232        endpoint: &'static str,
233        bucket: TradingBucket,
234        cost: u32,
235    ) -> Result<()> {
236        if cost == 0 {
237            return Err(Error::bad_request(format!(
238                "{endpoint} token cost must be positive"
239            )));
240        }
241
242        loop {
243            let wait = {
244                let mut state = self.state.lock().await;
245                let now = Instant::now();
246                state.refill(now);
247
248                let limits = state.tier.limits().bucket(bucket);
249                if cost > limits.burst {
250                    return Err(Error::BurstExceeded {
251                        endpoint,
252                        token_cost: cost,
253                        tier: state.tier.to_string(),
254                        bucket: bucket.to_string(),
255                        burst: limits.burst,
256                    });
257                }
258
259                let bucket = state.bucket_mut(bucket);
260                let blocked_for = bucket
261                    .blocked_until
262                    .and_then(|blocked_until| blocked_until.checked_duration_since(now))
263                    .unwrap_or_default();
264                let token_wait = if bucket.tokens >= f64::from(cost) {
265                    Duration::ZERO
266                } else {
267                    Duration::from_secs_f64((f64::from(cost) - bucket.tokens) / limits.rate)
268                };
269                let wait = blocked_for.max(token_wait);
270
271                if wait.is_zero() {
272                    bucket.tokens -= f64::from(cost);
273                    return Ok(());
274                }
275
276                wait
277            };
278
279            sleep(wait).await;
280        }
281    }
282
283    pub(crate) async fn burst(&self, bucket: TradingBucket) -> u32 {
284        self.state.lock().await.tier.limits().bucket(bucket).burst
285    }
286
287    pub(crate) async fn observe_response(
288        &self,
289        endpoint: &'static str,
290        bucket: TradingBucket,
291        request_cost: u32,
292        post_response_cost: u32,
293        headers: &RateLimitHeaders,
294        rejected: bool,
295    ) {
296        let effective_tier = {
297            let mut state = self.state.lock().await;
298            let now = Instant::now();
299            state.refill(now);
300
301            if let Some(tier) = headers.tier
302                && tier != state.tier
303            {
304                state.set_tier(tier);
305            }
306
307            let limits = state.tier.limits();
308            let bucket_limits = limits.bucket(bucket);
309            let bucket_state = state.bucket_mut(bucket);
310            bucket_state.tokens -= f64::from(post_response_cost);
311
312            if let Some(remaining) = headers.remaining {
313                bucket_state.tokens = bucket_state
314                    .tokens
315                    .min(remaining)
316                    .min(f64::from(bucket_limits.burst));
317            }
318
319            if bucket != TradingBucket::Cancel || !limits.negative_cancel_balance {
320                bucket_state.tokens = bucket_state.tokens.max(0.0);
321            }
322
323            if let Some(retry_after) = headers.retry_after {
324                bucket_state.block_for(now, retry_after);
325            }
326
327            if (rejected || bucket_state.tokens < 0.0)
328                && let Some(reset_wait) = reset_wait(headers.reset)
329            {
330                bucket_state.block_for(now, reset_wait);
331            }
332
333            state.tier
334        };
335
336        if headers.warning {
337            log::warn!(
338                "{}",
339                warning_message(
340                    endpoint,
341                    request_cost.saturating_add(post_response_cost),
342                    effective_tier,
343                    headers,
344                )
345            );
346        }
347    }
348}
349
350#[derive(Debug)]
351struct RateLimitState {
352    tier: RateLimitTier,
353    order: BucketState,
354    cancel: BucketState,
355}
356
357impl RateLimitState {
358    fn new() -> Self {
359        let now = Instant::now();
360        let limits = RateLimitTier::Standard.limits();
361        Self {
362            tier: RateLimitTier::Standard,
363            order: BucketState::new(limits.order, now),
364            cancel: BucketState::new(limits.cancel, now),
365        }
366    }
367
368    fn refill(&mut self, now: Instant) {
369        let limits = self.tier.limits();
370        self.order.refill(limits.order, now);
371        self.cancel.refill(limits.cancel, now);
372    }
373
374    fn set_tier(&mut self, tier: RateLimitTier) {
375        self.tier = tier;
376        let limits = tier.limits();
377        self.order.tokens = self.order.tokens.clamp(0.0, f64::from(limits.order.burst));
378        self.cancel.tokens = self.cancel.tokens.min(f64::from(limits.cancel.burst));
379        if !limits.negative_cancel_balance {
380            self.cancel.tokens = self.cancel.tokens.max(0.0);
381        }
382    }
383
384    fn bucket_mut(&mut self, bucket: TradingBucket) -> &mut BucketState {
385        match bucket {
386            TradingBucket::Order => &mut self.order,
387            TradingBucket::Cancel => &mut self.cancel,
388        }
389    }
390}
391
392#[derive(Debug)]
393struct BucketState {
394    tokens: f64,
395    updated_at: Instant,
396    blocked_until: Option<Instant>,
397}
398
399impl BucketState {
400    fn new(limits: BucketLimits, now: Instant) -> Self {
401        Self {
402            tokens: f64::from(limits.burst),
403            updated_at: now,
404            blocked_until: None,
405        }
406    }
407
408    fn refill(&mut self, limits: BucketLimits, now: Instant) {
409        let elapsed = now.saturating_duration_since(self.updated_at);
410        self.tokens =
411            (self.tokens + elapsed.as_secs_f64() * limits.rate).min(f64::from(limits.burst));
412        self.updated_at = now;
413
414        if self
415            .blocked_until
416            .is_some_and(|blocked_until| blocked_until <= now)
417        {
418            self.blocked_until = None;
419        }
420    }
421
422    fn block_for(&mut self, now: Instant, duration: Duration) {
423        let Some(blocked_until) = now.checked_add(duration) else {
424            return;
425        };
426        self.blocked_until = Some(
427            self.blocked_until
428                .map_or(blocked_until, |current| current.max(blocked_until)),
429        );
430    }
431}
432
433fn parse_number(
434    headers: &HashMap<String, String>,
435    name: &'static str,
436    allow_negative: bool,
437) -> Option<f64> {
438    let value = headers.get(name)?;
439    let parsed = value.parse::<f64>();
440    match parsed {
441        Ok(parsed) if parsed.is_finite() && (allow_negative || parsed >= 0.0) => Some(parsed),
442        _ => {
443            log::warn!("Invalid Polymarket rate-limit header {name}={value:?}");
444            None
445        }
446    }
447}
448
449fn parse_tier(headers: &HashMap<String, String>) -> Option<RateLimitTier> {
450    let value = headers.get(HEADER_RATE_LIMIT_TIER)?;
451    let tier = RateLimitTier::parse(value);
452    if tier.is_none() {
453        log::warn!("Invalid Polymarket rate-limit header {HEADER_RATE_LIMIT_TIER}={value:?}");
454    }
455    tier
456}
457
458fn parse_warning(headers: &HashMap<String, String>) -> bool {
459    let Some(value) = headers.get(HEADER_RATE_LIMIT_WARNING) else {
460        return false;
461    };
462
463    match value.trim().to_ascii_lowercase().as_str() {
464        "true" => true,
465        "false" => false,
466        _ => {
467            log::warn!(
468                "Invalid Polymarket rate-limit header {HEADER_RATE_LIMIT_WARNING}={value:?}"
469            );
470            false
471        }
472    }
473}
474
475fn parse_retry_after(headers: &HashMap<String, String>) -> Option<Duration> {
476    let seconds = parse_number(headers, HEADER_RETRY_AFTER, false)?;
477    let Ok(duration) = Duration::try_from_secs_f64(seconds) else {
478        log::warn!("Invalid Polymarket rate-limit header {HEADER_RETRY_AFTER}={seconds:?}");
479        return None;
480    };
481
482    if Instant::now().checked_add(duration).is_none() {
483        log::warn!("Invalid Polymarket rate-limit header {HEADER_RETRY_AFTER}={seconds:?}");
484        return None;
485    }
486    Some(duration)
487}
488
489fn reset_wait(reset: Option<f64>) -> Option<Duration> {
490    let reset = reset?;
491    let now = SystemTime::now()
492        .duration_since(UNIX_EPOCH)
493        .ok()?
494        .as_secs_f64();
495    let wait = reset - now;
496    if wait <= 0.0 {
497        return None;
498    }
499    let duration = Duration::try_from_secs_f64(wait).ok()?;
500    Instant::now().checked_add(duration).map(|_| duration)
501}
502
503fn warning_message(
504    endpoint: &str,
505    token_cost: u32,
506    tier: RateLimitTier,
507    headers: &RateLimitHeaders,
508) -> String {
509    let remaining = headers
510        .remaining
511        .map_or_else(|| "unknown".to_string(), |value| value.to_string());
512    let reset = headers
513        .reset
514        .map_or_else(|| "unknown".to_string(), |value| value.to_string());
515    format!(
516        "Polymarket rate limit warning: endpoint={endpoint}, token_cost={token_cost}, \
517         tier={tier}, remaining={remaining}, reset={reset}"
518    )
519}
520
521#[cfg(test)]
522mod tests {
523    use rstest::rstest;
524
525    use super::*;
526
527    fn header_map(entries: &[(&str, &str)]) -> HashMap<String, String> {
528        entries
529            .iter()
530            .map(|(name, value)| ((*name).to_string(), (*value).to_string()))
531            .collect()
532    }
533
534    #[rstest]
535    #[case::remaining(&[(HEADER_RATE_LIMIT_REMAINING, "0")], true)]
536    #[case::reset(&[(HEADER_RATE_LIMIT_RESET, "1")], true)]
537    #[case::tier(&[(HEADER_RATE_LIMIT_TIER, "Standard")], true)]
538    #[case::retry_after_only(&[(HEADER_RETRY_AFTER, "2")], false)]
539    #[case::warning_only(&[(HEADER_RATE_LIMIT_WARNING, "true")], false)]
540    #[case::empty(&[], false)]
541    fn test_has_signer_headers(#[case] entries: &[(&str, &str)], #[case] expected: bool) {
542        assert_eq!(
543            RateLimitHeaders::parse(&header_map(entries)).has_signer_headers(),
544            expected
545        );
546    }
547
548    #[rstest]
549    #[tokio::test(start_paused = true)]
550    async fn test_signer_limiters_share_state_and_keep_buckets_separate() {
551        let first = PolymarketRateLimiter::for_signer("0xSigner-Separation");
552        let second = PolymarketRateLimiter::for_signer("0xsigner-separation");
553        let other = PolymarketRateLimiter::for_signer("0xother-signer-separation");
554
555        first
556            .acquire("/orders", TradingBucket::Order, 10)
557            .await
558            .unwrap();
559
560        let state = second.state.lock().await;
561        assert!(Arc::ptr_eq(&first, &second));
562        assert!(!Arc::ptr_eq(&first, &other));
563        assert_eq!(state.order.tokens, 50.0);
564        assert_eq!(state.cancel.tokens, 120.0);
565        drop(state);
566        drop(first);
567        drop(second);
568
569        let replacement = PolymarketRateLimiter::for_signer("0xsigner-separation");
570        let replacement_state = replacement.state.lock().await;
571        assert_eq!(replacement_state.order.tokens, 50.0);
572        assert_eq!(replacement_state.cancel.tokens, 120.0);
573    }
574
575    #[rstest]
576    #[case::standard(RateLimitTier::Standard, 40, 60, 80, 120, true)]
577    #[case::copper(RateLimitTier::Copper, 60, 90, 120, 180, true)]
578    #[case::bronze(RateLimitTier::Bronze, 80, 120, 160, 240, true)]
579    #[case::silver(RateLimitTier::Silver, 200, 300, 400, 600, true)]
580    #[case::gold(RateLimitTier::Gold, 400, 600, 800, 1_200, true)]
581    #[case::platinum(RateLimitTier::Platinum, 450, 675, 900, 1_350, false)]
582    #[case::diamond(RateLimitTier::Diamond, 525, 787, 1_050, 1_575, false)]
583    #[case::elite(RateLimitTier::Elite, 600, 900, 1_200, 1_800, false)]
584    fn test_tier_limits_match_documented_contract(
585        #[case] tier: RateLimitTier,
586        #[case] order_rate: u32,
587        #[case] order_burst: u32,
588        #[case] cancel_rate: u32,
589        #[case] cancel_burst: u32,
590        #[case] negative_cancel_balance: bool,
591    ) {
592        let limits = tier.limits();
593
594        assert_eq!(limits.order.rate, f64::from(order_rate));
595        assert_eq!(limits.order.burst, order_burst);
596        assert_eq!(limits.cancel.rate, f64::from(cancel_rate));
597        assert_eq!(limits.cancel.burst, cancel_burst);
598        assert_eq!(limits.negative_cancel_balance, negative_cancel_balance);
599    }
600
601    #[rstest]
602    #[tokio::test(start_paused = true)]
603    async fn test_weighted_batches_debit_entry_counts() {
604        let limiter = PolymarketRateLimiter::for_signer("0xweighted-batches");
605
606        limiter
607            .acquire("/orders", TradingBucket::Order, 15)
608            .await
609            .unwrap();
610        limiter
611            .acquire("/orders", TradingBucket::Cancel, 7)
612            .await
613            .unwrap();
614
615        let state = limiter.state.lock().await;
616        assert_eq!(state.order.tokens, 45.0);
617        assert_eq!(state.cancel.tokens, 113.0);
618    }
619
620    #[rstest]
621    #[tokio::test(start_paused = true)]
622    async fn test_stale_remaining_header_cannot_credit_spent_tokens() {
623        let limiter = PolymarketRateLimiter::for_signer("0xstale-remaining");
624        limiter
625            .acquire("/order", TradingBucket::Order, 1)
626            .await
627            .unwrap();
628        limiter
629            .acquire("/order", TradingBucket::Order, 1)
630            .await
631            .unwrap();
632        let current = RateLimitHeaders::parse(&header_map(&[(HEADER_RATE_LIMIT_REMAINING, "58")]));
633        let stale = RateLimitHeaders::parse(&header_map(&[(HEADER_RATE_LIMIT_REMAINING, "59")]));
634
635        limiter
636            .observe_response("/order", TradingBucket::Order, 1, 0, &current, false)
637            .await;
638        limiter
639            .observe_response("/order", TradingBucket::Order, 1, 0, &stale, false)
640            .await;
641
642        assert_eq!(limiter.state.lock().await.order.tokens, 58.0);
643    }
644
645    #[rstest]
646    #[tokio::test(start_paused = true)]
647    async fn test_response_tier_updates_both_bucket_limits() {
648        let limiter = PolymarketRateLimiter::for_signer("0xtier-update");
649        limiter
650            .acquire("/order", TradingBucket::Order, 1)
651            .await
652            .unwrap();
653        let headers = RateLimitHeaders::parse(&header_map(&[
654            (HEADER_RATE_LIMIT_TIER, "Silver"),
655            (HEADER_RATE_LIMIT_REMAINING, "299"),
656        ]));
657
658        limiter
659            .observe_response("/order", TradingBucket::Order, 1, 0, &headers, false)
660            .await;
661
662        let state = limiter.state.lock().await;
663        assert_eq!(state.tier, RateLimitTier::Silver);
664        assert_eq!(state.tier.limits().order.rate, 200.0);
665        assert_eq!(state.tier.limits().order.burst, 300);
666        assert_eq!(state.tier.limits().cancel.rate, 400.0);
667        assert_eq!(state.tier.limits().cancel.burst, 600);
668        assert_eq!(state.order.tokens, 59.0);
669        drop(state);
670
671        assert_eq!(limiter.burst(TradingBucket::Cancel).await, 600);
672    }
673
674    #[rstest]
675    fn test_warning_headers_parse_and_log_all_required_fields() {
676        let headers = RateLimitHeaders::parse(&header_map(&[
677            (HEADER_RATE_LIMIT_REMAINING, "-2.5"),
678            (HEADER_RATE_LIMIT_RESET, "1786100000"),
679            (HEADER_RATE_LIMIT_TIER, "Gold"),
680            (HEADER_RATE_LIMIT_WARNING, "true"),
681            (HEADER_RETRY_AFTER, "1.2500001"),
682        ]));
683
684        assert_eq!(
685            RateLimitHeaders::names(),
686            vec![
687                HEADER_RATE_LIMIT_REMAINING.to_string(),
688                HEADER_RATE_LIMIT_RESET.to_string(),
689                HEADER_RATE_LIMIT_TIER.to_string(),
690                HEADER_RATE_LIMIT_WARNING.to_string(),
691                HEADER_RETRY_AFTER.to_string(),
692            ]
693        );
694        assert_eq!(headers.remaining, Some(-2.5));
695        assert_eq!(headers.reset, Some(1_786_100_000.0));
696        assert_eq!(headers.tier, Some(RateLimitTier::Gold));
697        assert!(headers.warning);
698        assert_eq!(
699            headers.retry_after,
700            Some(Duration::from_nanos(1_250_000_100))
701        );
702        assert_eq!(headers.retry_after_ms(), Some(1_251));
703        assert_eq!(
704            warning_message("/orders", 8, RateLimitTier::Gold, &headers),
705            "Polymarket rate limit warning: endpoint=/orders, token_cost=8, tier=Gold, \
706             remaining=-2.5, reset=1786100000"
707        );
708    }
709
710    #[rstest]
711    fn test_duration_header_overflow_is_ignored() {
712        let headers =
713            RateLimitHeaders::parse(&header_map(&[(HEADER_RETRY_AFTER, "18446744073709551616")]));
714
715        assert_eq!(headers.retry_after, None);
716        assert_eq!(reset_wait(Some(18_446_744_073_709_552_000.0)), None);
717    }
718
719    #[rstest]
720    #[tokio::test(start_paused = true)]
721    async fn test_negative_cancel_balance_blocks_until_refilled() {
722        let limiter = PolymarketRateLimiter::for_signer("0xnegative-cancel");
723        let headers = RateLimitHeaders::parse(&header_map(&[
724            (HEADER_RATE_LIMIT_REMAINING, "-5"),
725            (HEADER_RATE_LIMIT_TIER, "Standard"),
726        ]));
727        limiter
728            .observe_response("/cancel-all", TradingBucket::Cancel, 1, 0, &headers, false)
729            .await;
730
731        let limiter_clone = limiter.clone();
732        let acquire = tokio::spawn(async move {
733            limiter_clone
734                .acquire("/order", TradingBucket::Cancel, 1)
735                .await
736        });
737        tokio::task::yield_now().await;
738        assert!(!acquire.is_finished());
739
740        tokio::time::advance(Duration::from_millis(74)).await;
741        tokio::task::yield_now().await;
742        assert!(!acquire.is_finished());
743
744        tokio::time::advance(Duration::from_millis(1)).await;
745        acquire.await.unwrap().unwrap();
746    }
747
748    #[rstest]
749    #[tokio::test(start_paused = true)]
750    async fn test_post_cancel_debit_allows_standard_debt_and_floors_platinum() {
751        let standard = PolymarketRateLimiter::for_signer("0xpost-cancel-standard");
752        let empty = RateLimitHeaders::default();
753        standard
754            .acquire("/cancel-all", TradingBucket::Cancel, 1)
755            .await
756            .unwrap();
757        standard
758            .observe_response("/cancel-all", TradingBucket::Cancel, 1, 3, &empty, false)
759            .await;
760        assert_eq!(standard.state.lock().await.cancel.tokens, 116.0);
761
762        let zero_standard = RateLimitHeaders::parse(&header_map(&[
763            (HEADER_RATE_LIMIT_REMAINING, "0"),
764            (HEADER_RATE_LIMIT_TIER, "Standard"),
765        ]));
766        standard
767            .observe_response(
768                "/cancel-all",
769                TradingBucket::Cancel,
770                1,
771                0,
772                &zero_standard,
773                false,
774            )
775            .await;
776        standard
777            .observe_response("/cancel-all", TradingBucket::Cancel, 1, 3, &empty, false)
778            .await;
779
780        let platinum = PolymarketRateLimiter::for_signer("0xpost-cancel-platinum");
781        let zero_platinum = RateLimitHeaders::parse(&header_map(&[
782            (HEADER_RATE_LIMIT_REMAINING, "0"),
783            (HEADER_RATE_LIMIT_TIER, "Platinum"),
784        ]));
785        platinum
786            .observe_response(
787                "/cancel-market-orders",
788                TradingBucket::Cancel,
789                1,
790                2,
791                &zero_platinum,
792                false,
793            )
794            .await;
795
796        assert_eq!(standard.state.lock().await.cancel.tokens, -3.0);
797        assert_eq!(platinum.state.lock().await.cancel.tokens, 0.0);
798    }
799
800    #[rstest]
801    #[tokio::test(start_paused = true)]
802    async fn test_standard_order_burst_is_all_or_nothing() {
803        let limiter = PolymarketRateLimiter::for_signer("0xorder-burst");
804        limiter
805            .acquire("/orders", TradingBucket::Order, 60)
806            .await
807            .unwrap();
808
809        let limiter_clone = limiter.clone();
810        let acquire = tokio::spawn(async move {
811            limiter_clone
812                .acquire("/order", TradingBucket::Order, 1)
813                .await
814        });
815        tokio::task::yield_now().await;
816        assert!(!acquire.is_finished());
817
818        tokio::time::advance(Duration::from_millis(24)).await;
819        tokio::task::yield_now().await;
820        assert!(!acquire.is_finished());
821
822        tokio::time::advance(Duration::from_millis(1)).await;
823        acquire.await.unwrap().unwrap();
824
825        let error = limiter
826            .acquire("/orders", TradingBucket::Order, 61)
827            .await
828            .unwrap_err();
829        assert_eq!(
830            error.to_string(),
831            "bad request: /orders token cost 61 exceeds Standard tier order burst 60"
832        );
833    }
834
835    #[rstest]
836    #[tokio::test(start_paused = true)]
837    async fn test_retry_after_blocks_next_attempt_for_required_delay() {
838        let limiter = PolymarketRateLimiter::for_signer("0xretry-after");
839        let headers = RateLimitHeaders::parse(&header_map(&[
840            (HEADER_RATE_LIMIT_REMAINING, "60"),
841            (HEADER_RETRY_AFTER, "2"),
842        ]));
843        limiter
844            .observe_response("/order", TradingBucket::Order, 1, 0, &headers, true)
845            .await;
846
847        let limiter_clone = limiter.clone();
848        let acquire = tokio::spawn(async move {
849            limiter_clone
850                .acquire("/order", TradingBucket::Order, 1)
851                .await
852        });
853        tokio::task::yield_now().await;
854        assert!(!acquire.is_finished());
855
856        tokio::time::advance(Duration::from_millis(1_999)).await;
857        tokio::task::yield_now().await;
858        assert!(!acquire.is_finished());
859
860        tokio::time::advance(Duration::from_millis(1)).await;
861        acquire.await.unwrap().unwrap();
862    }
863}