From 1411ba241fb299c755099c8dccf6b8156e0d2781 Mon Sep 17 00:00:00 2001 From: Vlad Doloman Date: Mon, 22 Jun 2026 18:46:12 +0300 Subject: [PATCH] feat: LogParser state machine with full state transitions --- log_parser.py | 142 +++++++++++++++++++++++++++++++++ tests/test_log_parser.py | 165 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 307 insertions(+) diff --git a/log_parser.py b/log_parser.py index 9a564f8..1cceabc 100644 --- a/log_parser.py +++ b/log_parser.py @@ -60,3 +60,145 @@ class ServerState: "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 diff --git a/tests/test_log_parser.py b/tests/test_log_parser.py index 4e1b0d8..f08ec82 100644 --- a/tests/test_log_parser.py +++ b/tests/test_log_parser.py @@ -93,3 +93,168 @@ def test_server_state_to_dict_generating(): assert d["metrics"]["gen_speed"] == 72.2 assert d["metrics"]["n_decoded"] == 218 assert d["request_count"] == 5 + + +from log_parser import LogParser + +# Helper: feed a list of raw lines (with ts prefix), return final state name +def feed_lines(lines): + p = LogParser() + last = None + for line in lines: + result = p.feed(line) + if result is not None: + last = result + return p.state.state, p.state + + +# ── Router startup ─────────────────────────────────────────────────────────── + +def test_router_startup(): + state, _ = feed_lines([ + "[2026-06-22 14:29:49] 0.00.076.839 I srv llama_server: router server is listening on http://0.0.0.0:8082" + ]) + assert state == "starting" + + +# ── Model spawn + loading progress ────────────────────────────────────────── + +SPAWN_LINE = "[2026-06-22 14:31:59] 2.10.286.889 I srv load: spawning server instance with name=unsloth/gemma-4-31B-it:UD-Q8_K_XL-recommended on port 41803" +LOADING_START = '[2026-06-22 14:31:59] [41803] cmd_child_to_router:state:{"state":"loading","payload":{"stages":["text_model","mmproj_model"],"current":"text_model","value":0.0}}' +LOADING_MID = '[2026-06-22 14:32:05] [41803] cmd_child_to_router:state:{"state":"loading","payload":{"stages":["text_model","mmproj_model"],"current":"text_model","value":0.5}}' +LOADING_DONE = '[2026-06-22 14:32:08] [41803] cmd_child_to_router:state:{"state":"loading","payload":{"stages":["text_model","mmproj_model"],"current":"text_model","value":1.0}}' +WARMUP_LINE = "[2026-06-22 14:32:10] [41803] 0.11.107.736 I common_init_from_params: warming up the model with an empty run - please wait ..." +READY_LINE = '[2026-06-22 14:32:13] [41803] cmd_child_to_router:state:{"state":"ready","payload":{"id":"unsloth/gemma-4-31B-it:UD-Q8_K_XL-recommended","aliases":[],"tags":[],"object":"model","created":1782127933,"owned_by":"llamacpp","meta":{}}}' +IDLE_LINE = "[2026-06-22 14:32:13] [41803] 0.13.230.405 I srv update_slots: all slots are idle" + + +def test_model_loading_progress(): + p = LogParser() + p.feed(SPAWN_LINE) + p.feed(LOADING_START) + assert p.state.state == "loading" + assert p.state.loading.progress == 0.0 + assert p.state.loading.current == "text_model" + + p.feed(LOADING_MID) + assert p.state.loading.progress == 0.5 + + p.feed(LOADING_DONE) + assert p.state.loading.progress == 1.0 + + +def test_spawn_captures_model_name(): + p = LogParser() + p.feed(SPAWN_LINE) + assert p.state.model == "unsloth/gemma-4-31B-it:UD-Q8_K_XL-recommended" + + +def test_warmup_state(): + state, _ = feed_lines([SPAWN_LINE, LOADING_START, WARMUP_LINE]) + assert state == "warming_up" + + +def test_ready_transitions_to_idle(): + state, s = feed_lines([SPAWN_LINE, LOADING_START, WARMUP_LINE, READY_LINE]) + assert state == "idle" + assert s.model == "unsloth/gemma-4-31B-it:UD-Q8_K_XL-recommended" + + +def test_idle_line(): + state, _ = feed_lines([SPAWN_LINE, LOADING_START, READY_LINE, IDLE_LINE]) + assert state == "idle" + + +# ── Request lifecycle ──────────────────────────────────────────────────────── + +PROXY_LINE = "[2026-06-22 14:32:13] 2.23.870.146 I srv proxy_reques: proxying request to model unsloth/gemma-4-31B-it:UD-Q8_K_XL-recommended on port 41803" +LAUNCH_LINE = "[2026-06-22 14:34:20] [41803] 0.30.889.639 I slot launch_slot_: id 0 | task 0 | processing task, is_child = 0" +PROMPT_TIMING = "[2026-06-22 14:33:17] [41803] 0.56.093.911 I slot print_timing: id 0 | task 0 | prompt processing, n_tokens = 276, progress = 0.99, t = 4.90 s / 56.38 tokens per second" +GEN_TIMING = "[2026-06-22 14:33:20] [41803] 0.59.444.220 I slot print_timing: id 0 | task 0 | n_decoded = 156, tg = 51.94 t/s, tg_3s = 51.94 t/s" +RELEASE_LINE = "[2026-06-22 14:33:24] [41803] 1.02.743.373 I slot release: id 0 | task 0 | stop processing: n_tokens = 606, truncated = 0" + + +def test_proxy_sets_waiting(): + state, _ = feed_lines([READY_LINE, IDLE_LINE, PROXY_LINE]) + assert state == "waiting" + + +def test_launch_sets_processing_prompt(): + state, _ = feed_lines([IDLE_LINE, PROXY_LINE, LAUNCH_LINE]) + assert state == "processing_prompt" + + +def test_prompt_timing_updates_metrics(): + _, s = feed_lines([IDLE_LINE, PROXY_LINE, LAUNCH_LINE, PROMPT_TIMING]) + assert s.state == "processing_prompt" + assert abs(s.metrics.prompt_speed - 56.38) < 0.01 + assert s.metrics.prompt_tokens == 276 + + +def test_gen_timing_sets_generating(): + _, s = feed_lines([IDLE_LINE, PROXY_LINE, LAUNCH_LINE, PROMPT_TIMING, GEN_TIMING]) + assert s.state == "generating" + assert abs(s.metrics.gen_speed - 51.94) < 0.01 + assert s.metrics.n_decoded == 156 + + +# Short prompts never emit a live "prompt processing" line — only the end-of-slot +# summary fires, after generation has already started. +PROMPT_SUMMARY = "[2026-06-22 14:34:24] [41983] 0.34.388.097 I slot print_timing: id 0 | task 0 | prompt eval time = 425.83 ms / 11 tokens ( 38.71 ms per token, 25.83 tokens per second)" + + +def test_prompt_summary_captures_speed_during_generating(): + _, s = feed_lines([IDLE_LINE, PROXY_LINE, LAUNCH_LINE, GEN_TIMING, PROMPT_SUMMARY]) + assert s.state == "generating" # state name unchanged + assert abs(s.metrics.prompt_speed - 25.83) < 0.01 + assert s.metrics.prompt_tokens == 11 + assert abs(s.metrics.gen_speed - 51.94) < 0.01 # gen metrics preserved + + +def test_release_increments_request_count(): + _, s = feed_lines([IDLE_LINE, PROXY_LINE, LAUNCH_LINE, GEN_TIMING, RELEASE_LINE]) + assert s.state == "idle" + assert s.request_count == 1 + assert s.metrics.gen_speed is None # cleared on release + + +def test_release_clears_metrics(): + _, s = feed_lines([IDLE_LINE, PROXY_LINE, LAUNCH_LINE, GEN_TIMING, RELEASE_LINE]) + assert s.metrics.prompt_speed is None + assert s.metrics.gen_speed is None + assert s.metrics.n_decoded is None + + +# ── Sleeping ───────────────────────────────────────────────────────────────── + +SLEEPING_LINE = '[2026-06-22 15:23:11] [36417] cmd_child_to_router:state:{"state":"sleeping","payload":null}' + + +def test_sleeping_state(): + state, _ = feed_lines([IDLE_LINE, SLEEPING_LINE]) + assert state == "sleeping" + + +# ── Unloading ──────────────────────────────────────────────────────────────── + +UNLOAD_LINE = "[2026-06-22 14:32:19] 2.30.337.683 I srv unload_lru: models_max limit reached, removing LRU name=foo" + + +def test_unload_lru(): + state, _ = feed_lines([IDLE_LINE, UNLOAD_LINE]) + assert state == "unloading" + + +# ── Without ts prefix (single-instance mode) ──────────────────────────────── + +def test_no_ts_prefix(): + p = LogParser() + p.feed("0.00.075.237 I srv llama_server: router server is listening on http://0.0.0.0:8082") + assert p.state.state == "starting" + + +def test_child_no_ts_prefix(): + p = LogParser() + p.feed('[41803] cmd_child_to_router:state:{"state":"loading","payload":{"stages":["text_model"],"current":"text_model","value":0.3}}') + assert p.state.state == "loading" + assert p.state.loading.progress == 0.3