feat: LogParser state machine with full state transitions
This commit is contained in:
142
log_parser.py
142
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
|
||||
|
||||
Reference in New Issue
Block a user