Coverage for /home/admin/Documents/AI/applications/lexigram-dev/lexigram/experimental/ai/lexigram-ai-llm/src/lexigram/ai/llm/rate_limiting/core.py: 28%

47 statements  

« prev     ^ index     » next       coverage.py v7.15.4, created at 2026-08-25 07:19 +0800

1"""Rate limiter for LLM providers. 

2 

3Provides TPM (Tokens Per Minute) and RPM (Requests Per Minute) limiting. 

4""" 

5 

6from __future__ import annotations 

7 

8import asyncio 

9 

10from lexigram.logging import ( 

11 get_logger, 

12) 

13from lexigram.primitives import clock as ambient_clock 

14 

15logger = get_logger(__name__) 

16 

17 

18class TokenBucket: 

19 """Token bucket implementation for rate limiting.""" 

20 

21 def __init__(self, capacity: float, refill_rate: float): 

22 """Initialize bucket. 

23 

24 Args: 

25 capacity: Maximum bucket size 

26 refill_rate: Tokens added per second 

27 """ 

28 self.capacity = capacity 

29 self.refill_rate = refill_rate 

30 self.tokens = capacity 

31 self.last_update = ambient_clock.monotonic() 

32 self._lock = asyncio.Lock() 

33 

34 async def consume(self, amount: float = 1.0) -> bool: 

35 """Attempt to consume tokens from the bucket. 

36 

37 Args: 

38 amount: Number of tokens to consume 

39 

40 Returns: 

41 True if consumed, False if blocked 

42 """ 

43 async with self._lock: 

44 now = ambient_clock.monotonic() 

45 elapsed = now - self.last_update 

46 self.tokens = min(self.capacity, self.tokens + elapsed * self.refill_rate) 

47 self.last_update = now 

48 

49 if self.tokens >= amount: 

50 self.tokens -= amount 

51 return True 

52 return False 

53 

54 

55class RateLimiter: 

56 """Rate limiter for LLM requests (RPM and TPM). 

57 

58 Manages multiple buckets for different models and providers. 

59 """ 

60 

61 def __init__(self) -> None: 

62 """Initialize rate limiter.""" 

63 self._buckets: dict[str, TokenBucket] = {} 

64 self._lock = asyncio.Lock() 

65 

66 def _get_bucket_key(self, provider: str, model: str, limit_type: str) -> str: 

67 return f"{provider}:{model}:{limit_type}" 

68 

69 async def check( 

70 self, 

71 provider: str, 

72 model: str, 

73 tpm_limit: int | None = None, 

74 rpm_limit: int | None = None, 

75 estimated_tokens: int = 0, 

76 ) -> bool: 

77 """Check if request is allowed under current limits. 

78 

79 Args: 

80 provider: AI provider name 

81 model: Model name 

82 tpm_limit: Tokens Per Minute limit 

83 rpm_limit: Requests Per Minute limit 

84 estimated_tokens: Estimated tokens in request 

85 

86 Returns: 

87 True if allowed, False if blocked 

88 """ 

89 async with self._lock: 

90 # Initialize buckets if limits provided 

91 if rpm_limit: 

92 key = self._get_bucket_key(provider, model, "rpm") 

93 if key not in self._buckets: 

94 # RPM limit: 60 seconds interval 

95 self._buckets[key] = TokenBucket(rpm_limit, rpm_limit / 60.0) 

96 

97 if not await self._buckets[key].consume(1.0): 

98 logger.warning("RPM limit reached", provider=provider, model=model) 

99 return False 

100 

101 if tpm_limit: 

102 key = self._get_bucket_key(provider, model, "tpm") 

103 if key not in self._buckets: 

104 # TPM limit: 60 seconds interval 

105 self._buckets[key] = TokenBucket(tpm_limit, tpm_limit / 60.0) 

106 

107 if not await self._buckets[key].consume(float(estimated_tokens)): 

108 logger.warning( 

109 "TPM limit reached", 

110 provider=provider, 

111 model=model, 

112 tokens=estimated_tokens, 

113 ) 

114 return False 

115 

116 return True 

117 

118 def __repr__(self) -> str: 

119 return f"RateLimiter(buckets={len(self._buckets)})"