Skip to content

A Job runs a command to completion. Nodus picks the capacity, keeps its checkpoints and recovers it if the capacity is reclaimed; max_cost caps what it may spend.

examples/python/jobs/main.py
"""Upload this directory as a Job's source, follow its logs and download its output.
Run it with `python examples/python/jobs/main.py`.
"""
from pathlib import Path
import nodus
HERE = Path(__file__).parent
def main() -> None:
job = nodus.Job.run(
image="nodus/python:3.12",
command=["python", "train.py"],
source=HERE, # uploaded once as a content-addressed blob; .gitignore and .nodusignore apply
cpu=2,
memory="4Gi",
timeout="30m",
max_cost=1,
outputs={"report": "/nodus/outputs/report.txt"},
)
print("estimate:", job.estimate())
for line in job.logs(follow=True):
print(line, end="")
job.wait() # raises nodus.errors.JobFailed with the exit code and log tail
print("saved", job.outputs["report"].download(HERE / "report.txt"))
if __name__ == "__main__":
main()
job = nodus.Job.run(
name="finetune-llama", # optional; a generated name otherwise
image="nodus/pytorch:2.8-cuda12.8",
command=["python", "train.py", "--epochs", "3"],
source=".", # uploads the directory; or {"repo": "acme/trainer", "ref": "main"}
gpu="H100", secrets=["hf-token"],
max_cost=40, timeout="12h", expected_duration="6h",
checkpoint="/nodus/state", interruptible=True, region=["us", "eu"],
outputs={"adapter": "/nodus/outputs/adapter"},
)

The source upload skips what .gitignore and .nodusignore exclude and refuses more than 500 MiB compressed unless you pass allow_large_source=True. Everything written under /nodus/outputs is collected when the Job succeeds.

print(job.estimate()) # expected cost p50 and p90, start time and the hold
for line in job.logs(follow=True):
print(line, end="")
job.wait() # raises nodus.errors.JobFailed(job, exit_code, log_tail)
job.outputs["adapter"].download("./adapter") # checks the sha256, then renames into place

job.suspend(), job.resume() and job.cancel() change the Job’s state; job.attempts() lists every attempt with its placement; job.exec("nvidia-smi") runs a command in the running container.

ddp = nodus.Job.run(
name="ddp-smoke", image="nodus/pytorch:2.8-cuda12.8", command=["torchrun", "train.py"],
gpu="H100", max_cost=10,
distributed=nodus.Distributed(nodes=2, launcher="Torchrun", network="Global", transport="Auto"),
)
for line in ddp.logs(follow=True, rank="all"): # every rank, prefixed [r0], [r1]
print(line, end="")

Every rank gets the rendezvous environment (MASTER_ADDR, WORLD_SIZE, torchrun’s PET_* variables), and nodus.cluster.info() reads it for you. nodus.checkpoint.dcp.save(state, step) and .load(state) write and read torch.distributed.checkpoint shards to the gang’s checkpoint storage.