diff --git a/.agents/skills/live-logs/SKILL.md b/.agents/skills/live-logs/SKILL.md new file mode 100644 index 0000000000..567e624c50 --- /dev/null +++ b/.agents/skills/live-logs/SKILL.md @@ -0,0 +1,65 @@ +--- +name: live-logs +description: Open a live, per-node log viewer in the browser for a running multi-node sweep job, given a PR number, PR URL, workflow run URL or job URL. Resolves the GitHub job to its Slurm job on the cluster, then streams every worker (prefill / decode / agg), srtctl, Dynamo frontend, Mooncake / etcd service and aiperf log into resizable panes with phase tracking, error highlighting and node memory. Use when the user asks to watch, stream, tail or open the logs of a running sweep, canary or AgentX job. +--- + +# Live logs for a sweep job + +GitHub Actions buffers multi-node job output until the job ends. This skill streams the +srt-slurm logs straight from the cluster to a local page instead. + +## Run it + +```bash +python3 .agents/skills/live-logs/scripts/live_logs.py +``` + +- **PR number or PR URL:** uses the in-progress `Run Sweep` run on the PR head. If none is running, it falls back to the latest one. +- **Run URL:** uses every in-progress job with a runner. +- **Job URL:** uses just that job. +- **Result:** one local viewer per Slurm job, for example the benchmark job and the eval job of the same run. Each opens in the browser, and each URL is printed. +- **Stop:** `live_logs.py --stop` stops every viewer the script started. +- **Manual mode:** when you already know the job, pass `--host --job --logdir <.../outputs//logs>` instead of a target. +- **Faster start on very large logs:** pass `--history 4000` to fetch only the last 4000 lines of each file. +- **Requirements:** only the Python standard library, plus `gh` (authenticated) and `ssh` to the cluster login node. The viewer binds to `127.0.0.1` only. + +Job resolution relies on srtctl naming the Slurm job after the GitHub runner (for example `b300-dsxe_01`). Jobs named after the worker instead (for example `worker-2` on h200-dgxc) are found through the runner's `gharunnerNN` work directory in `squeue`. `squeue -n ` finds the job, and `sacct` covers a job that has just finished. The logs directory is taken from the job's `StdOut` (srtctl writes `sweep_.log` there), falling back to `/outputs//logs`. Single-node jobs have no such Slurm job, and the script says so. + +## Cluster login hosts: never put them in the repo + +Login addresses are infra details that stay out of this repo; see `$debug-runs`. The script reads the host for each cluster (the runner-name prefix, for example `b300-dsxe`) from either: + +- `~/.config/infx-live-logs/clusters.json`, for example `{"b300-dsxe": {"host": ""}}` +- or `INFX_LIVE_LOGS_HOST_B300_DSXE=`, with the cluster name upper-cased and `-` turned into `_`. + +If the host is missing, the script names the cluster and exits. Get the login address from the InferenceX Clusters canvas; ask the user for the link, or ask them for the SSH target. Add it to the local config, and never commit it. + +## What the page shows + +- **Header:** + - SSH connection state and Slurm state (state, elapsed time, nodes). + - Per-node free host memory and CPU load, refreshed every 20 s. The memory bar turns red below 200 GiB free, which is useful for Mooncake / EFA memory-registration OOMs. +- **Rows:** + 1. Engine workers (prefill / decode / agg). + 2. srtctl sweep, Dynamo frontend, Mooncake master, etcd. + 3. aiperf / benchmark logs. These appear on their own when the benchmark starts, because new files are discovered every 30 s. +- **Each pane header:** + - The latest phase: weights loaded → KV cache sized → EFA devices up → Mooncake segment registration → autotune → ✅ healthy. + - A count of real errors, excluding known noise such as NCCL `ibv_query_port_speed`, the pip resolver notice and the node-exporter TaskProlog message. + - Buttons: jump to start / end, ⬇ download that full log (fetched fresh from the cluster, whatever the display settings), copy, and maximize (Esc restores). +- **Controls:** + - Regex filter and "errors only". + - "Hide Mooncake metrics noise", on by default. It hides the client metric report blocks, throughput and latency summaries, and dynamo HTTP 200 spam. + - `show`: last 4000 lines (the default) or the entire log. The server keeps the whole file, so switching is instant. + - **Download all logs:** a zip of every log file (`*.out`, `*.log`, `*.txt`) under the job's `logs/` directory, freshly `tar`-ed on the cluster. Files are always complete, whatever the `show`, filter or `--history` settings. Tachometer metrics are left out. + - Theme and reset layout. +- **Resizing:** + - Drag the gutters between panes and rows; double-click a gutter to even out its two panes. + - Sizes are saved per layout shape. + - Scrolling up in a pane pauses follow mode (amber outline); scrolling back to the bottom resumes it. + +## Notes + +- One SSH `tail -F` runs per log file, so a huge log doesn't hold up the others. After an SSH drop, that file is re-read from the start, so nothing is lost. +- The server keeps the full logs, but on connect it sends only the last 4000 matching lines of each; noise and "errors only" are also filtered on the server. This keeps the page responsive with logs of hundreds of MB. "Entire log" asks for confirmation first. +- The viewer only reads logs. It never touches processes or files on shared hosts. diff --git a/.agents/skills/live-logs/agents/openai.yaml b/.agents/skills/live-logs/agents/openai.yaml new file mode 100644 index 0000000000..3f858fda1b --- /dev/null +++ b/.agents/skills/live-logs/agents/openai.yaml @@ -0,0 +1,4 @@ +interface: + display_name: "Live Logs" + short_description: "Stream a sweep job's per-node logs to the browser" + default_prompt: "Use $live-logs to open the live per-node logs for this PR or workflow run." diff --git a/.agents/skills/live-logs/scripts/index.html b/.agents/skills/live-logs/scripts/index.html new file mode 100644 index 0000000000..b42a69a05f --- /dev/null +++ b/.agents/skills/live-logs/scripts/index.html @@ -0,0 +1,254 @@ + + + + + +Live Job Logs + + + +
+

