Skip to content

Commit 0eda271

Browse files
committed
added compression completion check and fixed cottas docker error
1 parent 764d27d commit 0eda271

9 files changed

Lines changed: 813 additions & 25 deletions

‎COMPRESSION_CHANGELOG.md‎

Lines changed: 35 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -146,6 +146,17 @@ The Docker image installs `pycottas` in `/opt/pycottas-venv` and invokes the
146146
small adapter at `/opt/vcf-rdfizer/cottas_tool.py`. COTTAS is therefore part of
147147
the Docker workflow, not a dependency of the lightweight pip/conda wrapper.
148148

149+
The adapter runs every `pycottas.rdf2cottas` and `pycottas.cat` operation in a
150+
new temporary subdirectory of the container-local `/work` filesystem. pycottas
151+
uses a default `pycottas.duckdb` database in its current working directory;
152+
without isolation, a later chunk can reopen the prior database and fail while
153+
creating its existing `quads` table. Input and output paths are resolved before
154+
entering the temporary directory, so COTTAS artifacts remain in the designated
155+
workspace/output mount while the DuckDB database and related scratch files are
156+
removed immediately after each operation. `COTTAS_SCRATCH_DIR` can override
157+
`/work` for image-internal deployments, but the wrapper does not require users
158+
to set it.
159+
149160
## Metrics and Cleanup
150161

151162
Every chunk conversion, merge, index-generation, and final representation
@@ -162,6 +173,30 @@ default; `--remove-rdf-storage-output` makes aggregate removal explicit, while
162173
`--keep-rmlstreamer-rdf-output` retains the aggregate produced by RMLStreamer.
163174
The two flags are mutually exclusive.
164175

176+
### Representation Validation
177+
178+
Every generated base HDT or COTTAS representation is validated before any
179+
packaging stage and before raw RDF cleanup. The validation has two parts:
180+
181+
1. The artifact is opened and decoded with its native reader. HDT is loaded
182+
through `ensure_hdt_index.sh`/HDT Java and streamed through `hdt2rdf`;
183+
COTTAS is streamed through `pycottas.cottas2rdf`.
184+
2. The decoded N-Triples stream is counted and compared with the source count.
185+
186+
Full mode passes the `output_triples` count recorded by RMLStreamer whenever
187+
it is available. The partition planner already reads the source once, so its
188+
record count is reused for validation. Compression-only mode has no upstream
189+
conversion metric; in that case the validator counts plain or gzip-compressed
190+
N-Triples directly as a fallback. No decoded RDF file is written during this
191+
check.
192+
193+
A missing decoder, unreadable artifact, or count mismatch makes the
194+
compression stage fail and prevents RDF cleanup. Reports are stored in the
195+
method details of the raw compression JSON and expose `source_triples`,
196+
`decoded_triples`, `count_match`, and `valid`. The aggregate `metrics.csv`
197+
also records the source count, decoded count, and validation status for HDT
198+
and COTTAS.
199+
165200
## Compatibility and Limits
166201

167202
Full mode requires `--rdf-storage-mode plain` or

‎README.md‎

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -353,6 +353,16 @@ Compression metrics now include per-method:
353353
- `sys_seconds_*`
354354
- `max_rss_kb_*`
355355

356+
When HDT or COTTAS is selected, compression also validates the final base
357+
artifact before packaging or RDF cleanup. The validator reads the source
358+
triple count, streams the artifact back through the native decoder, and
359+
requires equal counts. HDT validation also initializes the `.hdt.index`
360+
sidecar. Validation results and `source_triples`/`decoded_triples` are stored
361+
in the per-run compression JSON and in the HDT/COTTAS columns of `metrics.csv`.
362+
Compression fails closed if the artifact cannot be decoded or the counts do
363+
not match. In compression-only mode, the source count is obtained by a
364+
streaming fallback when no upstream conversion metrics are available.
365+
356366
For partitioned HDT/COTTAS runs, the final method metric reports one
357367
sample-level result, while raw metrics also include a sample-scoped
358368
`__partitioned_compression__` artifact describing chunk conversion, merge
@@ -377,11 +387,20 @@ HDT index is generated after merging through HDT Java's indexed search loader,
377387
and COTTAS chunks are merged with `pycottas.cat`, which rebuilds the query
378388
indexes for the merged representation.
379389

