guilder_client_hyperliquid/
rate_limiter.rs1use std::collections::VecDeque;
18use std::time::{Duration, Instant};
19use tokio::sync::Mutex;
20
21const WINDOW: Duration = Duration::from_secs(60);
24pub const MAX_WEIGHT: u32 = 1200;
25
26pub const ADDR_INITIAL_BUFFER: u64 = 10_000;
29const ADDR_THROTTLE_INTERVAL: Duration = Duration::from_secs(10);
30
31#[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
46pub 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 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 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 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
138pub 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 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 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 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 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}