Skip to main content

guilder_client_hyperliquid/
rate_limiter.rs

1//! Rate limiters for the Hyperliquid REST API.
2//!
3//! ## IP-based (`RestRateLimiter`)
4//! Tracks an aggregated weight budget of **1 200 per minute** in a 60-second sliding window.
5//!
6//! ## Address-based (`AddressRateLimiter`)
7//! Tracks the per-address action budget. Each address starts with 10 000 requests and accrues
8//! 1 request per 1 USDC traded. When exhausted, 1 request is allowed every 10 seconds.
9//! Cancels use a higher limit: `min(limit + 100_000, limit * 2)`.
10//! A batch of n orders/cancels counts as n requests against this limit.
11//!
12//! ## Interface
13//! Both limiters share the same interface:
14//! - `acquire(...)` → `Result<(), RateLimitError>` — returns `Err` immediately if throttled.
15//! - `acquire_blocking(...)` — retries on `Err`, logging each wait via `tracing::warn`.
16
17use std::collections::VecDeque;
18use std::time::{Duration, Instant};
19use tokio::sync::Mutex;
20
21// ── IP-based constants ────────────────────────────────────────────────────────
22
23const WINDOW: Duration = Duration::from_secs(60);
24pub const MAX_WEIGHT: u32 = 1200;
25
26// ── Address-based constants ───────────────────────────────────────────────────
27
28pub const ADDR_INITIAL_BUFFER: u64 = 10_000;
29const ADDR_THROTTLE_INTERVAL: Duration = Duration::from_secs(10);
30
31// ── Error ─────────────────────────────────────────────────────────────────────
32
33#[derive(Debug, Clone)]
34pub struct RateLimitError {
35    pub retry_after: Duration,
36}
37
38impl std::fmt::Display for RateLimitError {
39    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
40        write!(f, "rate limited, retry after {:?}", self.retry_after)
41    }
42}
43
44impl std::error::Error for RateLimitError {}
45
46// ── IP-based limiter ──────────────────────────────────────────────────────────
47
48pub struct RestRateLimiter {
49    entries: Mutex<VecDeque<(Instant, u32)>>,
50    max_weight: u32,
51}
52
53impl Default for RestRateLimiter {
54    fn default() -> Self {
55        Self::new()
56    }
57}
58
59impl RestRateLimiter {
60    pub fn new() -> Self {
61        Self::new_with_budget(MAX_WEIGHT)
62    }
63
64    pub fn new_with_budget(max_weight: u32) -> Self {
65        RestRateLimiter {
66            entries: Mutex::new(VecDeque::new()),
67            max_weight,
68        }
69    }
70
71    /// Returns `Err(RateLimitError)` immediately if the weight budget is exhausted.
72    pub async fn acquire(&self, weight: u32) -> Result<(), RateLimitError> {
73        let mut entries = self.entries.lock().await;
74        let now = Instant::now();
75
76        while let Some(&(t, _)) = entries.front() {
77            if now.duration_since(t) >= WINDOW {
78                entries.pop_front();
79            } else {
80                break;
81            }
82        }
83
84        let used: u32 = entries.iter().map(|(_, w)| w).sum();
85        if used + weight <= self.max_weight {
86            entries.push_back((now, weight));
87            return Ok(());
88        }
89
90        let oldest = entries.front().unwrap().0;
91        let retry_after = WINDOW.saturating_sub(now.duration_since(oldest));
92        Err(RateLimitError { retry_after })
93    }
94
95    /// Retries on `Err`, logging each wait. Resolves once the request is accepted.
96    pub async fn acquire_blocking(&self, weight: u32, call: &str) {
97        loop {
98            match self.acquire(weight).await {
99                Ok(()) => return,
100                Err(e) => {
101                    let used = {
102                        let entries = self.entries.lock().await;
103                        entries.iter().map(|(_, w)| w).sum::<u32>()
104                    };
105                    tracing::warn!(
106                        call,
107                        weight,
108                        used,
109                        budget = self.max_weight,
110                        retry_after_ms = e.retry_after.as_millis(),
111                        "REST rate limited"
112                    );
113                    tokio::time::sleep(e.retry_after).await;
114                }
115            }
116        }
117    }
118
119    /// T4 (2026-09-30): current snapshot of the sliding window —
120    /// `(remaining_weight, max_weight, used_entries)`. Lets the host
121    /// application export the info-budget water level as a telemetry
122    /// gauge (exhaustion becomes a queryable trend, not a surprise).
123    pub async fn budget_snapshot(&self) -> (u32, u32, usize) {
124        let mut entries = self.entries.lock().await;
125        let now = Instant::now();
126        while let Some(&(t, _)) = entries.front() {
127            if now.duration_since(t) >= WINDOW {
128                entries.pop_front();
129            } else {
130                break;
131            }
132        }
133        let used: u32 = entries.iter().map(|(_, w)| w).sum();
134        (self.max_weight.saturating_sub(used), self.max_weight, entries.len())
135    }
136}
137
138// ── Address-based limiter ─────────────────────────────────────────────────────
139
140/// Tracks the per-address action budget.
141///
142/// Call [`record_fill`] when a fill is confirmed to grow the budget.
143/// Call [`acquire`] before every action (not info) request.
144/// Pass `is_cancel = true` for cancel actions to apply the higher cancel limit.
145/// Pass the actual batch size as `count` for batched requests.
146pub struct AddressRateLimiter {
147    inner: Mutex<AddrState>,
148}
149
150struct AddrState {
151    budget: u64,
152    consumed: u64,
153    last_throttled: Option<Instant>,
154}
155
156impl Default for AddressRateLimiter {
157    fn default() -> Self {
158        Self::new()
159    }
160}
161
162impl AddressRateLimiter {
163    pub fn new() -> Self {
164        Self::new_with_budget(ADDR_INITIAL_BUFFER)
165    }
166
167    pub fn new_with_budget(initial_budget: u64) -> Self {
168        AddressRateLimiter {
169            inner: Mutex::new(AddrState {
170                budget: initial_budget,
171                consumed: 0,
172                last_throttled: None,
173            }),
174        }
175    }
176
177    /// Increase the budget by `usdc_volume` requests (1 request per 1 USDC traded).
178    pub async fn record_fill(&self, usdc_volume: u64) {
179        let mut s = self.inner.lock().await;
180        s.budget = s.budget.saturating_add(usdc_volume);
181    }
182
183    fn cancel_limit(default: u64) -> u64 {
184        (default + 100_000).min(default * 2)
185    }
186
187    /// Returns `Err(RateLimitError)` immediately if the address budget is exhausted.
188    pub async fn acquire(&self, count: u64, is_cancel: bool) -> Result<(), RateLimitError> {
189        let mut s = self.inner.lock().await;
190        let effective_limit = if is_cancel {
191            Self::cancel_limit(s.budget)
192        } else {
193            s.budget
194        };
195
196        if s.consumed + count <= effective_limit {
197            s.consumed += count;
198            return Ok(());
199        }
200
201        // Budget exhausted — compute throttle wait.
202        let now = Instant::now();
203        let retry_after = match s.last_throttled {
204            Some(t) => ADDR_THROTTLE_INTERVAL.saturating_sub(now.duration_since(t)),
205            None => Duration::ZERO,
206        };
207
208        if retry_after.is_zero() {
209            if count > 1 {
210                return Err(RateLimitError {
211                    retry_after: ADDR_THROTTLE_INTERVAL,
212                });
213            }
214            s.last_throttled = Some(now);
215            s.consumed += 1;
216            return Ok(());
217        }
218
219        Err(RateLimitError { retry_after })
220    }
221
222    /// Retries on `Err`, logging each wait. Resolves once the request is accepted.
223    pub async fn acquire_blocking(&self, count: u64, is_cancel: bool) {
224        loop {
225            match self.acquire(count, is_cancel).await {
226                Ok(()) => return,
227                Err(e) => {
228                    tracing::warn!(
229                        retry_after_ms = e.retry_after.as_millis(),
230                        "address rate limited"
231                    );
232                    tokio::time::sleep(e.retry_after).await;
233                }
234            }
235        }
236    }
237}