Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion docs/cli-reference.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ artifacts, run metrics and logs, and hidden intermediates all live beneath it.

| Flag | Meaning |
| --- | --- |
| `-i, --input` | VCF file or directory. Only `*.vcf` and `*.vcf.gz` are recognised; a directory is enumerated one level deep and snapshotted at run start |
| `-i, --input` | VCF file or directory. Only `*.vcf` and `*.vcf.gz` are recognised; a directory is enumerated one level deep and snapshotted at run start. macOS AppleDouble sidecars (`._*.vcf`, which tar and copies from macOS add) are skipped with a notice. An input that is not UTF-8 text fails on its own at stage `input-encoding` ("not a UTF-8 text VCF") and the other inputs still convert |
| `--rdf` | RDF input for `compress`/`link`, or the artifact to check in `validation` |
| `-C, --compressed-input` | `.nt.gz`, `.nt.br`, `.hdt`, `.cottas`, `.cottas.gz`, `.cottas.br` |
| `-H, --hdt` / `--cottas` | Existing artifact for `--mode index` |
Expand Down
3 changes: 2 additions & 1 deletion src/vcf_as_tsv.sh
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,8 @@ if [ -f "$input_path" ]; then
elif [ -d "$input_path" ]; then
while IFS= read -r discovered_file; do
files+=("$discovered_file")
done < <(find "$input_path" -maxdepth 1 -type f \( -name '*.vcf' -o -name '*.vcf.gz' \) | sort)
# ._* are macOS AppleDouble sidecars: binary, never VCFs, despite the suffix.
done < <(find "$input_path" -maxdepth 1 -type f \( -name '*.vcf' -o -name '*.vcf.gz' \) ! -name '._*' | sort)
else
echo "Error: input path '$input_path' not found."
exit 1
Expand Down
138 changes: 138 additions & 0 deletions test/test_cohort_guard_unit.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@
"""

import csv
import os
import shutil
import sys
import tempfile
Expand All @@ -24,6 +25,7 @@

import vcf_rdfizer
from test.helpers import VerboseTestCase
from test.test_vcf_rdfizer_unit import invoke_main, latest_metrics_run_dir


def write_records_tsv(path: Path, *, sample_ids: list[str], records: int) -> Path:
Expand Down Expand Up @@ -95,6 +97,25 @@ def test_an_unreadable_tsv_reports_zero_samples_instead_of_raising(self):
path = Path(tmp) / "missing.tsv"
self.assertEqual(vcf_rdfizer.read_records_tsv_sample_count(path), 0)

def test_a_non_utf8_tsv_raises_an_input_error_not_a_decode_error(self):
"""The guard's reader must not leak UnicodeDecodeError past the per-input handler."""
with tempfile.TemporaryDirectory() as tmp:
path = Path(tmp) / "r.tsv"
path.write_bytes(APPLEDOUBLE_BYTES)
with self.assertRaises(vcf_rdfizer.InputEncodingError) as caught:
vcf_rdfizer.read_records_tsv_sample_count(path)
self.assertIn("not a UTF-8 text VCF", str(caught.exception))
with self.assertRaises(vcf_rdfizer.InputEncodingError):
vcf_rdfizer.check_records_tsv_is_utf8(path)

def test_the_encoding_probe_accepts_utf8_cut_at_the_probe_boundary(self):
"""A multi-byte character split by the bounded read is not a bad byte."""
with tempfile.TemporaryDirectory() as tmp:
path = Path(tmp) / "r.tsv"
limit = vcf_rdfizer.RECORDS_TSV_ENCODING_PROBE_BYTES
path.write_bytes(b"a" * (limit - 1) + "\u00e9".encode("utf-8"))
vcf_rdfizer.check_records_tsv_is_utf8(path)


class GuardDecisionTests(VerboseTestCase):
def _refusal(self, tmp, *, sample_ids, records, free_bytes, allow=False,
Expand Down Expand Up @@ -195,6 +216,123 @@ def test_the_guard_leaves_a_margin_rather_than_filling_the_volume(self):
)


#: The head of a real macOS AppleDouble sidecar: magic, version, then binary
#: entries. 0xa3 is the byte the vcf-bench-1 run died on.
APPLEDOUBLE_BYTES = (
b"\x00\x05\x16\x07\x00\x02\x00\x00Mac OS X \x00\x02"
+ b"\x00" * 110
+ b"\xa3\x9f\xff\xfe" * 16
)


class NonUtf8InputTests(VerboseTestCase):
"""A binary input fails on its own; it does not take the directory down.

vcf-bench-1, v3.1.0: a directory carrying macOS ``._P001.vcf`` sidecars
crashed ``--sample-representation expanded`` with a UnicodeDecodeError out
of the cohort guard, while condensed reported the same file as a failed
input and converted the rest.
"""

