From 6aa583eb1ab906ff01046a60adc616739234ab60 Mon Sep 17 00:00:00 2001 From: Xuepoo Date: Fri, 18 Sep 2026 03:45:12 +0800 Subject: [PATCH] [CTX-0055] fix(transport): enforce sustained refill and bounded burst --- crates/devtools-client/src/transport.rs | 214 +++++++++++++++++++++++- src/transport.ts | 43 ++++- tests/transport.test.ts | 124 +++++++++++++- 3 files changed, 365 insertions(+), 16 deletions(-) diff --git a/crates/devtools-client/src/transport.rs b/crates/devtools-client/src/transport.rs index e94ca21..3aaffe0 100644 --- a/crates/devtools-client/src/transport.rs +++ b/crates/devtools-client/src/transport.rs @@ -213,6 +213,8 @@ pub struct RateLimiter { timestamps: std::collections::VecDeque, limit_per_sec: u32, burst: u32, + credit: u64, + last_refill_ms: Option, } impl RateLimiter { @@ -224,19 +226,21 @@ impl RateLimiter { #[must_use] pub fn new(limit_per_sec: u32, burst: u32) -> Self { Self { - timestamps: std::collections::VecDeque::with_capacity(burst as usize), + timestamps: std::collections::VecDeque::new(), limit_per_sec, burst, + credit: u64::from(burst) * RC9_WINDOW_MS, + last_refill_ms: None, } } pub fn count_in_window(&mut self, now_ms: u64) -> usize { - self.evict_old(now_ms); + self.observe_time(now_ms); self.timestamps.len() } pub fn check(&mut self, now_ms: u64) -> Result<(), TransportError> { - self.evict_old(now_ms); + let now_ms = self.observe_time(now_ms); if (self.timestamps.len() as u32) >= self.burst { return Err(TransportError::RateLimited(format!( "{} requests in {}ms exceeds burst {}", @@ -245,11 +249,29 @@ impl RateLimiter { self.burst ))); } - let _ = self.limit_per_sec; + if self.credit < RC9_WINDOW_MS { + return Err(TransportError::RateLimited(format!( + "sustained limit {} requests per {}ms exhausted", + self.limit_per_sec, RC9_WINDOW_MS + ))); + } + self.credit -= RC9_WINDOW_MS; self.timestamps.push_back(now_ms); Ok(()) } + fn observe_time(&mut self, now_ms: u64) -> u64 { + let now_ms = now_ms.max(self.last_refill_ms.unwrap_or(now_ms)); + let elapsed_ms = now_ms - self.last_refill_ms.unwrap_or(now_ms); + self.credit = self + .credit + .saturating_add(elapsed_ms.saturating_mul(u64::from(self.limit_per_sec))) + .min(u64::from(self.burst) * RC9_WINDOW_MS); + self.last_refill_ms = Some(now_ms); + self.evict_old(now_ms); + now_ms + } + fn evict_old(&mut self, now_ms: u64) { while let Some(&front) = self.timestamps.front() { if now_ms.saturating_sub(front) >= RC9_WINDOW_MS { @@ -707,13 +729,187 @@ mod tests { } #[test] - fn rate_limiter_burst() { - let mut lim = RateLimiter::new(10, 5); - for _ in 0..5 { + fn rate_limiter_shared_timeline() { + let mut lim = RateLimiter::new(2, 2); + lim.check(0).unwrap(); + lim.check(0).unwrap(); + assert_eq!(lim.count_in_window(1000), 0); + lim.check(500).unwrap(); + lim.check(500).unwrap(); + assert_eq!(lim.count_in_window(1500), 2); + assert!(matches!( + lim.check(1500), + Err(TransportError::RateLimited(_)) + )); + assert_eq!(lim.count_in_window(0), 2); + assert_eq!(lim.count_in_window(2000), 0); + } + + #[test] + fn rate_limiter_count_before_check() { + let mut lim = RateLimiter::new(2, 2); + assert_eq!(lim.count_in_window(1000), 0); + lim.check(0).unwrap(); + assert_eq!(lim.count_in_window(1000), 1); + assert_eq!(lim.count_in_window(1999), 1); + assert_eq!(lim.count_in_window(2000), 0); + } + + #[test] + fn rate_limiter_integer_boundaries() { + for end in [(1_u64 << 53) - 1, u64::MAX] { + let mut lim = RateLimiter::new(2, 4); + for _ in 0..4 { + lim.check(end - 1500).unwrap(); + } + for _ in 0..2 { + lim.check(end - 500).unwrap(); + } + assert!(matches!( + lim.check(end - 1), + Err(TransportError::RateLimited(_)) + )); + lim.check(end).unwrap(); + assert!(matches!( + lim.check(end), + Err(TransportError::RateLimited(_)) + )); + } + let mut lim = RateLimiter::new(u32::MAX, u32::MAX); + assert!(lim.is_empty()); + assert_eq!(lim.timestamps.capacity(), 0); + lim.check(0).unwrap(); + lim.check(u64::MAX).unwrap(); + assert_eq!(lim.count_in_window(u64::MAX), 1); + } + + #[test] + fn rate_limiter_rate_above_burst() { + let mut lim = RateLimiter::new(u32::MAX, 1); + lim.check(0).unwrap(); + assert!(matches!( + lim.check(999), + Err(TransportError::RateLimited(_)) + )); + lim.check(1000).unwrap(); + lim.check((1_u64 << 53) - 1).unwrap(); + assert!(matches!( + lim.check((1_u64 << 53) - 1), + Err(TransportError::RateLimited(_)) + )); + } + + #[test] + fn rate_limiter_rc9_defaults() { + let mut lim = RateLimiter::rc9_default(); + for _ in 0..RC9_BURST_PER_SEC { lim.check(0).unwrap(); } - assert!(lim.check(0).is_err()); - assert!(lim.check(1000).is_ok()); + assert!(matches!(lim.check(0), Err(TransportError::RateLimited(_)))); + for _ in 0..RC9_REQ_PER_SEC { + lim.check(RC9_WINDOW_MS).unwrap(); + } + assert!(matches!( + lim.check(RC9_WINDOW_MS), + Err(TransportError::RateLimited(_)) + )); + } + + #[test] + fn rate_limiter_sustained_refill() { + let mut lim = RateLimiter::new(2, 4); + for _ in 0..4 { + lim.check(0).unwrap(); + } + assert!(matches!(lim.check(0), Err(TransportError::RateLimited(_)))); + for second in 1..=10 { + let now_ms = second * RC9_WINDOW_MS; + for _ in 0..2 { + lim.check(now_ms).unwrap(); + } + assert!(matches!( + lim.check(now_ms), + Err(TransportError::RateLimited(_)) + )); + assert_eq!(lim.count_in_window(now_ms), 2); + } + } + + #[test] + fn rate_limiter_fractional_refill() { + let mut lim = RateLimiter::new(3, 6); + for _ in 0..6 { + lim.check(0).unwrap(); + } + for _ in 0..3 { + lim.check(1000).unwrap(); + } + for now_ms in [1100, 1200, 1300, 1333] { + assert!(matches!( + lim.check(now_ms), + Err(TransportError::RateLimited(_)) + )); + } + assert_eq!(lim.count_in_window(1333), 3); + assert!(lim.check(1334).is_ok()); + assert!(matches!( + lim.check(1334), + Err(TransportError::RateLimited(_)) + )); + } + + #[test] + fn rate_limiter_idle_credit_and_burst_ceiling() { + let mut lim = RateLimiter::new(2, 4); + lim.check(0).unwrap(); + for _ in 0..4 { + lim.check(60_000).unwrap(); + } + for now_ms in [60_000, 60_500] { + assert!(matches!( + lim.check(now_ms), + Err(TransportError::RateLimited(_)) + )); + } + for _ in 0..2 { + lim.check(61_000).unwrap(); + } + assert!(matches!( + lim.check(61_000), + Err(TransportError::RateLimited(_)) + )); + } + + #[test] + fn rate_limiter_clock_regression() { + let mut lim = RateLimiter::new(2, 4); + for _ in 0..4 { + lim.check(1000).unwrap(); + } + for _ in 0..2 { + lim.check(2000).unwrap(); + } + for now_ms in [1500, 2000, 2499] { + assert!(matches!( + lim.check(now_ms), + Err(TransportError::RateLimited(_)) + )); + } + assert!(lim.check(2500).is_ok()); + } + + #[test] + fn rate_limiter_zero_limits() { + let mut lim = RateLimiter::new(0, 1); + lim.check(0).unwrap(); + assert!(matches!( + lim.check(60_000), + Err(TransportError::RateLimited(_)) + )); + assert!(matches!( + RateLimiter::new(2, 0).check(60_000), + Err(TransportError::RateLimited(_)) + )); } #[test] diff --git a/src/transport.ts b/src/transport.ts index 8ede9e8..a80c5a6 100644 --- a/src/transport.ts +++ b/src/transport.ts @@ -199,32 +199,69 @@ export class RateLimiter { // entry) instead of Array.shift() (O(K) memmove per entry). private timestamps: number[] = []; private head = 0; + private credit: number; + private lastRefillMs: number | undefined; constructor( private readonly limitPerSec: number = RC9_REQ_PER_SEC, private readonly burst: number = RC9_BURST_PER_SEC, - ) {} + ) { + for (const value of [limitPerSec, burst]) { + if (!Number.isInteger(value) || value < 0 || value > 2 ** 32 - 1) { + throw new RangeError("rate and burst must be unsigned 32-bit integers"); + } + } + this.credit = burst * RC9_WINDOW_MS; + } static rc9Default(): RateLimiter { return new RateLimiter(RC9_REQ_PER_SEC, RC9_BURST_PER_SEC); } countInWindow(nowMs: number): number { - this.evictOld(nowMs); + this.observeTime(nowMs); return this.timestamps.length - this.head; } check(nowMs: number): void { - this.evictOld(nowMs); + nowMs = this.observeTime(nowMs); if (this.timestamps.length - this.head >= this.burst) { throw new TransportError( "RateLimited", `rate limited: ${this.timestamps.length - this.head} requests in ${RC9_WINDOW_MS}ms exceeds burst ${this.burst}`, ); } + if (this.credit < RC9_WINDOW_MS) { + throw new TransportError( + "RateLimited", + `rate limited: sustained limit ${this.limitPerSec} requests per ${RC9_WINDOW_MS}ms exhausted`, + ); + } + this.credit -= RC9_WINDOW_MS; this.timestamps.push(nowMs); } + private observeTime(nowMs: number): number { + if (!Number.isSafeInteger(nowMs) || nowMs < 0) { + throw new RangeError( + "time must be nonnegative safe-integer milliseconds", + ); + } + nowMs = Math.max(nowMs, this.lastRefillMs ?? nowMs); + const elapsedMs = nowMs - (this.lastRefillMs ?? nowMs); + const capacity = this.burst * RC9_WINDOW_MS; + if (this.limitPerSec > 0) { + const refillMs = Math.ceil((capacity - this.credit) / this.limitPerSec); + this.credit = + elapsedMs >= refillMs + ? capacity + : this.credit + elapsedMs * this.limitPerSec; + } + this.lastRefillMs = nowMs; + this.evictOld(nowMs); + return nowMs; + } + private evictOld(nowMs: number): void { while (this.head < this.timestamps.length) { const front = this.timestamps[this.head]; diff --git a/tests/transport.test.ts b/tests/transport.test.ts index 0f5f4ef..be2c62e 100644 --- a/tests/transport.test.ts +++ b/tests/transport.test.ts @@ -12,6 +12,9 @@ import { MAX_FRAME_BYTES, RC9_PAYLOAD_CAP_BYTES, RC9_MAX_CONNECTIONS, + RC9_REQ_PER_SEC, + RC9_BURST_PER_SEC, + RC9_WINDOW_MS, } from "../src/transport.js"; import { peerCredentials } from "../src/auth.js"; @@ -46,11 +49,124 @@ describe("transport framing (phase 2, live runtime, bounded 256 KiB IPC / 1 MiB expect(new TextDecoder().decode(out[1]!.payload)).toBe("bb"); }); - test("rate limiter RC-9 100/s burst 200, window 1s", () => { - const lim = new RateLimiter(10, 5); - for (let i = 0; i < 5; i++) lim.check(0); + test("rate limiter defaults retain RC-9 sustained and burst values", () => { + const lim = RateLimiter.rc9Default(); + for (let i = 0; i < RC9_BURST_PER_SEC; i++) lim.check(0); expect(() => lim.check(0)).toThrow("rate limited"); - expect(() => lim.check(1000)).not.toThrow(); + for (let i = 0; i < RC9_REQ_PER_SEC; i++) lim.check(RC9_WINDOW_MS); + expect(() => lim.check(RC9_WINDOW_MS)).toThrow("rate limited"); + }); + + test("rate limiter refills at the sustained rate after the initial burst", () => { + const lim = new RateLimiter(2, 4); + for (let i = 0; i < 4; i++) lim.check(0); + expect(() => lim.check(0)).toThrow("rate limited"); + for (let second = 1; second <= 10; second++) { + const nowMs = second * 1000; + for (let i = 0; i < 2; i++) lim.check(nowMs); + expect(() => lim.check(nowMs)).toThrow("rate limited"); + expect(lim.countInWindow(nowMs)).toBe(2); + } + }); + + test("rate limiter preserves fractional refill across rejections", () => { + const fractional = new RateLimiter(3, 6); + for (let i = 0; i < 6; i++) fractional.check(0); + for (let i = 0; i < 3; i++) fractional.check(1000); + for (const nowMs of [1100, 1200, 1300, 1333]) { + expect(() => fractional.check(nowMs)).toThrow("rate limited"); + } + expect(fractional.countInWindow(1333)).toBe(3); + expect(() => fractional.check(1334)).not.toThrow(); + expect(() => fractional.check(1334)).toThrow("rate limited"); + }); + + test("rate limiter caps idle credit and retains the rolling burst ceiling", () => { + const lim = new RateLimiter(2, 4); + lim.check(0); + for (let i = 0; i < 4; i++) lim.check(60_000); + expect(() => lim.check(60_000)).toThrow("rate limited"); + expect(() => lim.check(60_500)).toThrow("rate limited"); + for (let i = 0; i < 2; i++) lim.check(61_000); + expect(() => lim.check(61_000)).toThrow("rate limited"); + }); + + test("rate limiter does not refill twice after a clock regression", () => { + const lim = new RateLimiter(2, 4); + for (let i = 0; i < 4; i++) lim.check(1000); + for (let i = 0; i < 2; i++) lim.check(2000); + expect(() => lim.check(1500)).toThrow("rate limited"); + expect(() => lim.check(2000)).toThrow("rate limited"); + expect(() => lim.check(2499)).toThrow("rate limited"); + expect(() => lim.check(2500)).not.toThrow(); + }); + + test("rate limiter zero rate never refills and zero burst never admits", () => { + const lim = new RateLimiter(0, 1); + lim.check(0); + expect(() => lim.check(60_000)).toThrow("rate limited"); + expect(() => new RateLimiter(2, 0).check(60_000)).toThrow("rate limited"); + }); + + test("rate limiter shares count and check time observations", () => { + const lim = new RateLimiter(2, 2); + lim.check(0); + lim.check(0); + expect(lim.countInWindow(1000)).toBe(0); + lim.check(500); + lim.check(500); + expect(lim.countInWindow(1500)).toBe(2); + expect(() => lim.check(1500)).toThrow("rate limited"); + expect(lim.countInWindow(0)).toBe(2); + expect(lim.countInWindow(2000)).toBe(0); + }); + + test("rate limiter counting before admission establishes the clock", () => { + const lim = new RateLimiter(2, 2); + expect(lim.countInWindow(1000)).toBe(0); + lim.check(0); + expect(lim.countInWindow(1000)).toBe(1); + expect(lim.countInWindow(1999)).toBe(1); + expect(lim.countInWindow(2000)).toBe(0); + }); + + test("rate limiter rejects unsupported numeric configuration", () => { + for (const value of [-1, 0.5, NaN, Infinity, -Infinity, 2 ** 32]) { + expect(() => new RateLimiter(value, 2)).toThrow(RangeError); + expect(() => new RateLimiter(2, value)).toThrow(RangeError); + } + expect(() => new RateLimiter(2 ** 32 - 1, 2 ** 32 - 1)).not.toThrow(); + }); + + test("rate limiter rejects invalid time without changing state", () => { + const lim = new RateLimiter(2, 2); + lim.check(0); + for (const nowMs of [-1, 0.5, NaN, Infinity, -Infinity, 2 ** 53]) { + expect(() => lim.check(nowMs)).toThrow(RangeError); + expect(() => lim.countInWindow(nowMs)).toThrow(RangeError); + expect(lim.countInWindow(0)).toBe(1); + } + lim.check(0); + expect(() => lim.check(0)).toThrow("rate limited"); + }); + + test("rate limiter preserves millisecond precision at the safe integer boundary", () => { + const lim = new RateLimiter(2, 4); + const end = Number.MAX_SAFE_INTEGER; + for (let i = 0; i < 4; i++) lim.check(end - 1500); + for (let i = 0; i < 2; i++) lim.check(end - 500); + expect(() => lim.check(end - 1)).toThrow("rate limited"); + lim.check(end); + expect(() => lim.check(end)).toThrow("rate limited"); + }); + + test("rate limiter caps wide elapsed refill and rate above burst", () => { + const lim = new RateLimiter(2 ** 32 - 1, 1); + lim.check(0); + expect(() => lim.check(999)).toThrow("rate limited"); + lim.check(1000); + lim.check(Number.MAX_SAFE_INTEGER); + expect(() => lim.check(Number.MAX_SAFE_INTEGER)).toThrow("rate limited"); }); test("rate limiter bulk eviction stays correct and bounded", () => {