Skip to content

Commit f7e42e8

Browse files
committed
fixed bug where old TSV files trigger unwanted RDF conversions
1 parent 83ac6ef commit f7e42e8

3 files changed

Lines changed: 168 additions & 37 deletions

File tree

‎README.md‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -98,9 +98,9 @@ Outputs:
9898
- `run_metrics/.wrapper_logs/wrapper-<timestamp>.log` stores detailed Docker/stdout/stderr command output
9999

100100
Small VCF fixtures for RDF size/inflation test runs:
101-
- `test_vcf_files/test-100.vcf` (100 total lines)
102-
- `test_vcf_files/test-1k.vcf` (1000 total lines)
103-
- `test_vcf_files/test-10k.vcf` (10000 total lines)
101+
- `test_vcf_files/infl100.vcf` (100 total lines)
102+
- `test_vcf_files/infl1k.vcf` (1000 total lines)
103+
- `test_vcf_files/infl10k.vcf` (10000 total lines)
104104

105105
Example inflation check:
106106
```bash
@@ -124,6 +124,7 @@ The wrapper validates:
124124
- Mode-specific required inputs are provided
125125
- Full mode input path exists and contains `.vcf` or `.vcf.gz`
126126
- Full mode rules file exists
127+
- Full mode converts only the VCF file(s) selected at pipeline start (ignores unrelated preexisting TSV intermediates)
127128
- Compression mode input is a `.nq` file
128129
- Decompression mode input is `.gz`, `.br`, or `.hdt`
129130
- Docker image exists or is built (if `--image-version` is set, it will attempt to pull that version and fail if missing)

‎test/test_vcf_rdfizer_unit.py‎

Lines changed: 60 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -412,7 +412,12 @@ def test_main_multiple_triplets_run_multiple_conversions_and_compress_all_output
412412
"""Multiple input triplets trigger per-sample conversion runs and all-output compression."""
413413
with tempfile.TemporaryDirectory() as td:
414414
tmp_path = Path(td)
415-
input_dir, rules_path = prepare_inputs(tmp_path)
415+
input_dir = tmp_path / "input"
416+
input_dir.mkdir()
417+
(input_dir / "sample_a.vcf").write_text("##fileformat=VCFv4.2\n#CHROM\tPOS\n1\t10\n")
418+
(input_dir / "sample_b.vcf").write_text("##fileformat=VCFv4.2\n#CHROM\tPOS\n1\t20\n")
419+
rules_path = tmp_path / "rules.ttl"
420+
rules_path.write_text("@prefix ex: <http://example.org/> .\n")
416421
commands = []
417422

418423
def fake_run(cmd, cwd=None, env=None):
@@ -449,10 +454,60 @@ def fake_run(cmd, cwd=None, env=None):
449454
os.chdir(old_cwd)
450455

451456
self.assertEqual(rc, 0)
452-
self.assertEqual(len(commands), 4)
453-
self.assertIn("OUT_NAME=sample_a", commands[1])
454-
self.assertIn("OUT_NAME=sample_b", commands[2])
455-
self.assertIn("OUT_NAME=", commands[3])
457+
self.assertEqual(len(commands), 5)
458+
self.assertIn("/opt/vcf-rdfizer/vcf_as_tsv.sh", commands[0])
459+
self.assertIn("/data/in/sample_a.vcf", commands[0])
460+
self.assertIn("/opt/vcf-rdfizer/vcf_as_tsv.sh", commands[1])
461+
self.assertIn("/data/in/sample_b.vcf", commands[1])
462+
self.assertIn("OUT_NAME=sample_a", commands[2])
463+
self.assertIn("OUT_NAME=sample_b", commands[3])
464+
self.assertIn("OUT_NAME=", commands[4])
465+
466+
def test_main_ignores_unrelated_existing_tsv_triplets(self):
467+
"""Wrapper converts only triplets that match the CLI-selected VCF snapshot."""
468+
with tempfile.TemporaryDirectory() as td:
469+
tmp_path = Path(td)
470+
input_dir, rules_path = prepare_inputs(tmp_path)
471+
commands = []
472+
473+
def fake_run(cmd, cwd=None, env=None):
474+
commands.append(cmd)
475+
return 0
476+
477+
triplets = [
478+
{
479+
"prefix": "sample",
480+
"records": Path("sample.records.tsv"),
481+
"headers": Path("sample.header_lines.tsv"),
482+
"metadata": Path("sample.file_metadata.tsv"),
483+
},
484+
{
485+
"prefix": "stale",
486+
"records": Path("stale.records.tsv"),
487+
"headers": Path("stale.header_lines.tsv"),
488+
"metadata": Path("stale.file_metadata.tsv"),
489+
},
490+
]
491+
492+
old_cwd = os.getcwd()
493+
os.chdir(tmp_path)
494+
try:
495+
with mock.patch.object(vcf_rdfizer, "run", side_effect=fake_run), mock.patch.object(
496+
vcf_rdfizer, "check_docker", return_value=True
497+
), mock.patch.object(
498+
vcf_rdfizer, "docker_image_exists", return_value=True
499+
), mock.patch.object(
500+
vcf_rdfizer, "discover_tsv_triplets", return_value=triplets
501+
):
502+
rc = invoke_main(["--input", str(input_dir), "--rules", str(rules_path), "--keep-tsv"])
503+
finally:
504+
os.chdir(old_cwd)
505+
506+
self.assertEqual(rc, 0)
507+
self.assertEqual(len(commands), 3)
508+
conversion_cmd = commands[1]
509+
self.assertIn("OUT_NAME=sample", conversion_cmd)
510+
self.assertNotIn("OUT_NAME=stale", conversion_cmd)
456511

457512
def test_main_uses_default_rules_when_flag_is_omitted(self):
458513
"""Wrapper uses repository default rules file when --rules is omitted."""

‎vcf_rdfizer.py‎

Lines changed: 104 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -105,6 +105,52 @@ def resolve_input(input_path: Path):
105105
raise ValueError("Input path must be a file or a directory")
106106

107107

108+
def vcf_output_prefix(path: Path) -> str:
109+
name = path.name
110+
if name.endswith(".vcf.gz"):
111+
return name[: -len(".vcf.gz")]
112+
if name.endswith(".vcf"):
113+
return name[: -len(".vcf")]
114+
return path.stem
115+
116+
117+
def unique_in_order(items):
118+
seen = set()
119+
ordered = []
120+
for item in items:
121+
if item in seen:
122+
continue
123+
seen.add(item)
124+
ordered.append(item)
125+
return ordered
126+
127+
128+
def resolve_input_snapshot(input_path: Path):
129+
if not input_path.exists():
130+
raise ValueError(f"Input path not found: {input_path}")
131+
132+
if input_path.is_file():
133+
if not is_vcf_file(input_path):
134+
raise ValueError("Input file must end with .vcf or .vcf.gz")
135+
mount_dir = input_path.parent
136+
container_inputs = [f"/data/in/{input_path.name}"]
137+
input_metrics_target = container_inputs[0]
138+
prefixes = [vcf_output_prefix(input_path)]
139+
return mount_dir, container_inputs, input_metrics_target, prefixes
140+
141+
if input_path.is_dir():
142+
snapshot_files = list_vcfs_in_dir(input_path)
143+
if not snapshot_files:
144+
raise ValueError("No .vcf or .vcf.gz files found in the input directory")
145+
mount_dir = input_path
146+
container_inputs = [f"/data/in/{p.name}" for p in snapshot_files]
147+
input_metrics_target = "/data/in"
148+
prefixes = [vcf_output_prefix(p) for p in snapshot_files]
149+
return mount_dir, container_inputs, input_metrics_target, prefixes
150+
151+
raise ValueError("Input path must be a file or a directory")
152+
153+
108154
def ensure_dir(path: Path):
109155
path.mkdir(parents=True, exist_ok=True)
110156

@@ -313,8 +359,10 @@ def ensure_image_available(
313359

314360
def run_full_mode(
315361
*,
316-
input_dir: Path,
317-
container_input: str,
362+
input_mount_dir: Path,
363+
container_inputs: list[str],
364+
input_metrics_target: str,
365+
expected_prefixes: list[str],
318366
rules_path: Path,
319367
out_dir: Path,
320368
tsv_dir: Path,
@@ -330,23 +378,26 @@ def run_full_mode(
330378
print("Step 3/5: Converting VCF to TSV")
331379
tsv_existed = tsv_dir.exists()
332380
ensure_dir(tsv_dir)
333-
tsv_cmd = [
334-
"sudo",
335-
"docker",
336-
"run",
337-
"--rm",
338-
"-v",
339-
f"{str(input_dir)}:/data/in:ro",
340-
"-v",
341-
f"{str(tsv_dir)}:/data/tsv",
342-
image_ref,
343-
"/opt/vcf-rdfizer/vcf_as_tsv.sh",
344-
container_input,
345-
"/data/tsv",
346-
]
347-
if run(tsv_cmd) != 0:
348-
eprint(f"Error: TSV conversion failed. See log: {wrapper_log_path}")
349-
return 1
381+
total_inputs = len(container_inputs)
382+
for idx, container_input in enumerate(container_inputs, start=1):
383+
print(f" - TSV conversion {idx}/{total_inputs}: {Path(container_input).name}")
384+
tsv_cmd = [
385+
"sudo",
386+
"docker",
387+
"run",
388+
"--rm",
389+
"-v",
390+
f"{str(input_mount_dir)}:/data/in:ro",
391+
"-v",
392+
f"{str(tsv_dir)}:/data/tsv",
393+
image_ref,
394+
"/opt/vcf-rdfizer/vcf_as_tsv.sh",
395+
container_input,
396+
"/data/tsv",
397+
]
398+
if run(tsv_cmd) != 0:
399+
eprint(f"Error: TSV conversion failed. See log: {wrapper_log_path}")
400+
return 1
350401

351402
print("Step 4/5: Running Conversion with RMLStreamer")
352403
ensure_dir(out_dir)
@@ -359,6 +410,23 @@ def run_full_mode(
359410
eprint(f"See log for details: {wrapper_log_path}")
360411
return 1
361412

413+
expected_order = unique_in_order(expected_prefixes)
414+
triplets_by_prefix = {triplet["prefix"]: triplet for triplet in tsv_triplets}
415+
missing_prefixes = [prefix for prefix in expected_order if prefix not in triplets_by_prefix]
416+
if missing_prefixes:
417+
eprint(
418+
"Error: TSV conversion did not produce expected triplets for: "
419+
+ ", ".join(missing_prefixes)
420+
)
421+
eprint(f"See log for details: {wrapper_log_path}")
422+
return 1
423+
424+
ignored_prefixes = sorted(set(triplets_by_prefix.keys()) - set(expected_order))
425+
if ignored_prefixes:
426+
print(" - Ignoring unrelated TSV triplets in output directory")
427+
428+
tsv_triplets = [triplets_by_prefix[prefix] for prefix in expected_order]
429+
362430
generated_rules_dir = metrics_dir / "_generated_rules"
363431
if generated_rules_dir.exists():
364432
shutil.rmtree(generated_rules_dir, ignore_errors=True)
@@ -407,13 +475,13 @@ def run_full_mode(
407475
f"OUT_NAME={output_name}",
408476
"-e",
409477
f"RUN_ID={run_id}",
410-
"-e",
411-
f"TIMESTAMP={timestamp}",
412-
"-e",
413-
f"IN_VCF={container_input}",
414-
"-e",
415-
"LOGDIR=/data/metrics",
416-
image_ref,
478+
"-e",
479+
f"TIMESTAMP={timestamp}",
480+
"-e",
481+
f"IN_VCF={input_metrics_target}",
482+
"-e",
483+
"LOGDIR=/data/metrics",
484+
image_ref,
417485
"/opt/vcf-rdfizer/run_conversion.sh",
418486
]
419487
if run(run_cmd) != 0:
@@ -693,7 +761,12 @@ def main():
693761
if args.input is None:
694762
raise ValueError("--input is required in --mode full")
695763
input_path = Path(args.input).expanduser().resolve()
696-
input_dir, container_input = resolve_input(input_path)
764+
(
765+
input_mount_dir,
766+
container_inputs,
767+
input_metrics_target,
768+
expected_prefixes,
769+
) = resolve_input_snapshot(input_path)
697770
if args.rules is None:
698771
rules_path = (repo_root / "rules" / "default_rules.ttl").resolve()
699772
else:
@@ -785,8 +858,10 @@ def main():
785858

786859
if mode == "full":
787860
return run_full_mode(
788-
input_dir=input_dir,
789-
container_input=container_input,
861+
input_mount_dir=input_mount_dir,
862+
container_inputs=container_inputs,
863+
input_metrics_target=input_metrics_target,
864+
expected_prefixes=expected_prefixes,
790865
rules_path=rules_path,
791866
out_dir=out_dir,
792867
tsv_dir=tsv_dir,

0 commit comments

Comments
 (0)