"""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`" )