"
+ for title, value in (
+ ("Complete settings and provenance", spec),
+ ("Input totals and categories", run["inputs"]),
+ ("Output totals and categories", run["outputs"]),
+ ("Container exit and OOM status", run["docker"]),
+ )
+ )
+ return (
+ f"
{escape(spec['label'])}
{'SUCCEEDED' if run['valid'] else 'FAILED / INCOMPLETE'}"
+ f" · elapsed: {elapsed} s · peak: {run['peak'] / 2**30:.3f} GiB
{failure}"
+ f"{memory_chart(run, xmax, ymax)}
Experiment settings
{settings}
"
+ f"
Population and outputs
Table
Input rows
Output rows
{''.join(counts)}
"
+ f"
Output households and persons are the realized sample.
{details}"
+ )
+
+
+def comparison_notes(runs):
+ """Surface workload/environment differences without claiming causal speedups."""
+ notes = []
+ fields = (
+ "profile_name",
+ "sources",
+ "households",
+ "multiprocess",
+ "processes",
+ "sharrow",
+ "config_sha256",
+ "input_files",
+ "docker",
+ )
+ for field in fields:
+ values = [json.dumps(run["spec"].get(field), sort_keys=True) for run in runs]
+ if len(set(values)) > 1:
+ notes.append(field.replace("_", " "))
+ if not notes:
+ return ""
+ return (
+ "
Comparison differences: "
+ + escape(", ".join(notes))
+ + ". Inspect provenance before attributing differences to a source revision.
"
+ )
+
+
+def report(directories, destination):
+ """Generate a portable, offline HTML report and normalized comparison data."""
+ runs = [load_run(path) for path in directories]
+ components = list(dict.fromkeys(name for run in runs for name in run["components"]))
+ headers = "".join(f"
{escape(r['spec']['label'])}
" for r in runs)
+ table = []
+ for name in components:
+ eligible = [
+ r["components"][name]["mean"]
+ for r in runs
+ if r["valid"] and name in r["components"]
+ ]
+ fastest = min(eligible, default=None)
+ cells = []
+ for run in runs:
+ c = run["components"].get(name)
+ if c is None:
+ cells.append("
—
")
+ continue
+ winner = (
+ run["valid"]
+ and fastest is not None
+ and math.isclose(c["mean"], fastest, rel_tol=1e-9)
+ )
+ cells.append(
+ f'
{c["mean"]:.3f} ± {c["sd"]:.3f} s n={c["n"]}; max={c["maximum"]:.3f} s
'
+ )
+ table.append(f"
{escape(name)}
{''.join(cells)}
")
+ xmax = (
+ max((row["elapsed_seconds"] for r in runs for row in r["memory"]), default=1)
+ or 1
+ )
+ ymax = (
+ max((row["current_bytes"] for r in runs for row in r["memory"]), default=1) or 1
+ )
+ cards = [experiment_card(run, xmax, ymax) for run in runs]
+ legend = " ".join(
+ f'■ {escape(r["spec"]["label"])}'
+ for i, r in enumerate(runs)
+ )
+ selectable = list(
+ dict.fromkeys(
+ components
+ + [
+ window["component"]
+ for run in runs
+ for window in run["component_windows"]
+ ]
+ )
+ )
+ options = "".join(
+ f'' for name in selectable
+ )
+ # Component names live only in escaped HTML attributes/text, never in JS.
+ selector_script = """
+
+"""
+ document = f"""abench report
+
+
abench report
Whole-container cgroup v2 memory counts shared pages once, including file cache, kernel memory and the supervisor. Swap is recorded separately in memory.csv. Cache preparation and post-run summaries are excluded. Peak is the kernel high-water mark sampled during the model lifetime, including container startup. Memory panels use identical axes.
+{comparison_notes(runs)}
+
+
+
Selection applies to every memory chart. Each translucent band is one worker execution; darker overlaps indicate concurrent workers. Gaps remain unshaded. Hover over a band for its worker and time range.
+
+
{"".join(cards)}
Component runtimes
Mean ± population standard deviation across worker executions, with observation count and maximum. These describe worker imbalance, not uncertainty across repeated experiments. Component timings exclude pipeline checkpoint writes; elapsed time includes startup, I/O and coordination. Parallel component times must not be summed to estimate wall time. Green cells mark the fastest successful experiment's mean; failed runs are excluded from winners.
Component
{headers}
{"".join(table)}
Runtime comparison
{legend}
{runtime_chart(runs, components)}
{selector_script}"""
+ destination.parent.mkdir(parents=True, exist_ok=True)
+ destination.write_text(document)
+ write_json(destination.with_suffix(".json"), runs)
diff --git a/src/abench/runtime/__init__.py b/src/abench/runtime/__init__.py
new file mode 100644
index 0000000..bdc50ef
--- /dev/null
+++ b/src/abench/runtime/__init__.py
@@ -0,0 +1 @@
+"""Standalone files copied into each experiment and mounted in its containers."""
diff --git a/src/abench/runtime/build_sources.py b/src/abench/runtime/build_sources.py
new file mode 100644
index 0000000..a531de6
--- /dev/null
+++ b/src/abench/runtime/build_sources.py
@@ -0,0 +1,109 @@
+"""Build exact source checkouts into wheels and resolve them together in Docker."""
+
+import email
+import importlib.metadata
+import json
+import re
+import subprocess
+import sys
+import zipfile
+from pathlib import Path
+
+
+def command(args):
+ """Use argument arrays throughout; profile values are never shell fragments."""
+ return subprocess.check_output(args, text=True).strip()
+
+
+def normalized(name):
+ return re.sub(r"[-_.]+", "-", name).lower()
+
+
+def install(manifest, directory=Path("/opt")):
+ """Verify Git objects and wheel identities before installing all overrides."""
+ directory.mkdir(parents=True, exist_ok=True)
+ constraints = directory / "constraints.txt"
+ constraints.write_text("\n".join(manifest.get("constraints", [])) + "\n")
+ wheels = []
+ provenance = []
+ for item in manifest["sources"]:
+ checkout = directory / "sources" / item["name"]
+ checkout.mkdir(parents=True)
+ command(["git", "init", str(checkout)])
+ url = f"https://github.com/{item['repository']}.git"
+ command(["git", "-C", str(checkout), "remote", "add", "origin", url])
+ command(
+ ["git", "-C", str(checkout), "fetch", "--tags", "origin", item["commit"]]
+ )
+ command(["git", "-C", str(checkout), "checkout", "--detach", "FETCH_HEAD"])
+ resolved = command(["git", "-C", str(checkout), "rev-parse", "HEAD"])
+ if resolved != item["commit"]:
+ raise ValueError(f"Git commit mismatch for {item['name']}")
+ package = (checkout / item.get("subdirectory", "")).resolve()
+ if not package.is_relative_to(checkout.resolve()):
+ raise ValueError("package subdirectory escapes checkout")
+ wheel_dir = directory / "wheels" / item["name"]
+ command(
+ [
+ sys.executable,
+ "-m",
+ "pip",
+ "wheel",
+ "--no-deps",
+ "--wheel-dir",
+ str(wheel_dir),
+ str(package),
+ ]
+ )
+ built = list(wheel_dir.glob("*.whl"))
+ if len(built) != 1:
+ raise ValueError(f"expected one wheel for {item['name']}")
+ with zipfile.ZipFile(built[0]) as archive:
+ metadata_path = next(
+ n for n in archive.namelist() if n.endswith(".dist-info/METADATA")
+ )
+ metadata = email.message_from_bytes(archive.read(metadata_path))
+ if normalized(metadata["Name"]) != item["name"]:
+ raise ValueError(
+ f"wrong distribution: expected {item['name']}, got {metadata['Name']}"
+ )
+ extras = "[" + ",".join(item["extras"]) + "]" if item.get("extras") else ""
+ wheels.append(str(built[0]) + extras)
+ provenance.append(
+ dict(item, version=metadata["Version"], resolved_commit=resolved)
+ )
+ # Supplying all local wheels in one transaction prevents dependency resolution
+ # from quietly replacing a requested source package with a registry release.
+ requirements = ["pyyaml", "pyarrow", *manifest.get("requirements", [])]
+ command(
+ [
+ sys.executable,
+ "-m",
+ "pip",
+ "install",
+ "--no-cache-dir",
+ "-c",
+ str(constraints),
+ *wheels,
+ *requirements,
+ ]
+ )
+ command([sys.executable, "-m", "pip", "check"])
+ for item in provenance:
+ dist = importlib.metadata.distribution(item["name"])
+ direct = json.loads(dist.read_text("direct_url.json") or "{}")
+ expected = (directory / "wheels" / item["name"]).resolve().as_uri() + "/"
+ if dist.version != item["version"] or not direct.get("url", "").startswith(
+ expected
+ ):
+ raise ValueError(f"installed source identity mismatch: {item['name']}")
+ (directory / "source-provenance.json").write_text(
+ json.dumps(provenance, indent=2) + "\n"
+ )
+ (directory / "pip-freeze.txt").write_text(
+ command([sys.executable, "-m", "pip", "freeze"]) + "\n"
+ )
+
+
+if __name__ == "__main__":
+ install(json.loads(Path(sys.argv[1]).read_text()))
diff --git a/src/abench/runtime/instrumentation.py b/src/abench/runtime/instrumentation.py
new file mode 100644
index 0000000..fdd7199
--- /dev/null
+++ b/src/abench/runtime/instrumentation.py
@@ -0,0 +1,61 @@
+"""Process-local component timing and a strict Sharrow disk-cache guard."""
+
+import json
+import multiprocessing
+import os
+import time
+from pathlib import Path
+
+
+def install(flow_cache=Path("/results/cache/flows")):
+ """Install in the model parent and every spawned worker, before model imports."""
+ from activitysim.core.workflow import runner
+
+ original = runner.run_named_step
+ records = Path(os.environ["BENCH_PHASE_DIR"])
+
+ def timed(name, context, **kwargs):
+ started = time.perf_counter()
+ succeeded = False
+ try:
+ result = original(name, context, **kwargs)
+ succeeded = True
+ return result
+ finally:
+ finished = time.perf_counter()
+ row = {
+ "component": name,
+ "seconds": finished - started,
+ "process": multiprocessing.current_process().name,
+ "pid": os.getpid(),
+ "succeeded": succeeded,
+ }
+ # Linux's monotonic clock is shared across processes. Use the
+ # supervisor's origin so worker windows align with memory samples.
+ if "BENCH_STARTED_MONOTONIC" in os.environ:
+ origin = float(os.environ["BENCH_STARTED_MONOTONIC"])
+ row.update(
+ start_seconds=started - origin, end_seconds=finished - origin
+ )
+ # Separate files avoid interleaving writes from concurrent workers.
+ with (records / f"components-{os.getpid()}.jsonl").open("a") as stream:
+ stream.write(json.dumps(row) + "\n")
+
+ runner.run_named_step = timed
+ if os.environ.get("BENCH_STRICT_CACHE") == "1":
+ from numba.core.dispatcher import _FunctionCompiler
+
+ compile_original = _FunctionCompiler.compile
+ flow_cache = flow_cache.resolve()
+
+ def compile_checked(self, *args, **kwargs):
+ # Numba reaches this method only after failing to load a compiled
+ # overload from disk. Ordinary ActivitySim/Numba JIT is still allowed.
+ filename = Path(self.py_func.__code__.co_filename).resolve()
+ if filename.is_relative_to(flow_cache):
+ with (records / f"cache-miss-{os.getpid()}.txt").open("a") as stream:
+ stream.write(str(filename) + "\n")
+ raise RuntimeError(f"Measured Sharrow flow cache miss: {filename}")
+ return compile_original(self, *args, **kwargs)
+
+ _FunctionCompiler.compile = compile_checked
diff --git a/src/abench/runtime/worker.py b/src/abench/runtime/worker.py
new file mode 100644
index 0000000..3cedcc6
--- /dev/null
+++ b/src/abench/runtime/worker.py
@@ -0,0 +1,255 @@
+"""Linux container supervisor; the model runs in a separate process tree."""
+
+import csv
+import json
+import os
+import subprocess
+import sys
+import time
+from pathlib import Path
+
+# multiprocessing's spawn imports this file again as __mp_main__.
+if os.environ.get("BENCH_MODEL") == "1":
+ from instrumentation import install
+
+ install()
+
+
+def write_json(path, value):
+ path.write_text(json.dumps(value, indent=2) + "\n")
+
+
+def make_state(
+ spec,
+ phase,
+ model_root=Path("/model"),
+ data_root=Path("/data"),
+ results_root=Path("/results"),
+):
+ """Resolve model configs < profile defaults < overlays < explicit CLI controls.
+
+ A generated inheriting config carries defaults through worker reconstruction;
+ direct settings overrides are reserved for the experiment's required controls.
+ """
+ import activitysim.abm # noqa: F401 -- register standard components
+ import yaml
+ from activitysim.core.workflow import State
+
+ profile = spec["profile"]
+ defaults = dict(profile.get("settings", {}))
+ if profile.get("models_from"):
+ base = yaml.safe_load((model_root / profile["models_from"]).read_text())
+ defaults["models"] = [
+ name
+ for name in base["models"]
+ if name not in profile.get("exclude_models", [])
+ ]
+ if spec["multiprocess"] and profile.get("mp_settings"):
+ mp = yaml.safe_load((model_root / profile["mp_settings"]).read_text())
+ defaults["multiprocess_steps"] = mp["multiprocess_steps"]
+ generated = phase / "profile-config"
+ generated.mkdir(exist_ok=True)
+ for name in ("settings.yaml", "settings_mp.yaml", "settings_mp_sharrow.yaml"):
+ (generated / name).write_text(
+ yaml.safe_dump(dict(defaults, inherit_settings=True))
+ )
+ configs = [
+ model_root / f"overlay-{i}" for i in range(len(spec.get("config_overlay", [])))
+ ]
+ configs += [generated]
+ if spec["multiprocess"]:
+ configs += [model_root / name for name in profile.get("mp_configs", [])]
+ configs += [model_root / name for name in profile["configs"]]
+ state = State.make_default(
+ working_dir=model_root,
+ configs_dir=configs,
+ data_dir=data_root,
+ output_dir=phase / "output",
+ cache_dir=results_root / "cache/model",
+ settings=dict(
+ households_sample_size=spec["households"],
+ multiprocess=spec["multiprocess"],
+ num_processes=spec["processes"],
+ sharrow="require" if spec["sharrow"] else False,
+ fail_fast=True,
+ ),
+ )
+ # Every sliced phase honors the requested count, including phases that have
+ # their own worker count in production configs or explicit chunk overlays.
+ for step in state.settings.multiprocess_steps or []:
+ if isinstance(step, dict):
+ if "slice" in step:
+ step["num_processes"] = spec["processes"]
+ elif getattr(step, "slice", None) is not None:
+ step.num_processes = spec["processes"]
+ state.filesystem.sharrow_cache_dir = results_root / "cache/flows"
+ state.settings.sharrow_cache_dir = str(state.filesystem.sharrow_cache_dir)
+ sys.path.insert(0, str(model_root))
+ state.set("imported_extensions", [])
+ for extension in profile.get("extensions", []):
+ state.import_extensions(extension)
+ if profile.get("adapter"):
+ import importlib
+
+ module, function = profile["adapter"].split(":")
+ getattr(importlib.import_module(module), function)(state, spec, phase)
+ state.set("run_timestamp", "benchmark")
+ state.set("run_id", str(state.tracing.run_id))
+ return state
+
+
+def run_model(spec, phase):
+ """Use identical samples, process layout and cache paths in both phases."""
+ state = make_state(spec, phase)
+ state.logging.config_logger()
+ write_json(
+ phase / "effective-settings.json", state.settings.model_dump(mode="json")
+ )
+ state.run.all(resume_after=None)
+ if not spec["multiprocess"]:
+ state.checkpoint.close_store()
+
+
+def table_summary(directory, prefix="", tables=None):
+ """Stream input/output tables so full-population summaries need bounded RAM."""
+ import pandas as pd
+ import pyarrow.parquet as pq
+
+ result = {}
+ tables = tables or {
+ name: {}
+ for name in (
+ "households",
+ "persons",
+ "land_use",
+ "tours",
+ "trips",
+ "joint_tour_participants",
+ "vehicles",
+ )
+ }
+ for name, options in tables.items():
+ path = directory / f"{prefix}{options.get('file', name)}.csv"
+ if path.exists():
+ chunks = pd.read_csv(path, chunksize=100_000)
+ else:
+ path = directory / f"{prefix}{options.get('file', name)}.parquet"
+ if not path.exists():
+ continue
+ chunks = (
+ batch.to_pandas() for batch in pq.ParquetFile(path).iter_batches()
+ )
+ info = {"rows": 0, "totals": {}, "categories": {}}
+ zones = set()
+ for chunk in chunks:
+ if name == "land_use":
+ for col in options.get("zone_columns", ["taz", "TAZ"]):
+ if col in chunk:
+ zones.update(chunk[col].dropna().unique().tolist())
+ info["rows"] += len(chunk)
+ for col in options.get(
+ "totals", ["TOTPOP", "TOTHH", "TOTEMP", "pop", "hh", "emp_total"]
+ ):
+ if col in chunk:
+ info["totals"][col] = info["totals"].get(col, 0) + float(
+ chunk[col].sum()
+ )
+ for col in options.get(
+ "categories",
+ (
+ "tour_category",
+ "tour_type",
+ "tour_mode",
+ "trip_mode",
+ "primary_purpose",
+ ),
+ ):
+ if col in chunk:
+ counts = info["categories"].setdefault(col, {})
+ for value, count in chunk[col].value_counts(dropna=False).items():
+ counts[str(value)] = counts.get(str(value), 0) + int(count)
+ if zones:
+ info["taz_count"] = len(zones)
+ result[name] = info
+ return result
+
+
+def sample_memory(root, elapsed):
+ """cgroup v2 charges shared pages once; file includes shmem, not vice versa."""
+ stat = dict(
+ line.split() for line in (root / "memory.stat").read_text().splitlines()
+ )
+ return {
+ "elapsed_seconds": elapsed,
+ "current_bytes": int((root / "memory.current").read_text()),
+ "peak_bytes": int((root / "memory.peak").read_text()),
+ "swap_bytes": int((root / "memory.swap.current").read_text()),
+ "anonymous_bytes": int(stat.get("anon", 0)),
+ "file_bytes": int(stat.get("file", 0)),
+ "shared_bytes": int(stat.get("shmem", 0)),
+ }
+
+
+def supervise(spec, phase):
+ """Measure only the model subprocess lifetime; summarize after sampling ends."""
+ root = Path("/sys/fs/cgroup")
+ if not (root / "memory.peak").exists():
+ raise RuntimeError("A private cgroup v2 with memory.peak is required")
+ (phase / "output").mkdir()
+ for name in ("flows", "model"):
+ Path(f"/results/cache/{name}").mkdir(parents=True, exist_ok=True)
+ env = dict(os.environ, BENCH_MODEL="1", BENCH_PHASE_DIR=str(phase))
+ env["BENCH_STRICT_CACHE"] = (
+ "1" if spec["sharrow"] and phase.name == "measured" else "0"
+ )
+ started = time.perf_counter()
+ env["BENCH_STARTED_MONOTONIC"] = str(started)
+ with (phase / "memory.csv").open("w", buffering=1) as stream:
+ writer = csv.DictWriter(stream, fieldnames=sample_memory(root, 0).keys())
+ writer.writeheader()
+ writer.writerow(sample_memory(root, 0))
+ process = subprocess.Popen(
+ [sys.executable, __file__, "model", str(phase)], env=env
+ )
+ while True:
+ writer.writerow(sample_memory(root, time.perf_counter() - started))
+ if process.poll() is not None:
+ break
+ try:
+ process.wait(timeout=spec["interval"])
+ except subprocess.TimeoutExpired:
+ pass
+ write_json(
+ phase / "status.json",
+ {
+ "returncode": process.returncode,
+ "elapsed_seconds": time.perf_counter() - started,
+ },
+ )
+ # This separate process keeps pandas/Arrow and summary allocations out of the
+ # measured lifetime; the host report uses only the recorded cgroup samples.
+ subprocess.run([sys.executable, __file__, "summary", str(phase)], check=True)
+ return process.returncode
+
+
+if __name__ == "__main__":
+ mode, phase_arg = sys.argv[1:]
+ phase = Path(phase_arg)
+ spec = json.loads(Path("/results/experiment.json").read_text())
+ if mode == "model":
+ run_model(spec, phase)
+ elif mode == "summary":
+ write_json(
+ phase / "input-summary.json",
+ table_summary(Path("/data"), tables=spec["profile"].get("input_tables")),
+ )
+ write_json(
+ phase / "output-summary.json",
+ table_summary(
+ phase / "output",
+ spec["profile"].get("output_prefix", "final_"),
+ spec["profile"].get("output_tables"),
+ ),
+ )
+ else:
+ sys.exit(supervise(spec, phase))
diff --git a/src/abench/sources.py b/src/abench/sources.py
new file mode 100644
index 0000000..d10ef0e
--- /dev/null
+++ b/src/abench/sources.py
@@ -0,0 +1,93 @@
+"""Normalize exact GitHub source overrides before executing any build commands."""
+
+import re
+from pathlib import PurePosixPath
+
+
+def canonical_name(name):
+ """Compare distribution names using Python packaging's normalization rules."""
+ return re.sub(r"[-_.]+", "-", name).lower()
+
+
+def source(value):
+ """Accept a CLI shorthand or a profile mapping, including extras/subdirectory."""
+ if isinstance(value, str):
+ match = re.fullmatch(
+ r"([A-Za-z0-9][A-Za-z0-9_.-]*)(?:\[([A-Za-z0-9_,.-]+)\])?=([A-Za-z0-9_.-]+/[A-Za-z0-9_.-]+)@([0-9a-fA-F]{40})(?:#subdirectory=([^\s]+))?",
+ value,
+ )
+ if not match:
+ raise ValueError(
+ "source must be distribution[extras]=organization/repository@full-SHA[#subdirectory=path]"
+ )
+ name, extras, repository, commit, subdirectory = match.groups()
+ value = dict(
+ name=name,
+ repository=repository,
+ commit=commit,
+ extras=extras.split(",") if extras else [],
+ subdirectory=subdirectory or "",
+ )
+ if not isinstance(value, dict) or set(value) - {
+ "name",
+ "repository",
+ "commit",
+ "extras",
+ "subdirectory",
+ }:
+ raise ValueError("invalid source fields")
+ result = dict(value)
+ for key, pattern in [
+ ("name", r"[A-Za-z0-9][A-Za-z0-9_.-]*"),
+ ("repository", r"[A-Za-z0-9_.-]+/[A-Za-z0-9_.-]+"),
+ ("commit", r"[0-9a-fA-F]{40}"),
+ ]:
+ if not isinstance(result.get(key), str) or not re.fullmatch(
+ pattern, result[key]
+ ):
+ raise ValueError(f"invalid source {key}: {result.get(key)!r}")
+ result["name"] = canonical_name(result["name"])
+ result["commit"] = result["commit"].lower()
+ result.setdefault("subdirectory", "")
+ result.setdefault("extras", [])
+ subdir = result["subdirectory"]
+ if (
+ not isinstance(subdir, str)
+ or PurePosixPath(subdir).is_absolute()
+ or ".." in PurePosixPath(subdir).parts
+ or "\\" in subdir
+ ):
+ raise ValueError("source subdirectory must stay inside the checkout")
+ if not isinstance(result["extras"], list) or any(
+ not isinstance(e, str) or not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9_.-]*", e)
+ for e in result["extras"]
+ ):
+ raise ValueError("invalid package extras")
+ result["extras"] = sorted(set(result["extras"]))
+ return result
+
+
+def resolve_sources(defaults, overrides, activitysim=None, sharrow=None):
+ """CLI entries replace profile pins; aliases cannot silently contradict them."""
+ selected = {}
+ for value in defaults:
+ item = source(value)
+ if item["name"] in selected:
+ raise ValueError(f"duplicate profile source: {item['name']}")
+ selected[item["name"]] = item
+ explicit = {}
+ for value in overrides:
+ item = source(value)
+ if item["name"] in explicit:
+ raise ValueError(f"duplicate CLI source: {item['name']}")
+ explicit[item["name"]] = item
+ for name, sha in [("activitysim", activitysim), ("sharrow", sharrow)]:
+ if sha:
+ item = source(f"{name}=ActivitySim/{name}@{sha}")
+ if name in explicit and explicit[name] != item:
+ raise ValueError(f"conflicting --source and --{name}-commit")
+ explicit[name] = item
+ selected.update(explicit)
+ if "activitysim" not in selected:
+ raise ValueError("provide an exact ActivitySim source commit")
+ return [selected[name] for name in sorted(selected)]
diff --git a/tests/fixtures/tiny/benchmark.yaml b/tests/fixtures/tiny/benchmark.yaml
new file mode 100644
index 0000000..59a5211
--- /dev/null
+++ b/tests/fixtures/tiny/benchmark.yaml
@@ -0,0 +1,20 @@
+schema_version: 1
+name: CI extension model
+configs: [configs]
+snapshot: [configs, tiny_extension.py]
+extensions: [tiny_extension]
+data_dir: data
+required_inputs: [households.csv]
+settings:
+ models: [bench_initialize, bench_compute, bench_export]
+ checkpoints: true
+ rng_base_seed: 0
+ chunk_size: 0
+ trace_hh_id: null
+ trace_od: null
+constraints: ['pandas<3', 'multimethod<2']
+input_tables:
+ households: {}
+output_tables:
+ households: {}
+adapter: tiny_extension:prepare
diff --git a/tests/fixtures/tiny/configs/settings.yaml b/tests/fixtures/tiny/configs/settings.yaml
new file mode 100644
index 0000000..052f22d
--- /dev/null
+++ b/tests/fixtures/tiny/configs/settings.yaml
@@ -0,0 +1,15 @@
+models: [bench_initialize, bench_compute, bench_export]
+multiprocess_steps:
+- name: mp_initialize
+ begin: bench_initialize
+- name: mp_households
+ begin: bench_compute
+ num_processes: 20
+ slice:
+ tables: [households]
+- name: mp_finalize
+ begin: bench_export
+input_table_list:
+- tablename: households
+ filename: households.csv
+ index_col: household_id
diff --git a/tests/fixtures/tiny/data/households.csv b/tests/fixtures/tiny/data/households.csv
new file mode 100644
index 0000000..21f2c9b
--- /dev/null
+++ b/tests/fixtures/tiny/data/households.csv
@@ -0,0 +1,5 @@
+household_id,value
+1,10
+2,20
+3,30
+4,40
diff --git a/tests/fixtures/tiny/tiny_extension.py b/tests/fixtures/tiny/tiny_extension.py
new file mode 100644
index 0000000..9ea406b
--- /dev/null
+++ b/tests/fixtures/tiny/tiny_extension.py
@@ -0,0 +1,54 @@
+"""Tiny real ActivitySim workflow: exercises registration, slicing, and caching."""
+
+import importlib
+import sys
+from pathlib import Path
+
+import pandas as pd
+from activitysim.core import workflow
+
+
+@workflow.step
+def bench_initialize(state: workflow.State):
+ """Seed four rows and a generated Numba module in the guarded flow directory."""
+ households = pd.read_csv(Path("/data/households.csv"), index_col="household_id")
+ state.add_table("households", households)
+ cache = Path(state.settings.sharrow_cache_dir)
+ cache.mkdir(parents=True, exist_ok=True)
+ module = cache / "tiny_generated.py"
+ if not module.exists():
+ module.write_text(
+ "from numba import njit\n@njit(cache=True)\ndef twice(x):\n return x * 2\n"
+ )
+
+
+@workflow.step
+def bench_compute(state: workflow.State, households: pd.DataFrame):
+ """Run independently in each slice; measured cache misses must be rejected."""
+ sys.path.insert(0, state.settings.sharrow_cache_dir)
+ twice = importlib.import_module("tiny_generated").twice
+ households = households.copy()
+ households["value"] = [twice(int(value)) for value in households["value"]]
+ state.add_table("households", households)
+
+
+@workflow.step
+def bench_export(state: workflow.State, households: pd.DataFrame):
+ """Export merged results so the harness validates the realized population."""
+ households.to_csv(Path(state.filesystem.output_dir) / "final_households.csv")
+
+
+def prepare(state, spec, phase):
+ """Supply tiny shared skims; this fixture has no destination/shadow models."""
+ state.set("shadow_pricing_info", None)
+ state.set("shadow_pricing_choice_info", None)
+ state.set("network_los_preload", None)
+ if spec["sharrow"]:
+ import numpy as np
+ import sharrow as sh
+
+ dataset = sh.Dataset({"distance": (("otaz", "dtaz"), np.ones((1, 1)))})
+ state.set(
+ "skim_dataset",
+ dataset.shm.to_shared_memory("skim_dataset", mode="r", load=True),
+ )
diff --git a/tests/requirements.txt b/tests/requirements.txt
new file mode 100644
index 0000000..99da3b9
--- /dev/null
+++ b/tests/requirements.txt
@@ -0,0 +1,6 @@
+# The measurement hook uses ActivitySim/Numba internals; pin real implementations
+# in CI so unrelated upstream changes do not silently change the test contract.
+activitysim @ git+https://github.com/ActivitySim/activitysim.git@5c6fae24a91a57a2d6dfc2e1dbe062a61d94545a
+sharrow @ git+https://github.com/ActivitySim/sharrow.git@fc175b27d8e0c5d202721c67d96b050e6117b235
+multimethod<2
+pandas<3
diff --git a/tests/test_builder.py b/tests/test_builder.py
new file mode 100644
index 0000000..932af91
--- /dev/null
+++ b/tests/test_builder.py
@@ -0,0 +1,82 @@
+"""Exercise source verification without network access or host package mutation."""
+
+import json
+import zipfile
+from types import SimpleNamespace
+
+import pytest
+
+from abench.runtime import build_sources
+from abench.sources import source
+
+SHA = "a" * 40
+
+
+def fake_commands(tmp_path, monkeypatch, *, resolved=SHA, metadata_name="addon"):
+ """Emulate Git/pip boundaries but pass real wheel metadata through the builder."""
+ calls = []
+
+ def command(args):
+ calls.append(args)
+ if args[-2:] == ["rev-parse", "HEAD"]:
+ return resolved
+ if "wheel" in args:
+ wheel_dir = tmp_path / "wheels/addon"
+ wheel_dir.mkdir(parents=True)
+ with zipfile.ZipFile(
+ wheel_dir / "addon-1.0-py3-none-any.whl", "w"
+ ) as archive:
+ archive.writestr(
+ "addon-1.0.dist-info/METADATA",
+ f"Name: {metadata_name}\nVersion: 1.0\n",
+ )
+ return ""
+
+ monkeypatch.setattr(build_sources, "command", command)
+ monkeypatch.setattr(
+ build_sources.importlib.metadata,
+ "distribution",
+ lambda name: SimpleNamespace(
+ version="1.0",
+ read_text=lambda key: json.dumps(
+ {"url": (tmp_path / "wheels/addon/addon-1.0-py3-none-any.whl").as_uri()}
+ ),
+ ),
+ )
+ return calls
+
+
+def test_builder_verifies_and_installs_wheels_together(tmp_path, monkeypatch):
+ calls = fake_commands(tmp_path, monkeypatch)
+ item = source(f"addon[fast]=org/addon@{SHA}")
+ build_sources.install(
+ dict(sources=[item], requirements=["numpy"], constraints=["numpy<3"]), tmp_path
+ )
+ install = next(c for c in calls if "install" in c)
+ assert any(value.endswith(".whl[fast]") for value in install)
+ assert "numpy" in install
+ assert (tmp_path / "constraints.txt").read_text() == "numpy<3\n"
+ assert (
+ json.loads((tmp_path / "source-provenance.json").read_text())[0][
+ "resolved_commit"
+ ]
+ == SHA
+ )
+
+
+@pytest.mark.parametrize(
+ "options,match",
+ [
+ ({"resolved": "b" * 40}, "Git commit mismatch"),
+ ({"metadata_name": "different"}, "wrong distribution"),
+ ],
+)
+def test_builder_rejects_wrong_identity_before_install(
+ tmp_path, monkeypatch, options, match
+):
+ calls = fake_commands(tmp_path, monkeypatch, **options)
+ with pytest.raises(ValueError, match=match):
+ build_sources.install(
+ dict(sources=[source(f"addon=org/addon@{SHA}")]), tmp_path
+ )
+ assert not any("install" in c for c in calls)
diff --git a/tests/test_docker.py b/tests/test_docker.py
new file mode 100644
index 0000000..add8287
--- /dev/null
+++ b/tests/test_docker.py
@@ -0,0 +1,61 @@
+"""Opt-in integration through the real Docker supervisor and ActivitySim workflow."""
+
+import json
+import os
+import shutil
+from pathlib import Path
+
+import pytest
+
+from abench import cli
+from abench.report import load_run
+
+pytestmark = [
+ pytest.mark.docker,
+ pytest.mark.skipif(
+ os.environ.get("ABENCH_DOCKER_TESTS") != "1",
+ reason="set ABENCH_DOCKER_TESTS=1 to run Docker integration",
+ ),
+]
+ACTIVITYSIM = "5c6fae24a91a57a2d6dfc2e1dbe062a61d94545a"
+SHARROW = "fc175b27d8e0c5d202721c67d96b050e6117b235"
+
+
+@pytest.mark.parametrize("multiprocess,sharrow", [(False, False), (True, True)])
+def test_tiny_model(tmp_path, multiprocess, sharrow):
+ """Build pinned sources, run a full warmup where enabled, then measure."""
+ root = tmp_path / "model"
+ shutil.copytree(Path(__file__).parent / "fixtures/tiny", root)
+ output = tmp_path / "experiment"
+ args = [
+ "run",
+ "--model-dir",
+ str(root),
+ "--source",
+ f"activitysim=ActivitySim/activitysim@{ACTIVITYSIM}",
+ "--source",
+ f"sharrow=ActivitySim/sharrow@{SHARROW}",
+ "--households",
+ "4",
+ "--memory",
+ "3g",
+ "--shm-size",
+ "256m",
+ "--output-dir",
+ str(output),
+ ]
+ args += (
+ ["--multiprocess", "--processes", "2"] if multiprocess else ["--single-process"]
+ )
+ args += ["--sharrow"] if sharrow else ["--no-sharrow"]
+ assert cli.main(args) == 0
+ run = load_run(output)
+ assert run["valid"]
+ assert run["components"]["bench_compute"]["n"] == (2 if multiprocess else 1)
+ assert run["outputs"]["households"]["rows"] == 4
+ assert run["memory"] and all(row["current_bytes"] > 0 for row in run["memory"])
+ assert (output / "warmup").exists() == sharrow
+ assert len(json.loads((output / "source-provenance.json").read_text())) == 2
+ assert not list((output / "measured").glob("cache-miss-*"))
+ lines = (output / "measured/output/final_households.csv").read_text().splitlines()
+ assert set(lines[1:]) == {"1,20", "2,40", "3,60", "4,80"}
diff --git a/tests/test_measurement.py b/tests/test_measurement.py
new file mode 100644
index 0000000..9dca0d7
--- /dev/null
+++ b/tests/test_measurement.py
@@ -0,0 +1,364 @@
+"""Measurement/report regression tests; no running Docker daemon required."""
+
+import csv
+import json
+import os
+import re
+import shutil
+import subprocess
+import sys
+import tempfile
+import time
+import unittest
+from html.parser import HTMLParser
+from pathlib import Path
+
+from abench import report as benchmark
+from abench.cli import commit, positive
+from abench.runtime import worker
+
+benchmark.commit = commit
+benchmark.positive = positive
+HERE = Path(worker.__file__).parent
+
+
+class BenchmarkTests(unittest.TestCase):
+ def make_run(self, root, label, times, success=True):
+ """Construct raw artifacts with the same schema emitted by workers."""
+ root.mkdir()
+ phase = root / "measured"
+ phase.mkdir()
+ benchmark.write_json(
+ root / "experiment.json", {"schema_version": 1, "label": label}
+ )
+ benchmark.write_json(
+ phase / "status.json", {"returncode": 0, "elapsed_seconds": 9}
+ )
+ benchmark.write_json(
+ phase / "docker-state.json", {"ExitCode": 0, "OOMKilled": not success}
+ )
+ for i, duration in enumerate(times):
+ (phase / f"components-{i}.jsonl").write_text(
+ json.dumps(
+ {
+ "component": "auto_ownership",
+ "seconds": duration,
+ "succeeded": True,
+ }
+ )
+ + "\n"
+ )
+ # A duplicate locutor CSV must never contribute another observation.
+ (phase / "timing_log.csv").write_text(
+ "model_name,seconds\nauto_ownership,999\n"
+ )
+ with (phase / "memory.csv").open("w") as stream:
+ writer = csv.DictWriter(
+ stream, fieldnames=["elapsed_seconds", "current_bytes", "peak_bytes"]
+ )
+ writer.writeheader()
+ writer.writerows(
+ [
+ {"elapsed_seconds": 0, "current_bytes": 100, "peak_bytes": 100},
+ {"elapsed_seconds": 9, "current_bytes": 200, "peak_bytes": 300},
+ ]
+ )
+ return root
+
+ def test_requested_sample_must_be_realized(self):
+ """Do not call a fallback to a smaller input population a valid trial."""
+ with tempfile.TemporaryDirectory() as tmp:
+ run = self.make_run(Path(tmp) / "sample", "sample", [1])
+ benchmark.write_json(
+ run / "experiment.json", {"schema_version": 1, "households": 500000}
+ )
+ benchmark.write_json(
+ run / "measured/output-summary.json", {"households": {"rows": 28365}}
+ )
+ self.assertFalse(benchmark.load_run(run)["valid"])
+
+ def test_worker_statistics_and_failed_run_winners(self):
+ with tempfile.TemporaryDirectory() as temporary:
+ root = Path(temporary)
+ a = self.make_run(root / "a", "A ", document, re.DOTALL)[1]
+ controller_test = r"""
+const assert = require('node:assert/strict');
+const vm = require('node:vm');
+const input = JSON.parse(require('node:fs').readFileSync(0, 'utf8'));
+let change;
+const selector = {value: '', addEventListener(type, fn) {assert.equal(type, 'change'); change = fn;}};
+const panels = input.panels.map(attrs => {
+ const bands = attrs.map(a => ({dataset: {component: a['data-component'], source: a['data-source']}, style: {display: 'none'}}));
+ const status = {textContent: ''};
+ return {bands, status, querySelectorAll() {return bands;}, querySelector() {return status;}};
+});
+vm.runInNewContext(input.script, {document: {
+ getElementById(id) {assert.equal(id, 'memory-component'); return selector;},
+ querySelectorAll() {return panels;}
+}});
+selector.value = input.name;
+change();
+assert.equal(panels[0].bands.filter(b => b.style.display === '').length, 3);
+assert.match(panels[0].status.textContent, /3 worker execution windows/);
+assert.match(panels[1].status.textContent, /No execution-window data/);
+selector.value = '';
+change();
+assert.ok(panels.every(p => p.bands.every(b => b.style.display === 'none')));
+"""
+ subprocess.run(
+ [shutil.which("node"), "-e", controller_test],
+ input=json.dumps(
+ {"panels": parsed.panels, "name": name, "script": script}
+ ),
+ text=True,
+ check=True,
+ )
+
+ def test_commit_and_interval_validation(self):
+ import argparse
+
+ for value in ("main", "1234567", "x" * 40):
+ with self.assertRaises(argparse.ArgumentTypeError):
+ benchmark.commit(value)
+ for value in ("nan", "inf", "0", "-1"):
+ with self.assertRaises(argparse.ArgumentTypeError):
+ benchmark.positive(value)
+
+ def test_table_summaries(self):
+ with tempfile.TemporaryDirectory() as temporary:
+ root = Path(temporary)
+ (root / "households.csv").write_text("HHID,PERSONS\n1,2\n2,1\n")
+ (root / "land_use.csv").write_text("TAZ,TOTPOP,TOTEMP\n1,20,10\n2,30,15\n")
+ (root / "trips.csv").write_text(
+ "trip_id,trip_mode\n1,WALK\n2,WALK\n3,SOV\n"
+ )
+ summary = worker.table_summary(root)
+ self.assertEqual(summary["households"]["rows"], 2)
+ self.assertEqual(summary["land_use"]["totals"]["TOTPOP"], 50)
+ self.assertEqual(
+ summary["trips"]["categories"]["trip_mode"], {"WALK": 2, "SOV": 1}
+ )
+
+ def test_real_numba_disk_hit_allowed_and_miss_blocked(self):
+ """A fresh interpreter must load warmed overloads but reject new ones."""
+ with tempfile.TemporaryDirectory() as temporary:
+ root = Path(temporary)
+ (root / "generated_flow.py").write_text(
+ "from numba import njit\n@njit(cache=True)\ndef f(x):\n return x + 1\n"
+ )
+ env = dict(
+ os.environ,
+ PYTHONPATH=os.pathsep.join([str(HERE), str(root)]),
+ BENCH_PHASE_DIR=str(root),
+ BENCH_STRICT_CACHE="1",
+ )
+ warm = subprocess.run(
+ [
+ sys.executable,
+ "-c",
+ "from generated_flow import f; assert f(1) == 2",
+ ],
+ env=env,
+ check=False,
+ capture_output=True,
+ text=True,
+ )
+ self.assertEqual(warm.returncode, 0, warm.stderr)
+ script = f"from pathlib import Path; from instrumentation import install; install(Path({str(root)!r})); from generated_flow import f; "
+ hit = subprocess.run(
+ [sys.executable, "-c", script + "assert f(1) == 2"],
+ env=env,
+ check=False,
+ capture_output=True,
+ text=True,
+ )
+ self.assertEqual(hit.returncode, 0, hit.stderr)
+ miss = subprocess.run(
+ [sys.executable, "-c", script + "f(1.5)"],
+ env=env,
+ check=False,
+ capture_output=True,
+ text=True,
+ )
+ self.assertNotEqual(miss.returncode, 0)
+ self.assertIn("Measured Sharrow flow cache miss", miss.stderr)
+ self.assertTrue(list(root.glob("cache-miss-*.txt")))
+
+ def test_spawned_workers_each_record_components(self):
+ """Exercise the same top-level hook installation used by MP workers."""
+ with tempfile.TemporaryDirectory() as temporary:
+ root = Path(temporary)
+ script = root / "spawn_check.py"
+ script.write_text(
+ "from activitysim.core.workflow import runner\n"
+ "runner.run_named_step = lambda name, context: 'ok'\n"
+ "import worker\n"
+ "import multiprocessing\n"
+ "def task():\n"
+ " assert runner.run_named_step('smoke', {}) == 'ok'\n"
+ "if __name__ == '__main__':\n"
+ " ctx = multiprocessing.get_context('spawn')\n"
+ " processes = [ctx.Process(target=task) for _ in range(2)]\n"
+ " for p in processes: p.start()\n"
+ " for p in processes:\n"
+ " p.join()\n"
+ " assert p.exitcode == 0\n"
+ )
+ env = dict(
+ os.environ,
+ PYTHONPATH=str(HERE),
+ BENCH_MODEL="1",
+ BENCH_PHASE_DIR=str(root),
+ BENCH_STRICT_CACHE="0",
+ BENCH_STARTED_MONOTONIC=str(time.perf_counter()),
+ )
+ result = subprocess.run(
+ [sys.executable, str(script)],
+ env=env,
+ capture_output=True,
+ text=True,
+ check=False,
+ )
+ self.assertEqual(result.returncode, 0, result.stderr)
+ records = [
+ json.loads(path.read_text()) for path in root.glob("components-*.jsonl")
+ ]
+ self.assertEqual(len({row["pid"] for row in records}), 2)
+ self.assertTrue(
+ all(row["succeeded"] and row["component"] == "smoke" for row in records)
+ )
+ for row in records:
+ self.assertGreater(row["start_seconds"], 0)
+ self.assertAlmostEqual(
+ row["end_seconds"] - row["start_seconds"], row["seconds"]
+ )
+
+
+if __name__ == "__main__":
+ unittest.main()
diff --git a/tests/test_profiles_sources.py b/tests/test_profiles_sources.py
new file mode 100644
index 0000000..5ba4a9e
--- /dev/null
+++ b/tests/test_profiles_sources.py
@@ -0,0 +1,165 @@
+"""Contracts that let a new model or source dependency reuse the same harness."""
+
+import json
+
+import pytest
+import yaml
+
+from abench import cli
+from abench.profiles import load_profile, validate_model
+from abench.runtime.worker import make_state
+from abench.sources import resolve_sources, source
+
+SHA = "a" * 40
+OTHER = "b" * 40
+
+
+def test_source_overrides_and_alias_conflicts():
+ defaults = [f"activitysim=ActivitySim/activitysim@{SHA}", f"My_Ext=org/old@{SHA}"]
+ resolved = resolve_sources(
+ defaults, [f"my-ext[fast]=org/new@{OTHER}#subdirectory=python/package"]
+ )
+ assert resolved[1] == dict(
+ name="my-ext",
+ repository="org/new",
+ commit=OTHER,
+ subdirectory="python/package",
+ extras=["fast"],
+ )
+ with pytest.raises(ValueError, match="conflicting"):
+ resolve_sources([], [f"activitysim=fork/activitysim@{SHA}"], SHA)
+ with pytest.raises(ValueError, match="duplicate"):
+ resolve_sources([], [f"activitysim=org/a@{SHA}"] * 2)
+
+
+@pytest.mark.parametrize(
+ "value",
+ [
+ "x=org/repo@main",
+ f"x=org/repo@{SHA}#subdirectory=../bad",
+ f"x=org/repo@{SHA}#subdirectory=/bad",
+ f"x=org/repo;echo@{SHA}",
+ dict(name="x", repository="org/repo", commit=SHA, extras=["x;echo"]),
+ ],
+)
+def test_invalid_source_rejected(value):
+ with pytest.raises(ValueError):
+ source(value)
+
+
+def make_model(root):
+ """A self-contained profile is independent of either example repository."""
+ (root / "configs").mkdir()
+ (root / "data").mkdir()
+ (root / "data/households.csv").write_text("household_id\n1\n2\n")
+ (root / "configs/settings.yaml").write_text(
+ "chunk_size: 100\nmodels: []\nmultiprocess_steps:\n- name: mp_households\n begin: test_step\n num_processes: 20\n slice:\n tables: [households]\n"
+ )
+ profile = dict(
+ schema_version=1,
+ name="tiny",
+ configs=["configs"],
+ snapshot=["configs"],
+ settings={"chunk_size": 200},
+ required_inputs=["households.csv"],
+ sources=[f"activitysim=ActivitySim/activitysim@{SHA}"],
+ )
+ (root / "benchmark.yaml").write_text(yaml.safe_dump(profile))
+ return profile
+
+
+def test_profiles_preflight_and_oversampling(tmp_path):
+ make_model(tmp_path)
+ profile = load_profile("benchmark.yaml", tmp_path)
+ validate_model(profile, tmp_path, tmp_path / "data", 2)
+ with pytest.raises(ValueError, match="only 2"):
+ validate_model(profile, tmp_path, tmp_path / "data", 3)
+ (tmp_path / "data/households.csv").unlink()
+ with pytest.raises(ValueError, match="missing required"):
+ validate_model(profile, tmp_path, tmp_path / "data", 0)
+ for name in ("mtc", "sandag"):
+ assert load_profile(name, tmp_path)["configs"]
+
+
+def test_config_overlay_and_process_precedence(tmp_path):
+ profile = make_model(tmp_path)
+ overlay = tmp_path / "overlay-0"
+ overlay.mkdir()
+ (overlay / "settings.yaml").write_text(
+ "inherit_settings: true\nchunk_size: 300\nchunk_training_mode: explicit\nnum_processes: 99\n"
+ )
+ phase = tmp_path / "measured"
+ phase.mkdir()
+ (phase / "output").mkdir()
+ spec = dict(
+ profile=profile,
+ households=2,
+ multiprocess=True,
+ processes=4,
+ sharrow=False,
+ config_overlay=[str(overlay)],
+ )
+ state = make_state(spec, phase, tmp_path, tmp_path / "data", tmp_path)
+ assert state.settings.chunk_size == 300
+ assert state.settings.chunk_training_mode == "explicit"
+ assert state.settings.num_processes == 4
+ assert state.settings.multiprocess_steps[0].num_processes == 4
+ assert "chunk_size: 100" in (tmp_path / "configs/settings.yaml").read_text()
+
+
+def test_validate_does_not_build_or_write(tmp_path, monkeypatch, capsys):
+ make_model(tmp_path)
+ calls = []
+
+ def command(args, log=None):
+ calls.append(args)
+ return json.dumps(dict(OSType="linux", CgroupVersion="2"))
+
+ monkeypatch.setattr(cli, "command", command)
+ assert (
+ cli.main(
+ [
+ "validate",
+ "--model-dir",
+ str(tmp_path),
+ "--no-sharrow",
+ "--households",
+ "2",
+ ]
+ )
+ == 0
+ )
+ assert calls == [["docker", "info", "--format", "{{json .}}"]]
+ assert json.loads(capsys.readouterr().out)["sources"][0]["commit"] == SHA
+ assert not (tmp_path / "measured").exists()
+
+
+def test_cli_failure_keeps_report(tmp_path, monkeypatch):
+ make_model(tmp_path)
+ output = tmp_path / "experiment"
+
+ def command(args, log=None):
+ if args[:2] == ["docker", "info"]:
+ return json.dumps(dict(OSType="linux", CgroupVersion="2"))
+ raise RuntimeError("deliberate build failure")
+
+ monkeypatch.setattr(cli, "command", command)
+ with pytest.raises(RuntimeError, match="deliberate"):
+ cli.main(
+ [
+ "run",
+ "--model-dir",
+ str(tmp_path),
+ "--no-sharrow",
+ "--households",
+ "2",
+ "--output-dir",
+ str(output),
+ ]
+ )
+ assert (output / "report.html").is_file()
+ assert (
+ json.loads((output / "experiment.json").read_text())["failure"]["phase"]
+ == "build"
+ )
+ assert (output / "runner/build_sources.py").is_file()
From 6a9c35f5024c5c479f538826ce64526f1b93fbb7 Mon Sep 17 00:00:00 2001
From: Jeff Newman
Date: Sun, 13 Sep 2026 15:13:00 -0500
Subject: [PATCH 2/8] Add named YAML experiment suites with shared defaults
Support variable expansion, per-run overrides, suite preflight, sequential execution, and comparison reports. Include a SANDAG example, tests, documentation, and expanded CI artifacts.
---
.github/workflows/ci.yml | 15 +-
MANIFEST.in | 1 +
README.md | 73 ++++++++++
examples/sandag-chunked.yaml | 30 ++++
src/abench/cli.py | 19 ++-
src/abench/experiments.py | 258 +++++++++++++++++++++++++++++++++++
tests/test_docker.py | 37 +++++
tests/test_experiments.py | 180 ++++++++++++++++++++++++
8 files changed, 604 insertions(+), 9 deletions(-)
create mode 100644 examples/sandag-chunked.yaml
create mode 100644 src/abench/experiments.py
create mode 100644 tests/test_experiments.py
diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml
index df4d8d4..23f2a24 100644
--- a/.github/workflows/ci.yml
+++ b/.github/workflows/ci.yml
@@ -58,12 +58,11 @@ jobs:
with:
name: docker-experiments
path: |
- experiments/ci/**/experiment/*.html
- experiments/ci/**/experiment/*.json
- experiments/ci/**/experiment/*.txt
- experiments/ci/**/experiment/*.log
- experiments/ci/**/experiment/measured/*.json*
- experiments/ci/**/experiment/measured/*.csv
- experiments/ci/**/experiment/measured/*.log
- experiments/ci/**/experiment/warmup/*.log
+ experiments/ci/**/*.html
+ experiments/ci/**/*.json
+ experiments/ci/**/*.jsonl
+ experiments/ci/**/*.txt
+ experiments/ci/**/*.log
+ experiments/ci/**/*.csv
+ experiments/ci/**/*.yaml
if-no-files-found: ignore
diff --git a/MANIFEST.in b/MANIFEST.in
index eb5bcdb..d930cc1 100644
--- a/MANIFEST.in
+++ b/MANIFEST.in
@@ -1,3 +1,4 @@
include LICENSE README.md .pre-commit-config.yaml
recursive-include tests *.py *.yaml *.csv *.txt
recursive-include .github *.yml
+recursive-include examples *.yaml
diff --git a/README.md b/README.md
index 17f1df6..141e758 100644
--- a/README.md
+++ b/README.md
@@ -26,6 +26,79 @@ Bookworm/Python 3.11. Current instrumentation requires ActivitySim's
`workflow.State` API (1.4-era or newer); arbitrary historical revisions are not
promised to work. Build/runtime failures retain diagnostics and a failure report.
+## Named experiment files
+
+Write common options once and override only what differs between runs:
+
+```yaml
+schema_version: 1
+vars:
+ model: /path/to/sandag-abm3-example
+output_root: ./results/sandag-${timestamp}
+defaults:
+ model_dir: ${model}
+ profile: sandag
+ data_dir: ${model}/benchmarking-data
+ config_overlay: ["${model}/configs_explicit_chunk"]
+ multiprocess: true
+ processes: 4
+ sharrow: true
+ households: 28365
+ memory: 80g
+ shm_size: 8g
+ sources:
+ - sharrow=ActivitySim/sharrow@fc175b27d8e0c5d202721c67d96b050e6117b235
+runs:
+ main:
+ sources:
+ - activitysim=ActivitySim/activitysim@5c6fae24a91a57a2d6dfc2e1dbe062a61d94545a
+ pr1110:
+ sources:
+ - activitysim=ActivitySim/activitysim@51e298a84276813946e1d623c9a5785e078e022f
+```
+
+Save it as `sandag.yaml`, then run:
+
+```bash
+abench sandag.yaml
+# Or check all runs without building images or running models:
+abench validate sandag.yaml
+```
+
+A ready-to-use [SANDAG chunked suite](examples/sandag-chunked.yaml) is included
+in the repository. Its paths assume abench and the SANDAG repository are siblings.
+The pinned `main` revision is the one used in the earlier trials, not a moving
+branch reference.
+
+- `defaults` accepts CLI options using underscores (`shm_size`, `config_overlay`,
+ etc.). Use `multiprocess: false` for serial execution and `sharrow: false` to
+ disable Sharrow. `sources` accepts the same strings/mappings as model profiles.
+- `runs` is an ordered mapping of names to overrides. Each run inherits defaults;
+ ordinary values and lists are replaced. **Sources merge by normalized package
+ name**, so changing ActivitySim does not discard the shared Sharrow pin.
+- `${name}` substitutes a reusable scalar from `vars`; terms can reference other
+ terms. A whole-value reference preserves its type, including numbers/booleans.
+ Undefined references and cycles are errors. No shell or environment expansion
+ is performed. `${timestamp}` is a built-in UTC launch identifier shared by all
+ runs, with microseconds to avoid reusing output directories.
+- All explicit paths in the suite are relative to the YAML file, independent of
+ the terminal's current directory. This includes overlays and custom profile
+ paths. Built-in `mtc`/`sandag` profile names retain their meaning. When omitted,
+ `model_dir` defaults to the YAML file's directory; the model profile still
+ supplies its usual default data/config paths.
+- The suite owns output locations: `output_root//`. Set `output_root`
+ once instead of `output_dir` in each run. Existing roots are rejected.
+- All runs are preflighted before the first starts, then run sequentially in file
+ order. Failure stops the suite and retains partial results. The combined report
+ is `output_root/comparison.html`; individual runs retain their own reports.
+ `experiments.yaml` and `suite.json` record the original file and expanded plan.
+- File invocations do not accept additional CLI overrides. Edit `defaults` or the
+ relevant run to keep the file a complete description of the experiment.
+
+This experiment file describes **which tests to run**. A model profile such as
+`benchmark.yaml` describes **how to configure a model**, and remains reusable
+across suites.
+
## Run controls
- `--single-process` (default), or `--multiprocess --processes N`. The count applies
diff --git a/examples/sandag-chunked.yaml b/examples/sandag-chunked.yaml
new file mode 100644
index 0000000..a808e23
--- /dev/null
+++ b/examples/sandag-chunked.yaml
@@ -0,0 +1,30 @@
+# Run from any directory: abench /path/to/abench/examples/sandag-chunked.yaml
+schema_version: 1
+vars:
+ model: ../../sandag-abm3-example
+ activitysim_main: 5c6fae24a91a57a2d6dfc2e1dbe062a61d94545a
+ activitysim_pr1110: 51e298a84276813946e1d623c9a5785e078e022f
+output_root: ../experiments/sandag-chunked-${timestamp}
+defaults:
+ model_dir: ${model}
+ profile: sandag
+ data_dir: ${model}/benchmarking-data
+ config_overlay: ["${model}/configs_explicit_chunk"]
+ multiprocess: true
+ processes: 4
+ sharrow: true
+ households: 28365
+ memory: 80g
+ shm_size: 8g
+ platform: linux/arm64
+ sources:
+ - sharrow=ActivitySim/sharrow@fc175b27d8e0c5d202721c67d96b050e6117b235
+runs:
+ main:
+ label: SANDAG main — explicit chunking
+ sources:
+ - activitysim=ActivitySim/activitysim@${activitysim_main}
+ pr1110:
+ label: SANDAG PR1110 — explicit chunking
+ sources:
+ - activitysim=ActivitySim/activitysim@${activitysim_pr1110}
diff --git a/src/abench/cli.py b/src/abench/cli.py
index de67962..1ae6f50 100644
--- a/src/abench/cli.py
+++ b/src/abench/cli.py
@@ -53,7 +53,10 @@ def positive(value):
def parser():
- p = argparse.ArgumentParser(description=__doc__)
+ p = argparse.ArgumentParser(
+ description=__doc__,
+ epilog="Named experiments: abench experiments.yaml; preflight: abench validate experiments.yaml",
+ )
p.add_argument("--model-dir", type=Path, default=Path.cwd())
p.add_argument("--profile", default="benchmark.yaml")
p.add_argument(
@@ -181,6 +184,20 @@ def container_phase(spec, output, data, image, phase_name):
def main(argv=None):
p = parser()
argv = list(sys.argv[1:] if argv is None else argv)
+ # A file invocation stays separate from model profiles and ordinary flags.
+ candidate = argv[1:] if argv and argv[0] in ("run", "validate") else argv
+ if (
+ candidate
+ and not candidate[0].startswith("-")
+ and candidate[0] not in ("run", "report", "validate")
+ ):
+ if len(candidate) != 1:
+ p.error(
+ "an experiment file cannot be mixed with command-line overrides; edit its defaults or runs"
+ )
+ from .experiments import run_suite
+
+ return run_suite(Path(candidate[0]), main, validate_only=argv[0] == "validate")
action = argv.pop(0) if argv and argv[0] in ("run", "report", "validate") else "run"
args = p.parse_args(argv)
if action == "report":
diff --git a/src/abench/experiments.py b/src/abench/experiments.py
new file mode 100644
index 0000000..b5d380a
--- /dev/null
+++ b/src/abench/experiments.py
@@ -0,0 +1,258 @@
+"""Named experiment suites with shared options and safe string substitution."""
+
+import io
+import re
+from contextlib import redirect_stdout
+from datetime import datetime, timezone
+from pathlib import Path
+
+import yaml
+
+from .common import write_json
+from .report import report
+from .sources import source
+
+OPTIONS = {
+ "model_dir",
+ "profile",
+ "sources",
+ "activitysim_commit",
+ "sharrow_commit",
+ "multiprocess",
+ "processes",
+ "sharrow",
+ "households",
+ "data_dir",
+ "config_overlay",
+ "cache_from",
+ "label",
+ "compare",
+ "interval",
+ "memory",
+ "shm_size",
+ "platform",
+}
+PATHS = {"model_dir", "data_dir", "cache_from", "config_overlay", "compare"}
+LISTS = {"config_overlay", "compare"}
+TOKEN = re.compile(r"\$\{([^{}]+)\}")
+
+
+class SuiteLoader(yaml.SafeLoader):
+ """Reject duplicate keys so an accidental repeated default cannot disappear."""
+
+
+def unique_mapping(loader, node):
+ result = {}
+ for key_node, value_node in node.value:
+ key = loader.construct_object(key_node)
+ if not isinstance(key, str):
+ raise ValueError("experiment configuration keys must be strings")
+ if key in result:
+ raise ValueError(f"duplicate configuration key: {key}")
+ result[key] = loader.construct_object(value_node)
+ return result
+
+
+SuiteLoader.add_constructor(
+ yaml.resolver.BaseResolver.DEFAULT_MAPPING_TAG, unique_mapping
+)
+
+
+def expand_variables(document, timestamp):
+ """Resolve named terms recursively; never execute shell code or read env vars."""
+ terms = document.get("vars", {})
+ if not isinstance(terms, dict):
+ raise ValueError("vars must be a mapping")
+ if "timestamp" in terms:
+ raise ValueError("timestamp is a reserved variable")
+ resolved = {"timestamp": timestamp}
+
+ def term(name, stack):
+ if name in resolved:
+ return resolved[name]
+ if name in stack:
+ raise ValueError(f"cyclic variable: {' -> '.join((*stack, name))}")
+ if name not in terms:
+ raise ValueError(f"undefined variable: {name}")
+ value = terms[name]
+ if not isinstance(value, (str, int, float, bool)):
+ raise ValueError(f"variable {name} must be a scalar")
+ resolved[name] = expand(value, (*stack, name))
+ return resolved[name]
+
+ def expand(value, stack=()):
+ if isinstance(value, str):
+ match = TOKEN.fullmatch(value)
+ if match:
+ return term(match[1], stack)
+ return TOKEN.sub(lambda m: str(term(m[1], stack)), value)
+ if isinstance(value, list):
+ return [expand(v, stack) for v in value]
+ if isinstance(value, dict):
+ return {k: expand(v, stack) for k, v in value.items()}
+ return value
+
+ for name in terms:
+ term(name, ())
+ return expand(document)
+
+
+def merge_options(defaults, overrides):
+ """Runs replace ordinary defaults; source pins merge by distribution name."""
+ if not isinstance(overrides, dict):
+ raise ValueError("defaults and each run must be option mappings")
+ unknown = set(overrides) - OPTIONS
+ if unknown:
+ raise ValueError(f"unknown experiment options: {sorted(unknown)}")
+ merged = dict(defaults, **overrides)
+ if "sources" in overrides:
+ if not isinstance(overrides["sources"], list):
+ raise ValueError("sources must be a list")
+ pins = {item["name"]: item for item in defaults.get("sources", [])}
+ seen = set()
+ for value in overrides["sources"]:
+ item = source(value)
+ if item["name"] in seen:
+ raise ValueError(f"duplicate source: {item['name']}")
+ seen.add(item["name"])
+ pins[item["name"]] = item
+ merged["sources"] = list(pins.values())
+ return merged
+
+
+def arguments(options, base):
+ """Translate typed YAML options into the existing CLI's validation interface."""
+ argv = []
+ options = dict(options)
+ options.setdefault("model_dir", str(base))
+ for key, value in options.items():
+ if value is None:
+ continue
+ if key in ("multiprocess", "sharrow"):
+ if not isinstance(value, bool):
+ raise ValueError(f"{key} must be a YAML boolean")
+ argv.append(
+ ("--multiprocess" if value else "--single-process")
+ if key == "multiprocess"
+ else ("--sharrow" if value else "--no-sharrow")
+ )
+ continue
+ if key == "sources":
+ for item in value:
+ extras = "[" + ",".join(item["extras"]) + "]" if item["extras"] else ""
+ subdir = (
+ "#subdirectory=" + item["subdirectory"]
+ if item["subdirectory"]
+ else ""
+ )
+ argv += [
+ "--source",
+ f"{item['name']}{extras}={item['repository']}@{item['commit']}{subdir}",
+ ]
+ continue
+ values = value if key in LISTS else [value]
+ if not isinstance(values, list) or any(
+ not isinstance(v, (str, int, float)) or isinstance(v, bool) for v in values
+ ):
+ raise ValueError(f"invalid value for {key}")
+ if not values:
+ continue
+ if key in PATHS or (key == "profile" and value not in ("mtc", "sandag")):
+ values = [str((base / Path(v).expanduser()).resolve()) for v in values]
+ argv += ["--" + key.replace("_", "-"), *map(str, values)]
+ return argv
+
+
+def load_suite(path):
+ """Expand an entire suite before creating output or running any experiments."""
+ path = path.expanduser().resolve()
+ try:
+ raw = path.read_text()
+ document = yaml.load(raw, Loader=SuiteLoader)
+ except yaml.YAMLError as error:
+ raise ValueError(f"invalid experiment YAML: {error}") from error
+ if not isinstance(document, dict) or document.get("schema_version") != 1:
+ raise ValueError("experiment file requires schema_version: 1")
+ unknown = set(document) - {
+ "schema_version",
+ "vars",
+ "defaults",
+ "runs",
+ "output_root",
+ }
+ if unknown:
+ raise ValueError(f"unknown experiment file fields: {sorted(unknown)}")
+ document = expand_variables(
+ document, datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S-%f")
+ )
+ defaults = merge_options({}, document.get("defaults", {}))
+ runs = document.get("runs")
+ if not isinstance(runs, dict) or not runs:
+ raise ValueError("runs must be a nonempty mapping of run names to options")
+ output = document.get("output_root")
+ if not isinstance(output, str) or not output:
+ raise ValueError("output_root must be a path")
+ output = (path.parent / Path(output).expanduser()).resolve()
+ if output.exists():
+ raise ValueError(
+ f"output_root already exists: {output}; use ${{timestamp}} for repeatable launches"
+ )
+ plan = []
+ for name, overrides in runs.items():
+ if not re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9_-]*", name):
+ raise ValueError(f"invalid run name: {name!r}")
+ options = merge_options(defaults, overrides)
+ options.setdefault("label", name)
+ argv = arguments(options, path.parent)
+ destination = output / name
+ plan.append(
+ {
+ "name": name,
+ "argv": argv + ["--output-dir", str(destination)],
+ "output_dir": str(destination),
+ }
+ )
+ return {
+ "source_file": str(path),
+ "original_yaml": raw,
+ "configuration": document,
+ "output_root": str(output),
+ "runs": plan,
+ }
+
+
+def run_suite(path, invoke, validate_only=False):
+ """Preflight every run, execute serially, and preserve partial failure reports."""
+ plan = load_suite(path)
+ root = Path(plan["output_root"])
+ # Validation uses the same CLI checks as individual runs and creates nothing.
+ # Run it for the whole suite first, so a typo in run two cannot waste run one.
+ for run in plan["runs"]:
+ with redirect_stdout(io.StringIO()):
+ code = invoke(["validate", *run["argv"]])
+ if code:
+ return code
+ if validate_only:
+ print(f"Validated {len(plan['runs'])} experiments; output root: {root}")
+ return 0
+ root.mkdir(parents=True)
+ (root / "experiments.yaml").write_text(plan["original_yaml"])
+ write_json(root / "suite.json", plan)
+ completed = []
+ try:
+ for run in plan["runs"]:
+ print(f"Running experiment {run['name']}…", flush=True)
+ try:
+ code = invoke(["run", *run["argv"]])
+ finally:
+ # Failed builds/models still have an experiment record and belong
+ # in the comparison, with failure status rather than winner badges.
+ if (Path(run["output_dir"]) / "experiment.json").is_file():
+ completed.append(Path(run["output_dir"]))
+ if code:
+ return code
+ finally:
+ if completed:
+ report(completed, root / "comparison.html")
+ print(f"Comparison: {root / 'comparison.html'}", flush=True)
+ return 0
diff --git a/tests/test_docker.py b/tests/test_docker.py
index add8287..e430bb0 100644
--- a/tests/test_docker.py
+++ b/tests/test_docker.py
@@ -59,3 +59,40 @@ def test_tiny_model(tmp_path, multiprocess, sharrow):
assert not list((output / "measured").glob("cache-miss-*"))
lines = (output / "measured/output/final_households.csv").read_text().splitlines()
assert set(lines[1:]) == {"1,20", "2,40", "3,60", "4,80"}
+
+
+def test_named_suite(tmp_path):
+ """Exercise file dispatch, shared defaults, serial/MP overrides, and comparison."""
+ import yaml
+
+ root = tmp_path / "model"
+ shutil.copytree(Path(__file__).parent / "fixtures/tiny", root)
+ path = tmp_path / "experiments.yaml"
+ path.write_text(
+ yaml.safe_dump(
+ dict(
+ schema_version=1,
+ output_root="results",
+ defaults=dict(
+ model_dir="model",
+ sharrow=True,
+ households=4,
+ memory="3g",
+ shm_size="256m",
+ sources=[
+ f"activitysim=ActivitySim/activitysim@{ACTIVITYSIM}",
+ f"sharrow=ActivitySim/sharrow@{SHARROW}",
+ ],
+ ),
+ runs={"serial": {}, "parallel": {"multiprocess": True, "processes": 2}},
+ ),
+ sort_keys=False,
+ )
+ )
+ assert cli.main([str(path)]) == 0
+ root = tmp_path / "results"
+ runs = json.loads((root / "comparison.json").read_text())
+ assert len(runs) == 2 and all(run["valid"] for run in runs)
+ assert [run["components"]["bench_compute"]["n"] for run in runs] == [1, 2]
+ assert (root / "experiments.yaml").read_text() == path.read_text()
+ assert (root / "suite.json").is_file()
diff --git a/tests/test_experiments.py b/tests/test_experiments.py
new file mode 100644
index 0000000..e2fc2b5
--- /dev/null
+++ b/tests/test_experiments.py
@@ -0,0 +1,180 @@
+"""Named suites reuse CLI validation/execution without requiring Docker in tests."""
+
+import json
+from pathlib import Path
+
+import pytest
+import yaml
+
+from abench import cli, experiments
+from abench.experiments import load_suite, run_suite
+
+SHA = "a" * 40
+OTHER = "b" * 40
+
+
+def suite(tmp_path, **updates):
+ """Write a portable two-run suite with common paths and source dependencies."""
+ document = dict(
+ schema_version=1,
+ vars={"model": "model", "sample": 4},
+ output_root="results-${timestamp}",
+ defaults=dict(
+ model_dir="${model}",
+ profile="sandag",
+ data_dir="${model}/data",
+ config_overlay=["${model}/chunks"],
+ households="${sample}",
+ multiprocess=True,
+ processes=2,
+ sharrow=True,
+ sources=[
+ f"sharrow=ActivitySim/sharrow@{SHA}",
+ f"activitysim=ActivitySim/activitysim@{SHA}",
+ ],
+ ),
+ runs={
+ "main": {},
+ "pr": {"sources": [f"activitysim=ActivitySim/activitysim@{OTHER}"]},
+ },
+ )
+ document.update(updates)
+ path = tmp_path / "experiment.yaml"
+ path.write_text(yaml.safe_dump(document, sort_keys=False))
+ return path
+
+
+def test_defaults_variables_and_source_overrides(tmp_path, monkeypatch):
+ path = suite(tmp_path)
+ monkeypatch.chdir("/")
+ plan = load_suite(path)
+ first, second = plan["runs"]
+ assert first["name"] == "main"
+ assert second["name"] == "pr"
+ for run in plan["runs"]:
+ args = cli.parser().parse_args(run["argv"])
+ assert args.households == 4
+ assert args.model_dir == tmp_path / "model"
+ assert args.data_dir == tmp_path / "model/data"
+ assert args.config_overlay == [tmp_path / "model/chunks"]
+ assert f"sharrow=ActivitySim/sharrow@{SHA}" in args.source
+ assert f"activitysim=ActivitySim/activitysim@{OTHER}" in second["argv"]
+ assert f"activitysim=ActivitySim/activitysim@{SHA}" not in second["argv"]
+ assert not Path(plan["output_root"]).exists()
+
+
+@pytest.mark.parametrize(
+ "updates,match",
+ [
+ ({"vars": {"a": "${b}", "b": "${a}"}}, "cyclic"),
+ ({"vars": {}}, "undefined"),
+ ({"vars": {"timestamp": "x"}}, "reserved"),
+ ({"runs": {"../bad": {}}}, "invalid run name"),
+ ({"runs": {}}, "nonempty"),
+ ({"runs": {"bad": {"household": 1}}}, "unknown experiment options"),
+ ({"runs": {"bad": {"multiprocess": "false"}}}, "YAML boolean"),
+ ({"runs": {"bad": {"output_dir": "somewhere"}}}, "unknown experiment options"),
+ ],
+)
+def test_bad_suite_rejected_before_execution(tmp_path, updates, match):
+ with pytest.raises(ValueError, match=match):
+ load_suite(suite(tmp_path, **updates))
+
+
+def test_duplicate_yaml_and_existing_root(tmp_path):
+ path = suite(tmp_path, output_root="results")
+ (tmp_path / "results").mkdir()
+ with pytest.raises(ValueError, match="already exists"):
+ load_suite(path)
+ path.write_text("schema_version: 1\nruns: {}\nruns: {}\n")
+ with pytest.raises(ValueError, match="duplicate"):
+ load_suite(path)
+
+
+def test_suite_preflights_all_then_runs_and_reports(tmp_path, monkeypatch):
+ path = suite(tmp_path)
+ calls, reports = [], []
+
+ def invoke(argv):
+ calls.append(argv)
+ if argv[0] == "run":
+ output = Path(argv[argv.index("--output-dir") + 1])
+ output.mkdir()
+ (output / "experiment.json").write_text("{}")
+ return 0
+
+ monkeypatch.setattr(
+ experiments,
+ "report",
+ lambda paths, destination: reports.append((paths, destination)),
+ )
+ assert run_suite(path, invoke) == 0
+ assert [args[0] for args in calls] == ["validate", "validate", "run", "run"]
+ assert len(reports[0][0]) == 2
+ root = reports[0][1].parent
+ assert (root / "experiments.yaml").read_text() == path.read_text()
+ assert len(json.loads((root / "suite.json").read_text())["runs"]) == 2
+
+
+def test_invalid_later_run_prevents_first_run(tmp_path):
+ calls = []
+
+ def invoke(argv):
+ calls.append(argv[0])
+ if len(calls) == 2:
+ raise ValueError("bad input")
+ return 0
+
+ with pytest.raises(ValueError, match="bad input"):
+ run_suite(suite(tmp_path), invoke)
+ assert calls == ["validate", "validate"]
+ assert not list(tmp_path.glob("results-*"))
+
+
+def test_model_failure_stops_suite_and_reports_partial_results(tmp_path, monkeypatch):
+ calls, reports = [], []
+
+ def invoke(argv):
+ calls.append(argv[0])
+ if argv[0] == "run":
+ output = Path(argv[argv.index("--output-dir") + 1])
+ output.mkdir()
+ (output / "experiment.json").write_text("{}")
+ raise RuntimeError("model failed")
+ return 0
+
+ monkeypatch.setattr(
+ experiments, "report", lambda paths, dest: reports.append(paths)
+ )
+ with pytest.raises(RuntimeError, match="model failed"):
+ run_suite(suite(tmp_path), invoke)
+ assert calls == ["validate", "validate", "run"]
+ assert len(reports[0]) == 1
+
+
+def test_cli_file_dispatch_and_validation(tmp_path, monkeypatch):
+ calls = []
+ monkeypatch.setattr(
+ experiments,
+ "run_suite",
+ lambda path, invoke, validate_only: calls.append((path, validate_only)) or 0,
+ )
+ path = tmp_path / "named.yaml"
+ assert cli.main([str(path)]) == 0
+ assert cli.main(["run", str(path)]) == 0
+ assert cli.main(["validate", str(path)]) == 0
+ assert calls == [(path, False), (path, False), (path, True)]
+ with pytest.raises(SystemExit):
+ cli.main([str(path), "--households", "5"])
+
+
+def test_shipped_sandag_suite():
+ path = Path(__file__).parents[1] / "examples/sandag-chunked.yaml"
+ plan = load_suite(path)
+ for run in plan["runs"]:
+ args = cli.parser().parse_args(run["argv"])
+ assert args.households == 28365
+ assert args.processes == 4
+ assert args.sharrow
+ assert args.data_dir.name == "benchmarking-data"
+ assert args.config_overlay[0].name == "configs_explicit_chunk"
From ed37819cfa239968a4f1ffb1b7a39690937df3ed Mon Sep 17 00:00:00 2001
From: Jeff Newman
Date: Sun, 13 Sep 2026 19:34:02 -0500
Subject: [PATCH 3/8] Use capped serial Sharrow warmups and improve failure
diagnostics
Add configurable warmup household limits while preserving measured settings and strict cache coverage checks. Record phase settings and missing flow signatures, and surface actionable failures in CLI output and reports.
---
README.md | 29 +++++--
examples/sandag-chunked.yaml | 6 +-
src/abench/cli.py | 25 +++++-
src/abench/experiments.py | 1 +
src/abench/failures.py | 82 ++++++++++++++++++++
src/abench/report.py | 23 +++++-
src/abench/runtime/instrumentation.py | 14 +++-
src/abench/runtime/worker.py | 63 ++++++++++++++-
tests/fixtures/tiny/tiny_extension.py | 2 +
tests/test_docker.py | 15 ++++
tests/test_experiments.py | 1 +
tests/test_failure_diagnostics.py | 106 ++++++++++++++++++++++++++
tests/test_profiles_sources.py | 26 +++++++
tests/test_warmup.py | 63 +++++++++++++++
14 files changed, 442 insertions(+), 14 deletions(-)
create mode 100644 src/abench/failures.py
create mode 100644 tests/test_failure_diagnostics.py
create mode 100644 tests/test_warmup.py
diff --git a/README.md b/README.md
index 141e758..db37522 100644
--- a/README.md
+++ b/README.md
@@ -33,6 +33,8 @@ Write common options once and override only what differs between runs:
```yaml
schema_version: 1
vars:
+ households: 28365
+ warmup_households: 5000
model: /path/to/sandag-abm3-example
output_root: ./results/sandag-${timestamp}
defaults:
@@ -43,7 +45,8 @@ defaults:
multiprocess: true
processes: 4
sharrow: true
- households: 28365
+ households: ${households}
+ warmup_households: ${warmup_households}
memory: 80g
shm_size: 8g
sources:
@@ -111,7 +114,7 @@ across suites.
- `--memory 16g`, `--shm-size 8g`, `--interval 0.5`, and optional `--platform`.
- `--output-dir` must be new. `--label` names an experiment, and `--compare` accepts
earlier experiment directories. `--cache-from` seeds compatible generated flows;
- a complete warmup still runs.
+ the small serial warmup still runs.
`abench validate` accepts the same experiment arguments without `--output-dir`.
It checks the profile, required inputs, CSV population size, source pin syntax,
@@ -227,9 +230,25 @@ model outputs.
## Measurement and reports
-Sharrow runs first complete a full matching warmup in a separate container. The
-measured run uses the same sample, seed, layout, and cache path. Numba compilation
-of generated flow overloads is rejected during measurement; ordinary non-flow
+Sharrow runs first execute the model in a separate **single-process** warmup,
+using **min(target households, 500)** households by default. For `--households 0`,
+the target is the full available population, so warmup uses at most 500 of those
+households. Set `--warmup-households N` (or `warmup_households: N` in experiment
+YAML) to change this positive cap. Warmup always uses one process; measured runs
+retain their requested sample and worker count. Model config directories, seed,
+chunk overlays, and flow cache path are retained from the target experiment.
+
+A smaller serial warmup may not exercise every flow signature required by the
+measured run. **Cache misses still invalidate the measured experiment**: abench
+does not silently compile, retry with compilation enabled, or accept incomplete
+cache coverage. Increase the warmup cap when needed.
+The CLI and failure report identify the phase, failed component, cache flow, and
+log path; `cache-miss-details-*.jsonl` records the exact required Numba signature.
+`phase-settings.json` and
+`effective-settings.json` preserve the actual settings for each phase, and the
+report includes the cache-build settings separately from the measured target.
+
+Numba compilation of generated flow overloads is rejected during measurement; ordinary non-flow
compilation and disk-cache loading remain included. This guard uses private
Numba internals and is covered by real disk-cache hit/miss tests.
diff --git a/examples/sandag-chunked.yaml b/examples/sandag-chunked.yaml
index a808e23..a4774e6 100644
--- a/examples/sandag-chunked.yaml
+++ b/examples/sandag-chunked.yaml
@@ -1,6 +1,7 @@
# Run from any directory: abench /path/to/abench/examples/sandag-chunked.yaml
schema_version: 1
vars:
+ households: 28365
model: ../../sandag-abm3-example
activitysim_main: 5c6fae24a91a57a2d6dfc2e1dbe062a61d94545a
activitysim_pr1110: 51e298a84276813946e1d623c9a5785e078e022f
@@ -13,7 +14,10 @@ defaults:
multiprocess: true
processes: 4
sharrow: true
- households: 28365
+ households: ${households}
+ # The 500-household cache build misses rare at-work flows and dtype signatures.
+ # This model-specific override keeps the same sample, but warmup is still serial.
+ warmup_households: 5000
memory: 80g
shm_size: 8g
platform: linux/arm64
diff --git a/src/abench/cli.py b/src/abench/cli.py
index 1ae6f50..d10f101 100644
--- a/src/abench/cli.py
+++ b/src/abench/cli.py
@@ -15,6 +15,7 @@
from . import __version__
from .common import read_json, write_json
+from .failures import BenchmarkFailure, describe_failure
from .profiles import load_profile, validate_model
from .report import load_run, report
from .sources import resolve_sources
@@ -76,6 +77,12 @@ def parser():
p.add_argument(
"--households", type=int, default=1000, help="0 means full population"
)
+ p.add_argument(
+ "--warmup-households",
+ type=int,
+ default=500,
+ help="maximum cache-build households (default: 500); warmup is always single-process",
+ )
p.add_argument("--data-dir", type=Path, default=None)
p.add_argument(
"--config-overlay",
@@ -241,6 +248,8 @@ def main(argv=None):
)
if args.sharrow and not args.sharrow_commit:
p.error("Sharrow enabled: provide a sharrow source override")
+ if args.warmup_households < 1:
+ p.error("--warmup-households must be positive")
if args.households < 0:
p.error("--households must be nonnegative")
if args.multiprocess and (args.processes is None or args.processes < 1):
@@ -424,17 +433,27 @@ def main(argv=None):
if seed:
if (seed / "pip-freeze.txt").read_text().strip() != freeze.strip():
raise ValueError("Cannot seed cache: installed dependencies differ")
- # Only flow artifacts are reused. Full warmup still checks coverage
- # under the new configuration before compilation-blocked measurement.
+ # Only flow artifacts are reused. A small serial warmup prepares flows;
+ # measurement remains responsible for rejecting missing signatures.
shutil.copytree(seed / "cache/flows", output / "cache/flows")
if args.sharrow:
stage = "warmup"
- print("Preparing Sharrow cache with a complete matching run…", flush=True)
+ print(
+ f"Preparing Sharrow cache in single process (up to {args.warmup_households} households)…",
+ flush=True,
+ )
container_phase(spec, output, data, image, "warmup")
stage = "measured"
print("Running measured model…", flush=True)
container_phase(spec, output, data, image, "measured")
+ if not load_run(output)["valid"]:
+ raise BenchmarkFailure(describe_failure(output, "measured"))
except (Exception, KeyboardInterrupt) as error:
+ if isinstance(error, subprocess.CalledProcessError):
+ error = BenchmarkFailure(describe_failure(output, stage))
+ spec["failure"] = {"phase": stage, "error": str(error)}
+ write_json(output / "experiment.json", spec)
+ raise error from None
spec["failure"] = {"phase": stage, "error": str(error)}
write_json(output / "experiment.json", spec)
raise
diff --git a/src/abench/experiments.py b/src/abench/experiments.py
index b5d380a..fc60461 100644
--- a/src/abench/experiments.py
+++ b/src/abench/experiments.py
@@ -22,6 +22,7 @@
"processes",
"sharrow",
"households",
+ "warmup_households",
"data_dir",
"config_overlay",
"cache_from",
diff --git a/src/abench/failures.py b/src/abench/failures.py
new file mode 100644
index 0000000..0307713
--- /dev/null
+++ b/src/abench/failures.py
@@ -0,0 +1,82 @@
+"""Translate container exit codes into useful model failure messages."""
+
+import json
+import re
+from pathlib import Path
+
+from .common import read_json
+
+
+class BenchmarkFailure(ValueError):
+ """An experiment failed with diagnostics retained beside its report."""
+
+
+def tail(path, limit=131072):
+ """Read a bounded log tail even when a full model log is many gigabytes."""
+ if not path.is_file():
+ return ""
+ with path.open("rb") as stream:
+ stream.seek(0, 2)
+ stream.seek(max(0, stream.tell() - limit))
+ return stream.read().decode("utf-8", errors="replace")
+
+
+def describe_failure(output, phase):
+ """Prefer explicit cache/OOM evidence over a generic worker exit exception."""
+ output = Path(output)
+ directory = output / phase
+ log = output / "build.log" if phase == "build" else directory / "console.log"
+ docker = read_json(directory / "docker-state.json", {})
+ components = set()
+ for path in directory.glob("components-*.jsonl"):
+ for line in tail(path).splitlines():
+ try:
+ row = json.loads(line)
+ except (ValueError, TypeError):
+ continue
+ if not row.get("succeeded", True):
+ components.add(row["component"])
+ where = f" in {', '.join(sorted(components))}" if components else ""
+ misses = sorted(
+ {
+ line.strip()
+ for path in directory.glob("cache-miss-*.txt")
+ for line in tail(path).splitlines()
+ if line.strip()
+ }
+ )
+ if misses:
+ reason = f"Sharrow flow cache miss{where}: " + ", ".join(
+ Path(name).parent.name for name in misses
+ )
+ remedy = "Results rejected; no flow compilation was allowed. The warmup did not cover the required flow/type signature. Increase warmup_households (up to the target sample) and start a new experiment."
+ elif docker.get("OOMKilled"):
+ reason = f"Docker killed the container for exceeding its memory limit{where}."
+ remedy = (
+ "Increase the container/VM memory budget or reduce component chunk sizes."
+ )
+ elif docker.get("ExitCode") == 0 and phase == "measured":
+ spec = read_json(output / "experiment.json", {})
+ summaries = read_json(directory / "output-summary.json", {})
+ actual = summaries.get("households", {}).get("rows")
+ requested = spec.get("households", 0)
+ if requested and actual != requested:
+ reason = f"Household sample mismatch: requested {requested}, output contains {actual if actual is not None else 'no household table'}."
+ else:
+ reason = "Required runtime or memory measurements are missing."
+ remedy = "Results rejected even though the container exited successfully."
+ else:
+ errors = re.findall(
+ r"^([\w.]+(?:Error|Exception): .+)$", tail(log), re.MULTILINE
+ )
+ useful = [e for e in errors if "SubprocessError: Process " not in e]
+ reason = (
+ useful
+ or errors
+ or [
+ docker.get("Error")
+ or f"Container exited with code {docker.get('ExitCode', 'unavailable')}"
+ ]
+ )[0]
+ remedy = "See the retained log for the complete traceback."
+ return f"{output.name}: {phase} failed{where if not misses else ''}. {reason}\n{remedy}\nLog: {log}\nReport: {output / 'report.html'}"
diff --git a/src/abench/report.py b/src/abench/report.py
index ced5be9..7e1b2cc 100644
--- a/src/abench/report.py
+++ b/src/abench/report.py
@@ -8,6 +8,7 @@
import statistics
from .common import read_json, write_json
+from .failures import describe_failure
COLORS = ("#0072b2", "#d55e00", "#009e73", "#cc79a7", "#e69f00", "#56b4e9")
@@ -142,6 +143,14 @@ def load_run(directory):
and not docker.get("OOMKilled")
and not list(phase.glob("cache-miss-*.txt"))
)
+ warmup_settings = read_json(directory / "warmup/phase-settings.json", {})
+ effective = read_json(directory / "warmup/effective-settings.json", {})
+ if effective:
+ warmup_settings = {
+ "households": effective.get("households_sample_size"),
+ "multiprocess": effective.get("multiprocess"),
+ "processes": effective.get("num_processes"),
+ }
return {
"spec": spec,
"components": components,
@@ -151,6 +160,12 @@ def load_run(directory):
"valid": valid,
"inputs": read_json(phase / "input-summary.json", {}),
"outputs": outputs,
+ "failure_reason": describe_failure(
+ directory, spec.get("failure", {}).get("phase", "measured")
+ )
+ if not valid
+ else None,
+ "warmup_settings": warmup_settings,
"peak": max((r["peak_bytes"] for r in memory), default=0),
"docker": docker,
"path": str(directory),
@@ -255,6 +270,7 @@ def experiment_card(run, xmax, ymax):
"processes",
"sharrow",
"households",
+ "warmup_households",
"data_dir",
"output_dir",
"memory",
@@ -293,11 +309,16 @@ def experiment_card(run, xmax, ymax):
counts.append(f"
TAZs
{taz_cells}
")
elapsed = run["status"].get("elapsed_seconds")
elapsed = f"{elapsed:.3f}" if elapsed is not None else "unavailable"
- failure = f"