Files
ROMFASTSQL/proxmox/lxc171-claude-agent/discord-bridge/stream.py
Claude Agent a89cf201ab fix(discord-bridge): zgomot in bot.log si doua tipuri de stream necunoscute
Trei lucruri observate in logurile de productie dupa restartul precedent.

1. Ecoul propriilor mesaje umplea bot.log cu WARNING. Propriile mesaje au si
   ele `author.bot == True`, iar verificarea generica de bot venea INAINTEA
   celei pe `self_id` — deci raspunsurile botului se jurnalizau ca "bot
   strain", la fiecare mesaj. Verificarea pe `self_id` trece prima (motivul e
   acum precis), iar refuzurile de rutina — propriile mesaje si ceilalti boti —
   merg la DEBUG. Guild / canal / utilizator strain si webhook raman WARNING:
   alea chiar sunt semnal de securitate si erau inecate in zgomot.

2. `tool_progress` (heartbeat la 30s cat timp o unealta ruleaza) devine
   eveniment `ToolProgress`. Mesajul live arata acum "⏳ ruleaza de 2m30s" sub
   unealta curenta — singurul semn ca un tur lung lucreaza si nu a inghetat.

3. `rate_limit_event` devine eveniment `RateLimit` si apare in `/status` la
   randul `utilizare`. Cum nu exista plafon de cost (abonament, nu API),
   fereastra de utilizare e singura limita reala; se avertizeaza in log o data
   per schimbare de stare, nu la fiecare eveniment.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01B29CApsP1JkSdjYaGaHpE7
2026-08-31 07:46:51 +00:00

329 lines
11 KiB
Python

