ca3a2a3462
- run_create_manifest gains --force; it previously overwrote an existing manifest (including holdout.manifest, which check_holdout_overlap exists specifically to protect) with no warning or backup on a second run. - _git_user_name only caught OSError, not subprocess.TimeoutExpired (a SubprocessError, not an OSError) — a slow/loaded shared portal machine could crash `dwarf bump-gen`/`bump-schema` instead of degrading to by=None as intended. - `dwarf convert --jobs`/`make-root --jobs` now warn (never block) when the requested count exceeds ~1/4 of the machine's CPUs, matching the same shared-machine etiquette check added to giant train in the previous commit. - The Conditioning enum was independently redefined in both giant/cli.py and scripts/dwarf.py; moved to a single giant.config.Conditioning both now import, removing the drift risk of a third conditioning mode being added to one but not the other. Each fix has a regression test. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
523 lines
17 KiB
Python
523 lines
17 KiB
Python
import os
|
|
import subprocess
|
|
|
|
from scripts import bump_dataset_version
|
|
|
|
plan_bump_gen = bump_dataset_version.plan_bump_gen
|
|
plan_bump_schema = bump_dataset_version.plan_bump_schema
|
|
apply_bump = bump_dataset_version.apply_bump
|
|
plan_update_manifest = bump_dataset_version.plan_update_manifest
|
|
apply_update_manifest = bump_dataset_version.apply_update_manifest
|
|
plan_create_manifest = bump_dataset_version.plan_create_manifest
|
|
apply_create_manifest = bump_dataset_version.apply_create_manifest
|
|
check_holdout_overlap = bump_dataset_version.check_holdout_overlap
|
|
|
|
|
|
def test_git_user_name_returns_none_on_timeout(monkeypatch):
|
|
def _raise_timeout(*args, **kwargs):
|
|
raise subprocess.TimeoutExpired(cmd=["git"], timeout=2)
|
|
|
|
monkeypatch.setattr(subprocess, "run", _raise_timeout)
|
|
assert bump_dataset_version._git_user_name() is None
|
|
|
|
|
|
def test_bump_gen_starts_at_gen1_when_none_exist(tmp_path):
|
|
dirs, log_line = plan_bump_gen(
|
|
tmp_path, "steps", "first generation", None, "2026-01-01"
|
|
)
|
|
assert dirs == [
|
|
tmp_path / "raw" / "steps" / "gen1",
|
|
tmp_path / "processed" / "steps" / "gen1" / "schema1",
|
|
]
|
|
assert "`gen1`" in log_line
|
|
assert "first generation" in log_line
|
|
|
|
|
|
def test_bump_gen_increments_past_existing(tmp_path):
|
|
(tmp_path / "raw" / "steps" / "gen1").mkdir(parents=True)
|
|
(tmp_path / "raw" / "steps" / "gen2").mkdir(parents=True)
|
|
dirs, _ = plan_bump_gen(tmp_path, "steps", "next gen", None, "2026-01-01")
|
|
assert dirs[0] == tmp_path / "raw" / "steps" / "gen3"
|
|
|
|
|
|
def test_bump_gen_checks_both_raw_and_processed_trees(tmp_path):
|
|
# processed/ is ahead of raw/ — next gen must still be past the max of both.
|
|
(tmp_path / "raw" / "steps" / "gen1").mkdir(parents=True)
|
|
(tmp_path / "processed" / "steps" / "gen4" / "schema1").mkdir(parents=True)
|
|
dirs, _ = plan_bump_gen(tmp_path, "steps", "next gen", None, "2026-01-01")
|
|
assert dirs[0] == tmp_path / "raw" / "steps" / "gen5"
|
|
|
|
|
|
def test_bump_gen_kinds_are_independent(tmp_path):
|
|
(tmp_path / "raw" / "steps" / "gen5").mkdir(parents=True)
|
|
dirs, _ = plan_bump_gen(tmp_path, "hits", "first hits gen", None, "2026-01-01")
|
|
assert dirs[0] == tmp_path / "raw" / "hits" / "gen1"
|
|
|
|
|
|
def test_bump_schema_starts_at_schema1_for_a_fresh_gen(tmp_path):
|
|
(tmp_path / "raw" / "steps" / "gen1").mkdir(parents=True)
|
|
dirs, log_line = plan_bump_schema(
|
|
tmp_path, "steps", "gen1", "added e_sec column", None, "2026-01-01"
|
|
)
|
|
assert dirs == [tmp_path / "processed" / "steps" / "gen1" / "schema1"]
|
|
assert "`gen1`/`schema1`" in log_line
|
|
|
|
|
|
def test_bump_schema_increments_within_its_gen(tmp_path):
|
|
(tmp_path / "processed" / "steps" / "gen1" / "schema1").mkdir(parents=True)
|
|
(tmp_path / "processed" / "steps" / "gen1" / "schema2").mkdir(parents=True)
|
|
dirs, _ = plan_bump_schema(
|
|
tmp_path, "steps", "gen1", "next schema", None, "2026-01-01"
|
|
)
|
|
assert dirs == [tmp_path / "processed" / "steps" / "gen1" / "schema3"]
|
|
|
|
|
|
def test_bump_schema_does_not_see_other_gens_schemas(tmp_path):
|
|
(tmp_path / "processed" / "steps" / "gen1" / "schema5").mkdir(parents=True)
|
|
(tmp_path / "raw" / "steps" / "gen2").mkdir(parents=True)
|
|
dirs, _ = plan_bump_schema(
|
|
tmp_path, "steps", "gen2", "fresh schema for gen2", None, "2026-01-01"
|
|
)
|
|
assert dirs == [tmp_path / "processed" / "steps" / "gen2" / "schema1"]
|
|
|
|
|
|
def test_bump_schema_rejects_nonexistent_gen(tmp_path):
|
|
try:
|
|
plan_bump_schema(tmp_path, "steps", "gen9", "oops", None, "2026-01-01")
|
|
assert False, "expected SystemExit"
|
|
except SystemExit:
|
|
pass
|
|
|
|
|
|
def test_bump_gen_to_specific_tag(tmp_path):
|
|
(tmp_path / "raw" / "steps" / "gen1").mkdir(parents=True)
|
|
dirs, log_line = plan_bump_gen(
|
|
tmp_path, "steps", "jump to gen5", None, "2026-01-01", target="gen5"
|
|
)
|
|
assert dirs[0] == tmp_path / "raw" / "steps" / "gen5"
|
|
assert "`gen5`" in log_line
|
|
|
|
|
|
def test_bump_gen_rejects_invalid_to_tag(tmp_path):
|
|
try:
|
|
plan_bump_gen(tmp_path, "steps", "bad tag", None, "2026-01-01", target="v5")
|
|
assert False, "expected SystemExit"
|
|
except SystemExit:
|
|
pass
|
|
|
|
|
|
def test_bump_schema_to_specific_tag(tmp_path):
|
|
(tmp_path / "raw" / "steps" / "gen1").mkdir(parents=True)
|
|
(tmp_path / "processed" / "steps" / "gen1" / "schema1").mkdir(parents=True)
|
|
dirs, log_line = plan_bump_schema(
|
|
tmp_path,
|
|
"steps",
|
|
"gen1",
|
|
"jump to schema5",
|
|
None,
|
|
"2026-01-01",
|
|
target="schema5",
|
|
)
|
|
assert dirs == [tmp_path / "processed" / "steps" / "gen1" / "schema5"]
|
|
assert "`schema5`" in log_line
|
|
|
|
|
|
def test_bump_schema_rejects_invalid_to_tag(tmp_path):
|
|
(tmp_path / "raw" / "steps" / "gen1").mkdir(parents=True)
|
|
try:
|
|
plan_bump_schema(
|
|
tmp_path, "steps", "gen1", "bad tag", None, "2026-01-01", target="v3"
|
|
)
|
|
assert False, "expected SystemExit"
|
|
except SystemExit:
|
|
pass
|
|
|
|
|
|
def test_apply_bump_creates_dirs_and_appends_log(tmp_path):
|
|
dirs, log_line = plan_bump_gen(tmp_path, "steps", "reason A", "alice", "2026-01-01")
|
|
apply_bump(tmp_path, dirs, log_line)
|
|
for d in dirs:
|
|
assert d.is_dir()
|
|
text = (tmp_path / "VERSIONS.md").read_text()
|
|
assert "reason A" in text
|
|
assert "alice" in text
|
|
|
|
|
|
def test_apply_bump_appends_without_clobbering_existing_log(tmp_path):
|
|
(tmp_path / "VERSIONS.md").write_text("# Dataset versions\n\n- existing entry\n")
|
|
dirs, log_line = plan_bump_gen(tmp_path, "steps", "reason B", None, "2026-01-02")
|
|
apply_bump(tmp_path, dirs, log_line)
|
|
text = (tmp_path / "VERSIONS.md").read_text()
|
|
assert "existing entry" in text
|
|
assert "reason B" in text
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# update-manifest
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def _make_parquet(path):
|
|
"""Create a zero-byte stand-in for a parquet file."""
|
|
path.parent.mkdir(parents=True, exist_ok=True)
|
|
path.touch()
|
|
|
|
|
|
def test_update_manifest_bumps_to_specified_schema(tmp_path):
|
|
parquet = (
|
|
tmp_path
|
|
/ "processed"
|
|
/ "steps"
|
|
/ "gen1"
|
|
/ "schema2"
|
|
/ "pbwo4"
|
|
/ "shard-000.parquet"
|
|
)
|
|
_make_parquet(parquet)
|
|
|
|
manifest_dir = tmp_path / "pools" / "pbwo4"
|
|
manifest_dir.mkdir(parents=True)
|
|
manifest = manifest_dir / "full.manifest"
|
|
old_rel = "../../processed/steps/gen1/schema1/pbwo4/shard-000.parquet"
|
|
manifest.write_text(old_rel + "\n")
|
|
|
|
lines, missing = plan_update_manifest(manifest, "schema2")
|
|
assert missing == []
|
|
# exactly one data line was updated
|
|
changed = [(old, new) for old, new in lines if new is not None]
|
|
assert len(changed) == 1
|
|
assert "schema2" in changed[0][1]
|
|
assert "schema1" not in changed[0][1]
|
|
|
|
|
|
def test_update_manifest_auto_detects_highest_schema(tmp_path):
|
|
for schema in ("schema1", "schema2", "schema3"):
|
|
d = tmp_path / "processed" / "steps" / "gen1" / schema / "pbwo4"
|
|
d.mkdir(parents=True)
|
|
parquet = (
|
|
tmp_path
|
|
/ "processed"
|
|
/ "steps"
|
|
/ "gen1"
|
|
/ "schema3"
|
|
/ "pbwo4"
|
|
/ "shard-000.parquet"
|
|
)
|
|
parquet.touch()
|
|
|
|
manifest_dir = tmp_path / "pools" / "pbwo4"
|
|
manifest_dir.mkdir(parents=True)
|
|
manifest = manifest_dir / "full.manifest"
|
|
manifest.write_text("../../processed/steps/gen1/schema1/pbwo4/shard-000.parquet\n")
|
|
|
|
lines, missing = plan_update_manifest(manifest, None)
|
|
assert missing == []
|
|
changed = [(old, new) for old, new in lines if new is not None]
|
|
assert "schema3" in changed[0][1]
|
|
|
|
|
|
def test_update_manifest_reports_missing_targets(tmp_path):
|
|
manifest_dir = tmp_path / "pools" / "pbwo4"
|
|
manifest_dir.mkdir(parents=True)
|
|
manifest = manifest_dir / "full.manifest"
|
|
manifest.write_text("../../processed/steps/gen1/schema1/pbwo4/shard-000.parquet\n")
|
|
|
|
# schema2 dir exists but the parquet file does not
|
|
(tmp_path / "processed" / "steps" / "gen1" / "schema2").mkdir(parents=True)
|
|
|
|
lines, missing = plan_update_manifest(manifest, "schema2")
|
|
assert len(missing) == 1
|
|
assert "schema2" in str(missing[0])
|
|
|
|
|
|
def test_update_manifest_skips_already_at_target(tmp_path):
|
|
parquet = (
|
|
tmp_path
|
|
/ "processed"
|
|
/ "steps"
|
|
/ "gen1"
|
|
/ "schema2"
|
|
/ "pbwo4"
|
|
/ "shard-000.parquet"
|
|
)
|
|
_make_parquet(parquet)
|
|
|
|
manifest_dir = tmp_path / "pools" / "pbwo4"
|
|
manifest_dir.mkdir(parents=True)
|
|
manifest = manifest_dir / "full.manifest"
|
|
manifest.write_text("../../processed/steps/gen1/schema2/pbwo4/shard-000.parquet\n")
|
|
|
|
lines, missing = plan_update_manifest(manifest, "schema2")
|
|
assert missing == []
|
|
# line is unchanged — new is None
|
|
assert all(new is None for _, new in lines)
|
|
|
|
|
|
def test_update_manifest_preserves_comments_and_blanks(tmp_path):
|
|
parquet = (
|
|
tmp_path
|
|
/ "processed"
|
|
/ "steps"
|
|
/ "gen1"
|
|
/ "schema2"
|
|
/ "pbwo4"
|
|
/ "shard-000.parquet"
|
|
)
|
|
_make_parquet(parquet)
|
|
|
|
manifest_dir = tmp_path / "pools" / "pbwo4"
|
|
manifest_dir.mkdir(parents=True)
|
|
manifest = manifest_dir / "full.manifest"
|
|
content = "# header\n\n../../processed/steps/gen1/schema1/pbwo4/shard-000.parquet\n"
|
|
manifest.write_text(content)
|
|
|
|
lines, _ = plan_update_manifest(manifest, "schema2")
|
|
assert lines[0] == ("# header", None)
|
|
assert lines[1] == ("", None)
|
|
assert lines[2][1] is not None # the data line was updated
|
|
|
|
|
|
def test_update_manifest_bumps_gen(tmp_path):
|
|
parquet = (
|
|
tmp_path
|
|
/ "processed"
|
|
/ "steps"
|
|
/ "gen2"
|
|
/ "schema1"
|
|
/ "pbwo4"
|
|
/ "shard-000.parquet"
|
|
)
|
|
_make_parquet(parquet)
|
|
|
|
manifest_dir = tmp_path / "pools" / "pbwo4"
|
|
manifest_dir.mkdir(parents=True)
|
|
manifest = manifest_dir / "full.manifest"
|
|
manifest.write_text("../../processed/steps/gen1/schema1/pbwo4/shard-000.parquet\n")
|
|
|
|
lines, missing = plan_update_manifest(manifest, None, target_gen="gen2")
|
|
assert missing == []
|
|
changed = [(old, new) for old, new in lines if new is not None]
|
|
assert len(changed) == 1
|
|
assert "gen2" in changed[0][1]
|
|
assert "gen1" not in changed[0][1]
|
|
|
|
|
|
def test_update_manifest_bumps_gen_and_schema(tmp_path):
|
|
parquet = (
|
|
tmp_path
|
|
/ "processed"
|
|
/ "steps"
|
|
/ "gen2"
|
|
/ "schema3"
|
|
/ "pbwo4"
|
|
/ "shard-000.parquet"
|
|
)
|
|
_make_parquet(parquet)
|
|
|
|
manifest_dir = tmp_path / "pools" / "pbwo4"
|
|
manifest_dir.mkdir(parents=True)
|
|
manifest = manifest_dir / "full.manifest"
|
|
manifest.write_text("../../processed/steps/gen1/schema1/pbwo4/shard-000.parquet\n")
|
|
|
|
lines, missing = plan_update_manifest(manifest, "schema3", target_gen="gen2")
|
|
assert missing == []
|
|
changed = [(old, new) for old, new in lines if new is not None]
|
|
assert len(changed) == 1
|
|
assert "gen2" in changed[0][1]
|
|
assert "schema3" in changed[0][1]
|
|
|
|
|
|
def test_apply_update_manifest_writes_file(tmp_path):
|
|
parquet = (
|
|
tmp_path
|
|
/ "processed"
|
|
/ "steps"
|
|
/ "gen1"
|
|
/ "schema2"
|
|
/ "pbwo4"
|
|
/ "shard-000.parquet"
|
|
)
|
|
_make_parquet(parquet)
|
|
|
|
manifest_dir = tmp_path / "pools" / "pbwo4"
|
|
manifest_dir.mkdir(parents=True)
|
|
manifest = manifest_dir / "full.manifest"
|
|
manifest.write_text("../../processed/steps/gen1/schema1/pbwo4/shard-000.parquet\n")
|
|
|
|
plan, _ = plan_update_manifest(manifest, "schema2")
|
|
apply_update_manifest(manifest, plan)
|
|
|
|
written = manifest.read_text()
|
|
assert "schema2" in written
|
|
assert "schema1" not in written
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# create-manifest
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def test_create_manifest_writes_relative_paths(tmp_path):
|
|
pq1 = (
|
|
tmp_path
|
|
/ "processed"
|
|
/ "steps"
|
|
/ "gen1"
|
|
/ "schema2"
|
|
/ "pbwo4"
|
|
/ "shard-000.parquet"
|
|
)
|
|
pq2 = (
|
|
tmp_path
|
|
/ "processed"
|
|
/ "steps"
|
|
/ "gen1"
|
|
/ "schema2"
|
|
/ "pbwo4"
|
|
/ "shard-001.parquet"
|
|
)
|
|
_make_parquet(pq1)
|
|
_make_parquet(pq2)
|
|
|
|
output = tmp_path / "pools" / "pbwo4" / "train.manifest"
|
|
lines, missing, resolved = plan_create_manifest(output, [pq1, pq2])
|
|
|
|
assert missing == []
|
|
assert len(lines) == 2
|
|
assert all("schema2" in line for line in lines)
|
|
assert all(not line.startswith("/") for line in lines)
|
|
assert resolved == [pq1.resolve(), pq2.resolve()]
|
|
|
|
apply_create_manifest(output, lines)
|
|
assert output.exists()
|
|
written = output.read_text().strip().splitlines()
|
|
assert len(written) == 2
|
|
|
|
|
|
def test_create_manifest_reports_missing_files(tmp_path):
|
|
ghost = tmp_path / "processed" / "gen1" / "schema2" / "shard-000.parquet"
|
|
output = tmp_path / "pools" / "full.manifest"
|
|
lines, missing, _ = plan_create_manifest(output, [ghost])
|
|
assert len(missing) == 1
|
|
assert missing[0] == ghost.resolve()
|
|
|
|
|
|
def test_create_manifest_creates_parent_dirs(tmp_path):
|
|
pq = tmp_path / "a.parquet"
|
|
pq.touch()
|
|
output = tmp_path / "deep" / "nested" / "pool.manifest"
|
|
lines, _, _ = plan_create_manifest(output, [pq])
|
|
apply_create_manifest(output, lines)
|
|
assert output.exists()
|
|
|
|
|
|
def test_run_create_manifest_refuses_to_overwrite_existing_output(tmp_path):
|
|
pq = tmp_path / "a.parquet"
|
|
pq.touch()
|
|
output = tmp_path / "pools" / "pbwo4" / "holdout.manifest"
|
|
output.parent.mkdir(parents=True)
|
|
output.write_text("original contents\n")
|
|
|
|
try:
|
|
bump_dataset_version.run_create_manifest(
|
|
[str(pq)], execute=True, output=str(output)
|
|
)
|
|
assert False, "expected SystemExit"
|
|
except SystemExit:
|
|
pass
|
|
assert output.read_text() == "original contents\n"
|
|
|
|
|
|
def test_run_create_manifest_force_overwrites_existing_output(tmp_path):
|
|
pq = tmp_path / "a.parquet"
|
|
pq.touch()
|
|
output = tmp_path / "pools" / "pbwo4" / "holdout.manifest"
|
|
output.parent.mkdir(parents=True)
|
|
output.write_text("original contents\n")
|
|
|
|
bump_dataset_version.run_create_manifest(
|
|
[str(pq)], execute=True, output=str(output), force=True
|
|
)
|
|
assert output.read_text() != "original contents\n"
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# check_holdout_overlap
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def test_no_overlap_check_when_no_holdout_involved(tmp_path):
|
|
pool_dir = tmp_path / "pools" / "pbwo4"
|
|
pool_dir.mkdir(parents=True)
|
|
pq = tmp_path / "a.parquet"
|
|
pq.touch()
|
|
# writing full.manifest, no holdout.manifest exists
|
|
output = pool_dir / "full.manifest"
|
|
overlaps = check_holdout_overlap(output, [pq.resolve()])
|
|
assert overlaps == []
|
|
|
|
|
|
def test_overlap_detected_when_creating_holdout(tmp_path):
|
|
pool_dir = tmp_path / "pools" / "pbwo4"
|
|
pool_dir.mkdir(parents=True)
|
|
pq = tmp_path / "a.parquet"
|
|
pq.touch()
|
|
|
|
# full.manifest already lists the same file
|
|
full = pool_dir / "full.manifest"
|
|
full.write_text(os.path.relpath(pq.resolve(), start=pool_dir) + "\n")
|
|
|
|
output = pool_dir / "holdout.manifest"
|
|
overlaps = check_holdout_overlap(output, [pq.resolve()])
|
|
assert len(overlaps) == 1
|
|
assert overlaps[0][0] == "full.manifest"
|
|
assert overlaps[0][1] == pq.resolve()
|
|
|
|
|
|
def test_overlap_detected_when_holdout_already_exists(tmp_path):
|
|
pool_dir = tmp_path / "pools" / "pbwo4"
|
|
pool_dir.mkdir(parents=True)
|
|
pq = tmp_path / "a.parquet"
|
|
pq.touch()
|
|
|
|
# holdout.manifest already lists the file
|
|
holdout = pool_dir / "holdout.manifest"
|
|
holdout.write_text(os.path.relpath(pq.resolve(), start=pool_dir) + "\n")
|
|
|
|
# now creating full.manifest with the same file
|
|
output = pool_dir / "full.manifest"
|
|
overlaps = check_holdout_overlap(output, [pq.resolve()])
|
|
assert len(overlaps) == 1
|
|
assert overlaps[0][0] == "holdout.manifest"
|
|
|
|
|
|
def test_dev_full_overlap_not_checked(tmp_path):
|
|
pool_dir = tmp_path / "pools" / "pbwo4"
|
|
pool_dir.mkdir(parents=True)
|
|
pq = tmp_path / "a.parquet"
|
|
pq.touch()
|
|
|
|
# dev.manifest lists the same file — should not trigger a warning for full
|
|
dev = pool_dir / "dev.manifest"
|
|
dev.write_text(os.path.relpath(pq.resolve(), start=pool_dir) + "\n")
|
|
|
|
output = pool_dir / "full.manifest"
|
|
overlaps = check_holdout_overlap(output, [pq.resolve()])
|
|
assert overlaps == []
|
|
|
|
|
|
def test_no_overlap_when_files_are_disjoint(tmp_path):
|
|
pool_dir = tmp_path / "pools" / "pbwo4"
|
|
pool_dir.mkdir(parents=True)
|
|
pq_holdout = tmp_path / "holdout.parquet"
|
|
pq_full = tmp_path / "full.parquet"
|
|
pq_holdout.touch()
|
|
pq_full.touch()
|
|
|
|
holdout = pool_dir / "holdout.manifest"
|
|
holdout.write_text(os.path.relpath(pq_holdout.resolve(), start=pool_dir) + "\n")
|
|
|
|
output = pool_dir / "full.manifest"
|
|
overlaps = check_holdout_overlap(output, [pq_full.resolve()])
|
|
assert overlaps == []
|