OMNIML-5128 Capture Docker experiment id (#1840)

## Summary
- Capture nemo-run experiment_id for Docker submit_job by redirecting
detached launcher output to a side-channel log and tailing it briefly.
- Reuse the launcher output parser for both Docker and Slurm submit
paths.
- Document Docker PID + experiment_id behavior and ignore the MCP local
uv.lock.

## Jira
OMNIML-5128

## Validation
- uv run --project tools/mcp pytest tools/mcp/tests -q


<!-- This is an auto-generated comment: release notes by coderabbit.ai
-->
## Summary by CodeRabbit

* **Bug Fixes**
* Improved Docker job submission to capture `experiment_id` from
launcher output during a brief tail window, while always returning the
detached `pid` (and `stdout_log` diagnostics if capture times out).
* Made Slurm identifier extraction more robust across varying launcher
output formats.
* **Documentation**
* Updated `submit_job` docs/module descriptions to clarify Docker
PID/`experiment_id` and `stdout_log` behavior on timeout.
* **Tests**
* Added/expanded tests for Docker success, timeout/no-id scenarios,
parsing edge cases, and log-side-channel failure handling.
* **Chores**
  * Updated local tooling ignore rules for `.venv/` and `uv.lock`.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->

Signed-off-by: Chenhan Yu <chenhany@nvidia.com>
This commit is contained in:
Chenhan D. Yu
2026-06-29 12:19:07 -07:00
committed by GitHub
parent c5e716701d
commit 248cbf2fd2
6 changed files with 283 additions and 76 deletions
+5
View File
@@ -0,0 +1,5 @@
# Virtual environment
.venv/
# uv lock (generated, not portable)
uv.lock
+2 -3
View File
@@ -21,7 +21,7 @@ Mode is determined by which args you pass, not by which tool you call. One tool,
|---|---|
| `list_examples` | Enumerate bundled launcher YAMLs under `tools/launcher/examples/` with model + description metadata extracted from each YAML. Discovery primitive — call this first when you don't know which YAML to launch. |
| `verify_setup(executor, ...)` | Fail-fast probe for the named executor. Docker: `docker info` (daemon up) + `docker info --format` runtime-registry check (looks for `"nvidia"` runtime registered by the NVIDIA Container Toolkit — no image pull, daemon-fast). Slurm: `ssh -o BatchMode=yes -o ConnectTimeout=5` to the cluster login node. Returns structured failure on auth / network / daemon issues — no exception. |
| `submit_job(yaml_path, hf_local? \| cluster_host?, ..., dry_run?, source_ref?, source_repo?)` | Submit a launcher YAML. Mode resolved from mutually-exclusive args. Before launching, materializes a managed Model-Optimizer checkout at `source_ref` (branch, tag, or SHA; default `main`) and initializes recursive submodules, then runs that checkout's launcher. Returns `experiment_id` (Slurm) or PID (Docker) immediately; the actual job runs detached. Auto-runs `verify_setup` first by default (skippable). **Pass `dry_run=True`** to validate the YAML via `launch.py --dryrun --yes` without contacting the cluster / spawning a container / running sbatch — returns `{ok, dry_run: True, validated: bool, diagnostic?, exit_code, stdout_tail, stderr_tail, argv, source_sha, source_root}` instead of `experiment_id`. Used by verify-task workflow stages (deployment_support, hidden_state_dump_support, mlm_eval, ...). |
| `submit_job(yaml_path, hf_local? \| cluster_host?, ..., dry_run?, source_ref?, source_repo?)` | Submit a launcher YAML. Mode resolved from mutually-exclusive args. Before launching, materializes a managed Model-Optimizer checkout at `source_ref` (branch, tag, or SHA; default `main`) and initializes recursive submodules, then runs that checkout's launcher. Returns `experiment_id` (Slurm) or PID plus captured `experiment_id` when available (Docker) immediately; if Docker's short tail times out, it still returns PID with `experiment_id=None` and a persistent `stdout_log` under `$NEMORUN_HOME/.modelopt-mcp/docker-submit-logs/` for diagnostics. The actual job runs detached. Auto-runs `verify_setup` first by default (skippable). **Pass `dry_run=True`** to validate the YAML via `launch.py --dryrun --yes` without contacting the cluster / spawning a container / running sbatch — returns `{ok, dry_run: True, validated: bool, diagnostic?, exit_code, stdout_tail, stderr_tail, argv, source_sha, source_root}` instead of `experiment_id`. Used by verify-task workflow stages (deployment_support, hidden_state_dump_support, mlm_eval, ...). |
| `job_status(experiment_id)` | Filesystem-based status from nemo_run's experiment dir (`_DONE`, `status_*.out`). Returns `done` / `failed` / `running` plus per-task statuses. No in-memory registry; survives MCP server restarts. |
| `job_logs(experiment_id, task?, tail?)` | Read `log_<task>.out` from the experiment dir. Per-task filtering + optional tail to truncate. |
| `wait_for_experiment(experiment_id, timeout_sec?, poll_interval_sec?)` | Block until `job_status` returns `done` / `failed`, or until `timeout_sec` elapses. Single tool call replaces the agent's `while True: status; sleep` loop — saves tool-call turns and avoids overshooting the poll interval. Returns the final status plus `waited_seconds`. |
@@ -192,8 +192,7 @@ Tracked under [OMNIML-5123](https://jirasw.nvidia.com/browse/OMNIML-5123) (Epic)
**Phase 1.5 — shipped in this PR:** `wait_for_experiment`, `provision_passwordless_ssh_dry_run`, `read_cluster_artifact`, `open_draft_pr`. Anchors: [OMNIML-5128](https://jirasw.nvidia.com/browse/OMNIML-5128) (partial: the three high-leverage tools), [OMNIML-5132](https://jirasw.nvidia.com/browse/OMNIML-5132) (full).
**Phase 2 — close the remaining `cell.md` simplification loop:**
* Capture `experiment_id` from Docker subprocess output (Phase 1 returns PID; nemo_run's id only appears in stdout after a few seconds — Phase 2 tails launcher output via a side-channel log file).
**Phase 2 — shipped:** Docker `submit_job` captures `experiment_id` from launcher subprocess output by tailing a side-channel log file without blocking the detached process.
**Phase 3 — NEL integration + checkpoint introspection:**
* [OMNIML-5133](https://jirasw.nvidia.com/browse/OMNIML-5133) — `nel_submit`, `nel_status`, `nel_run_eval`, `nel_export`, `nel_compare`, `nel_gate` (wraps `nemo-evaluator-launcher`)
+4 -3
View File
@@ -28,9 +28,10 @@ Tool surface (Phase 1):
* ``submit_job`` — submit a launcher YAML; mode determined by args
(``hf_local`` → Docker; ``cluster_host`` → Slurm). Returns
immediately: Slurm returns ``experiment_id`` (parsed from launch.py's
detach-mode stdout), Docker returns the background subprocess
``pid`` (Phase 2 will tail launcher output to capture the nemo_run
experiment_id for the Docker path too).
detach-mode stdout), Docker always returns the background subprocess
``pid`` plus ``experiment_id`` when it appears during the short
launcher-output tail; on timeout Docker returns ``experiment_id=None``
and a persistent ``stdout_log`` diagnostic path.
* ``job_status`` — filesystem-based status from nemo_run's experiment
dir (``_DONE``, ``status_*.out``). No in-memory registry; survives
MCP server restarts.
+132 -66
View File
@@ -38,6 +38,7 @@ import os
import re
import shutil
import subprocess # nosec B404 - fixed-argv CLI probes are required; shell=True is not used.
import tempfile
import time
from dataclasses import dataclass
from pathlib import Path
@@ -94,6 +95,92 @@ def _launcher_reported_error(stdout: str, stderr: str) -> bool:
return bool(_LAUNCHER_ERROR_RE.search(f"{stdout}\n{stderr}"))
def _parse_launcher_submission(text: str) -> tuple[str | None, str | None, str | None]:
"""Best-effort parse of launcher/nemo_run submission output."""
experiment_id = None
experiment_dir = None
slurm_job_id = None
# nemo_run prints "Experiment Status for <id>" and often also the
# reconstructable form `Experiment.from_id("<id>")`.
m = re.search(r'Experiment\.from_id\("([^"]+)"\)', text)
if m:
experiment_id = m.group(1)
else:
m = re.search(
r"Experiment Status for\s+(\S+)",
text,
re.IGNORECASE,
)
if m:
experiment_id = m.group(1)
if not experiment_id:
m = re.search(
r"experiment[_\s-]+id[:\s]+(\S+)",
text,
re.IGNORECASE,
)
if m:
experiment_id = m.group(1)
if not experiment_id:
# Fallback for older nemo_run output that lacked the explicit
# "id:" label. Accepts any path-safe id token following the
# word "experiment" — not just timestamp-style.
m = re.search(
r"experiment[_\s-]+([A-Za-z0-9_-]+)",
text,
re.IGNORECASE,
)
if m and m.group(1).lower() not in {"status", "dir", "directory"}:
experiment_id = m.group(1)
# Match any path containing `/experiments/<id>/` — don't anchor on
# cluster-specific filesystem roots (NVIDIA's /lustre, partner
# clusters' /scratch / /work / /data / /p / ...).
m = re.search(r"(?:experiment_dir[:=]\s*|(?<!\S))(\S+/experiments/[^\s/]+)", text)
if m:
experiment_dir = m.group(1)
m = re.search(r"Submitted batch job (\d+)", text)
if m:
slurm_job_id = m.group(1)
else:
m = re.search(r"Job id:\s*(\d+)", text, re.IGNORECASE)
if m:
slurm_job_id = m.group(1)
return experiment_id, experiment_dir, slurm_job_id
def _docker_experiment_id_capture_timeout() -> float:
"""Return how long Docker submit should tail launcher output for an id."""
raw = os.environ.get("MODELOPT_MCP_DOCKER_ID_TIMEOUT_SEC", "10")
try:
return max(0.0, float(raw))
except ValueError:
return 10.0
def _tail_docker_launch_log(log_path: Path, proc: subprocess.Popen) -> tuple[str | None, str]:
"""Tail detached Docker launcher output for an early experiment id."""
deadline = time.monotonic() + _docker_experiment_id_capture_timeout()
text = ""
while True:
try:
text = log_path.read_text(errors="replace")
except OSError:
text = ""
complete_text = text if text.endswith(("\n", "\r")) else text.rsplit("\n", 1)[0]
experiment_id, _, _ = _parse_launcher_submission(complete_text)
if experiment_id:
return experiment_id, _tail(text, 2000)
if time.monotonic() >= deadline or proc.poll() is not None:
experiment_id, _, _ = _parse_launcher_submission(text)
if experiment_id:
return experiment_id, _tail(text, 2000)
return None, _tail(text, 2000)
time.sleep(0.2)
def _validate_experiment_id(experiment_id: str) -> dict | None:
"""Reject experiment ids that could escape path joins or alter glob matching."""
if _SAFE_EXPERIMENT_ID_RE.fullmatch(experiment_id):
@@ -962,38 +1049,69 @@ def submit_job_impl(
child_env["SLURM_HOST"] = cluster_host or ""
if executor == "docker":
# Docker mode: spawn detached. Discard stdout/stderr to /dev/null —
# leaving them as Popen.PIPE without a reader fills the kernel's
# ~64 KB pipe buffer and BLOCKS the launcher's next write(), which
# would hang long-running PTQ jobs forever while the MCP server
# appears to have "succeeded".
# Docker mode: spawn detached. Redirect stdout/stderr to a side-channel
# log file, then tail it briefly for nemo_run's experiment id. This
# avoids PIPE deadlock while still giving callers the id needed for
# job_status/job_logs polling.
# `start_new_session=True` detaches from the MCP server's process
# group so an MCP server restart / SIGINT doesn't SIGHUP the
# in-flight launcher.
# B603 false positive — argv is a controlled list built above.
log_dir = Path(child_env["NEMORUN_HOME"]) / ".modelopt-mcp" / "docker-submit-logs"
try:
log_dir.mkdir(parents=True, exist_ok=True)
log_file = tempfile.NamedTemporaryFile(
prefix="submit-",
suffix=".log",
dir=log_dir,
delete=False,
mode="w+b",
)
except OSError as e:
return {
"ok": False,
"executor": "docker",
"reason": "docker_submit_log_unavailable",
"diagnostic": f"Unable to create Docker submit log under {log_dir}: {e}",
"argv": argv,
**_source_result_fields(checkout),
}
log_path = Path(log_file.name)
try:
proc = subprocess.Popen( # nosec B603 - fixed launcher argv list; no shell.
argv,
env=child_env,
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
stdout=log_file,
stderr=subprocess.STDOUT,
start_new_session=True,
)
except FileNotFoundError:
log_file.close()
log_path.unlink(missing_ok=True)
return _launcher_not_installed(argv)
finally:
log_file.close()
experiment_id, stdout_tail = _tail_docker_launch_log(log_path, proc)
return {
"ok": True,
"executor": "docker",
"pid": proc.pid,
"argv": argv,
"nemorun_home": child_env["NEMORUN_HOME"],
"experiment_id": None, # Phase 2: tail launcher's output
"experiment_id": experiment_id,
"stdout_log": str(log_path),
"stdout_tail": stdout_tail,
**_source_result_fields(checkout),
# via a side-channel log file to capture nemo_run's id
"spike_note": (
"Docker mode launched detached. Phase 1: experiment_id "
"is None — list under $NEMORUN_HOME/experiments/ or use "
"Phase 2's tail-based id capture."
"diagnostic": (
"Docker mode launched detached and experiment_id was captured from launcher output."
if experiment_id
else (
"Docker mode launched detached, but no experiment_id was "
"captured before the short output-tail timeout. Inspect "
"stdout_log or retry with MODELOPT_MCP_DOCKER_ID_TIMEOUT_SEC "
"set higher."
)
),
}
@@ -1049,59 +1167,7 @@ def submit_job_impl(
**_source_result_fields(checkout),
}
# Best-effort experiment_id + dir + slurm_job_id parse. nemo_run's
# output format may shift across versions; on parse miss, fields
# come back None and the caller still gets stdout_tail to inspect
# by hand.
experiment_id = None
experiment_dir = None
slurm_job_id = None
# nemo_run prints "Entering Experiment <title>_<id> with id: <id>" —
# match the trailing id directly so we don't have to encode the
# title prefix or hard-code timestamp width.
m = re.search(r'Experiment\.from_id\("([^"]+)"\)', stdout_tail)
if m:
experiment_id = m.group(1)
else:
m = re.search(
r"Experiment Status for\s+(\S+)",
stdout_tail,
re.IGNORECASE,
)
if m:
experiment_id = m.group(1)
if not experiment_id:
m = re.search(
r"experiment[_\s-]+id[:\s]+(\S+)",
stdout_tail,
re.IGNORECASE,
)
if m:
experiment_id = m.group(1)
if not experiment_id:
# Fallback for older nemo_run output that lacked the explicit
# "id:" label. Accepts any path-safe id token following the
# word "experiment" — not just timestamp-style.
m = re.search(
r"experiment[_\s-]+([A-Za-z0-9_-]+)",
stdout_tail,
re.IGNORECASE,
)
if m and m.group(1).lower() != "status":
experiment_id = m.group(1)
# Match any path containing `/experiments/<id>/` — don't anchor on
# cluster-specific filesystem roots (NVIDIA's /lustre, partner
# clusters' /scratch / /work / /data / /p / ...).
m = re.search(r"(\S+/experiments/[^\s/]+)", stdout_tail)
if m:
experiment_dir = m.group(1)
m = re.search(r"Submitted batch job (\d+)", stdout_tail)
if m:
slurm_job_id = m.group(1)
else:
m = re.search(r"Job id:\s*(\d+)", stdout_tail, re.IGNORECASE)
if m:
slurm_job_id = m.group(1)
experiment_id, experiment_dir, slurm_job_id = _parse_launcher_submission(stdout_tail)
if not experiment_id:
return {
+6 -4
View File
@@ -130,10 +130,12 @@ def _build_server() -> FastMCP:
"determined by mutually-exclusive args:\n"
" - hf_local=<path> → Docker (local GPU)\n"
" - cluster_host=<host> → Slurm (remote SSH)\n\n"
"Returns the experiment_id (Slurm) or PID (Docker, "
"experiment_id captured in Phase 2) immediately; the actual "
"job runs detached. Poll status via job_status, fetch "
"output via job_logs.\n\n"
"Returns PID for Docker, plus experiment_id when the id is "
"printed during the short launch-output tail; if the tail times "
"out, Docker returns experiment_id=None and stdout_log for "
"diagnostics. Slurm returns experiment_id. The actual job runs "
"detached. Poll status via job_status, fetch output via "
"job_logs.\n\n"
"Auto-verifies the executor first by default (skip_verify="
"False is recommended unless you just called verify_setup)."
),
+134
View File
@@ -394,6 +394,140 @@ def test_submit_job_source_checkout_failure_short_circuits(monkeypatch):
assert result["reason"] == "source_ref_not_found"
def test_submit_job_docker_captures_experiment_id_from_launcher_output(monkeypatch, tmp_path):
"""Docker submit tails launcher output long enough to return nemo_run's id."""
yaml_dir = tmp_path / "examples"
yaml_dir.mkdir()
yaml_path = yaml_dir / "config.yaml"
yaml_path.write_text("job_name: t\npipeline: []\n")
monkeypatch.setenv("MODELOPT_LAUNCHER_EXAMPLES_DIR", str(yaml_dir))
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path / "nemo"))
monkeypatch.setattr(bridge, "verify_docker_setup_impl", lambda: {"ok": True})
class FakePopen:
pid = 4242
def __init__(self, argv, **kwargs):
self.argv = argv
out = kwargs["stdout"]
out.write(
b"Experiment Status for docker_1782173197\n"
b'experiment = run.Experiment.from_id("docker_1782173197")\n'
)
out.flush()
def poll(self):
return None
monkeypatch.setattr(subprocess, "Popen", FakePopen)
result = bridge.submit_job_impl(
yaml_path="config.yaml",
hf_local="/tmp/hf",
cluster_host=None,
cluster_user=None,
identity=None,
job_dir=None,
job_name=None,
extra_overrides=None,
skip_verify=False,
)
assert result["ok"] is True
assert result["executor"] == "docker"
assert result["pid"] == 4242
assert result["experiment_id"] == "docker_1782173197"
assert "docker_1782173197" in result["stdout_tail"]
assert result["stdout_log"].endswith(".log")
def test_submit_job_docker_no_experiment_id_returns_pid_and_log(monkeypatch, tmp_path):
"""Docker submit stays detached when the id is not printed immediately."""
yaml_dir = tmp_path / "examples"
yaml_dir.mkdir()
yaml_path = yaml_dir / "config.yaml"
yaml_path.write_text("job_name: t\npipeline: []\n")
monkeypatch.setenv("MODELOPT_LAUNCHER_EXAMPLES_DIR", str(yaml_dir))
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path / "nemo"))
monkeypatch.setenv("MODELOPT_MCP_DOCKER_ID_TIMEOUT_SEC", "0")
monkeypatch.setattr(bridge, "verify_docker_setup_impl", lambda: {"ok": True})
class FakePopen:
pid = 4243
def __init__(self, argv, **kwargs):
kwargs["stdout"].write(b"launcher starting\n")
kwargs["stdout"].flush()
def poll(self):
return None
monkeypatch.setattr(subprocess, "Popen", FakePopen)
result = bridge.submit_job_impl(
yaml_path="config.yaml",
hf_local="/tmp/hf",
cluster_host=None,
cluster_user=None,
identity=None,
job_dir=None,
job_name=None,
extra_overrides=None,
skip_verify=False,
)
assert result["ok"] is True
assert result["executor"] == "docker"
assert result["pid"] == 4243
assert result["experiment_id"] is None
assert "launcher starting" in result["stdout_tail"]
assert "no experiment_id was captured" in result["diagnostic"]
def test_parse_launcher_submission_ignores_experiment_dir_as_id():
"""A path field must not be misread as a usable experiment id."""
experiment_id, experiment_dir, slurm_job_id = bridge._parse_launcher_submission(
"experiment_dir=/tmp/nemorun/experiments/docker_1782173197\n"
)
assert experiment_id is None
assert experiment_dir == "/tmp/nemorun/experiments/docker_1782173197"
assert slurm_job_id is None
def test_submit_job_docker_log_creation_failure_is_structured(monkeypatch, tmp_path):
"""Docker submit should not raise if the side-channel log cannot be created."""
yaml_dir = tmp_path / "examples"
yaml_dir.mkdir()
yaml_path = yaml_dir / "config.yaml"
yaml_path.write_text("job_name: t\npipeline: []\n")
monkeypatch.setenv("MODELOPT_LAUNCHER_EXAMPLES_DIR", str(yaml_dir))
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path / "nemo"))
monkeypatch.setattr(bridge, "verify_docker_setup_impl", lambda: {"ok": True})
def fail_named_temporary_file(*args, **kwargs):
raise OSError("disk full")
monkeypatch.setattr(bridge.tempfile, "NamedTemporaryFile", fail_named_temporary_file)
result = bridge.submit_job_impl(
yaml_path="config.yaml",
hf_local="/tmp/hf",
cluster_host=None,
cluster_user=None,
identity=None,
job_dir=None,
job_name=None,
extra_overrides=None,
skip_verify=False,
)
assert result["ok"] is False
assert result["executor"] == "docker"
assert result["reason"] == "docker_submit_log_unavailable"
assert "disk full" in result["diagnostic"]
def test_submit_job_slurm_zero_exit_without_ids_is_failure(monkeypatch, tmp_path):
"""Slurm submit must not report success when launcher emits no ids."""
yaml_dir = tmp_path / "examples"