Files
Llamagochi/log_parser.py

205 lines
8.2 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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
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
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,
"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+),.*/ ([\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')
# ── 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")
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 _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):
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,
)
return self._set(state="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):
# End-of-slot summary: update prompt speed without changing state.
# For short prompts this fires after generating starts (no live prompt line).
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,
)
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:
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")
return None
# ── Helpers ──────────────────────────────────────────────────────────────
def _set(self, **kwargs) -> ServerState:
for k, v in kwargs.items():
setattr(self._state, k, v)
return self._state