Add update-manifest and create-manifest subcommands to bump_dataset_version
update-manifest rewrites the schemaN component in existing manifest files to a specified or auto-detected highest schema, verifying all target files exist before writing. create-manifest builds a new manifest from explicit parquet file paths, supporting --pool/--type (full|holdout|dev) to derive the output path from root, and enforcing holdout isolation by checking for cross-manifest overlap whenever a holdout manifest is involved. Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
+289
-23
@@ -4,6 +4,7 @@ dataset tree (see scripts/migrate_geant_steps.py for the layout):
|
||||
|
||||
raw/<kind>/<gen>/<detector>/shard-NNN.root
|
||||
processed/<kind>/<gen>/<schema>/<detector>/shard-NNN.parquet
|
||||
pools/<detector>/<pool>.manifest
|
||||
|
||||
`gen` bumps when the underlying ROOT changes (geometry/physics-list/macro).
|
||||
`schema` bumps when the parquet export (steps_to_parquet.py or similar)
|
||||
@@ -15,11 +16,14 @@ to VERSIONS.md. Defaults to a dry run; pass --execute to apply.
|
||||
Usage:
|
||||
bump_dataset_version.py bump-gen --kind steps --reason "switched EM physics list"
|
||||
bump_dataset_version.py bump-schema --kind steps --gen gen1 --reason "added e_sec column"
|
||||
bump_dataset_version.py update-manifest pools/pbwo4/full.manifest [--schema schema2]
|
||||
bump_dataset_version.py create-manifest --output pools/pbwo4/train.manifest a.parquet b.parquet
|
||||
bump_dataset_version.py status
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import datetime as dt
|
||||
import os
|
||||
import re
|
||||
import subprocess
|
||||
from pathlib import Path
|
||||
@@ -131,12 +135,161 @@ def print_status(root: Path) -> None:
|
||||
print(f" {gen_tag}: {schema_str}")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# update-manifest
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _find_schema_idx(abs_path: Path) -> int | None:
|
||||
"""Return the index of the first schemaN component in abs_path.parts, or None."""
|
||||
for i, part in enumerate(abs_path.parts):
|
||||
if SCHEMA_RE.match(part):
|
||||
return i
|
||||
return None
|
||||
|
||||
|
||||
def plan_update_manifest(
|
||||
manifest_path: Path, target_schema: str | None
|
||||
) -> tuple[list[tuple[str, str | None]], list[Path]]:
|
||||
"""Parse a manifest and plan schema replacements for each data line.
|
||||
|
||||
Returns:
|
||||
lines: list of (original_line, new_relative_path_or_None)
|
||||
None means the line is unchanged (comment, blank, or already at target).
|
||||
missing: resolved absolute paths that don't exist on disk.
|
||||
"""
|
||||
manifest_path = manifest_path.resolve()
|
||||
if not manifest_path.exists():
|
||||
raise SystemExit(f"error: manifest not found: {manifest_path}")
|
||||
manifest_dir = manifest_path.parent
|
||||
raw_lines = manifest_path.read_text().splitlines()
|
||||
|
||||
result: list[tuple[str, str | None]] = []
|
||||
missing: list[Path] = []
|
||||
schema_cache: dict[Path, str] = {}
|
||||
|
||||
for raw in raw_lines:
|
||||
stripped = raw.strip()
|
||||
if not stripped or stripped.startswith("#"):
|
||||
result.append((raw, None))
|
||||
continue
|
||||
|
||||
old_abs = (manifest_dir / stripped).resolve()
|
||||
schema_idx = _find_schema_idx(old_abs)
|
||||
if schema_idx is None:
|
||||
result.append((raw, None))
|
||||
continue
|
||||
|
||||
parts = list(old_abs.parts)
|
||||
old_schema = parts[schema_idx]
|
||||
gen_dir = Path(*parts[:schema_idx])
|
||||
|
||||
if target_schema is not None:
|
||||
new_schema = target_schema
|
||||
else:
|
||||
if gen_dir not in schema_cache:
|
||||
n = _max_index(gen_dir, SCHEMA_RE)
|
||||
if n == 0:
|
||||
raise SystemExit(f"error: no schema dirs found under {gen_dir}")
|
||||
schema_cache[gen_dir] = f"schema{n}"
|
||||
new_schema = schema_cache[gen_dir]
|
||||
|
||||
if new_schema == old_schema:
|
||||
result.append((raw, None))
|
||||
continue
|
||||
|
||||
parts[schema_idx] = new_schema
|
||||
new_abs = Path(*parts)
|
||||
if not new_abs.exists():
|
||||
missing.append(new_abs)
|
||||
|
||||
new_rel = os.path.relpath(new_abs, start=manifest_dir)
|
||||
result.append((raw, new_rel))
|
||||
|
||||
return result, missing
|
||||
|
||||
|
||||
def apply_update_manifest(manifest_path: Path, lines: list[tuple[str, str | None]]) -> None:
|
||||
out = [replacement if replacement is not None else original for original, replacement in lines]
|
||||
manifest_path.write_text("\n".join(out) + "\n")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# create-manifest
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def _resolve_manifest_files(manifest_path: Path) -> list[Path]:
|
||||
"""Read a manifest and return its entries as resolved absolute paths."""
|
||||
files = []
|
||||
for line in manifest_path.read_text().splitlines():
|
||||
line = line.strip()
|
||||
if not line or line.startswith("#"):
|
||||
continue
|
||||
files.append((manifest_path.parent / line).resolve())
|
||||
return files
|
||||
|
||||
|
||||
def plan_create_manifest(
|
||||
output_path: Path, parquet_files: list[Path]
|
||||
) -> tuple[list[str], list[Path], list[Path]]:
|
||||
"""Return (relative_lines, missing_files, resolved_abs_paths)."""
|
||||
manifest_dir = output_path.resolve().parent
|
||||
lines: list[str] = []
|
||||
missing: list[Path] = []
|
||||
resolved: list[Path] = []
|
||||
for p in parquet_files:
|
||||
abs_p = p.resolve()
|
||||
resolved.append(abs_p)
|
||||
if not abs_p.exists():
|
||||
missing.append(abs_p)
|
||||
lines.append(os.path.relpath(abs_p, start=manifest_dir))
|
||||
return lines, missing, resolved
|
||||
|
||||
|
||||
def check_holdout_overlap(
|
||||
output_path: Path, resolved_new_files: list[Path]
|
||||
) -> list[tuple[str, Path]]:
|
||||
"""Return (other_manifest_name, file) pairs where new files clash with existing manifests.
|
||||
|
||||
The check is triggered when output_path is (or will be) holdout.manifest, or when a
|
||||
holdout.manifest already exists in the same directory — in either case holdout data
|
||||
must be strictly isolated from all other pools.
|
||||
"""
|
||||
output_resolved = output_path.resolve()
|
||||
manifest_dir = output_resolved.parent
|
||||
holdout_path = manifest_dir / "holdout.manifest"
|
||||
|
||||
if output_resolved.name != "holdout.manifest" and not holdout_path.exists():
|
||||
return []
|
||||
|
||||
new_set = set(resolved_new_files)
|
||||
overlaps: list[tuple[str, Path]] = []
|
||||
for existing in sorted(manifest_dir.glob("*.manifest")):
|
||||
if existing.resolve() == output_resolved:
|
||||
continue
|
||||
try:
|
||||
existing_files = set(_resolve_manifest_files(existing))
|
||||
except OSError:
|
||||
continue
|
||||
for f in sorted(new_set & existing_files):
|
||||
overlaps.append((existing.name, f))
|
||||
return overlaps
|
||||
|
||||
|
||||
def apply_create_manifest(output_path: Path, lines: list[str]) -> None:
|
||||
output_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
output_path.write_text("\n".join(lines) + "\n")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# CLI
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def main() -> None:
|
||||
parser = argparse.ArgumentParser(description=__doc__)
|
||||
parser.add_argument(
|
||||
"--root",
|
||||
default="/ceph/lbogner/geant_steps",
|
||||
help="Dataset root (default: /ceph/lbogner/geant_steps)",
|
||||
help="Dataset root for bump-gen/bump-schema/status (default: /ceph/lbogner/geant_steps)",
|
||||
)
|
||||
sub = parser.add_subparsers(dest="command", required=True)
|
||||
|
||||
@@ -157,40 +310,153 @@ def main() -> None:
|
||||
p_schema.add_argument("--date", default=None, help="Override date (default: today, ISO)")
|
||||
p_schema.add_argument("--execute", action="store_true", help="Apply (default: dry run)")
|
||||
|
||||
p_update = sub.add_parser(
|
||||
"update-manifest",
|
||||
help="Repoint manifest(s) to a new schema, verifying all target files exist",
|
||||
)
|
||||
p_update.add_argument(
|
||||
"manifests", nargs="+", metavar="MANIFEST", help="One or more .manifest files to update"
|
||||
)
|
||||
p_update.add_argument(
|
||||
"--schema",
|
||||
default=None,
|
||||
metavar="schemaN",
|
||||
help="Target schema tag (default: highest schema found in the same gen dir)",
|
||||
)
|
||||
p_update.add_argument("--execute", action="store_true", help="Write updated manifests (default: dry run)")
|
||||
|
||||
p_create = sub.add_parser(
|
||||
"create-manifest",
|
||||
help="Create a new manifest from a list of parquet files",
|
||||
)
|
||||
dest_group = p_create.add_mutually_exclusive_group(required=True)
|
||||
dest_group.add_argument(
|
||||
"--output", "-o", metavar="PATH", help="Explicit path for the new .manifest file"
|
||||
)
|
||||
dest_group.add_argument(
|
||||
"--pool", metavar="DETECTOR",
|
||||
help="Detector name; combined with --type and --root to form <root>/pools/<detector>/<type>.manifest",
|
||||
)
|
||||
p_create.add_argument(
|
||||
"--type", choices=["full", "holdout", "dev"],
|
||||
help="Pool type — full, holdout, or dev (required with --pool)",
|
||||
)
|
||||
p_create.add_argument(
|
||||
"files", nargs="+", metavar="FILE", help="Parquet files to include"
|
||||
)
|
||||
p_create.add_argument("--execute", action="store_true", help="Write the manifest (default: dry run)")
|
||||
|
||||
sub.add_parser("status", help="List existing gens/schemas per kind")
|
||||
|
||||
args = parser.parse_args()
|
||||
root = Path(args.root)
|
||||
if not root.is_dir():
|
||||
parser.error(f"{root} is not a directory")
|
||||
|
||||
# Commands that need --root
|
||||
if args.command in ("bump-gen", "bump-schema", "status"):
|
||||
root = Path(args.root)
|
||||
if not root.is_dir():
|
||||
parser.error(f"{root} is not a directory")
|
||||
|
||||
if args.command == "status":
|
||||
print_status(root)
|
||||
return
|
||||
|
||||
date = args.date or dt.date.today().isoformat()
|
||||
by = args.by if args.by is not None else _git_user_name()
|
||||
if args.command in ("bump-gen", "bump-schema"):
|
||||
date = args.date or dt.date.today().isoformat()
|
||||
by = args.by if args.by is not None else _git_user_name()
|
||||
if args.command == "bump-gen":
|
||||
new_dirs, log_line = plan_bump_gen(root, args.kind, args.reason, by, date)
|
||||
else:
|
||||
new_dirs, log_line = plan_bump_schema(root, args.kind, args.gen, args.reason, by, date)
|
||||
|
||||
if args.command == "bump-gen":
|
||||
new_dirs, log_line = plan_bump_gen(root, args.kind, args.reason, by, date)
|
||||
else:
|
||||
new_dirs, log_line = plan_bump_schema(
|
||||
root, args.kind, args.gen, args.reason, by, date
|
||||
)
|
||||
print(f"=== {'EXECUTING' if args.execute else 'DRY RUN'} ===")
|
||||
print("new directories:")
|
||||
for d in new_dirs:
|
||||
print(f" {d}")
|
||||
print("VERSIONS.md entry:")
|
||||
print(f" {log_line}")
|
||||
|
||||
print(f"=== {'EXECUTING' if args.execute else 'DRY RUN'} ===")
|
||||
print("new directories:")
|
||||
for d in new_dirs:
|
||||
print(f" {d}")
|
||||
print("VERSIONS.md entry:")
|
||||
print(f" {log_line}")
|
||||
|
||||
if not args.execute:
|
||||
print("\nDry run only — pass --execute to apply.")
|
||||
if not args.execute:
|
||||
print("\nDry run only — pass --execute to apply.")
|
||||
return
|
||||
apply_bump(root, new_dirs, log_line)
|
||||
print("\nDone.")
|
||||
return
|
||||
|
||||
apply_bump(root, new_dirs, log_line)
|
||||
print("\nDone.")
|
||||
if args.command == "update-manifest":
|
||||
if args.schema and not SCHEMA_RE.match(args.schema):
|
||||
parser.error(f"--schema must look like 'schemaN', got {args.schema!r}")
|
||||
|
||||
all_plans: list[tuple[Path, list[tuple[str, str | None]]]] = []
|
||||
all_missing: list[Path] = []
|
||||
|
||||
for raw in args.manifests:
|
||||
mp = Path(raw)
|
||||
plan, missing = plan_update_manifest(mp, args.schema)
|
||||
all_plans.append((mp.resolve(), plan))
|
||||
all_missing.extend(missing)
|
||||
|
||||
print(f"=== {'EXECUTING' if args.execute else 'DRY RUN'} ===")
|
||||
for mp, plan in all_plans:
|
||||
changes = [(old, new) for old, new in plan if new is not None]
|
||||
print(f"\n{mp} ({len(changes)} path(s) to update)")
|
||||
for old, new in changes:
|
||||
print(f" - {old.strip()}")
|
||||
print(f" + {new}")
|
||||
|
||||
if all_missing:
|
||||
print(f"\nMISSING ({len(all_missing)} file(s) — target paths do not exist):")
|
||||
for p in all_missing:
|
||||
print(f" {p}")
|
||||
if args.execute:
|
||||
raise SystemExit("error: refusing to write manifests with missing targets")
|
||||
|
||||
if not args.execute:
|
||||
print("\nDry run only — pass --execute to apply.")
|
||||
return
|
||||
|
||||
for mp, plan in all_plans:
|
||||
apply_update_manifest(mp, plan)
|
||||
print("\nDone.")
|
||||
return
|
||||
|
||||
if args.command == "create-manifest":
|
||||
if args.pool is not None and args.type is None:
|
||||
parser.error("--type is required when --pool is given")
|
||||
|
||||
if args.pool is not None:
|
||||
root = Path(args.root)
|
||||
output_path = root / "pools" / args.pool / f"{args.type}.manifest"
|
||||
else:
|
||||
output_path = Path(args.output)
|
||||
|
||||
parquet_files = [Path(f) for f in args.files]
|
||||
lines, missing, resolved = plan_create_manifest(output_path, parquet_files)
|
||||
overlaps = check_holdout_overlap(output_path, resolved)
|
||||
|
||||
print(f"=== {'EXECUTING' if args.execute else 'DRY RUN'} ===")
|
||||
print(f"manifest: {output_path.resolve()}")
|
||||
for line in lines:
|
||||
print(f" {line}")
|
||||
|
||||
if missing:
|
||||
print(f"\nMISSING ({len(missing)} file(s) do not exist):")
|
||||
for p in missing:
|
||||
print(f" {p}")
|
||||
|
||||
if overlaps:
|
||||
print(f"\nHOLDOUT OVERLAP ({len(overlaps)} file(s) appear in other manifests):")
|
||||
for name, f in overlaps:
|
||||
print(f" {f} (also in {name})")
|
||||
|
||||
if (missing or overlaps) and args.execute:
|
||||
raise SystemExit("error: refusing to write manifest (see above)")
|
||||
|
||||
if not args.execute:
|
||||
print("\nDry run only — pass --execute to apply.")
|
||||
return
|
||||
|
||||
apply_create_manifest(output_path, lines)
|
||||
print("\nDone.")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
|
||||
Reference in New Issue
Block a user