1use 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
53pub 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, ¤t, 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}