__TITLE__

+ ssh: … + slurm: … + +
+ + + + + + + +
+
+
+ + + diff --git a/.agents/skills/live-logs/scripts/live_logs.py b/.agents/skills/live-logs/scripts/live_logs.py new file mode 100755 index 0000000000..cf961b1b22 --- /dev/null +++ b/.agents/skills/live-logs/scripts/live_logs.py @@ -0,0 +1,198 @@ +#!/usr/bin/env python3 +"""Resolve a PR / workflow run / job to its live Slurm jobs and open a log viewer per job. + + live_logs.py 3523 + live_logs.py https://github.com/SemiAnalysisAI/InferenceX/pull/3523 + live_logs.py https://github.com/SemiAnalysisAI/InferenceX/actions/runs/36345723679 + live_logs.py https://github.com/SemiAnalysisAI/InferenceX/actions/runs/36345723679/job/108694606739 + live_logs.py --stop # stop every viewer this script started + +Cluster login hosts are NOT stored in this repo. They are read from +~/.config/infx-live-logs/clusters.json, e.g. {"b300-dsxe": {"host": "my-b300-login-alias"}}, +or from INFX_LIVE_LOGS_HOST_ (cluster upper-cased, '-' -> '_'). +""" +import argparse +import json +import os +import re +import signal +import socket +import subprocess +import sys +import time +import webbrowser + +REPO = os.environ.get("INFX_LIVE_LOGS_REPO", "SemiAnalysisAI/InferenceX") +CONFIG = os.path.expanduser("~/.config/infx-live-logs/clusters.json") +STATE_DIR = os.path.expanduser("~/.cache/infx-live-logs") +HERE = os.path.dirname(os.path.abspath(__file__)) +SSH = ["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=20"] +BANNER = ("Access to this system", "All activity is logged") + + +def gh(path): + out = subprocess.run(["gh", "api", path], capture_output=True, text=True, check=False) + if out.returncode: + sys.exit(f"gh api {path} failed: {out.stderr.strip()}") + return json.loads(out.stdout) + + +def resolve_jobs(target): + """Return [(run_id, job)] for in-progress, non-setup jobs of the target.""" + m = re.search(r"/actions/runs/(\d+)/job/(\d+)", target) + if m: + return [(m.group(1), gh(f"repos/{REPO}/actions/jobs/{m.group(2)}"))] + m = re.search(r"/actions/runs/(\d+)", target) + if m: + run_ids = [m.group(1)] + else: + m = re.search(r"(?:/pull/)?(\d+)\s*$", target.strip()) + if not m: + sys.exit(f"cannot parse {target!r}: pass a PR number, PR URL, run URL or job URL") + pr = gh(f"repos/{REPO}/pulls/{m.group(1)}") + runs = gh(f"repos/{REPO}/actions/runs?head_sha={pr['head']['sha']}&per_page=50")["workflow_runs"] + runs = [r for r in runs if r["name"].startswith("Run Sweep") and r["status"] != "completed"] or \ + [r for r in runs if r["name"].startswith("Run Sweep")][:1] + if not runs: + sys.exit(f"PR #{m.group(1)} has no Run Sweep run on head {pr['head']['sha'][:8]}") + run_ids = [str(r["id"]) for r in runs] + jobs = [] + for rid in run_ids: + for j in gh(f"repos/{REPO}/actions/runs/{rid}/jobs?per_page=100")["jobs"]: + if j["name"] != "setup" and j["status"] == "in_progress" and j.get("runner_name"): + jobs.append((rid, j)) + return jobs + + +def cluster_host(cluster): + env = os.environ.get("INFX_LIVE_LOGS_HOST_" + re.sub(r"[^A-Z0-9]", "_", cluster.upper())) + if env: + return env + try: + with open(CONFIG) as f: + return json.load(f)[cluster]["host"] + except (OSError, KeyError, ValueError): + return None + + +def remote(host, script): + out = subprocess.run(SSH + [host, script], capture_output=True, text=True, timeout=90, check=False).stdout + return [l for l in out.splitlines() if l and not l.startswith(BANNER)] + + +def slurm_job_for_runner(host, runner): + """srtctl names the Slurm job after the GitHub runner (e.g. b300-dsxe_01).""" + lines = remote(host, f"squeue -h -n {runner} -o '%i %T' | sort -n | tail -1; " + f"sacct -n -X -P --name {runner} -S now-1days -o JobID | grep -E '^[0-9]+$' | sort -n | tail -1") + ids = [l.split()[0] for l in lines if l.split()[0].isdigit()] + if not ids: + # Some launchers name the job after the worker (e.g. worker-2); fall back to the runner's + # gharunnerNN work directory. + suffix = runner.rsplit("_", 1)[-1] + lines = remote(host, f"squeue -h -o '%i %Z' | grep '/gharunner{suffix}/' | sort -n | tail -1") + ids = [l.split()[0] for l in lines if l.split()[0].isdigit()] + if not ids: + return None, None + job = ids[0] + # srtctl points the job's StdOut at //logs/sweep_.log, which is the most + # reliable way to find the logs (WorkDir can be a checkout next to outputs/, not its parent). + info = remote(host, f"scontrol show job {job} 2>/dev/null | grep -oE '(WorkDir|StdOut)=[^ ]+'") + kv = dict(x.split("=", 1) for x in info if "=" in x) + out = kv.get("StdOut", "") + if os.path.basename(out).startswith("sweep_"): + return job, os.path.dirname(out) + wd = kv.get("WorkDir") or next(iter(remote(host, f"sacct -n -X -P -j {job} -o WorkDir")), None) + return job, (f"{wd}/outputs/{job}/logs" if wd else None) + + +def free_port(start=8765): + for p in range(start, start + 200): + with socket.socket() as s: + if s.connect_ex(("127.0.0.1", p)): + return p + sys.exit("no free local port") + + +def start_viewer(host, job, logdir, title, history="all"): + os.makedirs(STATE_DIR, exist_ok=True) + port = free_port() + with open(os.path.join(STATE_DIR, f"viewer-{job}.log"), "w") as log: + p = subprocess.Popen([sys.executable, os.path.join(HERE, "server.py"), "--host", host, "--job", job, + "--logdir", logdir, "--title", title, "--port", str(port), "--history", history], + stdout=log, stderr=subprocess.STDOUT, start_new_session=True) + with open(os.path.join(STATE_DIR, "pids"), "a") as f: + f.write(f"{p.pid} {port} {job}\n") + url = f"http://127.0.0.1:{port}/" + for _ in range(50): # wait until the server accepts connections + with socket.socket() as s: + if not s.connect_ex(("127.0.0.1", port)): + break + time.sleep(0.1) + return url + + +def stop_all(): + path = os.path.join(STATE_DIR, "pids") + if not os.path.exists(path): + print("no viewers recorded") + return + with open(path) as f: + lines = f.read().split("\n") + for line in filter(None, lines): + pid, port, job = line.split() + try: + os.killpg(int(pid), signal.SIGTERM) + print(f"stopped viewer for Slurm {job} (port {port})") + except ProcessLookupError: + pass + os.remove(path) + + +def main(): + ap = argparse.ArgumentParser(description=__doc__, formatter_class=argparse.RawDescriptionHelpFormatter) + ap.add_argument("target", nargs="?", help="PR number, PR URL, workflow run URL or job URL") + ap.add_argument("--stop", action="store_true", help="stop every viewer started by this script") + ap.add_argument("--no-open", action="store_true", help="print URLs without opening a browser") + ap.add_argument("--host", help="manual mode: ssh target of the login node (with --job and --logdir)") + ap.add_argument("--job", help="manual mode: Slurm job id") + ap.add_argument("--logdir", help="manual mode: the job's logs directory on the cluster") + ap.add_argument("--history", default="all", help="'all' (default) or N: only fetch the last N lines of each log (faster on huge logs)") + args = ap.parse_args() + if args.stop: + return stop_all() + if args.host and args.job and args.logdir: + url = start_viewer(args.host, args.job, args.logdir, f"Slurm {args.job}", args.history) + print(url) + if not args.no_open: + webbrowser.open(url) + return None + if not args.target: + ap.error("target is required (or --host/--job/--logdir for manual mode)") + + jobs = resolve_jobs(args.target) + if not jobs: + sys.exit("no in-progress jobs with a runner found (queued, finished, or setup-only)") + missing = set() + for rid, j in jobs: + runner = j["runner_name"] + cluster = runner.rsplit("_", 1)[0] + host = cluster_host(cluster) + if not host: + missing.add(cluster) + continue + job, logdir = slurm_job_for_runner(host, runner) + if not job or not logdir: + print(f"- {j['name'][:90]}: no Slurm job named {runner} on {cluster} (single-node or not submitted yet)") + continue + title = f"{cluster} · Slurm {job} · {j['name'].split('|')[-1].strip()[:80]}" + url = start_viewer(host, job, logdir, title, args.history) + print(f"- {j['name'][:90]}\n runner {runner} → Slurm {job}\n {url} (GitHub job {j['html_url']})") + if not args.no_open: + webbrowser.open(url) + if missing: + sys.exit(f"no login host configured for: {', '.join(sorted(missing))}. Add it to {CONFIG} " + '(e.g. {"b300-dsxe": {"host": ""}}); see SKILL.md for where to find it.') + + +if __name__ == "__main__": + main() diff --git a/.agents/skills/live-logs/scripts/server.py b/.agents/skills/live-logs/scripts/server.py new file mode 100755 index 0000000000..dbbb554105 --- /dev/null +++ b/.agents/skills/live-logs/scripts/server.py @@ -0,0 +1,351 @@ +#!/usr/bin/env python3 +"""Stream one Slurm job's srt-slurm logs to a local browser page. + +One ssh `tail -F` per discovered log file feeds an in-memory buffer; the page +subscribes over Server-Sent Events. New log files (for example aiperf output once the benchmark +starts) are discovered every 30 s and get their own pane. +""" +import argparse +import collections +import json +import os +import queue +import re +import subprocess +import tarfile +import tempfile +import threading +import time +import zipfile +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from urllib.parse import parse_qs, urlparse + +ANSI = re.compile(r"\x1b\[[0-9;]*[A-Za-z]") +SSH = ["ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=20", "-o", "ServerAliveInterval=15"] +BANNER = ("Access to this system", "All activity is logged") +# Worker, orchestrator and benchmark logs; exporters and configs are left out on purpose. +FIND_EXPR = ( + "\\( -name '*_prefill_w*.out' -o -name '*_decode_w*.out' -o -name '*_agg_w*.out' " + "-o -name 'sweep_*.log' -o -name '*_frontend_*.out' -o -name 'service_mooncake-master.out' " + "-o -name 'service_etcd.out' -o -name 'aiperf.log' -o -name 'benchmark.log' -o -name 'benchmark.out' \\)" +) + + +def page_regex(name): + """Reuse the page's NOISE / ERR regex so server- and client-side filtering agree.""" + with open(os.path.join(os.path.dirname(os.path.abspath(__file__)), "index.html")) as f: + page = f.read() + m = re.search(rf"^const {name} = /(.*)/;$", page, re.MULTILINE) + return re.compile(m.group(1)) + + +NOISE_RE = page_regex("NOISE") +ERR_RE = page_regex("ERR") + + +def wanted(line, noise, erronly): + if erronly: + return bool(ERR_RE.search(line)) + return not (noise and NOISE_RE.search(line) and not ERR_RE.search(line)) + + +def last_matching(buf, n, noise, erronly): + """Newest n lines of buf that pass the filters (n=None means all), oldest first.""" + if n is None: + return [l for l in buf if wanted(l, noise, erronly)] + out = [] + for l in reversed(buf): + if wanted(l, noise, erronly): + out.append(l) + if len(out) >= n: + break + out.reverse() + return out + + +def group_of(rel): + if re.search(r"_(prefill|decode|agg)_w\d+\.out$", rel): + return 0 + if rel.endswith(("aiperf.log", "benchmark.log", "benchmark.out")): + return 2 + return 1 + + +def label_of(rel): + base = os.path.basename(rel) + m = re.match(r"(?:.*-)?(gpu-\d+|[a-z0-9]+-\d+)_(prefill|decode|agg)_w(\d+)\.out$", base) + if m: + return f"{m.group(2).capitalize()} w{m.group(3)} · {m.group(1)}" + m = re.match(r".*-(gpu-\d+|[a-z0-9]+-\d+)_frontend_(\d+)\.out$", base) + if m: + return f"Dynamo frontend · {m.group(1)}" + if base.startswith("sweep_"): + return "srtctl sweep" + if base.startswith("service_"): + return base[len("service_"):-len(".out")] + m = re.search(r"(aiperf|benchmark)\.(log|out)$", rel) + if m: + conc = re.search(r"conc_\d+", rel) + where = conc.group(0) if conc else (rel.split("/")[0] if "/" in rel else "") + return f"{m.group(1)}.{m.group(2)}" + (f" · {where}" if where else "") + parts = rel.split("/") + return "/".join(parts[-3:]) if len(parts) > 1 else base + + +class State: + def __init__(self): + self.files = [] # [rel, label, group] + self.buffers = {} + self.totals = collections.Counter() + self.status = {"ssh": "connecting", "squeue": "", "mem": {}} + self.subscribers = [] + self.lock = threading.Lock() + + def publish(self, ev): + with self.lock: + for q in list(self.subscribers): + try: + q.put_nowait(ev) + except queue.Full: + pass + + +def remote(host, script, timeout=60): + out = subprocess.run(SSH + [host, script], capture_output=True, text=True, timeout=timeout, check=False).stdout + return [l for l in out.splitlines() if l and not l.startswith(BANNER)] + + +def tail_batch(args, st, rels): + """Follow a fixed set of files from line 1; reconnect and re-read on ssh drops.""" + by_base = {} + for r in rels: + by_base[r] = r + while True: + cmd = SSH + [args.host, f"cd {args.logdir} && exec tail -n {args.tail_from} -F " + " ".join(rels)] + st.status["ssh"] = "connected" + st.publish({"t": "status", "s": st.status}) + p = subprocess.Popen(cmd, stdout=subprocess.PIPE, stderr=subprocess.DEVNULL, text=True, errors="replace", bufsize=1) + cur = rels[0] if len(rels) == 1 else None # tail prints "==> f <==" headers only for several files + for line in p.stdout: + line = ANSI.sub("", line.rstrip("\n")) + m = re.match(r"^==> (.+) <==$", line) + if m: + cur = by_base.get(m.group(1), cur) + continue + if cur is None or not line: + continue + st.buffers[cur].append(line) + st.totals[cur] += 1 + st.publish({"t": "line", "k": cur, "l": line}) + p.wait() + with st.lock: + for r in rels: + st.buffers[r].clear() + st.totals[r] = 0 + st.publish({"t": "reset", "k": rels}) + st.status["ssh"] = f"reconnecting (exit {p.returncode})" + st.publish({"t": "status", "s": st.status}) + time.sleep(5) + + +def discover_loop(args, st): + while True: + try: + found = remote(args.host, f"cd {args.logdir} 2>/dev/null && find . -maxdepth 5 -type f {FIND_EXPR} | sed 's#^./##' | sort") + new = [r for r in found if r not in st.buffers] + if new: + with st.lock: + for r in new: + st.buffers[r] = collections.deque(maxlen=args.max_lines) + st.files.append([r, label_of(r), group_of(r)]) + st.publish({"t": "files", "f": st.files}) + # one tail per file: a multi-file tail prints each file in full before the next, + # so a huge prefill log would hold every other pane empty until it finished + for r in new: + threading.Thread(target=tail_batch, args=(args, st, [r]), daemon=True).start() + except (OSError, subprocess.SubprocessError, ValueError) as e: # flaky network: keep going + st.status["ssh"] = f"discovery error: {e}" + time.sleep(30) + + +def status_loop(args, st): + while True: + try: + script = ( + f"squeue -j {args.job} -h -o '%T %M %N'; " + f"for n in $(scontrol show hostnames $(squeue -j {args.job} -h -o %N) 2>/dev/null); do " + "echo \"NODE $n $(scontrol show node $n | grep -oE 'FreeMem=[0-9]+|RealMemory=[0-9]+|CPULoad=[0-9.]+' | tr '\\n' ' ')\"; done" + ) + out = remote(args.host, script) + sq = [l for l in out if not l.startswith("NODE ")] + st.status["squeue"] = sq[0] if sq else "job not in queue (finished?)" + mem = {} + for l in out: + if l.startswith("NODE "): + _, name, *kvs = l.split() + mem[name] = dict(x.split("=", 1) for x in kvs if "=" in x) + st.status["mem"] = mem + st.status["totals"] = dict(st.totals) + st.publish({"t": "status", "s": st.status}) + except (OSError, subprocess.SubprocessError, ValueError) as e: + st.status["squeue"] = f"status poll error: {e}" + time.sleep(20) + + +def build_zip(args): + """Pull every text log under /logs (full files, no tachometer metrics) as a gzipped tar and repack it as a zip.""" + name = f"slurm-{args.job}" + pick = (f"cd {args.logdir} && find . -type f \\( -name '*.out' -o -name '*.log' -o -name '*.txt' \\) " + "! -path './tachometer/*' | sed 's#^./##' | tar czf - -T -") + p = subprocess.Popen(SSH + [args.host, pick], + stdout=subprocess.PIPE, stderr=subprocess.DEVNULL) + fd, out = tempfile.mkstemp(prefix=name + "-", suffix=".zip") + os.close(fd) + with tarfile.open(fileobj=p.stdout, mode="r|gz") as tar, zipfile.ZipFile(out, "w", zipfile.ZIP_DEFLATED) as zf: + for m in tar: + if not m.isfile(): + continue + src = tar.extractfile(m) + with zf.open(f"{name}-logs/{m.name}", "w") as dst: + while True: + chunk = src.read(1 << 20) + if not chunk: + break + dst.write(chunk) + p.wait() + return out + + +def make_handler(args, st): + page_path = os.path.join(os.path.dirname(os.path.abspath(__file__)), "index.html") + + class H(BaseHTTPRequestHandler): + def log_message(self, *a): + pass + + def send_file(self, rel): + """Stream one full log straight from the cluster; only files the viewer discovered are allowed.""" + if rel not in st.buffers: + self.send_error(404, "unknown log file") + return + p = subprocess.Popen(SSH + [args.host, f"cat {args.logdir}/{rel}"], stdout=subprocess.PIPE, stderr=subprocess.DEVNULL) + self.send_response(200) + self.send_header("Content-Type", "text/plain; charset=utf-8") + self.send_header("Content-Disposition", f'attachment; filename="slurm-{args.job}-{os.path.basename(rel)}"') + self.end_headers() + try: + while True: + chunk = p.stdout.read(1 << 20) + if not chunk: + break + self.wfile.write(chunk) + except (BrokenPipeError, ConnectionResetError): + p.kill() + finally: + p.wait() + + def do_GET(self): + if self.path == "/": + with open(page_path) as f: + page = f.read() + body = page.replace("__TITLE__", args.title).replace("__JOB__", json.dumps(str(args.job))).encode() + self.send_response(200) + self.send_header("Content-Type", "text/html; charset=utf-8") + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + return + if self.path.startswith("/download.zip"): + try: + path = build_zip(args) + except (OSError, tarfile.TarError, subprocess.SubprocessError) as e: + self.send_error(502, f"could not fetch logs: {e}") + return + try: + size = os.path.getsize(path) + self.send_response(200) + self.send_header("Content-Type", "application/zip") + self.send_header("Content-Disposition", f'attachment; filename="slurm-{args.job}-logs.zip"') + self.send_header("Content-Length", str(size)) + self.end_headers() + with open(path, "rb") as f: + while True: + chunk = f.read(1 << 20) + if not chunk: + break + self.wfile.write(chunk) + except (BrokenPipeError, ConnectionResetError): + pass + finally: + os.remove(path) + return + url = urlparse(self.path) + if url.path == "/file": + self.send_file(parse_qs(url.query).get("f", [""])[0]) + return + if url.path != "/events": + self.send_error(404) + return + qs = parse_qs(url.query) + n_arg = qs.get("n", ["4000"])[0] + n = None if n_arg == "all" else max(1, int(n_arg)) + noise = qs.get("noise", ["1"])[0] == "1" + erronly = qs.get("err", ["0"])[0] == "1" + self.send_response(200) + self.send_header("Content-Type", "text/event-stream") + self.send_header("Cache-Control", "no-cache") + self.end_headers() + q = queue.Queue(maxsize=50000) + with st.lock: + snap = {k: last_matching(v, n, noise, erronly) for k, v in st.buffers.items()} + st.status["totals"] = dict(st.totals) + files = [list(f) for f in st.files] + st.subscribers.append(q) + try: + self.wfile.write(f"data: {json.dumps({'t': 'snapshot', 'f': files, 'b': snap, 's': st.status})}\n\n".encode()) + self.wfile.flush() + while True: + try: + batch = [q.get(timeout=15)] + while len(batch) < 1000: + try: + batch.append(q.get_nowait()) + except queue.Empty: + break + batch = [e for e in batch if e["t"] != "line" or wanted(e["l"], noise, erronly)] + if not batch: + continue + self.wfile.write(f"data: {json.dumps({'t': 'batch', 'e': batch})}\n\n".encode()) + except queue.Empty: + self.wfile.write(b": ping\n\n") + self.wfile.flush() + except (BrokenPipeError, ConnectionResetError): + pass + finally: + with st.lock: + st.subscribers.remove(q) + + return H + + +def main(): + ap = argparse.ArgumentParser(description=__doc__) + ap.add_argument("--host", required=True, help="ssh target of the cluster login node") + ap.add_argument("--job", required=True, help="Slurm job id") + ap.add_argument("--logdir", required=True, help="/outputs//logs on the cluster") + ap.add_argument("--title", default="") + ap.add_argument("--port", type=int, default=8765) + ap.add_argument("--max-lines", type=int, default=1_000_000, help="lines kept per file") + ap.add_argument("--history", default="all", help="'all' (read every file from line 1) or N (only the last N lines of each file)") + args = ap.parse_args() + args.title = args.title or f"Slurm {args.job}" + args.tail_from = "+1" if args.history == "all" else str(int(args.history)) + st = State() + threading.Thread(target=discover_loop, args=(args, st), daemon=True).start() + threading.Thread(target=status_loop, args=(args, st), daemon=True).start() + print(f"http://127.0.0.1:{args.port}", flush=True) + ThreadingHTTPServer(("127.0.0.1", args.port), make_handler(args, st)).serve_forever() + + +if __name__ == "__main__": + main() diff --git a/.claude/skills/live-logs b/.claude/skills/live-logs new file mode 120000 index 0000000000..e4d73f168e --- /dev/null +++ b/.claude/skills/live-logs @@ -0,0 +1 @@ +../../.agents/skills/live-logs \ No newline at end of file