Implementeaza planul claude-master-plan-discord-bridge-20260830 (15 taskuri, 3 lane-uri paralele) — un bot subtire discord.py peste CLI-ul `claude`, cu proces persistent per fir alimentat pe stdin cu --input-format stream-json. Nucleu: runner (proces persistent + reaper 20min + respawn --resume), stream (parser tolerant), session_store (scriere atomica, lock per fir, detectare PID reuse, recovery), limits (max 4 procese, timeout tur, rate per user, plafon cost pe zi), render (un loop de editare per canal, interval adaptiv). Adaptor: allowlist guild/canal/user fail-closed cu respingerea webhook-urilor, comenzi !new/!cd/!model/!status/!stop/!cleanup, cost si model in subsolul fiecarui raspuns. Mesajul sosit in timpul unui tur devine steering, nu tur nou. Securitate: hook PreToolUse fail-closed care cere confirmare in Discord pentru operatiuni ireversibile, wrapper `infra` cu lista explicita de hosturi. Deny rules raman strat cosmetic, nu bariera (verificat: /usr/bin/ssh trece pe langa). Ops: alerte email pe conventia repo-ului, !cleanup pentru orfani, unit systemd user cu KillMode=control-group si limite de memorie, install.sh idempotent. Verificat: 275 teste fara retea/Discord/API (10.8s), identic cu si fara discord.py instalat; e2e pe CLI real confirma steering-ul mid-tur (mesaj la 6s intr-un tool call de 25s schimba raspunsul final). Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01B29CApsP1JkSdjYaGaHpE7
165 lines
5.9 KiB
Python
165 lines
5.9 KiB
Python
"""Limite: procese vii, coada per fir, timeout de tur, rate limit per user, plafon de cost.
|
|
|
|
T8. Motiv: 406 MB RSS per proces claude, pe un container cu istoric de OOM.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import collections
|
|
import contextlib
|
|
import logging
|
|
import time
|
|
from typing import Callable
|
|
|
|
import config
|
|
|
|
log = logging.getLogger("discord-bridge.limits")
|
|
|
|
try: # pragma: no cover - Lane C poate lipsi
|
|
import alerts # type: ignore
|
|
except ImportError: # pragma: no cover
|
|
class _NoAlerts:
|
|
@staticmethod
|
|
def alert(level: str, subject: str, body: str, dedup_key: str | None = None) -> None:
|
|
log.warning("alerta (%s) %s: %s", level, subject, body)
|
|
|
|
alerts = _NoAlerts() # type: ignore
|
|
|
|
|
|
class LimitError(RuntimeError):
|
|
"""Baza pentru refuzurile de limita (mesajul e afisabil in Discord)."""
|
|
|
|
|
|
class RateLimited(LimitError):
|
|
def __init__(self, retry_after: float):
|
|
self.retry_after = retry_after
|
|
super().__init__(f"prea multe mesaje; mai asteapta {retry_after:.0f}s")
|
|
|
|
|
|
class CostCapReached(LimitError):
|
|
def __init__(self, spent: float, cap: float):
|
|
self.spent, self.cap = spent, cap
|
|
super().__init__(f"plafonul de cost pe zi a fost atins ({spent:.2f} / {cap:.2f} USD)")
|
|
|
|
|
|
class Limits:
|
|
def __init__(
|
|
self,
|
|
store=None,
|
|
*,
|
|
max_procs: int | None = None,
|
|
turn_timeout: float | None = None,
|
|
rate_per_min: int | None = None,
|
|
cost_cap: float | None = None,
|
|
clock: Callable[[], float] = time.monotonic,
|
|
alerter: Callable | None = None,
|
|
):
|
|
self.store = store
|
|
self.max_procs = max_procs if max_procs is not None else config.get_int("MAX_PROCS", 4)
|
|
self.turn_timeout = (
|
|
turn_timeout if turn_timeout is not None else config.get_float("TURN_TIMEOUT_S", 900.0)
|
|
)
|
|
self.rate_per_min = (
|
|
rate_per_min if rate_per_min is not None else config.get_int("RATE_PER_USER_PER_MIN", 10)
|
|
)
|
|
self.cost_cap = cost_cap if cost_cap is not None else config.get_float("COST_CAP_USD_DAY", 10.0)
|
|
self.clock = clock
|
|
self._alert = alerter or alerts.alert
|
|
self._slots = asyncio.Semaphore(self.max_procs)
|
|
self._thread_locks: dict[str, asyncio.Lock] = {}
|
|
self._hits: dict[str, collections.deque[float]] = {}
|
|
self._local_cost = 0.0 # folosit doar daca nu avem store
|
|
self._cap_alerted = False
|
|
|
|
# ------------------------------------------------------------- procese
|
|
@property
|
|
def free_slots(self) -> int:
|
|
return self._slots._value # noqa: SLF001 - diagnostic pentru !status
|
|
|
|
@contextlib.asynccontextmanager
|
|
async def process_slot(self):
|
|
"""Maxim `max_procs` procese claude vii simultan; al 5-lea fir asteapta."""
|
|
await self._slots.acquire()
|
|
try:
|
|
yield
|
|
finally:
|
|
self._slots.release()
|
|
|
|
def thread_lock(self, thread_id: str) -> asyncio.Lock:
|
|
"""Coada per fir: un singur tur odata pe acelasi fir."""
|
|
key = str(thread_id)
|
|
lock = self._thread_locks.get(key)
|
|
if lock is None:
|
|
lock = self._thread_locks[key] = asyncio.Lock()
|
|
return lock
|
|
|
|
def queued(self, thread_id: str) -> bool:
|
|
return self.thread_lock(thread_id).locked()
|
|
|
|
# ---------------------------------------------------------- rate limit
|
|
def check_rate(self, user_id: str, now: float | None = None) -> None:
|
|
"""Fereastra glisanta de 60s per utilizator. Ridica RateLimited."""
|
|
now = self.clock() if now is None else now
|
|
q = self._hits.setdefault(str(user_id), collections.deque())
|
|
while q and now - q[0] >= 60.0:
|
|
q.popleft()
|
|
if len(q) >= self.rate_per_min:
|
|
raise RateLimited(60.0 - (now - q[0]))
|
|
q.append(now)
|
|
|
|
# ---------------------------------------------------------------- cost
|
|
def cost_today(self) -> float:
|
|
if self.store is not None:
|
|
with contextlib.suppress(Exception):
|
|
return float(self.store.cost_today())
|
|
return self._local_cost
|
|
|
|
def stopped(self) -> bool:
|
|
"""Plafonul e evaluat pe ziua curenta; reset-ul zilnic vine din `cost.day`."""
|
|
return self.cost_today() >= self.cost_cap
|
|
|
|
def cost_remaining(self) -> float:
|
|
return max(0.0, self.cost_cap - self.cost_today())
|
|
|
|
def check_cost(self) -> None:
|
|
if self.stopped():
|
|
raise CostCapReached(self.cost_today(), self.cost_cap)
|
|
|
|
def record_cost(self, usd: float, thread_id: str | None = None) -> float:
|
|
try:
|
|
usd = float(usd)
|
|
except (TypeError, ValueError):
|
|
usd = 0.0
|
|
if self.store is not None:
|
|
total = self.store.add_cost(thread_id, usd)
|
|
else:
|
|
self._local_cost = round(self._local_cost + usd, 6)
|
|
total = self._local_cost
|
|
if total >= self.cost_cap and not self._cap_alerted:
|
|
self._cap_alerted = True
|
|
with contextlib.suppress(Exception):
|
|
self._alert(
|
|
"CRITICAL",
|
|
"plafon de cost atins",
|
|
f"Cheltuit azi: {total:.2f} USD, plafon {self.cost_cap:.2f}. Puntea nu mai accepta tururi.",
|
|
"cost-cap",
|
|
)
|
|
elif total < self.cost_cap:
|
|
self._cap_alerted = False
|
|
return total
|
|
|
|
# --------------------------------------------------------------- admit
|
|
def admit(self, user_id: str) -> None:
|
|
"""Verificarile ieftine, inainte de a pune firul in coada. Ridica LimitError."""
|
|
self.check_cost()
|
|
self.check_rate(user_id)
|
|
|
|
@contextlib.asynccontextmanager
|
|
async def turn(self, thread_id: str, user_id: str):
|
|
"""Un tur complet: admis -> coada firului -> slot de proces."""
|
|
self.admit(user_id)
|
|
async with self.thread_lock(thread_id):
|
|
async with self.process_slot():
|
|
yield self.turn_timeout
|