390+
After each final HDT/COTTAS base artifact is produced, VCF-RDFizer performs a
391+
streaming decode/count check. This verifies both readability and that the
392+
decoded artifact contains exactly the number of source triples. The check is
393+
performed before `.hdt.gz`, `.hdt.br`, `.cottas.gz`, or `.cottas.br` packaging,
394+
and before raw RDF cleanup.
395+
380396
Partitioned HDT/COTTAS compression runs in an ephemeral Docker-managed
381397
workspace. Temporary RDF chunks, COTTAS/DuckDB scratch data, intermediate
382398
representations, and merge files are not written to the output directory.
383399
After a successful or failed run, the temporary workspace is removed; only
384400
the selected final artifacts and normal run metrics remain on the host.
401+
Each COTTAS conversion and merge also receives a fresh container-local DuckDB
402+
workspace, which is removed as soon as that operation completes. This prevents
403+
state from one chunk being reused by another and requires no user configuration.
385404

386405
HDT Java 3.0.10 does not provide a standalone `hdtGenerateIndex` executable.
387406
VCF-RDFizer sends an `exit` command to the supported `hdtSearch.sh` launcher;

‎src/cottas_tool.py‎

Lines changed: 41 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,29 @@
22
"""Small Docker-side adapter for the pycottas conversion and merge API."""
33

44
import argparse
5+
import os
56
import sys
7+
import tempfile
8+
from contextlib import contextmanager
9+
from pathlib import Path
10+
11+
12+
@contextmanager
13+
def cottas_scratch_workspace():
14+
"""Run one pycottas operation with an isolated DuckDB working directory."""
15+
scratch_root = Path(os.environ.get("COTTAS_SCRATCH_DIR", "/work")).resolve()
16+
scratch_root.mkdir(parents=True, exist_ok=True)
17+
original_working_directory = Path.cwd()
18+
# pycottas defaults to ``pycottas.duckdb`` in the current directory.
19+
# A fresh directory prevents one chunk from reusing another chunk's
20+
# database, whose ``quads`` table already exists.
21+
with tempfile.TemporaryDirectory(prefix="vcf-rdfizer-cottas-", dir=scratch_root) as directory:
22+
try:
23+
os.chdir(directory)
24+
yield
25+
finally:
26+
# Leave the directory before TemporaryDirectory removes it.
27+
os.chdir(original_working_directory)
628

729

830
def main() -> int:
@@ -28,22 +50,29 @@ def main() -> int:
2850
return 127
2951

3052
if args.command == "convert":
53+
rdf_path = str(Path(args.rdf_path).resolve())
54+
cottas_path = str(Path(args.cottas_path).resolve())
3155
# disk=True keeps parser/index construction from requiring the whole
3256
# RDF chunk in Python memory.
33-
pycottas.rdf2cottas(
34-
args.rdf_path,
35-
args.cottas_path,
36-
index=args.index,
37-
disk=True,
38-
)
57+
with cottas_scratch_workspace():
58+
pycottas.rdf2cottas(
59+
rdf_path,
60+
cottas_path,
61+
index=args.index,
62+
disk=True,
63+
)
3964
return 0
4065

41-
pycottas.cat(
42-
[args.left_path, args.right_path],
43-
args.cottas_path,
44-
index=args.index,
45-
remove_input_files=True,
46-
)
66+
left_path = str(Path(args.left_path).resolve())
67+
right_path = str(Path(args.right_path).resolve())
68+
cottas_path = str(Path(args.cottas_path).resolve())
69+
with cottas_scratch_workspace():
70+
pycottas.cat(
71+
[left_path, right_path],
72+
cottas_path,
73+
index=args.index,
74+
remove_input_files=True,
75+
)
4776
return 0
4877

4978

‎src/partitioned_compression.py‎

Lines changed: 94 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,12 @@
3131
)
3232

3333

34+
def is_triple_line(line: bytes) -> bool:
35+
"""Identify an N-Triples record without parsing or copying its terms."""
36+
stripped = line.strip()
37+
return bool(stripped) and not stripped.startswith(b"#") and stripped.endswith(b".")
38+
39+
3440
def parse_args() -> argparse.Namespace:
3541
parser = argparse.ArgumentParser(description="VCF-RDFizer partitioned compression runner")
3642
parser.add_argument("--source", required=True, help="plain or gzip-compressed N-Triples input")
@@ -40,6 +46,7 @@ def parse_args() -> argparse.Namespace:
4046
parser.add_argument("--target-chunk-bytes", required=True, type=int)
4147
parser.add_argument("--min-chunk-bytes", required=True, type=int)
4248
parser.add_argument("--max-chunk-bytes", required=True, type=int)
49+
parser.add_argument("--expected-triples", type=int)
4350
parser.add_argument("--result-path", required=True)
4451
return parser.parse_args()
4552

@@ -124,7 +131,8 @@ def close_chunk():
124131
handle.write(line)
125132
chunk_size += line_size
126133
logical_offset += line_size
127-
record_count += 1
134+
if is_triple_line(line):
135+
record_count += 1
128136
finally:
129137
close_chunk()
130138

