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
329 lines
11 KiB
Python
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`"
|
|
)
|