"""Parser tolerant pentru `claude --output-format stream-json --verbose`.
Reguli (T6):
- tip necunoscut -> se logeaza O SINGURA DATA per tip, cu versiunea CLI, si se ignora
(asa au fost descoperite `tool_progress` si `rate_limit_event` in claude 2.1.251;
cand un tip nou se dovedeste util, i se face un eveniment si iese din lista)
- linie non-JSON -> se logeaza si se ignora
- EOF inainte de `result` -> StreamEOFError, explicit, fara hang
Parser-ul NU arunca niciodata pe continut de stream; singura exceptie e EOF-ul de mai sus.
"""
from __future__ import annotations
import json
import logging
import subprocess
from dataclasses import dataclass, field
from typing import Any, AsyncIterator, Iterable
log = logging.getLogger("discord-bridge.stream")
_cli_version_cache: str | None = None
def cli_version() -> str:
"""`claude --version`, cu cache si tolerant la orice esec."""
global _cli_version_cache
if _cli_version_cache is None:
try:
out = subprocess.run(
["claude", "--version"], capture_output=True, text=True, timeout=10
)
_cli_version_cache = (out.stdout or out.stderr).strip() or "necunoscuta"
except Exception:
_cli_version_cache = "necunoscuta"
return _cli_version_cache
class StreamEOFError(RuntimeError):
"""Stream-ul s-a terminat inainte de evenimentul `result`."""
# ---------------------------------------------------------------- evenimente
@dataclass(frozen=True)
class SystemInit:
session_id: str | None
model: str | None
cwd: str | None
tools: tuple[str, ...] = ()
raw: dict = field(default_factory=dict, repr=False)
@dataclass(frozen=True)
class AssistantText:
text: str
session_id: str | None = None
raw: dict = field(default_factory=dict, repr=False)
@dataclass(frozen=True)
class ToolUse:
name: str
tool_id: str | None
input: dict = field(default_factory=dict)
session_id: str | None = None
raw: dict = field(default_factory=dict, repr=False)
@dataclass(frozen=True)
class ToolResult:
tool_id: str | None
text: str
is_error: bool = False
session_id: str | None = None
raw: dict = field(default_factory=dict, repr=False)
@dataclass(frozen=True)
class Result:
total_cost_usd: float
duration_ms: int
is_error: bool
num_turns: int
text: str = ""
session_id: str | None = None
raw: dict = field(default_factory=dict, repr=False)
@dataclass(frozen=True)
class ToolProgress:
"""Semn de viata pentru o unealta care ruleaza de mult (`tool_progress`).
CLI-ul il trimite periodic (heartbeat) cat timp o unealta e in executie. E
singurul semnal ca un tur lung inca lucreaza si nu a inghetat.
"""
tool_name: str
tool_id: str | None
elapsed_s: int
heartbeat: bool = False
session_id: str | None = None
raw: dict = field(default_factory=dict, repr=False)
@dataclass(frozen=True)
class RateLimit:
"""`rate_limit_event`: starea ferestrei de utilizare a abonamentului.
`status` e "allowed" cand totul e in regula; orice altceva inseamna ca
urmatoarele tururi pot fi franate sau refuzate.
"""
status: str
limit_type: str = ""
resets_at: int = 0
raw: dict = field(default_factory=dict, repr=False)
@property
def ok(self) -> bool:
return self.status == "allowed"
Event = SystemInit | AssistantText | ToolUse | ToolResult | ToolProgress | RateLimit | Result
def _blocks(msg: Any) -> list[dict]:
if not isinstance(msg, dict):
return []
content = msg.get("content")
if isinstance(content, str):
return [{"type": "text", "text": content}]
if isinstance(content, list):
return [b for b in content if isinstance(b, dict)]
return []
def _as_text(value: Any) -> str:
if isinstance(value, str):
return value
if isinstance(value, list):
parts = []
for b in value:
if isinstance(b, dict) and isinstance(b.get("text"), str):
parts.append(b["text"])
elif isinstance(b, str):
parts.append(b)
return "".join(parts)
if value is None:
return ""
return str(value)
class StreamParser:
"""Transforma linii JSONL in evenimente normalizate. Stateful doar pentru logare."""
def __init__(self, version_fn=None, logger: logging.Logger | None = None):
# rezolvat la apel, ca sa poata fi inlocuit in teste
self._version_fn = version_fn or (lambda: cli_version())
self._log = logger or log
self.unknown_types: set[str] = set()
self.bad_lines = 0
self.saw_result = False
self.session_id: str | None = None
# ---------------------------------------------------------------- feed
def feed_line(self, line: str) -> list[Event]:
"""Returneaza 0..n evenimente pentru o linie. Nu arunca niciodata."""
try:
return self._feed_line(line)
except Exception as exc: # plasa de siguranta: nimic din stream nu doboara botul
self._log.warning("stream: linie neasteptata ignorata: %r (%s)", line[:200], exc)
return []
def _feed_line(self, line: str) -> list[Event]:
text = line.strip()
if not text:
return []
try:
obj = json.loads(text)
except (ValueError, TypeError):
self.bad_lines += 1
self._log.warning("stream: linie non-JSON ignorata: %r", text[:200])
return []
if not isinstance(obj, dict):
self.bad_lines += 1
self._log.warning("stream: JSON care nu e obiect, ignorat: %r", text[:200])
return []
typ = obj.get("type")
sid = obj.get("session_id")
if isinstance(sid, str) and sid:
self.session_id = sid
if typ == "system":
if obj.get("subtype") == "init":
tools = obj.get("tools")
return [
SystemInit(
session_id=sid if isinstance(sid, str) else None,
model=obj.get("model"),
cwd=obj.get("cwd"),
tools=tuple(t for t in tools if isinstance(t, str)) if isinstance(tools, list) else (),
raw=obj,
)
]
# alte subtipuri de system (compact_boundary etc) -> zgomot, ignorat
return []
if typ == "assistant":
events: list[Event] = []
for b in _blocks(obj.get("message")):
if b.get("type") == "text" and b.get("text"):
events.append(AssistantText(text=str(b["text"]), session_id=self.session_id, raw=obj))
elif b.get("type") == "tool_use":
events.append(
ToolUse(
name=str(b.get("name") or "?"),
tool_id=b.get("id"),
input=b.get("input") if isinstance(b.get("input"), dict) else {},
session_id=self.session_id,
raw=obj,
)
)
return events
if typ == "user":
events = []
for b in _blocks(obj.get("message")):
if b.get("type") == "tool_result":
events.append(
ToolResult(
tool_id=b.get("tool_use_id"),
text=_as_text(b.get("content")),
is_error=bool(b.get("is_error")),
session_id=self.session_id,
raw=obj,
)
)
return events
if typ == "result":
self.saw_result = True
try:
cost = float(obj.get("total_cost_usd") or 0.0)
except (TypeError, ValueError):
cost = 0.0
try:
dur = int(obj.get("duration_ms") or 0)
except (TypeError, ValueError):
dur = 0
try:
turns = int(obj.get("num_turns") or 0)
except (TypeError, ValueError):
turns = 0
return [
Result(
total_cost_usd=cost,
duration_ms=dur,
is_error=bool(obj.get("is_error")) or obj.get("subtype") not in (None, "success"),
num_turns=turns,
text=_as_text(obj.get("result")),
session_id=self.session_id,
raw=obj,
)
]
if typ == "tool_progress":
try:
elapsed = int(obj.get("elapsed_time_seconds") or 0)
except (TypeError, ValueError):
elapsed = 0
return [
ToolProgress(
tool_name=str(obj.get("tool_name") or "?"),
tool_id=obj.get("tool_use_id"),
elapsed_s=elapsed,
heartbeat=bool(obj.get("heartbeat")),
session_id=self.session_id,
raw=obj,
)
]
if typ == "rate_limit_event":
info = obj.get("rate_limit_info")
info = info if isinstance(info, dict) else {}
try:
resets = int(info.get("resetsAt") or 0)
except (TypeError, ValueError):
resets = 0
return [
RateLimit(
status=str(info.get("status") or "necunoscut"),
limit_type=str(info.get("rateLimitType") or ""),
resets_at=resets,
raw=obj,
)
]
key = str(typ)
if key not in self.unknown_types:
self.unknown_types.add(key)
self._log.warning(
"stream: tip necunoscut %r ignorat (claude %s); exemplu: %r",
key, self._version_fn(), text[:200],
)
return []
# ------------------------------------------------------------ iteratoare
def feed_lines(self, lines: Iterable[str]) -> list[Event]:
out: list[Event] = []
for line in lines:
out.extend(self.feed_line(line))
return out
async def aiter_events(self, lines: AsyncIterator[str | bytes]) -> AsyncIterator[Event]:
"""Consuma un iterator asincron de linii; ridica StreamEOFError daca
stream-ul se termina inainte de `result`."""
self.saw_result = False
async for raw in lines:
if isinstance(raw, bytes):
raw = raw.decode("utf-8", "replace")
for ev in self.feed_line(raw):
yield ev
if isinstance(ev, Result):
return
raise StreamEOFError(
"stream-ul claude s-a terminat inainte de evenimentul `result`"
)