@@ -306,6 +314,12 @@ def main() -> int:
306314
)
307315
if not chunks:
308316
raise ValueError("RDF source contains no complete records")
317+
source_triples = int(plan["record_count"])
318+
if args.expected_triples is not None and source_triples != args.expected_triples:
319+
raise ValueError(
320+
"source triple count does not match the upstream conversion count: "
321+
f"source={source_triples}, expected={args.expected_triples}"
322+
)
309323
hdt_paths: list[Path] = []
310324
cottas_paths: list[Path] = []
311325
needs_hdt = any(method in HDT_METHODS for method in methods)
@@ -326,6 +340,65 @@ def main() -> int:
326340
if any(method in COTTAS_METHODS for method in methods) and not cottas_python:
327341
raise RuntimeError("Missing Python runtime for COTTAS")
328342

343+
def validate_artifact(
344+
*,
345+
name: str,
346+
artifact: Path,
347+
artifact_format: str,
348+
python_bin: str,
349+
skip_index_check: bool = False,
350+
) -> dict:
351+
"""Decode/count one final artifact and return its validation report."""
352+
validation_path = work_dir / f".{name}.validation.json"
353+
command = [
354+
python_bin,
355+
"/opt/vcf-rdfizer/validate_compression.py",
356+
"--source",
357+
str(source),
358+
"--artifact",
359+
str(artifact),
360+
"--format",
361+
artifact_format,
362+
"--source-triples",
363+
str(source_triples),
364+
"--result-path",
365+
str(validation_path),
366+
]
367+
if args.expected_triples is not None:
368+
command.extend(["--expected-triples", str(args.expected_triples)])
369+
if skip_index_check:
370+
command.append("--skip-index-check")
371+
stage = runner.run(name, command, validation_path)
372+
try:
373+
report = json.loads(validation_path.read_text(encoding="utf-8"))
374+
except (FileNotFoundError, json.JSONDecodeError) as exc:
375+
report = {
376+
"valid": False,
377+
"count_match": False,
378+
"error": f"validator did not produce a valid report: {exc}",
379+
}
380+
finally:
381+
validation_path.unlink(missing_ok=True)
382+
report.setdefault(
383+
"timing",
384+
{
385+
"wall_seconds": stage.get("wall_seconds"),
386+
"user_seconds": stage.get("user_seconds"),
387+
"sys_seconds": stage.get("sys_seconds"),
388+
"max_rss_kb": stage.get("max_rss_kb"),
389+
},
390+
)
391+
if (
392+
stage["exit_code"] != 0
393+
or not report.get("valid")
394+
or not report.get("count_match")
395+
):
396+
raise RuntimeError(
397+
f"{artifact_format.upper()} validation failed for {artifact.name}: "
398+
f"{report.get('error', 'decoded triple count mismatch')}"
399+
)
400+
return report
401+
329402
for index, chunk in enumerate(chunks):
330403
if needs_hdt:
331404
chunk_hdt = work_dir / f"chunk-{index:05d}.hdt"
@@ -369,6 +442,13 @@ def main() -> int:
369442
add_totals(hdt_total, index_stage)
370443
if index_stage["exit_code"] != 0:
371444
raise RuntimeError("final HDT index initialization failed")
445+
hdt_validation = validate_artifact(
446+
name="hdt-validate",
447+
artifact=output_hdt,
448+
artifact_format="hdt",
449+
python_bin=sys.executable,
450+
skip_index_check=True,
451+
)
372452
results["hdt"] = {
373453
**finalize_totals(hdt_total),
374454
"output_path": str(output_hdt),
@@ -379,6 +459,7 @@ def main() -> int:
379459
"merge_rounds": hdt_rounds,
380460
"index_path": str(output_index),
381461
"index_size_bytes": output_index.stat().st_size,
462+
"validation": hdt_validation,
382463
},
383464
}
384465

@@ -395,12 +476,23 @@ def main() -> int:
395476
raise RuntimeError("COTTAS merge failed")
396477
shutil.copyfile(final_cottas, output_cottas)
397478
final_cottas.unlink(missing_ok=True)
479+
cottas_validation = validate_artifact(
480+
name="cottas-validate",
481+
artifact=output_cottas,
482+
artifact_format="cottas",
483+
python_bin=cottas_python,
484+
)
398485
results["cottas"] = {
399486
**finalize_totals(cottas_total),
400487
"output_path": str(output_cottas),
401488
"output_size_bytes": output_cottas.stat().st_size,
402489
"source": "partitioned_generated",
403-
"details": {**plan, "merge_rounds": cottas_rounds, "index": "spo"},
490+
"details": {
491+
**plan,
492+
"merge_rounds": cottas_rounds,
493+
"index": "spo",
494+
"validation": cottas_validation,
495+
},
404496
}
405497

406498
for method in methods:

0 commit comments

Comments
 (0)