Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
40 changes: 26 additions & 14 deletions Snakefile
Original file line number Diff line number Diff line change
@@ -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 <repo>/Snakefile --configfile config.yaml \
--profile <repo>/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"))
Expand All @@ -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()

Expand Down
7 changes: 4 additions & 3 deletions pipeline/.env.example
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
26 changes: 15 additions & 11 deletions pipeline/config/config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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)
Expand Down
4 changes: 1 addition & 3 deletions pipeline/config/small.yaml
Original file line number Diff line number Diff line change
@@ -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
Expand Down
8 changes: 4 additions & 4 deletions pipeline/profiles/delta/config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
#

Expand All @@ -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:
Expand Down
20 changes: 10 additions & 10 deletions pipeline/profiles/ldg/config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
#

Expand Down Expand Up @@ -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 <workdir>`, 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:
Expand Down
13 changes: 0 additions & 13 deletions pipeline/profiles/local/config.yaml

This file was deleted.

64 changes: 63 additions & 1 deletion pipeline/resources.smk
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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.
Expand Down Expand Up @@ -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`.

Expand Down
3 changes: 2 additions & 1 deletion projects/export/export.smk
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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}"
Expand Down
7 changes: 5 additions & 2 deletions projects/infer/infer.smk
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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"
Expand Down
11 changes: 9 additions & 2 deletions projects/infer/scripts/start_triton.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -43,7 +50,7 @@
"uv",
"run",
"--directory",
"projects/infer",
params.infer_project,
"start-server",
"--model_repo_dir",
snakemake.input.model_repo,
Expand All @@ -54,7 +61,7 @@
"--model_version",
str(params.model_version),
"--gpus",
str(params.gpus),
gpus,
"--batch_size",
str(params.batch_size),
"--triton_image",
Expand Down
Loading
Loading