-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathvalidate_compression.py
More file actions
178 lines (155 loc) · 6.62 KB
/
Copy pathvalidate_compression.py
File metadata and controls
178 lines (155 loc) · 6.62 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
#!/usr/bin/env python3
"""Validate a generated HDT or COTTAS artifact without writing RDF output.
The validator deliberately compares decoded triple counts with the source
count. Decoding to ``/dev/stdout`` exercises the reader and keeps the full
N-Triples representation in a pipe rather than creating another large file.
"""
from __future__ import annotations
import argparse
import gzip
import json
import os
import subprocess
import sys
from pathlib import Path
def is_triple_line(line: bytes) -> bool:
"""Return whether a serialized line represents an N-Triples statement."""
stripped = line.strip()
return bool(stripped) and not stripped.startswith(b"#") and stripped.endswith(b".")
def count_nt(path: Path) -> int:
"""Count N-Triples records while transparently reading ``.nt.gz``."""
opener = gzip.open if path.name.endswith(".gz") else Path.open
with opener(path, "rb") as handle:
return sum(1 for line in handle if is_triple_line(line))
def count_decoded(command: list[str]) -> int:
"""Run an RDF exporter and count its streamed N-Triples output."""
# The context manager closes the stdout pipe even when counting raises, so a
# failed decode cannot leak a file descriptor into the rest of the run.
with subprocess.Popen(
command,
stdout=subprocess.PIPE,
stderr=subprocess.DEVNULL,
) as process:
assert process.stdout is not None
count = sum(1 for line in process.stdout if is_triple_line(line))
return_code = process.wait()
if return_code != 0:
raise RuntimeError(f"decoder exited with status {return_code}: {' '.join(command)}")
return count
def resolve_hdt2rdf() -> str:
candidates = (
os.environ.get("HDT2RDF_BIN", ""),
"/usr/local/bin/hdt2rdf",
"/opt/hdt-cpp/bin/hdt2rdf",
)
for candidate in candidates:
if candidate and Path(candidate).is_file() and os.access(candidate, os.X_OK):
return candidate
raise RuntimeError("Missing hdt2rdf binary in container")
def validate(args: argparse.Namespace) -> dict:
source = Path(args.source)
artifact = Path(args.artifact)
if not source.is_file():
raise FileNotFoundError(f"source RDF file not found: {source}")
if not artifact.is_file() or artifact.stat().st_size == 0:
raise FileNotFoundError(f"compression artifact is missing or empty: {artifact}")
source_triples = (
args.source_triples if args.source_triples is not None else count_nt(source)
)
if args.expected_triples is not None and source_triples != args.expected_triples:
return {
"valid": False,
"source_triples": source_triples,
"decoded_triples": None,
"expected_triples": args.expected_triples,
"count_match": False,
"error": (
"source triple count does not match the upstream conversion count: "
f"source={source_triples}, expected={args.expected_triples}"
),
}
if args.format == "hdt":
# The bundled Java-free helper streams the HDT and eagerly creates the
# query index. hdt2rdf below independently checks decoded readability.
if not args.skip_index_check:
index_helper = Path("/opt/vcf-rdfizer/ensure_hdt_index.sh")
if not index_helper.is_file():
raise RuntimeError(f"Missing HDT index helper: {index_helper}")
index_check = subprocess.run(
[str(index_helper), str(artifact)],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
check=False,
)
if index_check.returncode != 0:
raise RuntimeError(
"HDT index/readability check failed with status "
f"{index_check.returncode}"
)
# hdt2rdf uses "-" (not /dev/stdout) as its stdout sentinel. This
# keeps the decoded RDF in the pipe and avoids another large file.
decoded_triples = count_decoded([resolve_hdt2rdf(), str(artifact), "-"])
validator = "hdt2rdf"
else:
try:
import pycottas
except ImportError as exc:
raise RuntimeError(f"COTTAS dependency is unavailable: {exc}") from exc
# The adapter also preserves default/named graphs when decoding quads.
decoded_triples = count_decoded(
[sys.executable, str(Path(__file__).with_name("cottas_tool.py")),
"decompress", str(artifact), "/dev/stdout"]
)
validator = "cottas_tool.decompress"
return {
"valid": True,
"source_triples": source_triples,
"decoded_triples": decoded_triples,
"expected_triples": args.expected_triples,
"count_match": source_triples == decoded_triples,
"validator": validator,
} | ({
"error": (
"decoded triple count does not match the source: "
f"source={source_triples}, decoded={decoded_triples}"
)
} if source_triples != decoded_triples else {})
def main() -> int:
parser = argparse.ArgumentParser(description="Validate an HDT or COTTAS artifact")
parser.add_argument("--source", required=True, help="plain or gzip-compressed N-Triples/N-Quads")
parser.add_argument("--artifact", required=True, help="HDT or COTTAS artifact")
parser.add_argument("--format", required=True, choices=("hdt", "cottas"))
parser.add_argument("--expected-triples", type=int)
parser.add_argument(
"--source-triples",
type=int,
help="source count already collected by an upstream streaming pass",
)
parser.add_argument(
"--skip-index-check",
action="store_true",
help="skip HDT index initialization when the caller already performed it",
)
parser.add_argument("--result-path", required=True)
args = parser.parse_args()
result_path = Path(args.result_path)
try:
result = validate(args)
exit_code = 0 if result.get("valid") and result.get("count_match") else 1
except Exception as exc:
result = {
"valid": False,
"source_triples": None,
"decoded_triples": None,
"expected_triples": args.expected_triples,
"count_match": False,
"error": str(exc),
}
exit_code = 1
result_path.parent.mkdir(parents=True, exist_ok=True)
result_path.write_text(json.dumps(result, indent=2) + "\n", encoding="utf-8")
if exit_code:
print(result.get("error", "compression validation failed"), file=sys.stderr)
return exit_code
if __name__ == "__main__":
raise SystemExit(main())