Created
May 17, 2026 10:12
-
-
Save david-gang/3640356222a49d0face9907c0f4cd31a to your computer and use it in GitHub Desktop.
KAI-Scheduler PodGroup trace parser — streams a snapshot-tool log and emits a Markdown report for one pending PodGroup (related to kai-scheduler/KAI-Scheduler#1591)
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| #!/usr/bin/env -S uv run --script | |
| # /// script | |
| # requires-python = ">=3.11" | |
| # dependencies = [ | |
| # "tabulate>=0.9", | |
| # "rich>=13", | |
| # ] | |
| # /// | |
| """Parse a KAI-Scheduler snapshot log and emit a Markdown trace for one PodGroup. | |
| Diagnostic tool for issue #1591. Streams the log line-by-line, keys off stable | |
| file:line callsite anchors, and groups output by action -> scenario -> per-node | |
| attempt, plus a predicate-failure breakdown. | |
| Usage: | |
| uv run tools/log-analysis/parse_pg_trace.py [--log PATH] [--pg NAME] | |
| [--task-prefix PFX] | |
| [--out PATH] [--top-n N] | |
| """ | |
| import argparse | |
| import os | |
| import re | |
| import sys | |
| from collections import Counter, defaultdict | |
| from dataclasses import dataclass, field | |
| from typing import Optional | |
| from tabulate import tabulate | |
| from rich.console import Console | |
| from rich.table import Table | |
| REPO_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "..")) | |
| DEFAULT_LOG = os.path.abspath( | |
| os.path.join( | |
| REPO_ROOT, | |
| "..", | |
| "KAI-Scheduler-v0.14.1", | |
| "snapshots", | |
| "snapshot-output6.log", | |
| ) | |
| ) | |
| DEFAULT_PG = ( | |
| "pg-<your-podgroup-name>" | |
| ) | |
| DEFAULT_TASK_PREFIX = "<namespace>/<task-name-prefix>" | |
| DEFAULT_OUT = os.path.join(REPO_ROOT, "tools", "log-analysis", "report-1591.md") | |
| ANSI = re.compile(r"\x1b\[[0-9;]*m") | |
| # Match by message-prefix substring instead of `file:line` so the parser | |
| # survives line-number drift between KAI versions. | |
| MESSAGES = [ | |
| "Attempting to allocate job:", | |
| "Could not allocate resources for job:", | |
| "Attempting to consolidate running jobs", | |
| "Can't consolidate for job", | |
| "Sucesfully consolidated", | |
| "Attempting to reclaim for job:", | |
| "Didn't find a reclaim strategy", | |
| "Attempting to preempt for job:", | |
| "Didn't find a preemption strategy", | |
| "Trying to solve scenario:", | |
| "No solution found for", | |
| "Trying to solve scenario with potantial", | |
| "Looking for best node for task", | |
| "Failed statement allocate for task", | |
| "Statement evicted task:", | |
| "Predicates failed for task", | |
| ] | |
| ANCHOR_RE = re.compile("|".join(re.escape(m) for m in MESSAGES)) | |
| TS_RE = re.compile(r"^(\S+)\s+INFO\s+(\S+)\s+\[snapshot-runner\]\s+(?:\[(\w+)\]\s+)?(.*)$") | |
| DEMAND_RE = re.compile(r"resources: <([^>]+)>") | |
| PG_ANGLE_RE = re.compile(r"<default/([^>]+)>") | |
| PRIORITY_RE = re.compile(r"priority: <(\d+)>") | |
| QUEUE_RE = re.compile(r"queue: <([^>]+)>") | |
| SCENARIO_RE = re.compile( | |
| r"Trying to solve scenario:\s*Pending tasks:(?P<pending>.*?)\s*" | |
| r"Recorded victim jobs:(?P<recorded>.*?)\s*" | |
| r"Potential victim tasks:(?P<potential>.*?)\s*$" | |
| ) | |
| BY_POD_NODE_RE = re.compile(r"potantial victims from node:\s*(\S+)") | |
| EVICTED_RE = re.compile(r"evicted task:\s*<([^>]+)>\s*from node:\s*<([^>]+)>") | |
| PRED_FAIL_TASK_RE = re.compile( | |
| r"Predicates failed for task <([^>]+)> on node <([^>]+)>: .*?reasons:\s*(.*)$" | |
| ) | |
| NO_SOLUTION_RE = re.compile( | |
| r"No solution found for (\d+) tasks out of (\d+) tasks to allocate for (\S+)" | |
| ) | |
| PER_TASK_DEMAND_RE = re.compile( | |
| r"Looking for best node for task - Task: <([^>]+)>, init requested: <([^>]+)>" | |
| ) | |
| @dataclass | |
| class NodeAttempt: | |
| node: str | |
| evictions: list = field(default_factory=list) # list[(task, from_node)] | |
| rolled_back: bool = False | |
| @dataclass | |
| class Scenario: | |
| pending: str | |
| recorded_victims: str | |
| potential_victims: str | |
| node_attempts: list = field(default_factory=list) # list[NodeAttempt] | |
| no_solution: bool = False | |
| @dataclass | |
| class ActionRecord: | |
| name: str # allocate / consolidation / reclaim / preempt | |
| started_at: str | |
| demand: str = "" | |
| priority: str = "" | |
| queue: str = "" | |
| outcome: str = "(no terminal log)" | |
| closed_at: str = "" | |
| scenarios: list = field(default_factory=list) # list[Scenario] | |
| def parse_args(): | |
| p = argparse.ArgumentParser(description=__doc__) | |
| p.add_argument("--log", default=DEFAULT_LOG) | |
| p.add_argument("--pg", default=DEFAULT_PG) | |
| p.add_argument("--task-prefix", default=DEFAULT_TASK_PREFIX) | |
| p.add_argument("--out", default=DEFAULT_OUT) | |
| p.add_argument("--top-n", type=int, default=5) | |
| return p.parse_args() | |
| def parse(args): | |
| actions: list[ActionRecord] = [] | |
| current: Optional[ActionRecord] = None | |
| current_scenario: Optional[Scenario] = None | |
| current_node_attempt: Optional[NodeAttempt] = None | |
| predicate_failures = Counter() | |
| predicate_examples = defaultdict(set) | |
| lines_total = 0 | |
| lines_matched = 0 | |
| lines_for_our_pg = 0 | |
| gang_size: Optional[int] = None | |
| per_task_demand: Optional[str] = None | |
| probed_tasks: set = set() | |
| def close_action(): | |
| nonlocal current, current_scenario, current_node_attempt | |
| if current is not None: | |
| actions.append(current) | |
| current = None | |
| current_scenario = None | |
| current_node_attempt = None | |
| with open(args.log, "r", encoding="utf-8", errors="replace") as fh: | |
| for raw in fh: | |
| lines_total += 1 | |
| if not ANCHOR_RE.search(raw): | |
| continue | |
| line = ANSI.sub("", raw).rstrip("\n") | |
| m = TS_RE.match(line) | |
| if not m: | |
| continue | |
| lines_matched += 1 | |
| ts, callsite, action_tag, message = m.groups() | |
| # Action lifecycle — dispatch by message prefix (file path stable, line no. drifts) | |
| if "Attempting to allocate job:" in message and "allocate/allocate.go" in callsite: | |
| pg_match = PG_ANGLE_RE.search(message) | |
| if pg_match and pg_match.group(1) == args.pg: | |
| close_action() | |
| current = ActionRecord(name="allocate", started_at=ts) | |
| d = DEMAND_RE.search(message) | |
| current.demand = d.group(1) if d else "" | |
| q = QUEUE_RE.search(message) | |
| current.queue = q.group(1) if q else "" | |
| lines_for_our_pg += 1 | |
| elif "Could not allocate resources for job:" in message: | |
| pg_match = PG_ANGLE_RE.search(message) | |
| if pg_match and pg_match.group(1) == args.pg and current and current.name == "allocate": | |
| current.outcome = "Could not allocate resources for job" | |
| current.closed_at = ts | |
| close_action() | |
| lines_for_our_pg += 1 | |
| elif "Attempting to consolidate" in message: | |
| pg_match = PG_ANGLE_RE.search(message) | |
| if pg_match and pg_match.group(1) == args.pg: | |
| close_action() | |
| current = ActionRecord(name="consolidation", started_at=ts) | |
| d = DEMAND_RE.search(message) | |
| current.demand = d.group(1) if d else "" | |
| lines_for_our_pg += 1 | |
| elif "Can't consolidate for job" in message: | |
| pg_match = PG_ANGLE_RE.search(message) | |
| if pg_match and pg_match.group(1) == args.pg and current and current.name == "consolidation": | |
| current.outcome = "Can't consolidate, not enough allocatable GPUs in the cluster" | |
| current.closed_at = ts | |
| close_action() | |
| lines_for_our_pg += 1 | |
| elif "Sucesfully consolidated" in message: | |
| pg_match = PG_ANGLE_RE.search(message) | |
| if pg_match and pg_match.group(1) == args.pg and current and current.name == "consolidation": | |
| current.outcome = "Successfully consolidated" | |
| current.closed_at = ts | |
| close_action() | |
| lines_for_our_pg += 1 | |
| elif "Attempting to reclaim for job:" in message: | |
| pg_match = PG_ANGLE_RE.search(message) | |
| if pg_match and pg_match.group(1) == args.pg: | |
| close_action() | |
| current = ActionRecord(name="reclaim", started_at=ts) | |
| d = DEMAND_RE.search(message) | |
| current.demand = d.group(1) if d else "" | |
| q = QUEUE_RE.search(message) | |
| current.queue = q.group(1) if q else "" | |
| lines_for_our_pg += 1 | |
| elif "Didn't find a reclaim strategy" in message: | |
| if args.pg in message and current and current.name == "reclaim": | |
| current.outcome = "Didn't find a reclaim strategy" | |
| current.closed_at = ts | |
| close_action() | |
| lines_for_our_pg += 1 | |
| elif "Attempting to preempt for job:" in message: | |
| pg_match = PG_ANGLE_RE.search(message) | |
| if pg_match and pg_match.group(1) == args.pg: | |
| close_action() | |
| current = ActionRecord(name="preempt", started_at=ts) | |
| d = DEMAND_RE.search(message) | |
| current.demand = d.group(1) if d else "" | |
| pr = PRIORITY_RE.search(message) | |
| current.priority = pr.group(1) if pr else "" | |
| q = QUEUE_RE.search(message) | |
| current.queue = q.group(1) if q else "" | |
| lines_for_our_pg += 1 | |
| elif "Didn't find a preemption strategy" in message: | |
| if args.pg in message and current and current.name == "preempt": | |
| current.outcome = "Didn't find a preemption strategy" | |
| current.closed_at = ts | |
| close_action() | |
| lines_for_our_pg += 1 | |
| # Solver scenario open | |
| elif "Trying to solve scenario:" in message: | |
| if current is None: | |
| continue | |
| sm = SCENARIO_RE.search(message) | |
| if not sm: | |
| continue | |
| pending = sm.group("pending").strip() | |
| if not pending.startswith(args.task_prefix): | |
| continue | |
| current_scenario = Scenario( | |
| pending=pending, | |
| recorded_victims=sm.group("recorded").strip(), | |
| potential_victims=sm.group("potential").strip(), | |
| ) | |
| current.scenarios.append(current_scenario) | |
| current_node_attempt = None | |
| # Track every distinct task of our gang that gets probed | |
| for tok in pending.split(): | |
| if tok.startswith(args.task_prefix): | |
| probed_tasks.add(tok) | |
| lines_for_our_pg += 1 | |
| elif "Looking for best node for task" in message: | |
| ptm = PER_TASK_DEMAND_RE.search(message) | |
| if ptm and ptm.group(1).startswith(args.task_prefix): | |
| if per_task_demand is None: | |
| per_task_demand = ptm.group(2) | |
| probed_tasks.add(ptm.group(1)) | |
| lines_for_our_pg += 1 | |
| elif "No solution found for" in message: | |
| nm = NO_SOLUTION_RE.search(message) | |
| if nm and nm.group(3) == args.pg: | |
| if gang_size is None: | |
| gang_size = int(nm.group(2)) | |
| if current is not None and current_scenario is not None: | |
| current_scenario.no_solution = True | |
| lines_for_our_pg += 1 | |
| # Per-node attempt | |
| elif "Trying to solve scenario with potantial" in message: | |
| if current is None or current_scenario is None: | |
| continue | |
| bn = BY_POD_NODE_RE.search(message) | |
| if not bn: | |
| continue | |
| current_node_attempt = NodeAttempt(node=bn.group(1)) | |
| current_scenario.node_attempts.append(current_node_attempt) | |
| elif "Statement evicted task:" in message: | |
| if current_node_attempt is None: | |
| continue | |
| em = EVICTED_RE.search(message) | |
| if em: | |
| current_node_attempt.evictions.append((em.group(1), em.group(2))) | |
| elif "Failed statement allocate for task" in message: | |
| # rollback signal | |
| if current_node_attempt is not None and args.task_prefix in message: | |
| current_node_attempt.rolled_back = True | |
| lines_for_our_pg += 1 | |
| # Predicate failures for our task | |
| elif "Predicates failed for task" in message: | |
| pm = PRED_FAIL_TASK_RE.search(message) | |
| if not pm: | |
| continue | |
| task, node, reason = pm.group(1), pm.group(2), pm.group(3).strip() | |
| if not task.startswith(args.task_prefix): | |
| continue | |
| predicate_failures[reason] += 1 | |
| if len(predicate_examples[reason]) < 5: | |
| short_node = node.split(".")[0] | |
| predicate_examples[reason].add(short_node) | |
| lines_for_our_pg += 1 | |
| # Flush any open action without an explicit close | |
| if current is not None: | |
| actions.append(current) | |
| return { | |
| "actions": actions, | |
| "predicate_failures": predicate_failures, | |
| "predicate_examples": predicate_examples, | |
| "lines_total": lines_total, | |
| "lines_matched": lines_matched, | |
| "lines_for_our_pg": lines_for_our_pg, | |
| "gang_size": gang_size, | |
| "per_task_demand": per_task_demand, | |
| "probed_tasks": probed_tasks, | |
| } | |
| def short_task(t: str) -> str: | |
| # default/<name> -> just last 24 chars of <name> | |
| name = t.split("/", 1)[1] if "/" in t else t | |
| return name if len(name) <= 40 else "..." + name[-37:] | |
| def short_node(n: str) -> str: | |
| return n.split(".")[0] | |
| def rack_of(node: str) -> str: | |
| # <cluster>-r06-01.<domain> -> r06 | |
| n = node.split(".")[0] | |
| parts = n.split("-") | |
| for p in parts: | |
| if p.startswith("r") and p[1:].isdigit(): | |
| return p | |
| return n | |
| def md_table(rows, headers) -> str: | |
| return tabulate(rows, headers=headers, tablefmt="github") | |
| def gpus_in_demand(demand: str) -> Optional[int]: | |
| """Extract GPU count from either old human-readable format or new ResourceVector slice. | |
| Old format: 'CPU: 1540 (cores), memory: 0 (GB), Gpus: 44' | |
| New format: '[1.54e+06 0 44 11 0 0 11000 11000 0 11000 0 11000 0]' | |
| index 0=CPU(milli), 1=memory, 2=GPU, 3=Pods, rest=DRA/network | |
| """ | |
| if not demand: | |
| return None | |
| m = re.search(r"GPU(?:s)?:\s*([\d.]+)", demand) | |
| if m: | |
| return int(float(m.group(1))) | |
| # New ResourceVector slice format: take index 2 | |
| vec_match = re.match(r"\s*\[([^\]]+)\]\s*$", demand) | |
| if vec_match: | |
| parts = vec_match.group(1).split() | |
| if len(parts) >= 3: | |
| try: | |
| return int(float(parts[2])) | |
| except ValueError: | |
| pass | |
| return None | |
| def render_headline(args, result, per_node, rack_roll): | |
| """The 'why didn't it just preempt 11 nodes from r06' paragraph.""" | |
| gang = result["gang_size"] | |
| per_task = result["per_task_demand"] or "" | |
| per_task_gpus = gpus_in_demand(per_task) | |
| probed_n = len(result["probed_tasks"]) | |
| # Pick the largest rack the solver visited | |
| if rack_roll: | |
| biggest = max(rack_roll.items(), key=lambda kv: len(kv[1]["nodes"])) | |
| rack_name = biggest[0] | |
| rack_nodes = len(biggest[1]["nodes"]) | |
| rack_max_victims = biggest[1]["max_victims"] | |
| else: | |
| rack_name, rack_nodes, rack_max_victims = "n/a", 0, 0 | |
| # Find max victims evicted in any single per-node attempt across all racks | |
| max_victims_anywhere = max((st["max_victims"] for st in per_node.values()), default=0) | |
| needed_nodes = None | |
| rack_drain_gpus = None | |
| if gang and per_task_gpus: | |
| # 1 victim per GPU on each 4-GPU node, so each node yields 4 GPUs | |
| needed_nodes = gang # one pod per node (single-GPU victims, 4 GPUs per pod = whole node) | |
| rack_drain_gpus = rack_nodes * rack_max_victims | |
| lines = [ | |
| "## The puzzle: r06 alone had enough capacity — why didn't preempt take it?", | |
| "", | |
| ] | |
| summary_rows = [ | |
| ["Gang size (pods)", str(gang) if gang else "?"], | |
| ["Per-pod demand", per_task or "?"], | |
| ["Distinct pods of the gang ever probed by the solver", str(probed_n)], | |
| ["Max victims evicted in any single per-node statement", str(max_victims_anywhere)], | |
| [f"Largest rack visited (`{rack_name}`)", | |
| f"{rack_nodes} distinct nodes attempted, max {rack_max_victims} victims/node"], | |
| [f"`{rack_name}` capacity if fully drained", | |
| f"{rack_drain_gpus} GPUs ({rack_nodes} nodes × {rack_max_victims} GPUs)" if rack_drain_gpus else "?"], | |
| ["Gang demand", | |
| f"{gang * per_task_gpus} GPUs ({gang} pods × {per_task_gpus} GPUs)" | |
| if gang and per_task_gpus else "?"], | |
| ] | |
| lines.append(md_table(summary_rows, ["", ""])) | |
| lines.append("") | |
| if gang and per_task_gpus and probed_n == 1 and max_victims_anywhere <= per_task_gpus: | |
| lines += [ | |
| "### What the log actually shows", | |
| "", | |
| f"- The gang has **{gang} pods** but the solver only ever probed " | |
| f"**1 of them** (`{sorted(result['probed_tasks'])[0]}`) as the pending task — " | |
| "all 189 preempt scenarios use the same single pending task.", | |
| "", | |
| f"- Every per-node statement attempted to free at most **{max_victims_anywhere} " | |
| f"victims (= one 4-GPU node's occupants)**, then immediately tried to fit that " | |
| "single pending pod and rolled back when the gang-level placement still failed.", | |
| "", | |
| f"- Rack `{rack_name}` alone has **{rack_nodes} × {rack_max_victims} = " | |
| f"{rack_drain_gpus} GPUs** that could be freed — comfortably more than the " | |
| f"gang's **{gang * per_task_gpus} GPU** total demand.", | |
| "", | |
| "### Why it never combined 11 nodes into one statement", | |
| "", | |
| "The by-pod solver evaluates one pending task at a time. To satisfy a " | |
| "topology-required gang on a 100%-packed clique it would need to:", | |
| "", | |
| f"1. Open a **single statement** that evicts victims across **~{gang} different " | |
| f"nodes** in `{rack_name}`,", | |
| "2. Pipeline all 11 gang pods onto those freshly-freed nodes,", | |
| "3. Verify gang-level topology fit and commit.", | |
| "", | |
| "Instead, each statement covers **one node** (one pending pod), so the most it " | |
| "can free in a single committed statement is one node's worth of GPUs. The " | |
| "remaining 10 pods of the gang have no feasible node in the statement's view, " | |
| "the topology / gang-fit check rejects it, and the statement rolls back. " | |
| "Repeat ×189 across every candidate node and the action gives up with " | |
| "`Didn't find a preemption strategy`.", | |
| "", | |
| "This is exactly the NP-complete solver-coverage limitation noted in issue " | |
| "#1591 and the motivation for the solver refactor in PR #1574.", | |
| ] | |
| elif gang and probed_n < gang: | |
| lines += [ | |
| "### Observed", | |
| "", | |
| f"Gang of {gang} pods; only {probed_n} ever probed in solver scenarios.", | |
| ] | |
| return lines | |
| def aggregate_per_node(actions): | |
| per_node = defaultdict(lambda: {"visits": 0, "max_victims": 0, "rollbacks": 0, "rack": ""}) | |
| for a in actions: | |
| for s in a.scenarios: | |
| for na in s.node_attempts: | |
| key = short_node(na.node) | |
| per_node[key]["visits"] += 1 | |
| per_node[key]["max_victims"] = max(per_node[key]["max_victims"], len(na.evictions)) | |
| if na.rolled_back: | |
| per_node[key]["rollbacks"] += 1 | |
| per_node[key]["rack"] = rack_of(na.node) | |
| return per_node | |
| def render(args, result) -> str: | |
| per_node = aggregate_per_node(result["actions"]) | |
| rack_roll: dict = defaultdict(lambda: {"nodes": set(), "visits": 0, "rollbacks": 0, "max_victims": 0}) | |
| for node, st in per_node.items(): | |
| rack_roll[st["rack"]]["nodes"].add(node) | |
| rack_roll[st["rack"]]["visits"] += st["visits"] | |
| rack_roll[st["rack"]]["rollbacks"] += st["rollbacks"] | |
| rack_roll[st["rack"]]["max_victims"] = max(rack_roll[st["rack"]]["max_victims"], st["max_victims"]) | |
| out = [ | |
| f"# Why `{args.pg}` wasn't scheduled", | |
| "", | |
| f"_Generated from `{args.log}`._", | |
| "", | |
| ] | |
| # Headline puzzle box | |
| out += render_headline(args, result, per_node, rack_roll) | |
| out += ["", "## TL;DR — action outcomes", ""] | |
| tldr_rows = [ | |
| [a.name, "yes", a.demand or "-", a.outcome, a.started_at, a.closed_at or "-"] | |
| for a in result["actions"] | |
| ] or [["(no action ever attempted for this PG)", "", "", "", "", ""]] | |
| out.append(md_table(tldr_rows, ["Action", "Attempted", "Demand", "Outcome", "Started", "Closed"])) | |
| out += [ | |
| "", | |
| "## Scenarios summary", | |
| "", | |
| 'Each scenario is one `job_solver.go:109` "Trying to solve scenario" entry. ' | |
| "The by-pod solver greedily grows the victim set one task at a time, so many " | |
| "scenarios per action is expected.", | |
| "", | |
| ] | |
| scen_rows = [] | |
| for a in result["actions"]: | |
| if not a.scenarios: | |
| scen_rows.append([a.name, 0, "-", "-", "n/a"]) | |
| continue | |
| counts = [len(s.potential_victims.split()) if s.potential_victims else 0 for s in a.scenarios] | |
| solved = "no (all no_solution / no terminal record)" | |
| if any(not s.no_solution and s.node_attempts for s in a.scenarios) and not all(s.no_solution for s in a.scenarios): | |
| solved = "see per-node section" | |
| scen_rows.append([a.name, len(a.scenarios), min(counts), max(counts), solved]) | |
| out.append(md_table(scen_rows, ["Action", "Scenarios", "Min victims", "Max victims", "Any solved?"])) | |
| out += [ | |
| "", | |
| "## Per-node solver attempts (across all actions)", | |
| "", | |
| "Each unique node visited by `by_pod_solver`: how many times the solver tried it, " | |
| "max victim-set size, whether any attempt was rolled back.", | |
| "", | |
| ] | |
| if not per_node: | |
| out.append("(no per-node solver attempts captured)") | |
| else: | |
| out.append(f"Total unique nodes touched: **{len(per_node)}**.") | |
| out.append("") | |
| sorted_nodes = sorted(per_node.items(), key=lambda kv: (-kv[1]["visits"], kv[0])) | |
| rows = [ | |
| [node, st["rack"], st["visits"], st["max_victims"], st["rollbacks"]] | |
| for node, st in sorted_nodes | |
| ] | |
| out.append(md_table(rows, ["Node", "Rack", "Visits", "Max victims evicted", "Rolled-back attempts"])) | |
| if rack_roll: | |
| out += [ | |
| "", | |
| "## Rack roll-up", | |
| "", | |
| "Aggregated by `r<NN>` rack prefix. Confirms which topology domains the solver " | |
| "actually attempted evictions in.", | |
| "", | |
| ] | |
| rows = [ | |
| [rack, len(st["nodes"]), st["visits"], st["rollbacks"]] | |
| for rack, st in sorted(rack_roll.items(), key=lambda kv: (-len(kv[1]["nodes"]), kv[0])) | |
| ] | |
| out.append(md_table(rows, ["Rack", "Distinct nodes attempted", "Total solver visits", "Rolled-back attempts"])) | |
| out += [ | |
| "", | |
| "## Predicate failure breakdown (for tasks matching --task-prefix)", | |
| "", | |
| ] | |
| if not result["predicate_failures"]: | |
| out.append("(no `framework/session.go:223` lines matched the task prefix)") | |
| else: | |
| rows = [ | |
| [reason, count, ", ".join(sorted(result["predicate_examples"][reason])[:5])] | |
| for reason, count in result["predicate_failures"].most_common(args.top_n) | |
| ] | |
| out.append(md_table(rows, ["Reason", "Node-fail count", "Example nodes"])) | |
| preempt_actions = [a for a in result["actions"] if a.name == "preempt" and a.scenarios] | |
| if preempt_actions: | |
| a = preempt_actions[0] | |
| last = a.scenarios[-1] | |
| pv = last.potential_victims or "(none)" | |
| if len(pv) > 400: | |
| pv = pv[:397] + "..." | |
| n_last = len(last.potential_victims.split()) if last.potential_victims else 0 | |
| out += [ | |
| "", | |
| "## Preempt scenario growth (first / last)", | |
| "", | |
| "The by-pod solver grows the potential-victim list across scenarios; " | |
| "here are the first and last to show start vs. end state.", | |
| "", | |
| "**First scenario:**", | |
| f"- Pending: `{a.scenarios[0].pending}`", | |
| f"- Potential victims: `{a.scenarios[0].potential_victims or '(none)'}`", | |
| "", | |
| "**Last scenario:**", | |
| f"- Pending: `{last.pending}`", | |
| f"- Potential victims ({n_last}): `{pv}`", | |
| f"- Result: `{'no solution' if last.no_solution else '(no terminal record)'}`", | |
| ] | |
| out += [ | |
| "", | |
| "## Parse stats", | |
| "", | |
| f"- Lines total: **{result['lines_total']:,}**", | |
| f"- Lines matching one of the {len(MESSAGES)} message anchors: **{result['lines_matched']:,}**", | |
| f"- Lines attributed to our PG / task: **{result['lines_for_our_pg']:,}**", | |
| f"- Actions captured: **{len(result['actions'])}**", | |
| f"- Message anchors used: `{', '.join(MESSAGES)}`", | |
| "", | |
| ] | |
| return "\n".join(out) | |
| def print_console_summary(args, result): | |
| """Live rich-formatted summary printed to stderr while we work.""" | |
| con = Console(stderr=True) | |
| t = Table(title=f"Action outcomes — {args.pg}", show_lines=False) | |
| for col in ("Action", "Demand", "Outcome"): | |
| t.add_column(col) | |
| for a in result["actions"]: | |
| t.add_row(a.name, a.demand or "-", a.outcome) | |
| con.print(t) | |
| per_node = aggregate_per_node(result["actions"]) | |
| rack_roll = defaultdict(lambda: {"nodes": set(), "visits": 0}) | |
| for node, st in per_node.items(): | |
| rack_roll[st["rack"]]["nodes"].add(node) | |
| rack_roll[st["rack"]]["visits"] += st["visits"] | |
| rt = Table(title="Solver activity per rack") | |
| for col in ("Rack", "Nodes", "Visits"): | |
| rt.add_column(col, justify="right") | |
| for rack, st in sorted(rack_roll.items(), key=lambda kv: (-len(kv[1]["nodes"]), kv[0])): | |
| rt.add_row(rack, str(len(st["nodes"])), str(st["visits"])) | |
| con.print(rt) | |
| def main(): | |
| args = parse_args() | |
| if not os.path.exists(args.log): | |
| print(f"log not found: {args.log}", file=sys.stderr) | |
| return 2 | |
| print(f"Parsing {args.log} ...", file=sys.stderr) | |
| result = parse(args) | |
| print( | |
| f" lines_total={result['lines_total']:,} " | |
| f"matched={result['lines_matched']:,} " | |
| f"for_our_pg={result['lines_for_our_pg']:,} " | |
| f"actions={len(result['actions'])}", | |
| file=sys.stderr, | |
| ) | |
| print_console_summary(args, result) | |
| md = render(args, result) | |
| os.makedirs(os.path.dirname(args.out), exist_ok=True) | |
| with open(args.out, "w", encoding="utf-8") as fh: | |
| fh.write(md + "\n") | |
| print(f"Wrote {args.out}", file=sys.stderr) | |
| return 0 | |
| if __name__ == "__main__": | |
| sys.exit(main()) |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment