feat(tools/mcp): MCP server for ModelOpt launcher (OMNIML-5123) (#1701)

## Summary

`tools/mcp/` — a new MCP server exposing the existing
`tools/launcher/core.py` orchestration as **typed MCP tools** that codex
/ Claude Code agents can call directly, instead of shelling out to `uv
run launch.py --yaml ...` and parsing prose output.

Tracked under
[OMNIML-5123](https://jirasw.nvidia.com/browse/OMNIML-5123) (Epic).
Ships **Phase 1 + Phase 1.5** together: the core launcher surface plus
the four highest-leverage helpers from the `cell.md` simplification loop
([OMNIML-5128](https://jirasw.nvidia.com/browse/OMNIML-5128) partial,
[OMNIML-5132](https://jirasw.nvidia.com/browse/OMNIML-5132) full).

## Nine tools

**Phase 1 — core launcher surface:**

| Tool | Description |
|---|---|
| `list_examples` | Enumerate `tools/launcher/examples/` with model +
description metadata extracted from each YAML |
| `verify_setup` | Fail-fast probe for the named executor. Docker:
`docker info` (daemon up) + `docker info --format` runtime-registry
check for the `nvidia` runtime — no image pull, daemon-fast. Slurm: `ssh
-o BatchMode=yes -o ConnectTimeout=5` to the cluster login node. ~1 s
probe saves 30+ s of wasted submission on bad config |
| `submit_job` | Submit a launcher YAML. Mode is determined by
mutually-exclusive args: `hf_local` → Docker (local GPU), `cluster_host`
→ Slurm (remote SSH). Returns experiment_id immediately; the actual job
runs detached |
| `job_status` | Filesystem-based status from nemo_run's experiment dir
(`_DONE`, `status_*.out`) — no in-memory registry, survives MCP server
restarts |
| `job_logs` | Read `log_<task>.out` from experiment dir; per-task
filtering + optional tail |

**Phase 1.5 — `cell.md` simplification (OMNIML-5128 / 5132):**

| Tool | Description |
|---|---|
| `wait_for_experiment` | Replaces the agent's `while True: status;
sleep` poll loop with one tool call. Reuses `job_status_impl` so
terminal-state semantics stay identical. Returns final status plus
`waited_seconds`; on timeout returns structured `{ok: False, reason:
"wait_timeout", last_status: …}` |
| `provision_passwordless_ssh_dry_run` | No-side-effect inspection of
`~/.ssh/` that emits the exact `ssh-keygen` / `ssh-copy-id` commands the
operator should run to make `verify_setup(executor='slurm')` pass.
Closes the verify_setup "ssh_auth_failed → now what?" gap |
| `read_cluster_artifact` | Uses nemo_run's tunnel primitives, not a
reinvented SSH layer. `path=None` wraps `nemo experiment logs <id>
<job_idx>` (built-in log fetch); with a `path`, uses the experiment's
`Tunnel` to read the file. Structured failure on subprocess error /
timeout |
| `open_draft_pr` | `git push -u origin HEAD` + `gh pr create --draft
…`. Validates cwd is a git repo first; on gh failure after push
succeeds, reports `branch_pushed=True` so the operator can retry just
the PR-open step |

## Design constants

1. **Single `submit_job` with mode by args** (not separate
`submit_docker` / `submit_slurm` tools). Keeps the LLM tool catalog
compact; mutual-exclusion is a runtime check.
2. **Filesystem is the source of truth** for status + logs. No in-memory
registry. Survives MCP server restarts cleanly — important because
operators / agents kill + restart their hosts often.
3. **`verify_setup` is auto-called by `submit_job`** by default
(skippable when caller just probed). The probe is ~1 s; the cost of a
misconfigured submission is 30+ s of cluster timeout or container-pull.
Always-on verify pays back immediately.
4. **Delegate to nemo_run for tunnels.** `read_cluster_artifact` and
`wait_for_experiment` use nemo_run's existing `Experiment` / `Tunnel` /
`nemo experiment logs` primitives rather than reinventing SSH/rsync. One
source of truth for cluster I/O.

## Layout

```
tools/mcp/
├── pyproject.toml          # name: modelopt-mcp, console_script
├── modelopt_mcp/
│   ├── __init__.py
│   ├── server.py           # FastMCP entry; 9 tool definitions
│   └── bridge.py           # thin wrapper over launcher's core.py
│                           #   + filesystem status/log helpers
│                           #   + tunnel/PR helpers (Phase 1.5)
└── tests/
    └── test_bridge.py      # 34 unit tests, fully hermetic
                            # (mocked subprocess + tmp_path fixtures)
```

## Install

Two paths, both **from source via uv**. No PyPI wheel exists;
OMNIML-5123 opted for the uvx-from-git pattern to skip publication
overhead.

### End-user install (recommended)

`uvx` from the git subdirectory — single command, no manual clone:

```bash
# Claude Code
claude mcp add modelopt -- uvx --from \
  "git+https://github.com/NVIDIA/Model-Optimizer.git#subdirectory=tools/mcp" \
  modelopt-mcp

# Codex
codex mcp add modelopt -- uvx --from \
  "git+https://github.com/NVIDIA/Model-Optimizer.git#subdirectory=tools/mcp" \
  modelopt-mcp
```

Under the hood `uvx` clones the whole repo to its cache, installs
`tools/mcp/` as the entry, and resolves the sibling `modelopt-launcher`
dep via `[tool.uv.sources]` (`path = "../launcher"`) inside the cloned
tree.

### Dev install (local checkout)

```bash
uv pip install -e tools/launcher    # sibling dep first
uv pip install -e tools/mcp         # then this package
modelopt-mcp                         # entry on PATH
```

### Why no plain `pip install` today

Two specific reasons, worth flagging so reviewers know what's
intentional vs missing:

1. **Nothing on PyPI yet.** Neither `modelopt-mcp` nor
`modelopt-launcher` are published — this PR introduces the package but
doesn't add release machinery.
2. **`pip` doesn't read `[tool.uv.sources]`.** Even from a local
checkout, plain `pip install -e tools/mcp` fails because
`modelopt-launcher` is a bare name (no URL) and pip can't find it.
Sticking with `uv` / `uvx` is the practical path while we're git-only.

If we later want plain-pip support: publish to PyPI, or switch to a
PEP-440 direct URL (`"modelopt-launcher @
git+…#subdirectory=tools/launcher"`). Out of scope for this PR.

## Post-review changes

Addressed all CodeRabbit + claude[bot] review findings on the original
Phase-1 surface. See the inline replies for details; the substantive
bug-fix highlights:

* **Slurm `cluster_host`** — propagate via `env=child_env` (launch.py
reads SLURM_HOST, not a CLI arg)
* **`shlex.quote`** removed from nemo-run k=v overrides (subprocess
list-form doesn't shell-quote)
* **Docker `Popen`** now uses `stdout=DEVNULL, stderr=DEVNULL,
start_new_session=True` to avoid pipe-buffer blocking
* **`NEMORUN_HOME`** pinned in subprocess env so submit + status sides
agree
* **GPU verify** swapped from `docker run --gpus all` image-pull (slow +
flaky) to `docker info --format` runtime-registry check (daemon-fast)
* **Task-status word match** anchors on first word against a fixed
failure-word set (no more `"fail" in "succeeded after retry; previous
attempt failed"` false-positive)
* **`experiment_id` regex** generalized for non-NVIDIA cluster paths
* **`pyproject.toml`** dropped the unsatisfiable `modelopt-launcher`
bare-name dep (launcher is a file-layout sibling, not a Python import
dep)
* **`Field(ge=1)`** on `job_logs.tail`
* **Docstring contract** clarified (Docker returns `pid`, Slurm returns
`experiment_id`)

## Validation

- [x] `uv pip install -e .` succeeds (modelopt-launcher resolved
transitively)
- [x] 34/34 unit tests pass (`uv run python -m pytest tests/`)
- [x] stdio handshake works end-to-end; `tools/list` returns all 9 with
full schemas + descriptions
- [x] Mode-resolution: `submit_job` correctly rejects no-executor +
both-executors with structured `reason`
- [x] Filesystem status: correctly classifies `done` / `failed` /
`running` from `_DONE` + `status_*.out`
- [x] `wait_for_experiment` short-circuits on already-terminal
experiments; honors timeout without raising
- [x] `provision_passwordless_ssh_dry_run` distinguishes no-key /
key-only / key+pubkey cases
- [x] `read_cluster_artifact` handles subprocess timeout + non-zero exit
with structured reasons
- [x] `open_draft_pr` reports `branch_pushed=True` on
gh-failure-after-push so retries are cheap
- [x] Pre-commit clean: ruff, ruff-format, mypy, bandit, license-headers

## Acceptance criteria

**OMNIML-5123 (Phase 1):**

- [x] `list_examples` returns all bundled YAMLs with path and model name
- [x] `submit_job` with `hf_local` runs via Docker executor and returns
immediately (Phase 1: returns PID; experiment_id capture in Phase 2)
- [x] `submit_job` with `cluster_host`/`user` runs via Slurm executor
(`detach=True`) and returns experiment_id
- [x] `job_status` correctly reflects running / done / failed from
nemo_run filesystem
- [x] `job_logs` returns stdout for a completed job
- [x] `uvx --from git+...#subdirectory=tools/mcp modelopt-mcp --help`
resolves and starts
- [x] Existing launcher tests unaffected (no changes to
`tools/launcher/`)

**OMNIML-5128 (Phase 1.5, partial):**

- [x] `wait_for_experiment` blocks until terminal or timeout
- [x] `read_cluster_artifact` pulls remote artifacts via nemo_run tunnel
- [x] `open_draft_pr` opens a draft PR against a target repo
- [ ] Capture `experiment_id` from Docker subprocess output — deferred
to Phase 2

**OMNIML-5132 (Phase 1.5, full):**

- [x] `provision_passwordless_ssh_dry_run` emits operator-facing
commands without side effects

## Phase 2 (separate PR)

* Capture `experiment_id` from Docker subprocess output (tail until
nemo_run logs the id).
* Extract the verify + submit helpers into a shared lib that
[`nmm-sandbox-mcp`](https://gitlab-master.nvidia.com/omniml/integration/nmm-sandbox/-/tree/main/tools/mcp)
(companion server, separate repo) can consume for internal-ergonomics
tools — cluster short-name → factory lookup + GitLab CI dispatch.
* NEL integration
([OMNIML-5133](https://jirasw.nvidia.com/browse/OMNIML-5133)) +
checkpoint introspection
([OMNIML-5134](https://jirasw.nvidia.com/browse/OMNIML-5134)).


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

* **New Features**
* Added ModelOpt MCP server and console entrypoint; tools:
list_examples, verify_setup, submit_job, job_status, job_logs,
wait_for_experiment, provision_passwordless_ssh_dry_run,
read_cluster_artifact, open_draft_pr; Docker and Slurm support.

* **Documentation**
* Expanded README with install steps, design notes, end-to-end agent
example, roadmap, and repo layout.

* **Tests**
* Expanded unit tests covering bridge helpers, polling, SSH flows,
artifact reads, and PR automation.

* **Chores**
  * CI updated to run MCP tests; package/meta config added.
<!-- end of auto-generated comment: release notes by coderabbit.ai -->

---------

Signed-off-by: Chenhan Yu <chenhany@nvidia.com>
Co-authored-by: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
Chenhan D. Yu
2026-06-12 18:56:18 -07:00
committed by GitHub
co-authored by Claude Opus 4.7
parent 2640551515
commit 9f37fe1969
8 changed files with 2809 additions and 1 deletions
+27 -1
View File
@@ -12,6 +12,7 @@ on:
- "pyproject.toml"
- "tests/unit/**"
- "tools/launcher/**"
- "tools/mcp/**"
- ".agents/skills/**"
schedule:
- cron: "0 0 * * *" # Nightly
@@ -55,6 +56,7 @@ jobs:
pyproject.toml
tests/unit/**
tools/launcher/**
tools/mcp/**
.agents/skills/**
linux:
needs: [check-dco]
@@ -147,6 +149,29 @@ jobs:
uv venv .venv
uv pip install -e . pytest
uv run python3 -m pytest -v
mcp:
if: needs.check-file-changes.outputs.any_changed == 'true'
needs: [linux, check-file-changes]
runs-on: ubuntu-latest
timeout-minutes: 15
steps:
- uses: actions/checkout@v6
with:
submodules: recursive
- name: Run modelopt-mcp tests
working-directory: tools/mcp
run: |
curl -LsSf https://astral.sh/uv/install.sh | sh
export PATH="$HOME/.local/bin:$PATH"
uv venv .venv
# Install the sibling launcher package first; it's a runtime
# dep declared in tools/mcp/pyproject.toml as `modelopt-launcher`
# but uv resolves the source via [tool.uv.sources] to a local
# editable path. -e on both packages keeps the install cheap
# and matches the dev-mode install in tools/mcp/README.md.
uv pip install -e ../launcher
uv pip install -e . pytest
uv run python3 -m pytest -v
skills:
if: needs.check-file-changes.outputs.any_changed == 'true'
needs: [linux, check-file-changes]
@@ -167,7 +192,7 @@ jobs:
unit-pr-required-check:
# Run even if some jobs are skipped
if: ${{ github.event_name == 'pull_request' && always() }}
needs: [check-file-changes, linux, windows, multi-version, partial-install, launcher, skills]
needs: [check-file-changes, linux, windows, multi-version, partial-install, launcher, mcp, skills]
runs-on: ubuntu-latest
steps:
- name: Required unit tests did not succeed
@@ -177,6 +202,7 @@ jobs:
needs.multi-version.result != 'success' ||
needs.partial-install.result != 'success' ||
needs.launcher.result != 'success' ||
needs.mcp.result != 'success' ||
needs.skills.result != 'success'
)) }}
run: exit 1
+172
View File
@@ -0,0 +1,172 @@
# modelopt-mcp
MCP server exposing the ModelOpt launcher (`tools/launcher/`) as typed tools for codex / Claude Code agents.
Anchor design: [OMNIML-5123](https://jirasw.nvidia.com/browse/OMNIML-5123).
## What this is
A thin MCP wrapper around `tools/launcher/core.py`. Agents that want to submit ModelOpt jobs — PTQ, QAT, training, evaluation — call typed tools here instead of shelling out to `uv run launch.py --yaml ...` and parsing prose output.
Two executors, one tool surface:
* **Docker** (local GPU) — `submit_job(yaml_path, hf_local=...)`. The launcher runs the job in a container on the local machine.
* **Slurm** (remote cluster via SSH) — `submit_job(yaml_path, cluster_host=..., cluster_user=..., identity=...)`. The launcher tunnels in and submits via sbatch.
Mode is determined by which args you pass, not by which tool you call. One tool, two backends.
## Tool surface
| Tool | Description |
|---|---|
| `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?, ...)` | Submit a launcher YAML. Mode resolved from mutually-exclusive args. Returns `experiment_id` (Slurm) or PID (Docker) immediately; the actual job runs detached. Auto-runs `verify_setup` first by default (skippable). |
| `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`. |
| `provision_passwordless_ssh_dry_run(cluster_host, cluster_user, identity?)` | Operator UX helper. Inspects `~/.ssh/` and emits the exact `ssh-keygen` / `ssh-copy-id` commands the user should run to make `verify_setup(executor='slurm')` pass. No side effects — pure inspection + shell-command formatting. |
| `read_cluster_artifact(experiment_id, path?, job_idx?)` | Pull artifact content from the remote cluster via nemo_run's tunnel primitives. With `path=None`, wraps `nemo experiment logs <id> <job_idx>` (built-in log fetch). With a `path`, uses the experiment's `Tunnel` to `cat` the file. No reinvented SSH. |
| `open_draft_pr(target_repo, title, body, base_branch?, cwd?)` | Push the current branch and open a draft PR via `gh pr create --draft`. Validates `cwd` is a git repo before doing anything. On `gh` failure after a successful push, returns `branch_pushed=True` so the operator can retry just the PR-open step. |
## Install
Two paths, both **from source via uv**. No PyPI wheel — OMNIML-5123 deliberately picks `uvx`-from-git over publication overhead.
### End-user install (recommended)
```bash
# Claude Code
claude mcp add modelopt -- uvx --from \
"git+https://github.com/NVIDIA/Model-Optimizer.git#subdirectory=tools/mcp" \
modelopt-mcp
# Codex
codex mcp add modelopt -- uvx --from \
"git+https://github.com/NVIDIA/Model-Optimizer.git#subdirectory=tools/mcp" \
modelopt-mcp
```
`uvx` clones the whole repo to its cache, installs `tools/mcp/` as the entry point, and resolves the sibling `modelopt-launcher` dep via `[tool.uv.sources]` (path → `../launcher`) inside the cloned tree.
### Dev install (local checkout)
```bash
uv pip install -e tools/launcher # sibling dep first
uv pip install -e tools/mcp # then this package
modelopt-mcp # stdio server entry on PATH
```
Both packages share the launcher's `core.py` orchestrator. The dev path relies on `[tool.uv.sources]` to point `modelopt-launcher` at `../launcher`.
### Why no plain `pip install` today
`modelopt-mcp` and `modelopt-launcher` are not on PyPI. Plain `pip` doesn't read `[tool.uv.sources]`, so even from a local checkout, `pip install -e tools/mcp` fails to resolve the bare `modelopt-launcher` name. Stick with `uv` / `uvx` while we're git-only.
To enable `pip install` later, two options:
| Path | Tradeoff |
|---|---|
| **Publish to PyPI** — versioned wheels for both packages | Clean `pip install`, but requires release machinery + version cadence |
| **PEP-440 direct URL** — `"modelopt-launcher @ git+...#subdirectory=tools/launcher"` | Works with pip + uv, but double-clones the repo on install |
Out of scope for Phase 1.
## Example workflow
Agent picking a bundled example and running it on a remote cluster:
```python
# 1. Discover available YAMLs
examples = mcp__modelopt__list_examples()
# {"ok": True, "count": 47, "examples": [{"path": "launcher/examples/Qwen/Qwen3-8B/megatron_lm_ptq.yaml", "model": "Qwen/Qwen3-8B", ...}, ...]}
# 2. Pre-flight check before submission
mcp__modelopt__verify_setup(
executor="slurm",
cluster_host="cw-dfw-cs-001-login-01.nvidia.com",
cluster_user="alice",
)
# {"ok": True, "ssh_ok": True, "whoami": "alice", "remote_hostname": "cw-dfw-cs-001-login-01"}
# 3. Submit
result = mcp__modelopt__submit_job(
yaml_path="launcher/examples/Qwen/Qwen3-8B/megatron_lm_ptq.yaml",
cluster_host="cw-dfw-cs-001-login-01.nvidia.com",
cluster_user="alice",
identity="/home/alice/.ssh/id_ed25519",
skip_verify=True, # we just probed
)
# {"ok": True, "experiment_id": "cicd_1781240000", "slurm_job_id": "12345", ...}
# 4. Poll until done
while True:
status = mcp__modelopt__job_status(experiment_id="cicd_1781240000")
if status["status"] in ("done", "failed"):
break
# 5. Fetch logs
logs = mcp__modelopt__job_logs(
experiment_id="cicd_1781240000",
task="task_0",
tail=200,
)
```
For local Docker execution, drop `cluster_host`/`cluster_user`/`identity` and pass `hf_local=<path>` to `submit_job` instead.
## Required env vars
| Var | When | Notes |
|---|---|---|
| `NEMORUN_HOME` | submit + status + logs | Where the launcher writes experiment artifacts. Defaults to cwd if unset. `job_status` / `job_logs` search `$NEMORUN_HOME/experiments/<id>/`. |
| `MODELOPT_MCP_LOG` | (optional) server | Log level. Defaults to `INFO`. Logs go to stderr — stdout is the MCP wire. |
| `MODELOPT_MCP_SKIP_GPU_CHECK` | (optional) `verify_setup(executor='docker')` | Set to skip the `docker info --format` runtime-registry check. Useful for CI hosts where the daemon is up but the NVIDIA Container Toolkit isn't installed. |
| `MODELOPT_LAUNCHER_EXAMPLES_DIR` | (optional) `list_examples` | Override the examples directory location. Defaults to `../launcher/examples/` relative to this package. |
## Design principles
Three constants drive the surface here:
1. **Single `submit_job` with mode by args.** Not separate `submit_docker` / `submit_slurm` tools. The LLM tool catalog stays compact; the mutual-exclusion is a runtime check that returns structured failure when both or neither mode is specified.
2. **Filesystem is the source of truth.** Status + logs both read from nemo_run's experiment dir. No in-memory registry — survives MCP server restarts. The bridge module never holds per-job state across calls.
3. **`verify_setup` is auto-called by `submit_job` by default.** The probe takes ~1 second; the cost of a misconfigured submission is 30+ seconds of cluster timeout or container-pull. Always-on verification pays back immediately. Callers can pass `skip_verify=True` when they just probed.
## Internal companion (NVIDIA only)
For NVIDIA-internal users running on the in-house clusters, there's a companion server [`nmm-sandbox-mcp`](https://gitlab-master.nvidia.com/omniml/integration/nmm-sandbox/-/tree/main/tools/mcp) that adds:
* `resolve_cluster_factory(name)` — turn `"cw_dfw"` into the `{cluster_host, cluster_user, identity, account, ...}` dict that `submit_job` consumes, so internal users supply 1 arg instead of 6.
* `submit_via_gitlab_ci(...)` — alternative submission path that triggers the nmm-sandbox CI's `intern-step` pipeline. Useful when the operator doesn't have direct cluster access.
Both servers are registered in the operator's `.mcp.json` side by side. The agent threads results from one into calls to the other. No Python coupling between them.
## Layout
```text
tools/mcp/
├── pyproject.toml # name: modelopt-mcp, console_script
├── README.md # ← this file
├── modelopt_mcp/
│ ├── __init__.py
│ ├── server.py # FastMCP entry; 9 tool definitions
│ └── bridge.py # thin wrapper over launcher's core.py
│ # + filesystem status/log helpers
│ # + tunnel/PR helpers (Phase 1.5)
└── tests/
└── test_bridge.py # 32 unit tests, fully hermetic
# (mocked subprocess + tmp_path fixtures)
```
## Phase 2 & beyond
Tracked under [OMNIML-5123](https://jirasw.nvidia.com/browse/OMNIML-5123) (Epic). Highlights:
**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 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`)
* [OMNIML-5134](https://jirasw.nvidia.com/browse/OMNIML-5134) — `inspect_checkpoint`, `inspect_model`
+52
View File
@@ -0,0 +1,52 @@
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""ModelOpt launcher MCP server.
Anchor design: OMNIML-5123. Exposes the launcher's submit / status / logs
operations as typed MCP tools that codex / Claude Code agents can call
directly, instead of shelling out to ``uv run launch.py --yaml ...``.
Tool surface (Phase 1):
* ``list_examples`` — discover bundled example YAMLs under
``tools/launcher/examples/`` with their model + description metadata.
* ``verify_setup`` — fail-fast probe for the named executor (docker or
slurm) BEFORE the user burns minutes on a misconfigured submission.
* ``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).
* ``job_status`` — filesystem-based status from nemo_run's experiment
dir (``_DONE``, ``status_*.out``). No in-memory registry; survives
MCP server restarts.
* ``job_logs`` — read ``log_*.out`` per task, with optional tail.
Two design constants:
1. **Mode determined by args, not by tool choice.** Single
``submit_job`` rather than separate ``submit_docker`` / ``submit_slurm``.
The mutually-exclusive arg shape is a runtime check; the LLM catalog
stays compact.
2. **Filesystem is the source of truth.** Status + logs both read from
nemo_run's experiment dir. The MCP server carries no per-job state
across calls — survives restarts cleanly.
"""
from modelopt_mcp.server import main
__all__ = ["main"]
File diff suppressed because it is too large Load Diff
+478
View File
@@ -0,0 +1,478 @@
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""modelopt-mcp server entry point.
Stdio transport; codex / Claude Code launch this as a subprocess and
talk to it over stdin/stdout. See OMNIML-5123 for the design.
Phase 1 tool surface:
* list_examples — discover bundled launcher YAMLs
* verify_setup — fail-fast probe (docker or slurm)
* submit_job — submit a launcher YAML; mode by args
* job_status — filesystem-based status
* job_logs — filesystem-based logs
All tools return JSON-friendly dicts with explicit ``ok`` / ``reason`` /
``diagnostic`` fields so the calling LLM can route on structured
outcomes instead of free-form prose.
"""
from __future__ import annotations
import logging
import os
import sys
from typing import Annotated, Literal
from mcp.server.fastmcp import FastMCP
from pydantic import Field
from modelopt_mcp import bridge
logger = logging.getLogger("modelopt_mcp")
def _build_server() -> FastMCP:
"""Construct the MCP server with Phase 1 tools registered.
Factored out so tests can build an isolated instance without
stdio plumbing.
"""
mcp = FastMCP("modelopt")
@mcp.tool(
name="list_examples",
description=(
"List all bundled launcher YAML examples under "
"tools/launcher/examples/, with model + description "
"metadata extracted from each YAML. Use BEFORE submit_job "
"when you don't know which YAML to launch — this gives the "
"agent the discovery surface needed to pick one."
),
)
def list_examples() -> dict:
return bridge.list_examples_impl()
@mcp.tool(
name="verify_setup",
description=(
"Probe whether the named executor is reachable from THIS "
"host. Run BEFORE submit_job to fail fast — the actual "
"submission burns 30+ seconds (slurm) or starts a Docker "
"container (docker) before discovering setup is broken.\n\n"
"Docker mode: checks `docker info` (daemon up) + GPU "
"passthrough (`docker run --gpus all nvidia-smi`). Set "
"MODELOPT_MCP_SKIP_GPU_CHECK=1 in the env to skip the GPU "
"check on CPU-only hosts.\n\n"
"Slurm mode: ssh -o BatchMode=yes -o ConnectTimeout=5 "
"to the cluster login node. Refuses to prompt for "
"password, so key-auth failure is detected immediately."
),
)
def verify_setup(
executor: Annotated[
Literal["docker", "slurm"],
Field(
description=(
"Which executor to probe: 'docker' for local GPU or 'slurm' for remote cluster."
)
),
],
cluster_host: Annotated[
str | None,
Field(description=("Slurm cluster login hostname. Required when executor='slurm'.")),
] = None,
cluster_user: Annotated[
str | None, Field(description=("SSH user for the cluster. None uses ssh's default."))
] = None,
identity: Annotated[
str | None,
Field(
description=("SSH identity file (-i) override. None uses default key / ssh-agent.")
),
] = None,
) -> dict:
if executor == "docker":
return bridge.verify_docker_setup_impl()
if executor == "slurm":
if not cluster_host:
return {
"ok": False,
"executor": "slurm",
"reason": "missing_cluster_host",
"diagnostic": ("executor='slurm' requires cluster_host=<hostname>."),
}
return bridge.verify_slurm_setup_impl(
cluster_host=cluster_host,
cluster_user=cluster_user,
identity=identity,
)
# Pydantic Literal already constrains; this is a defensive fallback.
return {"ok": False, "reason": "unknown_executor"}
@mcp.tool(
name="submit_job",
description=(
"Submit a ModelOpt launcher YAML for execution. Mode is "
"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"
"Auto-verifies the executor first by default (skip_verify="
"False is recommended unless you just called verify_setup)."
),
)
def submit_job(
yaml_path: Annotated[
str,
Field(
description=(
"Launcher YAML to submit. Pass an absolute path, a path "
"relative to tools/launcher/examples/, or one of the paths "
"returned by list_examples."
)
),
],
hf_local: Annotated[
str | None,
Field(
description=(
"Local HF cache directory — when set, dispatches via "
"Docker. Mutually exclusive with cluster_host."
)
),
] = None,
cluster_host: Annotated[
str | None,
Field(
description=(
"Slurm cluster login hostname — when set, dispatches via "
"remote SSH. Mutually exclusive with hf_local."
)
),
] = None,
cluster_user: Annotated[
str | None,
Field(description=("SSH user for the cluster. None uses launcher's default.")),
] = None,
identity: Annotated[
str | None,
Field(description=("SSH identity file (-i). None uses ssh-agent / default key.")),
] = None,
job_dir: Annotated[
str | None,
Field(
description=(
"Override the experiment output directory. None uses the "
"launcher's per-mode default."
)
),
] = None,
job_name: Annotated[
str | None,
Field(description=("Override the job_name in the YAML. None uses the YAML's default.")),
] = None,
extra_overrides: Annotated[
dict[str, str] | None,
Field(
description=(
"Additional nemo-run-style overrides as a flat dict, e.g. "
"{'task.slurm_config.nodes': '2'}."
)
),
] = None,
skip_verify: Annotated[
bool,
Field(
description=(
"If True, skip the verify_setup probe before submission. "
"Default False — the probe takes ~1s and saves you from "
"30+s of wasted submission time on bad config."
)
),
] = False,
) -> dict:
return bridge.submit_job_impl(
yaml_path=yaml_path,
hf_local=hf_local,
cluster_host=cluster_host,
cluster_user=cluster_user,
identity=identity,
job_dir=job_dir,
job_name=job_name,
extra_overrides=extra_overrides,
skip_verify=skip_verify,
)
@mcp.tool(
name="job_status",
description=(
"Read filesystem-based status from a nemo_run experiment "
"dir. Returns 'done' / 'failed' / 'running' based on "
"presence of _DONE and contents of status_*.out files. "
"Per-task statuses also surfaced for multi-task pipelines."
),
)
def job_status(
experiment_id: Annotated[
str,
Field(
description=(
"The experiment id returned by submit_job (Slurm) or the "
"name nemo_run assigned to the experiment dir."
)
),
],
) -> dict:
return bridge.job_status_impl(experiment_id)
@mcp.tool(
name="job_logs",
description=(
"Read log_<task>.out files from the experiment dir. If "
"task=None, returns logs for all tasks. If tail=N, returns "
"only the last N lines per task."
),
)
def job_logs(
experiment_id: Annotated[str, Field(description=("The experiment id."))],
task: Annotated[
str | None,
Field(description=("Specific task name to filter logs by. None returns all.")),
] = None,
tail: Annotated[
int | None,
Field(
ge=1,
description=(
"Return only the last N lines per task. Must be >= 1 when set; "
"None returns the full log."
),
),
] = None,
) -> dict:
return bridge.job_logs_impl(experiment_id, task, tail)
@mcp.tool(
name="wait_for_experiment",
description=(
"Block until an experiment reaches a terminal status "
"('done' or 'failed') or the timeout elapses. Returns the "
"same shape as job_status plus a `waited_seconds` field. "
"Use this instead of writing your own polling while-loop "
"around job_status."
),
)
def wait_for_experiment(
experiment_id: Annotated[
str,
Field(description="The experiment id from submit_job."),
],
timeout_sec: Annotated[
int,
Field(
ge=1,
description=(
"Max seconds to wait before returning "
"`{ok: False, reason: 'wait_timeout'}`. Default 7200 "
"(2 hours) — large PTQ runs need this. Cap at your "
"agent's own deadline."
),
),
] = 7200,
poll_interval_sec: Annotated[
int,
Field(
ge=1,
description=(
"Seconds between status polls. Default 30 — "
"filesystem-based status is cheap so the poll "
"doesn't have to be slow."
),
),
] = 30,
) -> dict:
return bridge.wait_for_experiment_impl(
experiment_id,
timeout_sec,
poll_interval_sec,
)
@mcp.tool(
name="provision_passwordless_ssh_dry_run",
description=(
"Emit the exact commands the operator should run to set up "
"passwordless SSH to a slurm cluster. Does NOT execute "
"them and does NOT handle passwords — the MCP wire is "
"unsafe for cluster credentials. Two-phase:\n\n"
"1. If the SSH private key is missing → emit "
"`ssh-keygen` command. Operator runs it, re-invokes this "
"tool.\n"
"2. If the key exists → emit `ssh-copy-id` command + "
"public-key content. Operator runs it (this is the only "
"step that needs a cluster password, prompted by ssh).\n\n"
"After both steps complete, the `next_check` field "
"recommends calling `verify_setup(executor='slurm', ...)` "
"to confirm key-auth now works. Use this when "
"`verify_setup` returns `ssh_auth_failed` and you want to "
"tell the operator how to fix it."
),
)
def provision_passwordless_ssh_dry_run(
cluster_host: Annotated[
str,
Field(description="Slurm cluster login hostname."),
],
cluster_user: Annotated[
str | None,
Field(
description=("SSH user for the cluster. None uses the local user."),
),
] = None,
identity: Annotated[
str | None,
Field(
description=(
"Path to the SSH private key to inspect. None "
"uses $IDENTITY env var, then ~/.ssh/id_ed25519."
),
),
] = None,
) -> dict:
return bridge.provision_passwordless_ssh_dry_run_impl(
cluster_host,
cluster_user,
identity,
)
@mcp.tool(
name="read_cluster_artifact",
description=(
"Read an artifact from a remote experiment via nemo_run's "
"tunnel. nemo_run already knows the cluster host + user + "
"identity from the executor metadata stored alongside the "
"experiment — this tool does NOT take cluster credentials.\n\n"
"Two modes:\n"
"* `path=None, job_idx=N` → fetch the job's log via "
"`nemo experiment logs <id> <N>`.\n"
"* `path='<rel>'` → read the named relative path inside "
"the remote experiment dir via the executor's tunnel.\n\n"
"Returns the file content truncated to 8 KB (same cap as "
"the launcher's `log_excerpt`). Use this for files like "
"`specbench_results.json` that the launcher writes on the "
"cluster's lustre but the MCP server (running locally) can't "
"directly read."
),
)
def read_cluster_artifact(
experiment_id: Annotated[
str,
Field(description="The experiment id from submit_job."),
],
path: Annotated[
str | None,
Field(
description=(
"Relative path inside the experiment dir. None = "
"log-fetch mode via `nemo experiment logs`."
),
),
] = None,
job_idx: Annotated[
int,
Field(
ge=0,
description=(
"Job index for log-fetch mode. Default 0 = the first task in the pipeline."
),
),
] = 0,
) -> dict:
return bridge.read_cluster_artifact_impl(experiment_id, path, job_idx)
@mcp.tool(
name="open_draft_pr",
description=(
"Push the agent's current branch + open a draft PR on the "
"named target repo. Preconditions enforced by the caller "
"(NOT this tool): the agent's working tree is at the branch "
"it wants to PR, commits exist, and any required DCO "
"`Signed-off-by:` trailer is in place.\n\n"
"Steps internally: `git push -u origin HEAD` + "
"`gh pr create --draft --repo <target> --title --body --base`. "
"Returns `pr_url` parsed from gh's stdout on success, or "
"structured `reason` on failure (git_push_failed, "
"gh_pr_create_failed, etc.)."
),
)
def open_draft_pr(
target_repo: Annotated[
str,
Field(
description=("GitHub repo slug, e.g. 'NVIDIA/Model-Optimizer'."),
),
],
title: Annotated[
str,
Field(description="PR title."),
],
body: Annotated[
str,
Field(description="PR description body (Markdown)."),
],
base_branch: Annotated[
str,
Field(
description=("Base branch to PR against. Default 'main'."),
),
] = "main",
cwd: Annotated[
str | None,
Field(
description=(
"Path to the git checkout. None uses the MCP "
"server's cwd (usually the agent's working dir)."
),
),
] = None,
) -> dict:
return bridge.open_draft_pr_impl(
target_repo,
title,
body,
base_branch,
cwd,
)
return mcp
def main() -> None:
"""Entry point for the `modelopt-mcp` console_script."""
logging.basicConfig(
stream=sys.stderr,
level=os.environ.get("MODELOPT_MCP_LOG", "INFO"),
format="%(asctime)s %(name)s %(levelname)s %(message)s",
)
mcp = _build_server()
mcp.run() # stdio by default
if __name__ == "__main__":
main()
+48
View File
@@ -0,0 +1,48 @@
[project]
name = "modelopt-mcp"
version = "0.1.0"
description = "MCP server exposing ModelOpt launcher operations (submit, status, logs) as typed tools for codex / Claude Code agents"
requires-python = ">=3.10"
dependencies = [
"mcp>=1.0",
# NOTE on modelopt-launcher: `tools/launcher/pyproject.toml` declares
# the package name as `modelopt-launcher` but configures
# `py-modules = []` — there is NO importable `modelopt_launcher`
# Python package on disk. bridge.py invokes the launcher via
# `uv run launch.py` (subprocess) from `<repo>/tools/launcher/` as
# a sibling directory; it does NOT `import modelopt_launcher`.
# Declaring the bare name here would add an unsatisfiable PyPI
# dependency for end users installing via
# `uvx --from "git+...#subdirectory=tools/mcp" modelopt-mcp`. So we
# do NOT declare it. The install relationship is documented in
# README.md as a sibling-checkout layout requirement instead. The
# uvx-from-git path satisfies this naturally because uvx clones
# the whole repo, putting tools/launcher and tools/mcp next to
# each other on disk.
"pyyaml",
"pydantic>=2.0",
]
[project.scripts]
# Console entry. Codex / Claude Code launch this as a stdio subprocess.
# See OMNIML-5123 for the design + acceptance criteria.
modelopt-mcp = "modelopt_mcp.server:main"
[build-system]
requires = ["setuptools>=64"]
build-backend = "setuptools.build_meta"
[tool.setuptools.packages.find]
where = ["."]
include = ["modelopt_mcp*"]
# No [tool.uv.sources] for the launcher — bridge.py uses it via
# `subprocess.run(["uv", "run", "launch.py", ...], cwd=<repo>/tools/launcher/)`,
# so the launcher is a file-layout dependency, not a Python import
# dependency. The uvx-from-git path clones the whole repo so the
# sibling tools/launcher/ ends up on disk automatically. For dev:
# uv pip install -e .
# # then run from a clone where ../launcher exists.
[tool.pytest.ini_options]
testpaths = ["tests"]
+16
View File
@@ -0,0 +1,16 @@
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Unit tests for the modelopt-mcp package."""
+717
View File
@@ -0,0 +1,717 @@
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0
#
# Licensed under the Apache License, Version 2.0 (the "License");
# you may not use this file except in compliance with the License.
# You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
"""Unit tests for modelopt-mcp's bridge module — subprocess + filesystem interactions mocked."""
from __future__ import annotations
import subprocess
import pytest
# Skip the whole module if mcp / pydantic aren't installed (the [mcp]
# extra is opt-in).
pytest.importorskip("mcp")
pytest.importorskip("pydantic")
from modelopt_mcp import bridge
# ---------------------------------------------------------------------------
# list_examples
# ---------------------------------------------------------------------------
def test_list_examples_returns_structured_metadata(tmp_path, monkeypatch):
"""Drop two YAMLs into a fake examples dir and verify metadata extraction (model, description) and path shape."""
examples = tmp_path / "examples"
(examples / "Qwen").mkdir(parents=True)
(examples / "Qwen" / "ptq.yaml").write_text(
"job_name: qwen-ptq\nmodel: Qwen/Qwen3-8B\ndescription: PTQ test\n"
)
(examples / "moonshotai").mkdir(parents=True)
(examples / "moonshotai" / "train.yaml").write_text(
"job_name: kimi-train\nbase_model: moonshotai/Kimi-K2\n"
)
monkeypatch.setenv("MODELOPT_LAUNCHER_EXAMPLES_DIR", str(examples))
result = bridge.list_examples_impl()
assert result["ok"] is True
assert result["count"] == 2
by_model = {e["model"]: e for e in result["examples"]}
assert "Qwen/Qwen3-8B" in by_model
assert by_model["Qwen/Qwen3-8B"]["description"] == "PTQ test"
assert "moonshotai/Kimi-K2" in by_model
def test_list_examples_missing_dir(monkeypatch, tmp_path):
"""When examples dir can't be located, return a structured failure — no exception."""
monkeypatch.setenv("MODELOPT_LAUNCHER_EXAMPLES_DIR", str(tmp_path / "ghost"))
result = bridge.list_examples_impl()
assert result["ok"] is False
assert result["reason"] == "examples_dir_not_found"
def test_list_examples_tolerates_malformed_yaml(tmp_path, monkeypatch):
"""A single malformed YAML doesn't crash list_examples — it lands with model=None."""
examples = tmp_path / "examples"
examples.mkdir()
(examples / "good.yaml").write_text("job_name: g\nmodel: ok\n")
(examples / "bad.yaml").write_text("not: [unbalanced\n")
monkeypatch.setenv("MODELOPT_LAUNCHER_EXAMPLES_DIR", str(examples))
result = bridge.list_examples_impl()
assert result["ok"] is True
assert result["count"] == 2
by_path = {e["path"]: e for e in result["examples"]}
assert any("bad.yaml" in p for p in by_path)
bad = next(e for e in result["examples"] if "bad.yaml" in e["path"])
assert bad["model"] is None
# ---------------------------------------------------------------------------
# verify_docker_setup
# ---------------------------------------------------------------------------
def test_verify_docker_daemon_unavailable(monkeypatch):
"""When `docker info` exits non-zero, verify returns docker_daemon_unavailable."""
def fake_run(argv, **kwargs):
return subprocess.CompletedProcess(
args=argv,
returncode=1,
stdout="",
stderr="Cannot connect to the Docker daemon at unix:///var/run/docker.sock",
)
monkeypatch.setattr(subprocess, "run", fake_run)
monkeypatch.setenv("MODELOPT_MCP_SKIP_GPU_CHECK", "1")
result = bridge.verify_docker_setup_impl()
assert result["ok"] is False
assert result["reason"] == "docker_daemon_unavailable"
def test_verify_docker_daemon_not_installed(monkeypatch):
"""When `docker` is not on PATH, verify returns docker_not_installed."""
def fake_run(argv, **kwargs):
raise FileNotFoundError("docker")
monkeypatch.setattr(subprocess, "run", fake_run)
result = bridge.verify_docker_setup_impl()
assert result["ok"] is False
assert result["reason"] == "docker_not_installed"
def test_verify_docker_skip_gpu_when_env_set(monkeypatch):
"""MODELOPT_MCP_SKIP_GPU_CHECK lets CI hosts without GPUs report ok after the daemon check passes."""
def fake_run(argv, **kwargs):
# Daemon check passes; GPU check is skipped — so only one call.
assert argv[:2] == ["docker", "info"], f"only `docker info` should run; got {argv}"
return subprocess.CompletedProcess(
args=argv,
returncode=0,
stdout="",
stderr="",
)
monkeypatch.setattr(subprocess, "run", fake_run)
monkeypatch.setenv("MODELOPT_MCP_SKIP_GPU_CHECK", "1")
result = bridge.verify_docker_setup_impl()
assert result["ok"] is True
assert result["gpu_check_skipped"] is True
def test_verify_docker_gpu_unavailable(monkeypatch):
"""GPU passthrough container exits non-zero → gpu_unavailable + install-toolkit pointer."""
call_count = {"n": 0}
def fake_run(argv, **kwargs):
call_count["n"] += 1
if call_count["n"] == 1:
# Daemon check: ok
return subprocess.CompletedProcess(
args=argv,
returncode=0,
stdout="",
stderr="",
)
# GPU check: failed
return subprocess.CompletedProcess(
args=argv,
returncode=125,
stdout="",
stderr='could not select device driver "" with capabilities: [[gpu]]',
)
monkeypatch.setattr(subprocess, "run", fake_run)
monkeypatch.delenv("MODELOPT_MCP_SKIP_GPU_CHECK", raising=False)
result = bridge.verify_docker_setup_impl()
assert result["ok"] is False
assert result["reason"] == "gpu_unavailable"
assert "NVIDIA Container Toolkit" in result["diagnostic"]
# ---------------------------------------------------------------------------
# verify_slurm_setup
# ---------------------------------------------------------------------------
def test_verify_slurm_ssh_success(monkeypatch):
"""Mocked ssh probe returns whoami + hostname; verify returns ok."""
def fake_run(argv, **kwargs):
assert argv[0] == "ssh"
assert "BatchMode=yes" in argv
return subprocess.CompletedProcess(
args=argv,
returncode=0,
stdout="chenhany\ncluster-login-01\n",
stderr="",
)
monkeypatch.setattr(subprocess, "run", fake_run)
result = bridge.verify_slurm_setup_impl(
cluster_host="cw-dfw-cs-001-login-01.nvidia.com",
cluster_user="chenhany",
)
assert result["ok"] is True
assert result["whoami"] == "chenhany"
assert result["remote_hostname"] == "cluster-login-01"
def test_verify_slurm_auth_failed(monkeypatch):
"""Ssh -o BatchMode=yes exit 255 → ssh_auth_failed with diagnostic."""
def fake_run(argv, **kwargs):
return subprocess.CompletedProcess(
args=argv,
returncode=255,
stdout="",
stderr="Permission denied (publickey).",
)
monkeypatch.setattr(subprocess, "run", fake_run)
result = bridge.verify_slurm_setup_impl(
cluster_host="ghost-cluster.nvidia.com",
)
assert result["ok"] is False
assert result["reason"] == "ssh_auth_failed"
# ---------------------------------------------------------------------------
# submit_job mode resolution
# ---------------------------------------------------------------------------
def test_submit_job_rejects_no_executor():
"""Neither hf_local nor cluster_host → no_executor_specified."""
result = bridge.submit_job_impl(
yaml_path="examples/test.yaml",
hf_local=None,
cluster_host=None,
cluster_user=None,
identity=None,
job_dir=None,
job_name=None,
extra_overrides=None,
skip_verify=True,
)
assert result["ok"] is False
assert result["reason"] == "no_executor_specified"
def test_submit_job_rejects_both_executors():
"""Both hf_local AND cluster_host → ambiguous_executor."""
result = bridge.submit_job_impl(
yaml_path="examples/test.yaml",
hf_local="/tmp/hf",
cluster_host="cluster.nvidia.com",
cluster_user=None,
identity=None,
job_dir=None,
job_name=None,
extra_overrides=None,
skip_verify=True,
)
assert result["ok"] is False
assert result["reason"] == "ambiguous_executor"
def test_submit_job_yaml_not_found(monkeypatch, tmp_path):
"""yaml_path that doesn't resolve to an existing file → yaml_not_found."""
monkeypatch.setenv("MODELOPT_LAUNCHER_EXAMPLES_DIR", str(tmp_path))
result = bridge.submit_job_impl(
yaml_path="does/not/exist.yaml",
hf_local="/tmp/hf",
cluster_host=None,
cluster_user=None,
identity=None,
job_dir=None,
job_name=None,
extra_overrides=None,
skip_verify=True,
)
assert result["ok"] is False
assert result["reason"] == "yaml_not_found"
# ---------------------------------------------------------------------------
# job_status / job_logs — filesystem-based
# ---------------------------------------------------------------------------
def test_job_status_done_success(tmp_path, monkeypatch):
"""_DONE marker + all task statuses succeeded → status='done'."""
exp = tmp_path / "experiments" / "exp_1781000000"
exp.mkdir(parents=True)
(exp / "_DONE").touch()
(exp / "status_task_0.out").write_text("succeeded\n")
(exp / "status_task_1.out").write_text("succeeded\n")
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
result = bridge.job_status_impl("exp_1781000000")
assert result["ok"] is True
assert result["status"] == "done"
assert result["task_statuses"] == {"task_0": "succeeded", "task_1": "succeeded"}
def test_job_status_failed_task(tmp_path, monkeypatch):
"""_DONE marker + at least one task status contains 'fail' → status='failed'."""
exp = tmp_path / "experiments" / "exp_1781000001"
exp.mkdir(parents=True)
(exp / "_DONE").touch()
(exp / "status_task_0.out").write_text("succeeded\n")
(exp / "status_task_1.out").write_text("failed (rc=1)\n")
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
result = bridge.job_status_impl("exp_1781000001")
assert result["ok"] is True
assert result["status"] == "failed"
assert "failed" in result["task_statuses"]["task_1"]
def test_job_status_running(tmp_path, monkeypatch):
"""No _DONE marker → running."""
exp = tmp_path / "experiments" / "exp_1781000002"
exp.mkdir(parents=True)
(exp / "status_task_0.out").write_text("running\n")
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
result = bridge.job_status_impl("exp_1781000002")
assert result["ok"] is True
assert result["status"] == "running"
assert result["has_done_marker"] is False
def test_job_status_unknown_id(tmp_path, monkeypatch):
"""No experiment dir matching the id → experiment_dir_not_found."""
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
result = bridge.job_status_impl("does_not_exist")
assert result["ok"] is False
assert result["reason"] == "experiment_dir_not_found"
def test_job_logs_all_tasks(tmp_path, monkeypatch):
"""task=None returns logs for every log_*.out under the experiment dir."""
exp = tmp_path / "experiments" / "exp_1781000003"
exp.mkdir(parents=True)
(exp / "log_task_0.out").write_text("hello\nworld\n")
(exp / "log_task_1.out").write_text("done\n")
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
result = bridge.job_logs_impl("exp_1781000003", task=None, tail=None)
assert result["ok"] is True
assert set(result["logs"].keys()) == {"task_0", "task_1"}
assert "hello" in result["logs"]["task_0"]
def test_job_logs_with_tail(tmp_path, monkeypatch):
"""tail=N returns only the last N lines per task."""
exp = tmp_path / "experiments" / "exp_1781000004"
exp.mkdir(parents=True)
(exp / "log_task_0.out").write_text("line1\nline2\nline3\nline4\n")
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
result = bridge.job_logs_impl("exp_1781000004", task="task_0", tail=2)
assert result["ok"] is True
body = result["logs"]["task_0"]
assert body.splitlines() == ["line3", "line4"]
def test_job_logs_missing_task(tmp_path, monkeypatch):
"""Requested task name has no log file → task_log_not_found."""
exp = tmp_path / "experiments" / "exp_1781000005"
exp.mkdir(parents=True)
(exp / "log_task_0.out").write_text("only task 0\n")
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
result = bridge.job_logs_impl("exp_1781000005", task="task_99", tail=None)
assert result["ok"] is False
assert result["reason"] == "task_log_not_found"
# ---------------------------------------------------------------------------
# wait_for_experiment
# ---------------------------------------------------------------------------
def test_wait_for_experiment_returns_terminal_immediately(tmp_path, monkeypatch):
"""If the experiment is already terminal, return without polling."""
exp = tmp_path / "experiments" / "exp_already_done"
exp.mkdir(parents=True)
(exp / "_DONE").touch()
(exp / "status_task_0.out").write_text("succeeded\n")
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
result = bridge.wait_for_experiment_impl(
"exp_already_done",
timeout_sec=10,
poll_interval_sec=1,
)
assert result["ok"] is True
assert result["status"] == "done"
assert result["waited_seconds"] < 1 # didn't actually wait
def test_wait_for_experiment_polls_until_done(tmp_path, monkeypatch):
"""Spin through running → done."""
exp = tmp_path / "experiments" / "exp_in_flight"
exp.mkdir(parents=True)
(exp / "status_task_0.out").write_text("running\n")
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
# Flip the marker after 2 polls via a counter side-effect
call_count = {"n": 0}
real_status = bridge.job_status_impl
def fake_status(experiment_id):
call_count["n"] += 1
if call_count["n"] >= 2:
(exp / "_DONE").touch()
return real_status(experiment_id)
monkeypatch.setattr(bridge, "job_status_impl", fake_status)
result = bridge.wait_for_experiment_impl(
"exp_in_flight",
timeout_sec=10,
poll_interval_sec=0,
)
assert result["ok"] is True
assert result["status"] == "done"
assert call_count["n"] >= 2
def test_wait_for_experiment_timeout(tmp_path, monkeypatch):
"""Never reaches terminal → wait_timeout with last_status."""
exp = tmp_path / "experiments" / "exp_stuck"
exp.mkdir(parents=True)
(exp / "status_task_0.out").write_text("running\n")
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
result = bridge.wait_for_experiment_impl(
"exp_stuck",
timeout_sec=1,
poll_interval_sec=0,
)
assert result["ok"] is False
assert result["reason"] == "wait_timeout"
assert result["last_status"]["status"] == "running"
def test_wait_for_experiment_passes_through_dir_not_found(tmp_path, monkeypatch):
"""If the experiment dir doesn't exist, don't spin to timeout."""
monkeypatch.setenv("NEMORUN_HOME", str(tmp_path))
result = bridge.wait_for_experiment_impl(
"does_not_exist",
timeout_sec=60,
poll_interval_sec=0,
)
assert result["ok"] is False
assert result["reason"] == "experiment_dir_not_found"
# Bail out fast — not the full timeout
assert result["waited_seconds"] < 1
# ---------------------------------------------------------------------------
# provision_passwordless_ssh_dry_run
# ---------------------------------------------------------------------------
def test_provision_ssh_no_key_emits_keygen(tmp_path, monkeypatch):
"""Missing private key → keygen command, step='keygen_required'."""
fake_home = tmp_path / "home"
(fake_home / ".ssh").mkdir(parents=True)
monkeypatch.setenv("HOME", str(fake_home))
monkeypatch.delenv("IDENTITY", raising=False)
result = bridge.provision_passwordless_ssh_dry_run_impl(
cluster_host="cw-dfw.example.com",
cluster_user="alice",
identity=None,
)
assert result["ok"] is True
assert result["step"] == "keygen_required"
assert "ssh-keygen" in result["commands"][0]
assert "ed25519" in result["commands"][0]
assert result["identity_path"].endswith(".ssh/id_ed25519")
def test_provision_ssh_key_present_emits_copy_id(tmp_path, monkeypatch):
"""Key + pubkey present → ssh-copy-id command + pubkey content."""
fake_home = tmp_path / "home"
ssh = fake_home / ".ssh"
ssh.mkdir(parents=True)
(ssh / "id_ed25519").write_text("PRIVKEY")
pubkey_content = "ssh-ed25519 AAAAC3NzaC... alice@host"
(ssh / "id_ed25519.pub").write_text(pubkey_content + "\n")
monkeypatch.setenv("HOME", str(fake_home))
monkeypatch.delenv("IDENTITY", raising=False)
result = bridge.provision_passwordless_ssh_dry_run_impl(
cluster_host="cw-dfw.example.com",
cluster_user="alice",
identity=None,
)
assert result["ok"] is True
assert result["step"] == "ssh_copy_id_required"
assert "ssh-copy-id" in result["commands"][0]
assert "alice@cw-dfw.example.com" in result["commands"][0]
assert result["pubkey"] == pubkey_content
assert "verify_setup" in result["next_check"]
def test_provision_ssh_priv_without_pub_surfaces_failure(tmp_path, monkeypatch):
"""Private key but no .pub → pubkey_missing with recovery hint."""
fake_home = tmp_path / "home"
ssh = fake_home / ".ssh"
ssh.mkdir(parents=True)
(ssh / "id_ed25519").write_text("PRIVKEY")
monkeypatch.setenv("HOME", str(fake_home))
monkeypatch.delenv("IDENTITY", raising=False)
result = bridge.provision_passwordless_ssh_dry_run_impl(
cluster_host="cw-dfw.example.com",
cluster_user="alice",
identity=None,
)
assert result["ok"] is False
assert result["reason"] == "pubkey_missing"
assert "ssh-keygen" in result["diagnostic"]
def test_provision_ssh_explicit_identity_overrides_default(tmp_path, monkeypatch):
"""Explicit identity arg wins over $IDENTITY and ~/.ssh/id_ed25519."""
explicit = tmp_path / "custom_key"
explicit.write_text("CUSTOM")
(tmp_path / "custom_key.pub").write_text("ssh-ed25519 AAAA alice\n")
monkeypatch.setenv("IDENTITY", "/wrong/path") # should be ignored
result = bridge.provision_passwordless_ssh_dry_run_impl(
cluster_host="cw-dfw.example.com",
cluster_user="alice",
identity=str(explicit),
)
assert result["ok"] is True
assert result["identity_path"] == str(explicit)
# ---------------------------------------------------------------------------
# read_cluster_artifact — logs mode (subprocess mocked)
# ---------------------------------------------------------------------------
def test_read_cluster_artifact_logs_mode_ok(monkeypatch):
"""path=None → wraps `nemo experiment logs <id> <job_idx>`."""
captured = {}
def fake_run(argv, **kwargs):
captured["argv"] = argv
return subprocess.CompletedProcess(
args=argv,
returncode=0,
stdout="line 1\nline 2\nline 3\n",
stderr="",
)
monkeypatch.setattr(subprocess, "run", fake_run)
result = bridge.read_cluster_artifact_impl(
experiment_id="cicd_42",
path=None,
job_idx=0,
)
assert result["ok"] is True
assert result["mode"] == "logs"
assert "line 1" in result["content"]
# Verify the wrapped command
assert "nemo" in captured["argv"]
assert "experiment" in captured["argv"]
assert "logs" in captured["argv"]
assert "cicd_42" in captured["argv"]
def test_read_cluster_artifact_logs_mode_subprocess_failed(monkeypatch):
"""Nemo cli non-zero → structured logs_fetch_failed."""
def fake_run(argv, **kwargs):
return subprocess.CompletedProcess(
args=argv,
returncode=1,
stdout="",
stderr="experiment cicd_42 not found\n",
)
monkeypatch.setattr(subprocess, "run", fake_run)
result = bridge.read_cluster_artifact_impl(
experiment_id="cicd_42",
path=None,
job_idx=0,
)
assert result["ok"] is False
assert result["reason"] == "logs_fetch_failed"
assert result["exit_code"] == 1
def test_read_cluster_artifact_logs_mode_timeout(monkeypatch):
"""Hanging tunnel → structured logs_fetch_timeout, no exception."""
def fake_run(argv, **kwargs):
raise subprocess.TimeoutExpired(cmd=argv, timeout=60)
monkeypatch.setattr(subprocess, "run", fake_run)
result = bridge.read_cluster_artifact_impl(
experiment_id="cicd_42",
path=None,
job_idx=0,
)
assert result["ok"] is False
assert result["reason"] == "logs_fetch_timeout"
# ---------------------------------------------------------------------------
# open_draft_pr — subprocess mocked
# ---------------------------------------------------------------------------
def test_open_draft_pr_happy_path(monkeypatch, tmp_path):
"""Git push ok → gh pr create ok → returns parsed pr_url."""
(tmp_path / ".git").mkdir() # pretend it's a git repo
call_log = []
def fake_run(argv, **kwargs):
call_log.append(argv[0])
if argv[0] == "git":
return subprocess.CompletedProcess(
args=argv,
returncode=0,
stdout="",
stderr="",
)
# gh pr create — return a typical PR URL on stdout
return subprocess.CompletedProcess(
args=argv,
returncode=0,
stdout=(
"Warning: 1 uncommitted change\n"
"https://github.com/NVIDIA/Model-Optimizer/pull/9999\n"
),
stderr="",
)
monkeypatch.setattr(subprocess, "run", fake_run)
result = bridge.open_draft_pr_impl(
target_repo="NVIDIA/Model-Optimizer",
title="test pr",
body="body text",
base_branch="main",
cwd=str(tmp_path),
)
assert result["ok"] is True
assert result["pr_url"] == "https://github.com/NVIDIA/Model-Optimizer/pull/9999"
assert call_log == ["git", "gh"]
def test_open_draft_pr_not_a_git_repo(tmp_path):
"""Cwd without .git → structured not_a_git_repo failure, no subprocess."""
result = bridge.open_draft_pr_impl(
target_repo="NVIDIA/Model-Optimizer",
title="x",
body="x",
base_branch="main",
cwd=str(tmp_path),
)
assert result["ok"] is False
assert result["reason"] == "not_a_git_repo"
def test_open_draft_pr_git_push_failed(monkeypatch, tmp_path):
"""Git push non-zero → structured git_push_failed."""
(tmp_path / ".git").mkdir()
def fake_run(argv, **kwargs):
if argv[0] == "git":
return subprocess.CompletedProcess(
args=argv,
returncode=128,
stdout="",
stderr="rejected: non-fast-forward\n",
)
pytest.fail(f"gh should not be called when git push fails; got {argv}")
monkeypatch.setattr(subprocess, "run", fake_run)
result = bridge.open_draft_pr_impl(
target_repo="NVIDIA/Model-Optimizer",
title="x",
body="x",
base_branch="main",
cwd=str(tmp_path),
)
assert result["ok"] is False
assert result["reason"] == "git_push_failed"
def test_open_draft_pr_gh_failed_but_branch_pushed(monkeypatch, tmp_path):
"""Gh pr create non-zero → reports branch_pushed=True so the operator can retry just the PR-open step."""
(tmp_path / ".git").mkdir()
def fake_run(argv, **kwargs):
if argv[0] == "git":
return subprocess.CompletedProcess(
args=argv,
returncode=0,
stdout="",
stderr="",
)
return subprocess.CompletedProcess(
args=argv,
returncode=1,
stdout="",
stderr="resource not found: NVIDIA/Model-Optimizer\n",
)
monkeypatch.setattr(subprocess, "run", fake_run)
result = bridge.open_draft_pr_impl(
target_repo="NVIDIA/Model-Optimizer",
title="x",
body="x",
base_branch="main",
cwd=str(tmp_path),
)
assert result["ok"] is False
assert result["reason"] == "gh_pr_create_failed"
assert result["branch_pushed"] is True