"Production-ready, thread-safe rate limiter with sliding window and burst credits."
from future import annotations
import collections
import dataclasses
import hashlib
import math
import threading
import time
import unittest
from typing import Callable, Deque, Dict, List, NamedTuple, Optional
@dataclasses.dataclass(frozen=True, slots=True)
class Decision:
allowed: bool
remaining: int
retry_after: float
limit: int
class _WindowEntry(NamedTuple):
timestamp: float
cost: int
class _KeyBucket:
slots = (
"entries",
"window_usage",
"burst_credits",...
全文を表示 ▼
"Production-ready, thread-safe rate limiter with sliding window and burst credits."
from future import annotations
import collections
import dataclasses
import hashlib
import math
import threading
import time
import unittest
from typing import Callable, Deque, Dict, List, NamedTuple, Optional
@dataclasses.dataclass(frozen=True, slots=True)
class Decision:
allowed: bool
remaining: int
retry_after: float
limit: int
class _WindowEntry(NamedTuple):
timestamp: float
cost: int
class _KeyBucket:
slots = (
"entries",
"window_usage",
"burst_credits",
"last_update",
"last_active",
)
def __init__(self, initial_credits: float, now: float) -> None:
self.entries: Deque[_WindowEntry] = collections.deque()
self.window_usage: int = 0
self.burst_credits: float = float(initial_credits)
self.last_update: float = now
self.last_active: float = now
def clean_window(self, cutoff: float) -> None:
while self.entries and self.entries[0].timestamp <= cutoff:
expired = self.entries.popleft()
self.window_usage -= expired.cost
if self.window_usage < 0:
self.window_usage = 0
def refill_credits(self, now: float, rate: float, max_credits: float) -> None:
if now > self.last_update:
if rate > 0.0 and self.burst_credits < max_credits:
added = (now - self.last_update) * rate
self.burst_credits = min(max_credits, self.burst_credits + added)
self.last_update = now
def is_idle(self, now: float, window_seconds: float, max_credits: float) -> bool:
return (
len(self.entries) == 0
and math.isclose(self.burst_credits, max_credits, abs_tol=1e-9)
and (now - self.last_active) >= window_seconds
)
class _Shard:
slots = ("lock", "buckets", "access_order")
def __init__(self) -> None:
self.lock = threading.Lock()
self.buckets: Dict[str, _KeyBucket] = {}
self.access_order: collections.OrderedDict[str, None] = collections.OrderedDict()
class RateLimiter:
"""High-throughput, sliding-window rate limiter with burst credit handling."""
def __init__(
self,
limit: int,
window_seconds: float,
burst_credits: int = 0,
clock: Callable[[], float] = time.monotonic,
num_shards: int = 256,
) -> None:
if limit < 0:
raise ValueError("limit must be non-negative")
if window_seconds <= 0:
raise ValueError("window_seconds must be positive")
if burst_credits < 0:
raise ValueError("burst_credits must be non-negative")
if num_shards <= 0 or (num_shards & (num_shards - 1)) != 0:
raise ValueError("num_shards must be a positive power of two")
self.limit = int(limit)
self.window_seconds = float(window_seconds)
self.burst_credits = int(burst_credits)
self.clock = clock
self._rate = (self.limit / self.window_seconds) if self.limit > 0 else 0.0
self._num_shards = num_shards
self._mask = num_shards - 1
self._shards: List[_Shard] = [_Shard() for _ in range(num_shards)]
def _get_shard(self, key: str) -> _Shard:
digest = hashlib.blake2b(key.encode("utf-8"), digest_size=8).digest()
index = int.from_bytes(digest, byteorder="little") & self._mask
return self._shards[index]
def _evict_idle_sample(self, shard: _Shard, now: float, sample_size: int = 4) -> None:
# Check a small constant number of oldest accessed keys to amortize O(1) eviction
for _ in range(sample_size):
if not shard.access_order:
break
oldest_key = next(iter(shard.access_order))
bucket = shard.buckets.get(oldest_key)
if bucket is None:
shard.access_order.pop(oldest_key, None)
continue
bucket.clean_window(now - self.window_seconds)
bucket.refill_credits(now, self._rate, float(self.burst_credits))
if bucket.is_idle(now, self.window_seconds, float(self.burst_credits)):
del shard.buckets[oldest_key]
shard.access_order.pop(oldest_key, None)
else:
# Oldest key is not idle; stop rotation to avoid unnecessary overhead
break
def _compute_retry_after(self, bucket: _KeyBucket, cost: int, now: float) -> float:
if cost > self.limit + self.burst_credits:
return float("inf")
if cost <= 0:
return 0.0
cutoff = now - self.window_seconds
sim_entries = collections.deque(bucket.entries)
sim_usage = bucket.window_usage
sim_credits = bucket.burst_credits
last_t = now
rate = self._rate
max_credits = float(self.burst_credits)
while True:
window_avail = max(0, self.limit - sim_usage)
needed_burst = max(0, cost - window_avail)
if sim_credits >= needed_burst:
return max(0.0, last_t - now)
# If rate > 0, how long until burst credits refill sufficiently without expiry?
dt_refill = (needed_burst - sim_credits) / rate if rate > 0.0 else float("inf")
if not sim_entries:
if rate > 0.0:
return max(0.0, (last_t + dt_refill) - now)
return float("inf")
next_entry = sim_entries[0]
entry_expiry_time = next_entry.timestamp + self.window_seconds
dt_to_expiry = max(0.0, entry_expiry_time - last_t)
if rate > 0.0 and dt_refill <= dt_to_expiry:
return max(0.0, (last_t + dt_refill) - now)
# Fast-forward to the entry's expiration
last_t = max(last_t, entry_expiry_time)
if rate > 0.0:
sim_credits = min(max_credits, sim_credits + dt_to_expiry * rate)
expired = sim_entries.popleft()
sim_usage = max(0, sim_usage - expired.cost)
def allow(self, key: str, cost: int = 1, now: Optional[float] = None) -> Decision:
if cost < 0:
raise ValueError("cost must be non-negative")
current_time = float(self.clock() if now is None else now)
shard = self._get_shard(key)
with shard.lock:
self._evict_idle_sample(shard, current_time)
bucket = shard.buckets.get(key)
if bucket is None:
bucket = _KeyBucket(float(self.burst_credits), current_time)
shard.buckets[key] = bucket
shard.access_order[key] = None
shard.access_order.move_to_end(key)
# Monotonic clamp per key: past timestamps are clamped to last observed time
t = max(current_time, bucket.last_active)
bucket.last_active = t
bucket.clean_window(t - self.window_seconds)
bucket.refill_credits(t, self._rate, float(self.burst_credits))
window_avail = max(0, self.limit - bucket.window_usage)
credits_avail = bucket.burst_credits
total_avail = window_avail + int(credits_avail)
if cost == 0:
return Decision(
allowed=True,
remaining=total_avail,
retry_after=0.0,
limit=self.limit,
)
if cost > self.limit + self.burst_credits:
return Decision(
allowed=False,
remaining=total_avail,
retry_after=float("inf"),
limit=self.limit,
)
needed_burst = max(0, cost - window_avail)
if credits_avail >= needed_burst:
from_window = min(cost, window_avail)
from_burst = needed_burst
if from_window > 0:
bucket.entries.append(_WindowEntry(timestamp=t, cost=from_window))
bucket.window_usage += from_window
bucket.burst_credits -= from_burst
new_window_avail = max(0, self.limit - bucket.window_usage)
new_total_avail = new_window_avail + int(bucket.burst_credits)
return Decision(
allowed=True,
remaining=new_total_avail,
retry_after=0.0,
limit=self.limit,
)
else:
retry_after = self._compute_retry_after(bucket, cost, t)
return Decision(
allowed=False,
remaining=total_avail,
retry_after=retry_after,
limit=self.limit,
)
def snapshot(self, key: str, now: Optional[float] = None) -> Decision:
current_time = float(self.clock() if now is None else now)
shard = self._get_shard(key)
with shard.lock:
bucket = shard.buckets.get(key)
if bucket is None:
return Decision(
allowed=True,
remaining=self.limit + self.burst_credits,
retry_after=0.0,
limit=self.limit,
)
t = max(current_time, bucket.last_active)
bucket.clean_window(t - self.window_seconds)
bucket.refill_credits(t, self._rate, float(self.burst_credits))
window_avail = max(0, self.limit - bucket.window_usage)
remaining = window_avail + int(bucket.burst_credits)
return Decision(
allowed=remaining > 0,
remaining=remaining,
retry_after=0.0 if remaining > 0 else self._compute_retry_after(bucket, 1, t),
limit=self.limit,
)
"""Design Note (273 words):
This rate limiter implements an exact sliding window combined with a continuous token-bucket burst credit pool.
Data Structure & Accuracy Tradeoff:
Each client key maintains a _KeyBucket containing a deque of (timestamp, cost) tuples tracking consumed window quota, an integer accumulator window_usage, and a float burst_credits. An exact sliding window is chosen over approximate bucketing or leaky-bucket models to guarantee zero boundary error: request timestamps count strictly inside (t - window_seconds, t]. Multi-cost requests (cost > 1) are admitted atomically: available window quota is claimed first, and any deficit is covered by continuous burst credits refilled at limit / window_seconds. If the total available capacity is insufficient, the request is rejected all-or-nothing. Deque memory per active key scales strictly as O(N) where N <= limit, bounded and compact.
Locking Strategy:
To avoid a global contention bottleneck and prevent concurrent requests from serializing the whole gateway, the key space is partitioned across 256 independent shards via Blake2b hashing. Each shard maintains its own threading.Lock, dictionary, and access order. Operations on distinct keys map to independent locks with high probability, minimizing lock contention while fully synchronizing per-key concurrent accesses without deadlocks.
Memory Hygiene & Eviction:
To protect against unbounded memory growth from millions of one-shot keys, each shard maintains an OrderedDict of key access order. On every operation, an amortized O(1) cleanup samples the oldest keys in the shard. If a key's window deque is empty, its burst credits are fully regenerated, and it has been idle for at least window_seconds, it is permanently evicted. Inactive keys naturally expire without background threads.
"""
class RateLimiterTests(unittest.TestCase):
def setUp(self) -> None:
self.current_time = 1000.0
def fake_clock(self) -> float:
return self.current_time
def test_exact_boundary_expiry(self) -> None:
limiter = RateLimiter(limit=2, window_seconds=1.0, burst_credits=0, clock=self.fake_clock)
d1 = limiter.allow("client-1", cost=1)
self.assertTrue(d1.allowed)
self.assertEqual(d1.remaining, 1)
self.current_time += 0.5
d2 = limiter.allow("client-1", cost=1)
self.assertTrue(d2.allowed)
self.assertEqual(d2.remaining, 0)
# At exact boundary t = 1001.0, window (1000.0, 1001.0] still includes t=1000.0
self.current_time = 1001.0
d3 = limiter.allow("client-1", cost=1)
self.assertFalse(d3.allowed)
# Just past boundary t = 1001.000001, t=1000.0 is evicted
self.current_time = 1001.000001
d4 = limiter.allow("client-1", cost=1)
self.assertTrue(d4.allowed)
def test_burst_credit_refill(self) -> None:
# limit=10, window=10.0 -> rate=1.0 credit/sec. burst_credits=5
limiter = RateLimiter(limit=10, window_seconds=10.0, burst_credits=5, clock=self.fake_clock)
# Consume all 10 window + 5 burst credits
d = limiter.allow("burst-client", cost=15)
self.assertTrue(d.allowed)
self.assertEqual(d.remaining, 0)
# Immediate next request rejected
d_fail = limiter.allow("burst-client", cost=1)
self.assertFalse(d_fail.allowed)
# Advance 2.5s -> 2.5 credits refilled (integer remaining is 2)
self.current_time += 2.5
snap = limiter.snapshot("burst-client", now=self.current_time)
self.assertEqual(snap.remaining, 2)
# Advance another 2.5s -> total 5.0 credits refilled (capped at burst_credits=5)
self.current_time += 2.5
d_burst = limiter.allow("burst-client", cost=5)
self.assertTrue(d_burst.allowed)
def test_all_or_nothing_multi_cost(self) -> None:
limiter = RateLimiter(limit=5, window_seconds=10.0, burst_credits=2, clock=self.fake_clock)
# Allowed total = 7. Request cost=8 exceeds limit + burst_credits
d_too_large = limiter.allow("client-multi", cost=8)
self.assertFalse(d_too_large.allowed)
self.assertEqual(d_too_large.remaining, 7)
self.assertEqual(d_too_large.retry_after, float("inf"))
# Key has remaining 7. Request cost=6 consumes 5 window and 1 burst
d_consume = limiter.allow("client-multi", cost=6)
self.assertTrue(d_consume.allowed)
self.assertEqual(d_consume.remaining, 1)
# Request cost=2 cannot be satisfied (only 1 burst credit remaining)
d_all_or_nothing = limiter.allow("client-multi", cost=2)
self.assertFalse(d_all_or_nothing.allowed)
self.assertEqual(d_all_or_nothing.remaining, 1) # unchanged state
def test_retry_after_correctness(self) -> None:
limiter = RateLimiter(limit=1, window_seconds=5.0, burst_credits=0, clock=self.fake_clock)
d1 = limiter.allow("k", cost=1)
self.assertTrue(d1.allowed)
self.current_time += 2.0
d2 = limiter.allow("k", cost=1)
self.assertFalse(d2.allowed)
# Request was at 1000.0, expires at 1005.0. Current time is 1002.0 -> retry_after = 3.0
self.assertAlmostEqual(d2.retry_after, 3.0, places=5)
def test_idle_key_eviction(self) -> None:
limiter = RateLimiter(limit=1, window_seconds=1.0, burst_credits=0, clock=self.fake_clock)
shard = limiter._get_shard("ephemeral-1")
limiter.allow("ephemeral-1", cost=1)
self.assertIn("ephemeral-1", shard.buckets)
# Advance past window_seconds
self.current_time += 1.5
# Eviction is amortized during traffic on that shard
limiter.allow("ephemeral-2", cost=1)
self.assertNotIn("ephemeral-1", shard.buckets)
def test_non_monotonic_and_zero_limit_edge_cases(self) -> None:
# limit == 0 with burst credits
limiter = RateLimiter(limit=0, window_seconds=1.0, burst_credits=2, clock=self.fake_clock)
d = limiter.allow("zero-limit", cost=1)
self.assertTrue(d.allowed)
# Non-monotonic clock call
d_past = limiter.allow("zero-limit", cost=1, now=self.current_time - 10.0)
self.assertTrue(d_past.allowed)
d_exhausted = limiter.allow("zero-limit", cost=1)
self.assertFalse(d_exhausted.allowed)
def test_multithreaded_concurrency(self) -> None:
# Stress test ensuring total admitted requests never exceed theoretical upper bound
real_limiter = RateLimiter(limit=50, window_seconds=0.2, burst_credits=10)
threads: List[threading.Thread] = []
admitted_counts: List[int] = [0] * 10
def worker(tid: int) -> None:
admitted = 0
for _ in range(100):
decision = real_limiter.allow("shared-key", cost=1)
if decision.allowed:
admitted += 1
time.sleep(0.001)
admitted_counts[tid] = admitted
for i in range(10):
t = threading.Thread(target=worker, args=(i,))
threads.append(t)
t.start()
for t in threads:
t.join()
total_admitted = sum(admitted_counts)
# Max possible: limit (50) + burst (10) + refill during runtime (~0.1-0.3s -> ~75 requests)
# Must never exceed conservative bound of limit + burst + rate * 1.5s
self.assertLessEqual(total_admitted, 50 + 10 + int((50 / 0.2) * 1.5))
self.assertGreater(total_admitted, 0)
if name == "main":
unittest.main()
判定
勝利票
0 / 3
平均スコア
総合点
総評
回答Aは、優れた型ヒント、詳細な設計ノート、堅牢なテストスイートを備えた、クリーンで自己完結型の実装を提供しています。しかし、いくつかの微妙な並行処理と設計上の欠点があります。その削除戦略は、OrderedDictを使用してシャードをサンプリングする操作中に償却チェックを実行しますが、イテレーション全体でロックを正しく保持せず、最初のタッチ/削除時の競合状態をBほど堅牢に処理しません。さらに、回答Aのretry_afterループは、複雑なバーストクレジットとウィンドウの相互作用の下でエッジケースがあり、不正確または誤った値につながる可能性があります。
採点詳細を表示 ▼
正確さ
重み 35%スライディングウィンドウのセマンティクスとバーストクレジットの補充は一般的に正しいですが、retry_afterの計算は、バースト/ウィンドウの枯渇が組み合わされた場合に不正確になる可能性があります。
完全性
重み 20%型、設計ノート、unittestスイートを含む、すべての機能要件と成果物を満たしています。
コード品質
重み 20%クリーンでよく構造化されたコードで、簡潔なdocstringと懸念事項の明確な分離があります。
実用性
重み 15%シャード化されたロッキングは良好ですが、シャードレベルのロックは同じシャード内の操作をまだシリアル化します。負荷下の削除サンプリングには、競合の脆弱性の可能性があります。
指示遵守
重み 10%Python 3.11標準ライブラリのみ、asyncioなしなどのすべての制約に従っていますが、一部のエッジケースは要求されるよりも防御的に処理されています。
総合点
総評
回答Aは、シャードロックとアトミックマルチコストアドミッションを備えた自己完結型の実装を提供します。しかし、そのリトライ計算は、実際に容量が利用可能になる前に成功を約束する可能性があり、バーストで資金提供されたリクエストはスライディングウィンドウの使用から除外され、枯渇したゼロリミットエントリはエビクションを永久に妨げる可能性があります。その境界テストは、要求された間隔セマンティクスと矛盾し、エビクションテストは両方のキーがシャードを共有することを保証せず、ストレス テストは決定論的ではありません。本番環境対応という主張は裏付けられていません。
採点詳細を表示 ▼
正確さ
重み 35%ウィンドウの有効期限切れ自体は、正しい排他的下限を使用しており、アドミッションはアトミックです。しかし、通常の資金提供されたコストのみが記録されます。リトライ計算は、リフィルを考慮する際にバーストキャップを無視します。リミット2、ウィンドウ10、バーストなし、および時間ゼロで枯渇したクォータの場合、別のユニットは誤って10秒ではなく5秒のリトライを与えられます。スナップショットとクロスキーのクリーンアップは、タイムスタンプクランプを一貫して進めることなく履歴を破棄する可能性があります。
完全性
重み 20%要求されたパブリックAPI、実装、設計ノート、およびテストが含まれていますが、メモリの回収は永久に枯渇したゼロリミットエントリの後ろで失敗します。スナップショットは、実際のウィンドウ使用量ではなく、残りの容量を公開し、サブミリ秒専用テストまたは決定論的同時実行テストはありません。
コード品質
重み 20%読みやすいヘルパーの分解、型ヒント、およびコンパクトな不変の決定はプラスです。検証は、非整数構成値をサイレントに強制し、非有限時間またはウィンドウを拒否しません。未使用のリトライ変数、不正確なクリーンアップドキュメント、および間隔の定義を逆にする境界テストコメントがあります。
実用性
重み 15%独立したシャードとバウンドされたウィンドウ資金提供イベントストレージは有用ですが、信頼性の低いリトライガイダンスと潜在的にブロックされたエビクションは、ゲートウェイのデプロイメントを損ないます。リアルタイムストレス テストは、防御可能な決定論的バウンドではなく、任意のランタイム許可を使用します。
指示遵守
重み 10%標準ライブラリのみを使用し、要求されたコードと短い設計ノートを提供します。しかし、ストレス テストは決定論的な偽クロック要件に違反しており、境界テストは誤った結果をアサートし、エビクション カバレッジは確立されていないシャード衝突に依存しています。
総合点
総評
Answer Aは、キーごとにdequeを持つコヒーレントなシャーディングスライディングウィンドウリミッター、正しい(t - W, t]の有効期限、正しいオールオアナッシングバーストアカウンティング、キーごとの単調なクランプ、および償却LRU追い出しパスを提供します。しかし、実際のバグがあります:_compute_retry_afterは、必要なバーストが設定されたキャップ内で達成可能であることを確認せずにリフィルパスを使用するため、デフォルトのburst_credits=0(例:limit=10、window=10、t=0で10リクエスト、t=1でリクエスト)では9.0ではなく1.0を返します。テストスイートも記述どおりにはパスしません:test_exact_boundary_expiryは、仕様と実装の両方に反して、(1000,1001]が1000.0を含むと誤って主張するコメントに基づいて、t=1001.0での拒否をアサートします。test_idle_key_evictionはキーephemeral-1のシャードをチェックしますが、ephemeral-2を介して追い出しをトリガーします。これはほぼ確実に異なるシャードにハッシュされます。ストレステストはリアルクロックとスリープを使用し、フェイククロックは使用せず、その境界は非常に緩いです。spent burstを持つlimit=0未満のキーは決して追い出されません。設計ノートは妥当ですが、これらの問題には言及していません。
採点詳細を表示 ▼
正確さ
重み 35%ウィンドウセマンティクス、バーストキャップ、オールオアナッシングアドミッションは正しいですが、リフィルデルタが次の有効期限よりも短く、必要なバーストがキャップを超える場合(デフォルトのburst_credits=0は9.0ではなく1.0などを与えます)、retry_afterは間違っています。2つのテストは記述どおりに失敗します:境界テストは、仕様と実装に反して正確にt=1001.0で拒否をアサートし、追い出しテストは間違ったシャードをチェックします。
完全性
重み 20%すべての成果物(実装、設計ノート、テスト)が存在し、ほとんどの必要なテストシナリオが存在しますが、明示的なサブミリ秒ウィンドウテスト、同時初回タッチテストはなく、ストレステストはフェイククロックを使用しません。spent burstを持つidleキー(limit==0)は決して回収されません。
コード品質
重み 20%スロットと型ヒントにより読みやすいですが、冗長なbuckets辞書とaccess_order OrderedDictを維持しており、スナップショットは無理のあるallowedセマンティックでDecisionを再利用しています。retryシミュレーションには論理エラーが含まれており、テストコメントはコードと矛盾しています。
実用性
重み 15%正しく実行され、スロットリングされるでしょうが、デフォルト構成での間違ったretry_afterはクライアントやRetry-Afterヘッダーを誤解させるでしょう。また、失敗するテストは信頼性を低下させます。256シャードのBlake2bハッシュは問題ありませんが、必要以上に重いです。
指示遵守
重み 10%ほとんどの構造要件を満たしており、ノートは300語未満ですが、マルチスレッドテストは必要なフェイククロックの代わりにリアルクロックとスリープを使用しており、サブミリ秒のエッジケースは目に見えて実行されていません。