def _run_full(self, tmp_path: Path, vcf_names: dict[str, bytes], representation: str):
input_dir = tmp_path / "input"
input_dir.mkdir()
tsv_dir = tmp_path / "tsv_fixture"
tsv_dir.mkdir()
triplets = []
for name, payload in vcf_names.items():
(input_dir / name).write_bytes(payload)
prefix = vcf_rdfizer.vcf_output_prefix(Path(name))
records = tsv_dir / f"{prefix}.records.tsv"
if payload.startswith(b"##fileformat"):
write_records_tsv(records, sample_ids=["NA1", "NA2"], records=3)
else:
# vcf_as_tsv.sh passes bytes through, so the binary survives
# into the records TSV that the Python side then decodes.
records.write_bytes(payload)
(tsv_dir / f"{prefix}.header_lines.tsv").write_text(
f"SOURCE_FILE\tLINE\n{name}\t##fileformat=VCFv4.2\n", encoding="utf-8"
)
(tsv_dir / f"{prefix}.file_metadata.tsv").write_text(
"SOURCE_FILE\tKEY\tVALUE\n", encoding="utf-8"
)
triplets.append(
{
"prefix": prefix,
"records": records,
"headers": tsv_dir / f"{prefix}.header_lines.tsv",
"metadata": tsv_dir / f"{prefix}.file_metadata.tsv",
}
)
out_dir = tmp_path / "out"
stdout, stderr = StringIO(), StringIO()
old_cwd = os.getcwd()
os.chdir(tmp_path)
try:
with (
mock.patch.object(vcf_rdfizer, "run", return_value=0),
mock.patch.object(vcf_rdfizer, "check_docker", return_value=True),
mock.patch.object(vcf_rdfizer, "docker_image_exists", return_value=True),
mock.patch.object(vcf_rdfizer, "discover_tsv_triplets", return_value=triplets),
redirect_stdout(stdout),
redirect_stderr(stderr),
):
rc = invoke_main(
[
"--input", str(input_dir),
"--sample-representation", representation,
"--rdf-storage-mode", "plain",
"--compression", "none",
"--out", str(out_dir),
]
)
finally:
os.chdir(old_cwd)
return rc, out_dir, stdout.getvalue(), stderr.getvalue()

def test_a_directory_with_an_appledouble_sidecar_converts_the_vcf_and_reports_the_sidecar(self):
with tempfile.TemporaryDirectory() as td:
rc, out_dir, stdout, stderr = self._run_full(
Path(td),
{"P001.vcf": VALID_VCF, "._P001.vcf": APPLEDOUBLE_BYTES},
"expanded",
)
self.assertNotIn("Traceback", stderr)
self.assertNotIn("UnicodeDecodeError", stderr)
self.assertEqual(rc, 0, stderr)
self.assertIn("Skipping 1 macOS AppleDouble file(s)", stdout)
self.assertIn("._P001.vcf", stdout)
self.assertIn("Input 1/1: P001.vcf", stdout)
self.assertTrue((out_dir / "P001" / "P001.nt").exists())

def test_a_binary_vcf_fails_as_one_input_and_the_others_still_convert(self):
"""Not every binary is an AppleDouble sidecar; the encoding probe catches the rest."""
for representation in ("expanded", "condensed"):
with self.subTest(representation=representation), tempfile.TemporaryDirectory() as td:
rc, out_dir, stdout, stderr = self._run_full(
Path(td),
{"P001.vcf": VALID_VCF, "x.vcf": APPLEDOUBLE_BYTES},
representation,
)
self.assertNotIn("Traceback", stderr)
self.assertEqual(rc, 1, stderr)
self.assertIn("(x.vcf) failed at input-encoding: not a UTF-8 text VCF", stderr)
self.assertTrue((out_dir / "P001" / "P001.nt").exists())
report = (
latest_metrics_run_dir(out_dir / "run_metrics")
/ "reports" / "failed_inputs.csv"
)
with report.open(newline="", encoding="utf-8") as handle:
failures = list(csv.DictReader(handle))
self.assertEqual(
[(f["expected_prefix"], f["stage"]) for f in failures],
[("x", "input-encoding")],
)


VALID_VCF = b"##fileformat=VCFv4.2\n#CHROM\tPOS\n1\t10\n"


class CliTests(VerboseTestCase):
def test_the_override_flag_is_accepted_by_the_cli(self):
"""It must parse; without it the escape hatch is unreachable."""
Expand Down
99 changes: 90 additions & 9 deletions vcf_rdfizer.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@
"""

import argparse
import codecs
import csv
import gzip
import importlib.resources as importlib_resources
Expand Down Expand Up @@ -392,11 +393,55 @@ def count_records_tsv_rows(records_tsv: Path) -> int:
return max(0, total - 1)


class InputEncodingError(ValueError):
"""An input whose bytes do not decode as UTF-8, so it is not a text VCF."""


#: How much of a records TSV the encoding probe decodes. A binary file posing
#: as a VCF -- a macOS AppleDouble ``._x.vcf`` is the usual one -- is caught in
#: its first bytes, and a bounded read keeps the probe free at cohort scale.
RECORDS_TSV_ENCODING_PROBE_BYTES = 1024 * 1024


def _not_utf8_message(exc: UnicodeDecodeError) -> str:
bad_byte = exc.object[exc.start : exc.start + 1].hex() or "??"
return (
f"not a UTF-8 text VCF (byte 0x{bad_byte} does not decode as UTF-8; "
f"is this a binary file, such as a macOS ._ AppleDouble file?)"
)


def check_records_tsv_is_utf8(records_tsv: Path) -> None:
"""Raise InputEncodingError when the start of a records TSV is not UTF-8.

