Add shared rate_limiter.py module\n\nModular rate limiter any .py script can use:\n from rate_limiter import rate_limit_wait\n rate_limit_wait(agent)\n\nToken bucket per agent: 1 op/3s sustained, burst 5, max 20/min.\nState in /tmp/netvm-rate-limit.json.
This commit is contained in:
@@ -0,0 +1,93 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Shared rate limiter for NetVM scripts.
|
||||
|
||||
Any .py script can use this to avoid hitting Muse rate limits:
|
||||
from rate_limiter import rate_limit_wait
|
||||
rate_limit_wait("opm") # blocks until allowed
|
||||
|
||||
Uses a token bucket per agent stored in /tmp (persists across invocations,
|
||||
not across reboots). Default: 1 op per 3s sustained, burst of 5, max ~20/min.
|
||||
"""
|
||||
import json
|
||||
import time
|
||||
import os
|
||||
|
||||
RATE_LIMIT_FILE = "/tmp/netvm-rate-limit.json"
|
||||
RATE_LIMIT_INTERVAL = 3.0 # seconds between ops
|
||||
RATE_LIMIT_BURST = 5
|
||||
RATE_LIMIT_MAX_PER_MIN = 20
|
||||
|
||||
def _load_state():
|
||||
try:
|
||||
with open(RATE_LIMIT_FILE) as f:
|
||||
return json.load(f)
|
||||
except:
|
||||
return {}
|
||||
|
||||
def _save_state(state):
|
||||
try:
|
||||
# Atomic write via temp file
|
||||
tmp = RATE_LIMIT_FILE + ".tmp"
|
||||
with open(tmp, "w") as f:
|
||||
json.dump(state, f)
|
||||
os.rename(tmp, RATE_LIMIT_FILE)
|
||||
except:
|
||||
pass
|
||||
|
||||
def rate_limit_wait(agent, interval=RATE_LIMIT_INTERVAL, burst=RATE_LIMIT_BURST):
|
||||
"""
|
||||
Block until the agent is allowed to perform an operation.
|
||||
Returns the time waited in seconds (0 if no wait needed).
|
||||
"""
|
||||
state = _load_state()
|
||||
now = time.time()
|
||||
|
||||
# Prune old entries
|
||||
recent_key = f"{agent}_recent"
|
||||
recent = state.get(recent_key, [])
|
||||
recent = [t for t in recent if now - t < 60]
|
||||
|
||||
# Check burst limit (max per minute)
|
||||
if len(recent) >= RATE_LIMIT_MAX_PER_MIN:
|
||||
wait = 60 - (now - recent[0]) + 1
|
||||
if wait > 0:
|
||||
time.sleep(wait)
|
||||
# Refresh after wait
|
||||
state = _load_state()
|
||||
recent = state.get(recent_key, [])
|
||||
recent = [t for t in recent if time.time() - t < 60]
|
||||
|
||||
# Check interval limit
|
||||
last = state.get(agent, 0)
|
||||
now = time.time()
|
||||
if now - last < interval:
|
||||
wait = interval - (now - last)
|
||||
time.sleep(wait)
|
||||
|
||||
# Record this operation
|
||||
now = time.time()
|
||||
state[agent] = now
|
||||
recent.append(now)
|
||||
state[recent_key] = recent
|
||||
_save_state(state)
|
||||
return 0
|
||||
|
||||
def rate_limit_check(agent):
|
||||
"""
|
||||
Non-blocking check. Returns (allowed: bool, wait_seconds: float).
|
||||
"""
|
||||
state = _load_state()
|
||||
now = time.time()
|
||||
recent = state.get(f"{agent}_recent", [])
|
||||
recent = [t for t in recent if now - t < 60]
|
||||
|
||||
if len(recent) >= RATE_LIMIT_MAX_PER_MIN:
|
||||
wait = 60 - (now - recent[0]) + 1
|
||||
return False, max(0, wait)
|
||||
|
||||
last = state.get(agent, 0)
|
||||
if now - last < RATE_LIMIT_INTERVAL:
|
||||
return False, RATE_LIMIT_INTERVAL - (now - last)
|
||||
|
||||
return True, 0
|
||||
Reference in New Issue
Block a user