import re import json from dataclasses import dataclass, field, asdict from datetime import datetime from typing import Optional # ── Preprocessing ──────────────────────────────────────────────────────────── _TS_RE = re.compile(r'^\[\d{4}-\d{2}-\d{2} \d{2}:\d{2}:\d{2}\]\s*') _CHILD_RE = re.compile(r'^\[(\d+)\]\s*') def strip_ts(line: str) -> str: """Remove optional `ts`-injected timestamp prefix from a log line.""" return _TS_RE.sub('', line.strip()) def classify_line(line: str) -> tuple[str, Optional[int], str]: """ Return (kind, port, content) where kind is 'router' or 'child'. `line` must already have the ts prefix stripped. """ m = _CHILD_RE.match(line) if m: return ('child', int(m.group(1)), line[m.end():]) return ('router', None, line) # ── Data model ─────────────────────────────────────────────────────────────── @dataclass class LoadingInfo: stages: list current: str progress: float # 0.0 – 1.0 @dataclass class Metrics: prompt_speed: Optional[float] = None # tokens/sec during prompt eval prompt_tokens: Optional[int] = None # total prompt tokens prompt_progress: Optional[float] = None # 0.0–1.0 from live prompt timing lines gen_speed: Optional[float] = None # tokens/sec during generation n_decoded: Optional[int] = None # tokens decoded so far @dataclass class ServerState: state: str = "offline" model: Optional[str] = None loading: Optional[LoadingInfo] = None metrics: Metrics = field(default_factory=Metrics) request_count: int = 0 sleeping_since: Optional[str] = None # ISO timestamp when sleeping started def to_dict(self) -> dict: return { "state": self.state, "model": self.model, "loading": asdict(self.loading) if self.loading else None, "metrics": asdict(self.metrics), "request_count": self.request_count, "sleeping_since": self.sleeping_since, "timestamp": datetime.now().isoformat(timespec="seconds"), } # ── Log patterns ───────────────────────────────────────────────────────────── _SPAWN_RE = re.compile(r'spawning server instance with name=(\S+)') _CMD_STATE_RE = re.compile(r'^cmd_child_to_router:state:(.+)$') _ROUTER_LISTENING_RE = re.compile(r'llama_server: router server is listening') _WARMUP_RE = re.compile(r'common_init_from_params: warming up') _IDLE_RE = re.compile(r'update_slots: all slots are idle') _PROXY_RE = re.compile(r'proxy_reques: proxying request') _LAUNCH_RE = re.compile(r'slot launch_slot_:.*processing task') _PROMPT_TIMING_RE = re.compile( r'slot print_timing:.*prompt processing, n_tokens =\s+(\d+), progress = ([\d.]+),.*/ ([\d.]+) tokens per second' ) _GEN_TIMING_RE = re.compile( r'slot print_timing:.*n_decoded =\s+(\d+), tg =\s+([\d.]+) t/s' ) _PROMPT_SUMMARY_RE = re.compile( r'slot print_timing:.*prompt eval time =.*/ \s*(\d+) tokens.*?([\d.]+) tokens per second' ) _RELEASE_RE = re.compile(r'slot\s+release:.*stop processing') _UNLOAD_RE = re.compile(r'unload_lru:|unload: stopping model instance') _LOAD_ERROR_RE = re.compile(r'exiting due to model loading error') _EXIT_ERROR_RE = re.compile(r'operator\(\):.*exited with status 1') _CLIENT_CANCEL_RE = re.compile(r'http client error: Connection handling canceled') _CANCEL_TASK_RE = re.compile(r'stop: cancel task') # ── State machine ───────────────────────────────────────────────────────────── class LogParser: def __init__(self): self._state = ServerState() self._active_port: Optional[int] = None @property def state(self) -> ServerState: return self._state def feed(self, raw_line: str) -> Optional[ServerState]: """ Process one raw log line (with or without ts prefix). Returns the updated ServerState if something meaningful changed, else None. """ line = strip_ts(raw_line) if not line: return None kind, port, content = classify_line(line) if kind == "child": return self._handle_child(port, content) return self._handle_router(content) # ── Router-level lines ─────────────────────────────────────────────────── def _handle_router(self, content: str) -> Optional[ServerState]: if _ROUTER_LISTENING_RE.search(content): return self._set(state="starting") if m := _SPAWN_RE.search(content): self._state.model = m.group(1) return None # wait for child's loading signal if _PROXY_RE.search(content): return self._set(state="waiting") if _UNLOAD_RE.search(content): return self._set(state="unloading") if _EXIT_ERROR_RE.search(content): return self._set(state="error") if _CLIENT_CANCEL_RE.search(content): return self._set(state="idle", loading=None, metrics=Metrics()) return None # ── Child-process lines ────────────────────────────────────────────────── def _handle_child(self, port: int, content: str) -> Optional[ServerState]: self._active_port = port if m := _CMD_STATE_RE.match(content): return self._handle_cmd_state(m.group(1)) if _LOAD_ERROR_RE.search(content): return self._set(state="error") if _CANCEL_TASK_RE.search(content): return self._set(state="idle", loading=None, metrics=Metrics()) if _WARMUP_RE.search(content): return self._set(state="warming_up") if _IDLE_RE.search(content): return self._set(state="idle", loading=None, metrics=Metrics()) if _LAUNCH_RE.search(content): return self._set(state="processing_prompt") if m := _PROMPT_TIMING_RE.search(content): progress = float(m.group(2)) done = progress >= 1.0 metrics = Metrics( prompt_speed=float(m.group(3)), prompt_tokens=int(m.group(1)), prompt_progress=None if done else progress, ) return self._set(state="generating" if done else "processing_prompt", metrics=metrics) if m := _GEN_TIMING_RE.search(content): metrics = Metrics( prompt_speed=self._state.metrics.prompt_speed, prompt_tokens=self._state.metrics.prompt_tokens, gen_speed=float(m.group(2)), n_decoded=int(m.group(1)), ) return self._set(state="generating", metrics=metrics) if m := _PROMPT_SUMMARY_RE.search(content): metrics = Metrics( prompt_speed=float(m.group(2)), prompt_tokens=int(m.group(1)), gen_speed=self._state.metrics.gen_speed, n_decoded=self._state.metrics.n_decoded, ) if self._state.state == "processing_prompt": # Summary fires right after prompt eval finishes, before first token. # Advance to generating now rather than waiting for the first tg line. # For short prompts the summary arrives after generating has already # started, so the guard prevents stepping back. return self._set(state="generating", metrics=metrics) return self._set(metrics=metrics) if _RELEASE_RE.search(content): return self._set( state="idle", request_count=self._state.request_count + 1, loading=None, metrics=Metrics(), ) return None # ── cmd_child_to_router:state JSON handler ─────────────────────────────── def _handle_cmd_state(self, json_str: str) -> Optional[ServerState]: try: data = json.loads(json_str) except json.JSONDecodeError: return None s = data.get("state") payload = data.get("payload") if s == "loading" and payload: # Some transition messages only carry {"stage":"name"} with no # progress info — skip them; the next message has the full payload. if not all(k in payload for k in ("stages", "current", "value")): return None loading = LoadingInfo( stages=payload["stages"], current=payload["current"], progress=payload["value"], ) return self._set(state="loading", loading=loading) if s == "ready": if payload and (model_id := payload.get("id")): self._state.model = model_id return self._set(state="idle", loading=None) if s == "sleeping": return self._set( state="sleeping", sleeping_since=datetime.now().isoformat(timespec="seconds"), ) return None # ── Helpers ────────────────────────────────────────────────────────────── def _set(self, **kwargs) -> ServerState: if "state" in kwargs and kwargs["state"] != "sleeping": self._state.sleeping_since = None for k, v in kwargs.items(): setattr(self._state, k, v) return self._state