Every downstream reader opens the TSV as UTF-8 text, so a binary input
otherwise surfaces as a bare ``UnicodeDecodeError`` from whichever reader
touches it first. An unreadable file is left for the pipeline to report.
"""
try:
with records_tsv.open("rb") as handle:
head = handle.read(RECORDS_TSV_ENCODING_PROBE_BYTES)
except OSError:
return
try:
# final=False: a multi-byte character cut by the probe boundary is not
# an error, only a byte that can never start or continue one is.
codecs.getincrementaldecoder("utf-8")().decode(head, final=False)
except UnicodeDecodeError as exc:
raise InputEncodingError(_not_utf8_message(exc)) from exc


def read_records_tsv_sample_count(records_tsv: Path) -> int:
"""Sample columns declared by a records TSV, or 0 when it has none."""
"""Sample columns declared by a records TSV, or 0 when it has none.

Raises InputEncodingError for a TSV that is not UTF-8 text: that is a bad
input to report, not an unreadable one to wave through.
"""
try:
with SampleRecordStream(records_tsv) as stream:
return len(stream.columns)
except UnicodeDecodeError as exc:
raise InputEncodingError(_not_utf8_message(exc)) from exc
except (OSError, csv.Error, StopIteration):
return 0

Expand Down Expand Up @@ -1145,15 +1190,37 @@ def is_vcf_file(path: Path) -> bool:
return name.endswith(".vcf") or name.endswith(".vcf.gz")


def is_appledouble_file(path: Path) -> bool:
"""Return True for a macOS AppleDouble sidecar (``._name``).

tar and copies from macOS to non-HFS volumes add one ``._x.vcf`` binary
resource fork per file. It carries the VCF suffix but is never a VCF.
"""
return path.name.startswith("._")


def list_vcfs_in_dir(path: Path):
"""List VCF inputs in a stable order for deterministic processing."""
"""List VCF inputs in a stable order for deterministic processing.

AppleDouble ``._*`` sidecars are left out rather than queued as inputs
that can only fail; ``list_appledouble_in_dir`` names what was left out.
"""
files = []
for item in sorted(path.iterdir()):
if item.is_file() and is_vcf_file(item):
if item.is_file() and is_vcf_file(item) and not is_appledouble_file(item):
files.append(item)
return files


def list_appledouble_in_dir(path: Path):
"""The ``._*.vcf`` / ``._*.vcf.gz`` sidecars ``list_vcfs_in_dir`` skips."""
return [
item
for item in sorted(path.iterdir())
if item.is_file() and is_vcf_file(item) and is_appledouble_file(item)
]


def vcf_output_prefix(path: Path) -> str:
"""Derive stable sample prefix from VCF filename."""
name = path.name
Expand Down Expand Up @@ -1196,6 +1263,12 @@ def resolve_input_snapshot(input_path: Path):

if input_path.is_dir():
snapshot_files = list_vcfs_in_dir(input_path)
skipped = list_appledouble_in_dir(input_path)
if skipped:
print(
f" Skipping {len(skipped)} macOS AppleDouble file(s), which "
f"are never VCFs: {', '.join(p.name for p in skipped)}"
)
if not snapshot_files:
raise ValueError("No .vcf or .vcf.gz files found in the input directory")
mount_dir = input_path
Expand Down Expand Up @@ -7444,12 +7517,20 @@ def fail_current(stage: str, message: str):
# This sits ahead of the helper TSVs deliberately: building those is
# already records x samples work, so a guard after them has let the
# failure mode start.
refusal = cohort_scale_refusal(
records_tsv=triplet["records"],
sample_representation=sample_workflow.representation,
out_dir=out_dir,
allow=allow_cohort_expansion,
)
# The encoding probe runs first and for every representation: the
# guard and all later readers decode the TSV as UTF-8, and a binary
# input must fail here, as this input, not crash the whole run.
try:
check_records_tsv_is_utf8(triplet["records"])
refusal = cohort_scale_refusal(
records_tsv=triplet["records"],
sample_representation=sample_workflow.representation,
out_dir=out_dir,
allow=allow_cohort_expansion,
)
except InputEncodingError as exc:
fail_current("input-encoding", str(exc))
continue
if refusal is not None:
fail_current("cohort-scale-guard", refusal)
continue
Expand Down
Loading