diff --git a/Snakefile b/Snakefile index fa22badd8..c0613626b 100644 --- a/Snakefile +++ b/Snakefile @@ -1,28 +1,35 @@ """Top-level Snakefile -Usage, from the root directory: - snakemake -n # dry-run - snakemake --profile pipeline/profiles/ldg # HTCondor - snakemake --profile pipeline/profiles/local # local execution (dev) +Run from a run directory made by `aframe-init snakemake`, whose run.sh +calls, from that directory: -The pipeline config is loaded from pipeline/config/config.yaml by -default. To use a different config, copy the original, make -modifications, and run: + snakemake --snakefile /Snakefile --configfile config.yaml \ + --profile /pipeline/profiles/ldg - snakemake --configfile my_run.yaml --profile pipeline/profiles/ldg +Each run has its own .snakemake/ directory and run_dir defaults to the +working directory. Settings that the run's config doesn't give come from +pipeline/config/config.yaml. A relative train_config is relative to the +repo, like the presets in pipeline/config/. """ +import os from pathlib import Path +REPO = Path(workflow.basedir) -configfile: "pipeline/config/config.yaml" +configfile: str(REPO / "pipeline" / "config" / "config.yaml") -if "run_dir" not in config: - raise WorkflowError( - "'run_dir' must be set in your config. " - "Pass it via --configfile or --config run_dir=..." - ) + +if config["run_dir"] is None: + # Prevent runs from writing into the repo + if Path.cwd().resolve() == REPO.resolve(): + raise WorkflowError( + "Set run_dir (--config run_dir=...) or run from a run " + "directory made by `aframe-init snakemake`" + ) + config["run_dir"] = os.getcwd() +config["train_config"] = str(REPO / config["train_config"]) config.setdefault("background_dir", str(Path(config["run_dir"]) / "data")) config.setdefault("waveforms_dir", str(Path(config["run_dir"]) / "waveforms")) @@ -40,7 +47,12 @@ include: "projects/infer/infer.smk" include: "projects/plots/plots.smk" +# Outside of `onstart` so that dry-runs do the check +check_gpus() + + onstart: + set_container_binds() check_images() check_triton_image() diff --git a/pipeline/.env.example b/pipeline/.env.example index 706de9d44..4ceb23be5 100644 --- a/pipeline/.env.example +++ b/pipeline/.env.example @@ -10,9 +10,10 @@ # Directory holding the images built by scripts/build_containers.py # export AFRAME_CONTAINER_ROOT=${HOME}/aframe/images -# Comma-separated list of directories bound into every container. -# Must cover run_dir, background_dir, and waveforms_dir. -# export AFRAME_DATA_DIRS=${HOME}/aframe/run_dir,${HOME}/aframe/data_dir +# Optional: comma-separated extra directories to bind into every container, +# The run's own directories (run_dir, background_dir, waveforms_dir, log_dir) +# are bound automatically. +# export AFRAME_DATA_DIRS= # --- S3 (training on Nautilus; optional) -------------------------------------- # export AWS_ENDPOINT_URL=https://s3-west.nrp-nautilus.io diff --git a/pipeline/config/config.yaml b/pipeline/config/config.yaml index 795c3f075..5ab577ea8 100644 --- a/pipeline/config/config.yaml +++ b/pipeline/config/config.yaml @@ -7,9 +7,9 @@ # ============================================================================= # --- Output directories ------------------------------------------------------ -# REQUIRED: experiment-specific outputs (train/, export/, infer/, plots/) are -# written here. There is no default, so set this in your run config. -run_dir: /path/to/run +# Experiment-specific outputs (train/, export/, infer/, plots/) are written +# here. null for the directory snakemake runs from. +run_dir: null # Background strain data directory, containing train/ and test/ # subdirectories of fetched HDF5 files. Defaults to {run_dir}/data if @@ -139,16 +139,22 @@ resources: # Slurm GPU partition for training. train_partition: gpuA40x4 -# GPUs for training, which override trainer.devices in the train config. -# On a shared node, pin specific devices with train_gpus (comma-separated -# IDs, e.g. "2" or "2,3"), since Lightning would otherwise take every -# visible GPU. Leave it null under slurm, which assigns train_num_gpus. -train_gpus: null +# Which and how many GPUs for training, which overrides trainer.devices in +# the train config. On a shared node, pin train_num_gpus devices with +# train_gpus (comma-separated IDs, e.g. "2" or "2,3"), or use `auto` for the +# least-used GPUs when training starts, since Lightning would otherwise take +# every visible GPU. Leave it null under slurm, which assigns the GPUs. +train_gpus: auto train_num_gpus: 1 # Slurm GPU partition for export and inference. inference_partition: gpuA40x4 +# GPUs for export and the Triton server. Pass inference_num_gpus +# comma-separated IDs, or `auto` for the least-used GPUs when each starts. +inference_gpus: auto +inference_num_gpus: 1 + # Condor GPU matchmaking for in-process inference. # gpu_min_memory_mb: null sets no memory minimum. gpu_min_capability: 7.0 @@ -165,6 +171,7 @@ epnfs: false # --- Training ---------------------------------------------------------------- # Lightning CLI YAML config for `train fit`. +# Relative to the repo if not absolute. train_config: projects/train/train.yaml # Submit training to Nautilus via train-remote instead of @@ -225,9 +232,6 @@ model_name: aframe-stream model_version: -1 # Resolved against $AFRAME_CONTAINER_ROOT unless given as an absolute path. triton_image: tritonserver_25.06.sif -# Comma-separated GPU IDs for export and inference server only. -# Training uses train_gpus. -gpus: "0" # Snapshotter instances hosted per GPU. Baked into the model at export time # and used to size the concurrent-inference pool (streams_per_gpu * num_gpus) diff --git a/pipeline/config/small.yaml b/pipeline/config/small.yaml index 30890bab8..fc6378610 100644 --- a/pipeline/config/small.yaml +++ b/pipeline/config/small.yaml @@ -1,7 +1,5 @@ # A small-scale run that exercises every stage, for checking that the -# pipeline runs. Layer it over config.yaml: -# -# snakemake --configfile pipeline/config/small.yaml --config run_dir=... +# pipeline runs. Start one with `aframe-init snakemake --preset small`. # Continuous H1+L1 open data in O3a, four background files per split train_start: 1250917000 diff --git a/pipeline/profiles/delta/config.yaml b/pipeline/profiles/delta/config.yaml index d8bf3822e..874b37b78 100644 --- a/pipeline/profiles/delta/config.yaml +++ b/pipeline/profiles/delta/config.yaml @@ -5,8 +5,8 @@ # Requirements on the submit node: # - The aframe-smk conda environment (pipeline/envs/snakemake.yaml) # - AFRAME_CONTAINER_ROOT pointing at the built .sif images -# - AFRAME_DATA_DIRS: comma-separated directories to bind into -# containers (run_dir, background_dir, waveforms_dir parents) +# - AFRAME_DATA_DIRS (optional): comma-separated extra directories to +# bind into containers; the run's own directories are added at startup # - WANDB_API_KEY if using W&B logging # @@ -33,8 +33,8 @@ software-deployment-method: apptainer # The local repo is bound to /opt/aframe apptainer-args: >- --nv - --bind ${AFRAME_DATA_DIRS:?must be set, see pipeline/.env.example} - --bind $PWD:/opt/aframe + --bind ${AFRAME_DATA_DIRS:?set at startup by set_container_binds in pipeline/resources.smk} + --bind ${AFRAME_REPO:?set at startup by set_container_binds in pipeline/resources.smk}:/opt/aframe --home $HOME default-resources: diff --git a/pipeline/profiles/ldg/config.yaml b/pipeline/profiles/ldg/config.yaml index 23dbe864e..861627697 100644 --- a/pipeline/profiles/ldg/config.yaml +++ b/pipeline/profiles/ldg/config.yaml @@ -7,8 +7,8 @@ # - A valid LIGO SciToken (htgettoken -a vault.ligo.org) # - LIGO_GROUP and LIGO_USERNAME set for accounting # - AFRAME_CONTAINER_ROOT pointing at the built .sif images -# - AFRAME_DATA_DIRS: comma-separated directories to bind into -# containers (run_dir, background_dir, waveforms_dir parents) +# - AFRAME_DATA_DIRS (optional): comma-separated extra directories to +# bind into containers; the run's own directories are added at startup # - WANDB_API_KEY if using W&B logging # @@ -49,18 +49,18 @@ local-cores: 16 software-deployment-method: apptainer -# $AFRAME_DATA_DIRS binds the run/data/waveform dirs into the container. -# If it's unset, the rule fails with a message saying so, rather than -# apptainer misreading the next flag as the bind path. -# The local repo ($PWD, since snakemake runs from the repo root) is bound to -# /opt/aframe, where the images install the projects in editable mode, so -# code always comes from the local repo. +# $AFRAME_DATA_DIRS binds the run's directories into the container. +# $AFRAME_REPO, the local repo, is bound to /opt/aframe, where the images +# install the projects in editable mode, so code always comes from the local +# repo. The workflow sets both at startup via `set_container_binds` in +# pipeline/resources.smk. If one is somehow unset, the rule fails with a +# message so that apptainer doesn't misread the next flag. # Snakemake invokes apptainer with `--home `, which sets $HOME # to the repo root and creates caches there. A `--home $HOME` overrides it. apptainer-args: >- --nv - --bind ${AFRAME_DATA_DIRS:?must be set, see pipeline/.env.example} - --bind $PWD:/opt/aframe + --bind ${AFRAME_DATA_DIRS:?set at startup by set_container_binds in pipeline/resources.smk} + --bind ${AFRAME_REPO:?set at startup by set_container_binds in pipeline/resources.smk}:/opt/aframe --home $HOME default-resources: diff --git a/pipeline/profiles/local/config.yaml b/pipeline/profiles/local/config.yaml deleted file mode 100644 index e27cd5355..000000000 --- a/pipeline/profiles/local/config.yaml +++ /dev/null @@ -1,13 +0,0 @@ -# Snakemake local execution profile. -# -# snakemake --profile pipeline/profiles/local -# -# Runs every rule as a subprocess locally. Should be use with -# the aframe-dev conda environment (pipeline/envs/dev.yaml) -# with the uv dependencies installed so that all project CLIs -# are available. -# -# Intended for dry-runs (-n) or single-rule testing (--until ). - -cores: all -rerun-incomplete: true diff --git a/pipeline/resources.smk b/pipeline/resources.smk index b9d8481ef..00511e1df 100644 --- a/pipeline/resources.smk +++ b/pipeline/resources.smk @@ -6,11 +6,14 @@ and reads `htcondor_request_mem_mb`, `allowed_execute_duration`, and `request_gpus` plus its GPU matchmaking keys. Each helper returns both slurm and htcondor sets. -Also resolves each project's container image. +Also resolves each project's container image, checks the image hash +against the source code's, and sets up run's directories to be +bound into the image. """ import os import subprocess +from pathlib import Path from snakemake.exceptions import WorkflowError from snakemake.logging import logger @@ -35,6 +38,39 @@ def container(project): return os.path.join(os.getenv("AFRAME_CONTAINER_ROOT", ""), f"{project}.sif") +def set_container_binds(): + """Create the run's directories and set what the profiles' + `apptainer-args` bind into each container. + + `AFRAME_REPO` is the local repo, which bound over the code in the image. + `AFRAME_DATA_DIRS` is `run_dir`, `background_dir`, `waveforms_dir`, + `log_dir` and possible `rnp_frame_dir`, plus whatever the variable + already held. + """ + # REPO is defined in Snakefile + os.environ["AFRAME_REPO"] = str(REPO) + run_dirs = [ + config[key] for key in ("run_dir", "background_dir", "waveforms_dir", "log_dir") + ] + for path in run_dirs: + os.makedirs(path, exist_ok=True) + dirs = os.getenv("AFRAME_DATA_DIRS", "").split(",") + run_dirs + if config.get("rnp_frame_dir"): + dirs.append(config["rnp_frame_dir"]) + + # Absolute paths, sorted by their number of parts + paths = sorted( + {Path(os.path.abspath(d)) for d in dirs if d}, key=lambda p: len(p.parts) + ) + binds = [] + for path in paths: + # Skip if the parent will be bound + if not any(path.is_relative_to(parent) for parent in binds): + binds.append(path) + os.environ["AFRAME_DATA_DIRS"] = ",".join(str(path) for path in binds) + logger.info(f"Binding {os.environ['AFRAME_DATA_DIRS']} into containers") + + def check_images(projects=("data", "train", "export", "infer", "plots")): """Fail on missing published images, and warn about local images built for a different environment than the local repo's. @@ -71,6 +107,32 @@ def check_images(projects=("data", "train", "export", "infer", "plots")): ) +def gpu_env(gpus, num_gpus): + """Shell command that sets CUDA_VISIBLE_DEVICES for a rule. Either uses + `gpus` as given, or picks the `num_gpus` least-used GPUs with `auto` via + scripts/free_gpus.py, Empty for null, which uses all visible GPUs. + """ + if gpus is None: + return "" + if gpus == "auto": + gpus = f"$(python /opt/aframe/scripts/free_gpus.py {num_gpus})" + return f"CUDA_VISIBLE_DEVICES={gpus}; export CUDA_VISIBLE_DEVICES; " + + +def check_gpus(): + """Check that pinned GPU lists name as many GPUs as their counts.""" + for prefix in ("train", "inference"): + gpus = config[f"{prefix}_gpus"] + num_gpus = config[f"{prefix}_num_gpus"] + if gpus is None or gpus == "auto": + continue + if len(str(gpus).split(",")) != num_gpus: + raise WorkflowError( + f"{prefix}_gpus ({gpus}) should list {prefix}_num_gpus " + f"({num_gpus}) GPUs" + ) + + def rule_resources(name): """Memory and walltime for rule `name`, from the config's `resources`. diff --git a/projects/export/export.smk b/projects/export/export.smk index 4085948d2..e2487f10d 100644 --- a/projects/export/export.smk +++ b/projects/export/export.smk @@ -50,6 +50,7 @@ snakemake is invoked. psd_length=config["psd_length"], highpass=config["highpass"], streams_per_gpu=config["streams_per_gpu"], + gpu_env=gpu_env(config["inference_gpus"], 1), weights=( (config["remote_run_dir"] + "/model_exported.pt2") if remote_train @@ -61,7 +62,7 @@ snakemake is invoked. else str(train_out / "batch.hdf5") ), shell: - "python -m export" + "{params.gpu_env}python -m export" " --weights {params.weights}" " --batch_file {params.batch_file}" " --repository_directory {output.model_repo}" diff --git a/projects/infer/infer.smk b/projects/infer/infer.smk index 040bf1fc1..84407e32b 100644 --- a/projects/infer/infer.smk +++ b/projects/infer/infer.smk @@ -311,7 +311,7 @@ _GROUP_SHELL_SUFFIX = ( if INFERENCE_MODE == "triton": - num_gpus = len(str(config["gpus"]).split(",")) + num_gpus = config["inference_num_gpus"] streams_per_gpu = config["streams_per_gpu"] workflow.global_resources["triton_streams"] = streams_per_gpu * num_gpus @@ -341,9 +341,12 @@ if INFERENCE_MODE == "triton": logfile=str(triton_dir / "server.log"), model_name=config["model_name"], model_version=config["model_version"], - gpus=config["gpus"], + gpus=config["inference_gpus"], + num_gpus=num_gpus, + free_gpus=str(REPO / "scripts" / "free_gpus.py"), batch_size=config["inference_batch_size"], triton_image=TRITON_IMAGE, + infer_project=str(REPO / "projects" / "infer"), idle_timeout=config["triton_idle_timeout"], script: "scripts/start_triton.py" diff --git a/projects/infer/scripts/start_triton.py b/projects/infer/scripts/start_triton.py index 4951178f5..f4ddd022b 100644 --- a/projects/infer/scripts/start_triton.py +++ b/projects/infer/scripts/start_triton.py @@ -33,6 +33,13 @@ Path(params.output_dir).mkdir(parents=True, exist_ok=True) +gpus = str(params.gpus) +if gpus == "auto": + gpus = subprocess.check_output( + [sys.executable, params.free_gpus, str(params.num_gpus)], text=True + ).strip() +logging.info(f"Serving on GPUs {gpus}") + # Clear stale sentinel files before launching a fresh server. ip_file.unlink(missing_ok=True) Path(params.stop_sentinel).unlink(missing_ok=True) @@ -43,7 +50,7 @@ "uv", "run", "--directory", - "projects/infer", + params.infer_project, "start-server", "--model_repo_dir", snakemake.input.model_repo, @@ -54,7 +61,7 @@ "--model_version", str(params.model_version), "--gpus", - str(params.gpus), + gpus, "--batch_size", str(params.batch_size), "--triton_image", diff --git a/projects/train/train.smk b/projects/train/train.smk index ba8ee92d3..8a3a96266 100644 --- a/projects/train/train.smk +++ b/projects/train/train.smk @@ -21,15 +21,11 @@ train_log_dir = log_dir / "train" TRAIN_CONTAINER = container("train") -# GPUs for local training. `train_gpus` pins specific devices on a shared -# node; otherwise use `train_num_gpus` of whatever is visible, which under -# slurm is the allocation. -if config["train_gpus"] is not None: - TRAIN_GPU_ENV = f"CUDA_VISIBLE_DEVICES={config['train_gpus']} " - TRAIN_NUM_GPUS = len(str(config["train_gpus"]).split(",")) -else: - TRAIN_GPU_ENV = "" - TRAIN_NUM_GPUS = config["train_num_gpus"] +# GPUs for training. Uses `train_num_gpus` of them, determined by `train_gpus` +# on a shared node or the least used with `auto`. `null` uses whatever is +# visible, which for slurm is the allocation. +TRAIN_GPU_ENV = gpu_env(config["train_gpus"], config["train_num_gpus"]) +TRAIN_NUM_GPUS = config["train_num_gpus"] def _train_waveform_inputs(wildcards): diff --git a/scripts/aframe_init.py b/scripts/aframe_init.py index 428537729..897030b3f 100644 --- a/scripts/aframe_init.py +++ b/scripts/aframe_init.py @@ -8,6 +8,14 @@ root = Path(__file__).resolve().parent.parent +# Presets for `aframe-init snakemake --preset` +PRESETS = { + "small": ( + root / "pipeline" / "config" / "small.yaml", + root / "pipeline" / "config" / "small_train.yaml", + ), +} + ONLINE_CONFIGS = [ root / "projects" / "online" / "config.yaml", root / "projects" / "online" / "prior.yaml", @@ -40,13 +48,15 @@ def write_content(content: str, path: Path): return content -def create_snakemake_runfile(path: Path, profile: str): - config = path / "config.yaml" - cmd = f"snakemake --configfile {config} --profile {profile}" +def create_snakemake_runfile(path: Path, profile: Path): + cmd = ( + f"snakemake --snakefile {root}/Snakefile" + f" --configfile config.yaml --profile {profile}" + ) content = f""" #!/bin/bash - cd {root} - source pipeline/.env + cd {path} + [ -f {root}/pipeline/.env ] && source {root}/pipeline/.env {cmd} """ runfile = path / "run.sh" @@ -168,7 +178,16 @@ def main(): "--profile", type=str, default="pipeline/profiles/ldg", - help="Path to the snakemake profile directory", + help="Path to the snakemake profile directory, relative to the repo " + "if not absolute", + ) + snakemake_parser.add_argument( + "--preset", + type=str | None, + default=None, + choices=[None, *PRESETS], + help="Start from a preset run, e.g. `small` to check that the " + "pipeline runs", ) # online subcommand @@ -200,14 +219,22 @@ def main(): if subcommand == "snakemake": train_yaml = root / "projects" / "train" / "train.yaml" + overrides = "" + if args.preset is not None: + preset_config, train_yaml = PRESETS[args.preset] + # the preset's own train_config is replaced by the run's copy + lines = preset_config.read_text().splitlines(keepends=True) + overrides = f"\n# From the {args.preset} preset\n" + "".join( + line for line in lines if not line.startswith("train_config:") + ) shutil.copy(train_yaml, directory / "train.yaml") run_config = directory / "config.yaml" run_config.write_text( f"# Overrides for pipeline/config/config.yaml.\n" f"run_dir: {directory}\n" - f"train_config: {directory / 'train.yaml'}\n" + f"train_config: {directory / 'train.yaml'}\n" + overrides ) - create_snakemake_runfile(directory, args.profile) + create_snakemake_runfile(directory, root / args.profile) elif subcommand == "online": copy_configs(directory, ONLINE_CONFIGS) diff --git a/scripts/free_gpus.py b/scripts/free_gpus.py new file mode 100644 index 000000000..fb7d55ae2 --- /dev/null +++ b/scripts/free_gpus.py @@ -0,0 +1,44 @@ +"""Print the IDs of the least-used GPUs on this node + +Orders GPUs by memory in use, then utilization, each summed over several +samples since both fluctuate. Chosen when a job starts. + +Standard library only, so that it runs in any image and on the AP. +""" + +import subprocess +import sys +import time +from collections import defaultdict + +NUM_SAMPLES = 10 +INTERVAL = 0.5 # seconds between samples + + +def free_gpus(num_gpus: int) -> list[str]: + usage = defaultdict(lambda: [0.0, 0.0]) + for i in range(NUM_SAMPLES): + if i > 0: + time.sleep(INTERVAL) + output = subprocess.check_output( + [ + "nvidia-smi", + "--query-gpu=index,memory.used,utilization.gpu", + "--format=csv,noheader,nounits", + ], + text=True, + ) + for line in output.strip().splitlines(): + index, memory, utilization = line.split(", ") + usage[index][0] += float(memory) + usage[index][1] += float(utilization) + + if len(usage) < num_gpus: + raise RuntimeError( + f"Asked for {num_gpus} GPUs, but found {len(usage)}" + ) + return sorted(usage, key=usage.get)[:num_gpus] + + +if __name__ == "__main__": + print(",".join(free_gpus(int(sys.argv[1])))) # noqa: T201