feat(cli): scope each run's work directory per corpus

Runs are namespaced <work>/<pipeline>/<corpus-key>/ so different targets
through the same pipeline no longer share a sample store or report. The
corpus key is <basename>-<8 hex of sha256(resolved path)>, so a rerun of
the same target resumes in place while distinct targets stay isolated.
status and report resolve to the newest corpus run; --clean purges only
that corpus's run directory.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
416rehman
2026-07-26 11:49:00 -06:00
co-authored by Claude Opus 4.8
parent 9dc0fb388c
commit 8f0e973ab2
4 changed files with 142 additions and 15 deletions
+31 -2
View File
@@ -84,6 +84,21 @@ class _ShortNameFormatter(logging.Formatter):
return f"[{color}]\\[{short}][/{color}] {msg_escaped}"
def _latest_run_dir(base: Path) -> Path | None:
"""Pick the most recent corpus run directory under a pipeline's base dir.
Runs are namespaced per corpus (``<base>/<corpus-key>/``), so ``status`` and
``report`` with only ``--pipeline`` resolve to the newest run. A ``run.json``
directly in ``base`` (e.g. ``--work-dir`` pointed straight at a run) also works.
"""
if not base.exists():
return None
runs = [d for d in base.iterdir() if d.is_dir() and (d / "run.json").exists()]
if runs:
return max(runs, key=lambda d: (d / "run.json").stat().st_mtime)
return base if (base / "run.json").exists() else None
def _setup_logging(verbose: bool) -> None:
level = logging.DEBUG if verbose else logging.INFO
handler = RichHandler(
@@ -186,6 +201,10 @@ def run(
console.print(f"[bold red]X ERROR[/]: {e}")
raise SystemExit(1)
# scope the run to this corpus so different targets through the same pipeline
# get isolated, independently resumable work directories under one root.
pipeline_def.bind_corpus(target_path)
import os
import shutil
import threading
@@ -320,7 +339,12 @@ def report(
except ValueError as e:
console.print(f"[bold red]X ERROR[/]: {e}")
raise SystemExit(1)
work_path = pipeline_def.work_dir
work_path = _latest_run_dir(pipeline_def.base_work_dir)
if work_path is None:
console.print(
f"[bold red]X ERROR[/]: no runs found under {pipeline_def.base_work_dir}"
)
raise SystemExit(1)
report_cfg = pipeline_def.report
else:
console.print("[bold red]X ERROR[/]: pass --pipeline or --work-dir")
@@ -362,7 +386,12 @@ def status(pipeline: str | None, work_dir: str | None, verbose: bool):
console.print(f"[bold red]X ERROR[/]: {e}")
raise SystemExit(1)
work_path = pipeline_def.work_dir
work_path = _latest_run_dir(pipeline_def.base_work_dir)
if work_path is None:
console.print(
f"[yellow]no runs found under {pipeline_def.base_work_dir}[/]"
)
raise SystemExit(1)
else:
console.print("[red]specify --pipeline or --work-dir[/]")
raise SystemExit(1)
+38 -1
View File
@@ -1,7 +1,9 @@
from __future__ import annotations
import hashlib
import logging
import os
import re
from pathlib import Path
from typing import Any
@@ -22,6 +24,23 @@ from deepzero.engine.stage import (
log = logging.getLogger("deepzero.pipeline")
def corpus_segment(target: Path) -> str:
"""A stable, filesystem-safe directory segment identifying a corpus (the run
target). Two different corpora, or the same corpus at two paths, never
collide; the same resolved path always maps to the same segment, so a rerun
resumes in place.
Shape: ``<readable-basename>-<8 hex of sha256(resolved-abspath)>``.
"""
resolved = target.expanduser().resolve()
# case-fold on case-insensitive filesystems so C:\X and c:\x agree
key_src = str(resolved).lower() if os.name == "nt" else str(resolved)
digest = hashlib.sha256(key_src.encode("utf-8")).hexdigest()[:8]
base = resolved.name or "root"
base = re.sub(r"[^A-Za-z0-9._-]", "_", base)[:48].strip("._-") or "corpus"
return f"{base}-{digest}"
class PipelineDefinition:
# a loaded and validated pipeline ready for execution
@@ -53,14 +72,32 @@ class PipelineDefinition:
self.ingest_processor: IngestProcessor | None = None
self.stages: list[tuple[StageSpec, Processor]] = []
# set once the run target is known (see bind_corpus); each corpus gets
# its own isolated, resumable run directory under the pipeline.
self.corpus_key: str | None = None
@property
def work_dir(self) -> Path:
def base_work_dir(self) -> Path:
"""The pipeline's root: ``<work_dir setting>/<pipeline name>``. Holds one
subdirectory per corpus."""
raw = self.settings.get("work_dir", "work")
p = Path(raw)
if not p.is_absolute():
p = Path.cwd() / p
return p / self.name
@property
def work_dir(self) -> Path:
"""The active run directory. Corpus-scoped once bound, so different
targets through the same pipeline never share a sample store or report."""
base = self.base_work_dir
return base / self.corpus_key if self.corpus_key else base
def bind_corpus(self, target: Path) -> None:
"""Scope this run's work directory to a specific corpus (run target).
Must be called before ``work_dir`` is used for a run."""
self.corpus_key = corpus_segment(target)
@property
def max_workers(self) -> int:
return int(self.settings.get("max_workers", min(4, os.cpu_count() or 1)))
+15 -11
View File
@@ -113,10 +113,12 @@ class TestRunCommand:
pipeline_file = tmp_path / "pipeline.yaml"
pipeline_file.write_text(yaml_content)
# trash dirs live beside the corpus run dirs under <work>/<pipeline>/, and
# GC sweeps that directory on every run.
work_root = tmp_path / "work"
pipeline_work_dir = work_root / "test"
pipeline_work_dir.mkdir(parents=True)
trash_dir = pipeline_work_dir.with_name("trash_test_123")
trash_dir = pipeline_work_dir / "trash_test_123"
trash_dir.mkdir()
target = tmp_path / "test.sys"
@@ -150,16 +152,18 @@ class TestRunCommand:
pipeline_file = tmp_path / "pipeline.yaml"
pipeline_file.write_text(yaml_content)
work_root = tmp_path / "work"
pipeline_work_dir = work_root / "test"
pipeline_work_dir.mkdir(parents=True)
# Inject dummy file to prove deletion
(pipeline_work_dir / "dummy.txt").write_text("keep")
target = tmp_path / "test.sys"
target.write_bytes(b"MZ")
# runs are scoped per corpus: <work>/<pipeline>/<corpus-key>/. --clean
# purges this corpus's prior run dir, so pre-create it with a marker file.
from deepzero.engine.pipeline import corpus_segment
work_root = tmp_path / "work"
corpus_dir = work_root / "test" / corpus_segment(target)
corpus_dir.mkdir(parents=True)
(corpus_dir / "dummy.txt").write_text("keep")
runner = CliRunner()
# Pass --clean flag to initiate force reset
result = runner.invoke(
@@ -171,9 +175,9 @@ class TestRunCommand:
assert result.exit_code == 0
assert "purging" in result.output
# The clean flag moves the original workdir to a trash directory which is then recursively wiped
# Since we ran synchronously via mock, the directory should be completely empty and devoid of dummy.txt
assert not (pipeline_work_dir / "dummy.txt").exists()
# The clean flag moves the corpus run dir to a trash directory which is then
# recursively wiped; run synchronously via mock, so dummy.txt is gone.
assert not (corpus_dir / "dummy.txt").exists()
class TestServeCommand:
+58 -1
View File
@@ -1,8 +1,65 @@
import re
from pathlib import Path
import pytest
from deepzero.engine.pipeline import load_pipeline, validate_pipeline
from deepzero.engine.pipeline import corpus_segment, load_pipeline, validate_pipeline
def _pipe(tmp_path: Path, work_dir: str = "work"):
import deepzero.stages # noqa: F401 (populate the processor registry)
(tmp_path / "pipeline.yaml").write_text(
f"name: p\nsettings:\n work_dir: {work_dir}\n"
"stages:\n - name: discover\n processor: file_discovery\n"
)
return load_pipeline(str(tmp_path))
def test_corpus_segment_deterministic_and_safe(tmp_path: Path):
a = tmp_path / "DP_Vendor_26061"
a.mkdir()
seg1 = corpus_segment(a)
seg2 = corpus_segment(Path(str(a))) # same path, different Path object
assert seg1 == seg2 # stable -> a rerun resumes in place
assert seg1.startswith("DP_Vendor_26061-") # readable basename kept
assert re.fullmatch(r"[A-Za-z0-9._-]+-[0-9a-f]{8}", seg1) # filesystem-safe
b = tmp_path / "DP_Misc_26053"
b.mkdir()
assert corpus_segment(b) != seg1 # different corpora never collide
def test_corpus_segment_distinguishes_same_basename(tmp_path: Path):
(tmp_path / "x" / "drivers").mkdir(parents=True)
(tmp_path / "y" / "drivers").mkdir(parents=True)
# identical basename, different paths -> distinct segments (no silent sharing)
assert corpus_segment(tmp_path / "x" / "drivers") != corpus_segment(
tmp_path / "y" / "drivers"
)
def test_work_dir_scoped_by_corpus(tmp_path: Path, monkeypatch):
monkeypatch.chdir(tmp_path)
pipe = _pipe(tmp_path)
base = pipe.base_work_dir
assert base.name == "p"
# unbound: work_dir is the pipeline root itself
assert pipe.work_dir == base
target = tmp_path / "corpusA"
target.mkdir()
pipe.bind_corpus(target)
assert pipe.work_dir == base / corpus_segment(target)
assert pipe.work_dir.parent == base # corpus dir sits directly under the base
# a second corpus through the same pipeline gets its own isolated dir
pipe2 = _pipe(tmp_path)
other = tmp_path / "corpusB"
other.mkdir()
pipe2.bind_corpus(other)
assert pipe2.work_dir != pipe.work_dir
assert pipe2.base_work_dir == base
def test_load_pipeline_valid(tmp_path: Path):