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
210 lines
7.0 KiB
Python
210 lines
7.0 KiB
Python
"""Randare pentru Discord: chunker + un SINGUR loop de editare per CANAL.
|
|
|
|
T10. Bucket-ul de rate limit al Discord e per canal, deci loop-ul e per canal, coalescent,
|
|
cu interval adaptiv 1s -> 5s. Modulul NU importa discord.py: I/O-ul il face adaptorul.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import contextlib
|
|
import logging
|
|
import re
|
|
import time
|
|
from dataclasses import dataclass, field
|
|
from typing import Any, Awaitable, Callable
|
|
|
|
log = logging.getLogger("discord-bridge.render")
|
|
|
|
MAX_MSG = 2000
|
|
ATTACH_OVER = 6000
|
|
_FENCE = re.compile(r"^```([A-Za-z0-9_+.\-]*)\s*$")
|
|
|
|
|
|
@dataclass
|
|
class Attachment:
|
|
filename: str
|
|
content: str
|
|
preview: str
|
|
|
|
kind: str = "attachment"
|
|
|
|
|
|
@dataclass
|
|
class Chunks:
|
|
parts: list[str] = field(default_factory=list)
|
|
kind: str = "chunks"
|
|
|
|
|
|
def split_message(text: str, limit: int = MAX_MSG, attach_over: int = ATTACH_OVER,
|
|
filename: str = "raspuns.md") -> Chunks | Attachment:
|
|
"""Chunker-ul, intr-o singura functie.
|
|
|
|
- peste `attach_over` caractere -> obiect Attachment (adaptorul il urca in Discord)
|
|
- altfel bucati de cel mult `limit` caractere, fara sa rupa un fence ``` peste granita:
|
|
fence-ul se inchide la finalul bucatii si se redeschide cu acelasi limbaj in urmatoarea.
|
|
"""
|
|
text = text or ""
|
|
if len(text) > attach_over:
|
|
head = text[: limit - 200]
|
|
return Attachment(
|
|
filename=filename,
|
|
content=text,
|
|
preview=head + "\n...\n(raspuns lung, atasat integral)",
|
|
)
|
|
|
|
parts: list[str] = []
|
|
cur: list[str] = []
|
|
cur_len = 0
|
|
fence_lang: str | None = None # limbajul fence-ului deschis in bucata curenta
|
|
|
|
def flush(reopen: bool) -> None:
|
|
nonlocal cur, cur_len
|
|
if not cur:
|
|
return
|
|
body = "\n".join(cur)
|
|
if fence_lang is not None:
|
|
body += "\n```"
|
|
parts.append(body)
|
|
cur = [f"```{fence_lang}"] if (reopen and fence_lang is not None) else []
|
|
cur_len = sum(len(x) + 1 for x in cur)
|
|
|
|
for line in text.split("\n"):
|
|
# o linie singura mai lunga decat limita se taie dur
|
|
pieces = [line] if len(line) <= limit - 10 else [
|
|
line[i: i + limit - 10] for i in range(0, len(line), limit - 10)
|
|
]
|
|
for piece in pieces:
|
|
need = len(piece) + 1 + (4 if fence_lang is not None else 0)
|
|
if cur and cur_len + need > limit:
|
|
flush(reopen=True)
|
|
cur.append(piece)
|
|
cur_len += len(piece) + 1
|
|
m = _FENCE.match(piece.strip())
|
|
if m:
|
|
fence_lang = None if fence_lang is not None else (m.group(1) or "")
|
|
flush(reopen=False)
|
|
if not parts:
|
|
parts = [""]
|
|
return Chunks(parts=parts)
|
|
|
|
|
|
class ChannelEditLoop:
|
|
"""Un loop de editare per canal. Coalescent: conteaza doar ultimul text per tinta."""
|
|
|
|
def __init__(
|
|
self,
|
|
channel_id: str,
|
|
edit_fn: Callable[[Any, str], Awaitable[None]],
|
|
*,
|
|
min_interval: float = 1.0,
|
|
max_interval: float = 5.0,
|
|
clock: Callable[[], float] = time.monotonic,
|
|
):
|
|
self.channel_id = str(channel_id)
|
|
self.edit_fn = edit_fn
|
|
self.min_interval = min_interval
|
|
self.max_interval = max_interval
|
|
self.interval = min_interval
|
|
self.clock = clock
|
|
self.pending: dict[Any, str] = {}
|
|
self.edits = 0
|
|
self.coalesced = 0
|
|
self.errors = 0
|
|
self._wake = asyncio.Event()
|
|
self._task: asyncio.Task | None = None
|
|
self._stop = False
|
|
|
|
# ------------------------------------------------------------ interfata
|
|
def queue(self, target: Any, text: str) -> None:
|
|
"""Ultima versiune castiga; actualizarile intermediare se pierd intentionat."""
|
|
if target in self.pending:
|
|
self.coalesced += 1
|
|
self.pending[target] = text
|
|
self._wake.set()
|
|
|
|
def note_rate_limited(self) -> None:
|
|
"""Adaptorul semnaleaza un 429: urcam direct la intervalul maxim."""
|
|
self.interval = self.max_interval
|
|
|
|
def start(self) -> asyncio.Task:
|
|
if self._task is None or self._task.done():
|
|
self._stop = False
|
|
self._task = asyncio.create_task(self._run())
|
|
return self._task
|
|
|
|
async def stop(self, flush: bool = True) -> None:
|
|
self._stop = True
|
|
self._wake.set()
|
|
if self._task is not None:
|
|
with contextlib.suppress(asyncio.CancelledError, Exception):
|
|
await asyncio.wait_for(self._task, 5.0)
|
|
self._task = None
|
|
if flush:
|
|
await self.flush()
|
|
|
|
async def flush(self) -> None:
|
|
"""Trimite imediat tot ce e in asteptare (finalul unui tur)."""
|
|
while self.pending:
|
|
target, text = next(iter(self.pending.items()))
|
|
del self.pending[target]
|
|
await self._edit(target, text)
|
|
|
|
# -------------------------------------------------------------- interne
|
|
async def _edit(self, target: Any, text: str) -> None:
|
|
try:
|
|
await self.edit_fn(target, text)
|
|
self.edits += 1
|
|
except Exception:
|
|
self.errors += 1
|
|
log.warning("editare esuata pe canalul %s", self.channel_id, exc_info=True)
|
|
|
|
async def _run(self) -> None:
|
|
while not self._stop:
|
|
if not self.pending:
|
|
self._wake.clear()
|
|
with contextlib.suppress(asyncio.TimeoutError):
|
|
await asyncio.wait_for(self._wake.wait(), self.max_interval)
|
|
# liniste -> coboram spre intervalul minim
|
|
self.interval = max(self.min_interval, self.interval / 1.5)
|
|
continue
|
|
batch = self.pending
|
|
self.pending = {}
|
|
for target, text in batch.items():
|
|
await self._edit(target, text)
|
|
if len(batch) > 1:
|
|
self.interval = min(self.max_interval, self.interval * 1.5)
|
|
await asyncio.sleep(self.interval)
|
|
|
|
|
|
class RenderManager:
|
|
"""Un singur ChannelEditLoop per canal."""
|
|
|
|
def __init__(self, edit_fn: Callable[[Any, str], Awaitable[None]], **kw):
|
|
self.edit_fn = edit_fn
|
|
self.kw = kw
|
|
self.loops: dict[str, ChannelEditLoop] = {}
|
|
|
|
def loop_for(self, channel_id: str) -> ChannelEditLoop:
|
|
key = str(channel_id)
|
|
loop = self.loops.get(key)
|
|
if loop is None:
|
|
loop = self.loops[key] = ChannelEditLoop(key, self.edit_fn, **self.kw)
|
|
loop.start()
|
|
return loop
|
|
|
|
def queue(self, channel_id: str, target: Any, text: str) -> None:
|
|
self.loop_for(channel_id).queue(target, text)
|
|
|
|
async def stop_all(self) -> None:
|
|
for loop in list(self.loops.values()):
|
|
await loop.stop()
|
|
self.loops.clear()
|
|
|
|
|
|
def footer(model: str, duration_ms: int, cost_usd: float, thread_total: float) -> str:
|
|
"""Subsolul cerut la T11: model, durata, cost tur, cost cumulat pe fir."""
|
|
return (
|
|
f"-# {model} · {duration_ms / 1000:.1f}s · ${cost_usd:.4f} tur · ${thread_total:.4f} fir"
|
|
)
|