From 76667026aa6f46eee34121f71563490ff8f797f1 Mon Sep 17 00:00:00 2001 From: ecrum19 Date: Fri, 25 Sep 2026 09:25:36 +0200 Subject: [PATCH 1/2] Support multiple COTTAS index orders with shared chunk conversion --- .github/workflows/tests.yml | 2 + Dockerfile | 1 + README.md | 20 +- docs/cli-reference.md | 1 + docs/representations.md | 17 +- pyproject.toml | 2 +- src/cottas_tool.py | 176 +++++++++----- src/partitioned_compression.py | 183 +++++++-------- test/README.md | 17 ++ test/cottas_index_roundtrip.py | 54 +++++ test/test_cottas_indexes_unit.py | 382 +++++++++++++++++++++++++++++++ test/test_cottas_tool.py | 9 +- vcf_rdfizer.py | 60 ++++- vcf_rdfizer_cottas.py | 24 ++ 14 files changed, 781 insertions(+), 167 deletions(-) create mode 100644 test/cottas_index_roundtrip.py create mode 100644 test/test_cottas_indexes_unit.py create mode 100644 vcf_rdfizer_cottas.py diff --git a/.github/workflows/tests.yml b/.github/workflows/tests.yml index 6990355..c6923d1 100644 --- a/.github/workflows/tests.yml +++ b/.github/workflows/tests.yml @@ -111,7 +111,9 @@ jobs: run: | python -m pip install "pip==${{ env.PIP_VERSION }}" coverage python -m pip install -e . + python -m pip install pycottas==1.1.0 duckdb==1.5.5 pyarrow==22.0.0 coverage run -m unittest discover -s test -p "test_*_unit.py" + coverage run --append -m unittest test.test_cottas_tool coverage xml -o coverage.xml - name: Upload coverage to Codecov diff --git a/Dockerfile b/Dockerfile index ef995bb..4d868a4 100644 --- a/Dockerfile +++ b/Dockerfile @@ -254,6 +254,7 @@ COPY src/validation/ /opt/vcf-rdfizer/validation/ # Shared, dependency-free helpers used by both the host CLI and the in-container # runners, so they live at the repository root rather than in src/. COPY vcf_rdfizer_gzip.py /opt/vcf-rdfizer/ +COPY vcf_rdfizer_cottas.py /opt/vcf-rdfizer/ # The vocabulary terms, the VCF-version model and the lexical parsers. The # validation runner imports this as a sibling module, so the oracle and the # emitters share one description of what the graph should contain instead of diff --git a/README.md b/README.md index 7473ae6..e5d0366 100644 --- a/README.md +++ b/README.md @@ -269,6 +269,7 @@ host filesystem. - `space-optimized`: gzip each part into one `.nt.gz` aggregate and delete the source part immediately - `--rdf-compression {gzip,brotli,none}` raw RDF artifacts to retain - `--representations {hdt,cottas,none}` queryable primary representations +- `--cottas-indexes` one or more of `spo,sop,pso,pos,osp,ops`, or `all` (default `spo`) - `--artifact-compression {gzip,brotli,none}` optional packaging applied to each selected representation - `--hdt-strategy {auto,partitioned,single}` - `auto`: in full mode, build smaller HDT chunks and merge them with native `hdtc` @@ -431,11 +432,12 @@ compatible with condensed emission. - `-H, --hdt` existing `.hdt` file; creates or regenerates its sibling sidecar - `--cottas` existing `.cottas` file; rebuilds its embedded query index in place +- `--cottas-indexes` selects the replacement order and optional additional index copies - Exactly one of `--hdt` or `--cottas` is required. - HDT indexing is Java-free: the image uses `hdtc` 1.1.0 to generate the canonical v1-1 sidecar named `.hdt.index.v1-1`. -- COTTAS indexes are stored inside the Parquet-based `.cottas` file, so the - file is rewritten atomically; no separate COTTAS index file is expected. +- Each COTTAS copy stores its index inside the Parquet-based `.cottas` file; + the primary is rewritten atomically and there is no index sidecar. - Existing indexes are intentionally replaced. Use this mode when an HDT sidecar is missing/stale or when a COTTAS file needs its query ordering and zone-map metadata rebuilt. No VCF conversion, RDF conversion, packaging, or @@ -644,10 +646,18 @@ buffers and may increase temporary I/O; higher values can improve performance when memory is available. Temporary files live in the container's `/work` area and are removed after the attempt. +Choose COTTAS orders with `--cottas-indexes spo,pso,pos` or `--cottas-indexes all`. +The first order keeps `sample.cottas`; additional whole-graph copies are named +`sample.pso.cottas`, `sample.pos.cottas`, etc. Query one copy at a time. Each RDF +chunk is parsed once; extra orders sort its deduplicated Parquet data. Each order +gets its own streaming merge, round-trip check and requested gzip/Brotli packages. +JSON metrics list every index; existing CSV sizes refer to the primary copy. +See [COTTAS representations](docs/representations.md#5-cottas) for details. + COTTAS avoids both a global in-memory `DISTINCT` and a global external sort. -Each chunk is already written in `spo` order, so the final stage performs a +Each chunk is written in each requested order, so the final stage performs a bounded k-way Parquet merge: it holds one small batch from each chunk, writes -one copy of each adjacent equal triple, and preserves the `spo` index. Its +one copy of each adjacent equal triple, and preserves that index. Its memory use is controlled by `COTTAS_MERGE_BATCH_ROWS` (default `2048`), not by the total RDF graph size or a temporary DuckDB sort area. Override it only to tune the memory/throughput tradeoff, for example: @@ -822,7 +832,7 @@ COTTAS chunk conversion uses `pycottas.rdf2cottas(..., disk=True)`. The final COTTAS merge deliberately does **not** call `pycottas.cat`: version 1.1.0 runs its global `DISTINCT` plus `ORDER BY` through an unbounded in-memory DuckDB connection, which can be killed on large condensed graphs. VCF-RDFizer instead -uses a PyArrow k-way merge of the already `spo`-sorted Parquet chunks. It keeps +uses a PyArrow k-way merge of the already index-sorted Parquet chunks. It keeps only a configurable batch from each input, drops adjacent duplicate triples, and writes the final COTTAS file incrementally—no graph-wide DuckDB hash table or external-sort spill directory is created. The merge emits processed-source diff --git a/docs/cli-reference.md b/docs/cli-reference.md index f760610..a6f8b91 100644 --- a/docs/cli-reference.md +++ b/docs/cli-reference.md @@ -80,6 +80,7 @@ See [Data linking](datalinking.md) for commands and current limits. There is no | `--rdf-storage-mode` | `plain`, `space-optimized` | `space-optimized` | | `--rdf-compression` | `gzip`, `brotli`, `none` | `gzip,brotli` | | `--representations` | `hdt`, `cottas`, `none` | `hdt` | +| `--cottas-indexes` | `spo,sop,pso,pos,osp,ops` (one or more), or `all`; also applies to COTTAS index mode | `spo` | | `--artifact-compression` | `gzip`, `brotli`, `none` | `none` | | `--hdt-strategy` | `auto`, `partitioned`, `single` | `auto` | | `--chunk-target-bytes` | bytes | 512 MiB | diff --git a/docs/representations.md b/docs/representations.md index 676a774..734f2b5 100644 --- a/docs/representations.md +++ b/docs/representations.md @@ -120,12 +120,23 @@ the raw RDF is retained so it can be repaired later with `--mode index`. COTTAS is a Parquet-based representation with its index **inside** the artifact; there is no sidecar. Chunk conversion uses `pycottas.rdf2cottas(..., disk=True)` with a fresh container-local DuckDB -workspace per operation. +workspace per operation. `--cottas-indexes` accepts `spo` (default), `sop`, +`pso`, `pos`, `osp`, `ops`, a comma-separated selection, or `all`; case is +ignored and repeated orders are built once. The first order uses `name.cottas`; +additional orders use `name..cottas`. Each file contains the whole graph +in one order. Query one copy at a time; loading them together repeats the graph. + +The RDF chunk is parsed and deduplicated once. Additional orders sort that +chunk's Parquet data using DuckDB, without repeating RDF parsing or `DISTINCT`. +Indexes merge sequentially, so adding orders increases disk use and work without +multiplying merge memory. Gzip/Brotli packaging and round-trip checks cover each +copy. Per-index paths, sizes and validation are in JSON `details.indexes`; the +existing CSV size columns describe the primary copy, with timings for all orders. The final merge deliberately does **not** call `pycottas.cat`. In version 1.1.0 that runs a global `DISTINCT` plus `ORDER BY` through an unbounded in-memory DuckDB connection, which gets OOM-killed on large condensed graphs. VCF-RDFizer -instead performs a **k-way PyArrow merge** of the already `spo`-sorted Parquet +instead performs a **k-way PyArrow merge** of the already index-sorted Parquet chunks: it holds at most `COTTAS_MERGE_BATCH_ROWS` (default 2048) rows per input, drops adjacent duplicate triples, and writes the result incrementally. There is no graph-wide hash table and no external-sort spill directory; memory @@ -160,7 +171,7 @@ deliberate exception to "never overwrite a planned artifact". | Input | Behaviour | | --- | --- | | `--hdt file.hdt` | Existing versioned sidecars are moved aside, regenerated, and restored if indexing fails; incomplete replacements are removed first | -| `--cottas file.cottas` | The artifact is rewritten atomically through the same bounded streaming Parquet rewrite with the default `spo` index; the original stays in place if it fails | +| `--cottas file.cottas` | `--cottas-indexes` selects the replacement order and optional additional copies; default `spo`. Changed orders use bounded sorted batches and streaming merges (at most 128 runs per merge). The primary is replaced only after all indexes build successfully; existing additional outputs are refused | No conversion, packaging, or decompression output is produced. Standalone index mode is **strict** — unlike the in-run degradation above, a failure is a failure. diff --git a/pyproject.toml b/pyproject.toml index c224c5c..0db051c 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -45,7 +45,7 @@ dev = [ ] [tool.setuptools] -py-modules = ["vcf_rdfizer", "vcf_rdfizer_gzip", "vcf_rdfizer_rules", "vcf_rdfizer_link", "vcf_rdfizer_vocab"] +py-modules = ["vcf_rdfizer", "vcf_rdfizer_gzip", "vcf_rdfizer_rules", "vcf_rdfizer_link", "vcf_rdfizer_vocab", "vcf_rdfizer_cottas"] packages = ["vcf_rdfizer_data", "vcf_rdfizer_data.rules", "vcf_rdfizer_data.linkers", "vcf_rdfizer_linking"] include-package-data = true diff --git a/src/cottas_tool.py b/src/cottas_tool.py index a119227..49b9168 100644 --- a/src/cottas_tool.py +++ b/src/cottas_tool.py @@ -10,10 +10,68 @@ from contextlib import contextmanager from pathlib import Path +from vcf_rdfizer_cottas import COTTAS_INDEXES, cottas_index_paths, parse_cottas_indexes + DEFAULT_COTTAS_MERGE_BATCH_ROWS = 2048 COTTAS_OUTPUT_BATCH_ROWS = 16 * 1024 COTTAS_MERGE_PROGRESS_ROWS = 250_000 +COTTAS_REINDEX_FAN_IN = 128 + + +def sort_cottas_chunk(source: Path, outputs: dict[str, Path]) -> None: + """Build extra orders from one parsed chunk, with disk-backed sorting.""" + import duckdb + + with duckdb.connect("index-sort.duckdb") as connection: + connection.execute("SET preserve_insertion_order = false") + for index, output in outputs.items(): + if index not in COTTAS_INDEXES: + raise ValueError("COTTAS index must be a permutation of spo") + connection.execute( + f"COPY (SELECT s, p, o FROM read_parquet($source) ORDER BY {', '.join(index)}) " + f"TO $target (FORMAT PARQUET, COMPRESSION ZSTD, COMPRESSION_LEVEL 22, " + f"PARQUET_VERSION v2, KV_METADATA {{index: '{index}'}})", + {"source": str(source), "target": str(output)}, + ) + + +def reindex_cottas(source: Path, outputs: dict[str, Path]) -> None: + """Reorder bounded Parquet batches, then merge each requested index.""" + import pyarrow as pa + import pyarrow.parquet as pq + + with pq.ParquetFile(source) as parquet: + if len(outputs) == 1 and cottas_file_index(parquet) == next(iter(outputs)): + index, output = next(iter(outputs.items())) + streaming_cottas_merge([str(source)], str(output), index=index, remove_input_files=False) + return + # Sorting batches keeps reindexing bounded even for a cohort-sized file. + with tempfile.TemporaryDirectory(prefix="reindex-runs-", dir=Path.cwd()) as directory: + runs = {index: [] for index in outputs} + for number, batch in enumerate(parquet.iter_batches(batch_size=COTTAS_OUTPUT_BATCH_ROWS, columns=["s", "p", "o"])): + table = pa.Table.from_batches([batch]) + for index in outputs: + path = Path(directory) / f"{number}.{index}.cottas" + ordered = table.sort_by([(column, "ascending") for column in index]) + pq.write_table(ordered.replace_schema_metadata({b"index": index.encode()}), path, compression="zstd") + runs[index].append(str(path)) + for index, output in outputs.items(): + if runs[index]: + paths = runs[index] + round_number = 0 + while len(paths) > COTTAS_REINDEX_FAN_IN: + merged = [] + for start in range(0, len(paths), COTTAS_REINDEX_FAN_IN): + path = Path(directory) / f"merge-{index}-{round_number}-{start}.cottas" + streaming_cottas_merge(paths[start:start + COTTAS_REINDEX_FAN_IN], str(path), index=index, remove_input_files=True) + merged.append(str(path)) + paths = merged + round_number += 1 + streaming_cottas_merge(paths, str(output), index=index, remove_input_files=True) + else: + schema = pa.schema([parquet.schema_arrow.field(name) for name in ("s", "p", "o")], metadata={b"index": index.encode()}) + pq.write_table(pa.Table.from_batches([], schema=schema), output, compression="zstd") @contextmanager @@ -94,34 +152,38 @@ def __init__(self, parquet_module, path: Path, index: str, batch_rows: int): self.path = path self.index = index self.parquet_file = parquet_module.ParquetFile(path) - self.schema = self.parquet_file.schema_arrow - missing_columns = {"s", "p", "o"} - set(self.schema.names) - if missing_columns: - raise RuntimeError( - f"COTTAS input {path.name} is missing columns: {', '.join(sorted(missing_columns))}" - ) - source_index = cottas_file_index(self.parquet_file) - if source_index != index: - found = source_index or "missing" - raise RuntimeError( - f"COTTAS input {path.name} is indexed as {found!r}, not {index!r}; " - "a streaming merge requires every input to use the requested index" + try: + self.schema = self.parquet_file.schema_arrow + missing_columns = {"s", "p", "o"} - set(self.schema.names) + if missing_columns: + raise RuntimeError( + f"COTTAS input {path.name} is missing columns: {', '.join(sorted(missing_columns))}" + ) + source_index = cottas_file_index(self.parquet_file) + if source_index != index: + found = source_index or "missing" + raise RuntimeError( + f"COTTAS input {path.name} is indexed as {found!r}, not {index!r}; " + "a streaming merge requires every input to use the requested index" + ) + self.fields = tuple(self.schema.field(name) for name in ("s", "p", "o")) + positions = {"s": 0, "p": 1, "o": 2} + self._sort_positions = tuple(positions[column] for column in index) + self._batches = self.parquet_file.iter_batches( + batch_size=batch_rows, + columns=["s", "p", "o"], + use_threads=False, ) - self.fields = tuple(self.schema.field(name) for name in ("s", "p", "o")) - positions = {"s": 0, "p": 1, "o": 2} - self._sort_positions = tuple(positions[column] for column in index) - self._batches = self.parquet_file.iter_batches( - batch_size=batch_rows, - columns=["s", "p", "o"], - use_threads=False, - ) - self._values: tuple[list, list, list] | None = None - self._row = 0 - self.current: tuple | None = None - self.sort_key: tuple | None = None - self._previous_sort_key: tuple | None = None - self.exhausted = False - self.advance() + self._values: tuple[list, list, list] | None = None + self._row = 0 + self.current: tuple | None = None + self.sort_key: tuple | None = None + self._previous_sort_key: tuple | None = None + self.exhausted = False + self.advance() + except Exception: + self.close() + raise @property def row_count(self) -> int: @@ -200,10 +262,8 @@ def streaming_cottas_merge( temporary_path: Path | None = None writer = None try: - streams = [ - CottasTripleStream(pq, Path(path), normalized_index, batch_rows) - for path in input_paths - ] + for path in input_paths: + streams.append(CottasTripleStream(pq, Path(path), normalized_index, batch_rows)) fields = streams[0].fields expected_types = tuple(field.type for field in fields) for stream in streams[1:]: @@ -332,13 +392,13 @@ def main() -> int: convert = subparsers.add_parser("convert", help="convert one RDF file to COTTAS") convert.add_argument("rdf_path") convert.add_argument("cottas_path") - convert.add_argument("index", nargs="?", default="spo") + convert.add_argument("index", nargs="?", default="spo", type=parse_cottas_indexes) merge = subparsers.add_parser("merge", help="merge two COTTAS files") merge.add_argument("left_path") merge.add_argument("right_path") merge.add_argument("cottas_path") - merge.add_argument("index", nargs="?", default="spo") + merge.add_argument("index", nargs="?", default="spo", type=str.lower, choices=COTTAS_INDEXES) merge_many = subparsers.add_parser( "merge-many", @@ -351,7 +411,7 @@ def main() -> int: help="COTTAS inputs to merge", ) merge_many.add_argument("--output-cottas-file", required=True) - merge_many.add_argument("--index", default="spo") + merge_many.add_argument("--index", default="spo", type=str.lower, choices=COTTAS_INDEXES) merge_many.add_argument( "--progress-path", help="optional JSONL sidecar for bounded streaming-merge progress", @@ -362,7 +422,7 @@ def main() -> int: help="rebuild the embedded COTTAS query index in place", ) reindex.add_argument("cottas_path") - reindex.add_argument("index", nargs="?", default="spo") + reindex.add_argument("index", nargs="?", default="spo", type=parse_cottas_indexes) decompress = subparsers.add_parser("decompress", help="convert COTTAS to RDF") decompress.add_argument("cottas_path") @@ -384,9 +444,13 @@ def main() -> int: pycottas.rdf2cottas( rdf_path, cottas_path, - index=args.index, + index=args.index[0], disk=True, ) + outputs = cottas_index_paths(Path(cottas_path), args.index) + outputs.pop(args.index[0]) + if outputs: + sort_cottas_chunk(Path(cottas_path), outputs) return 0 if args.command == "decompress": @@ -408,29 +472,29 @@ def main() -> int: # files. Rebuild into a temporary file in the same directory, then # replace the original only after the streaming Parquet rewrite # completes successfully. - temporary_path = None + outputs = cottas_index_paths(cottas_path, args.index) + for output in list(outputs.values())[1:]: + if output.exists(): + raise FileExistsError(f"COTTAS index output already exists: {output}") + temporary_paths = {} try: - file_descriptor, temporary_name = tempfile.mkstemp( - prefix=f".{cottas_path.name}.reindex-", - suffix=".cottas", - dir=str(cottas_path.parent), - ) - os.close(file_descriptor) - temporary_path = Path(temporary_name) - temporary_path.unlink() - with cottas_scratch_workspace(): - streaming_cottas_merge( - [str(cottas_path)], - str(temporary_path), - index=args.index, - remove_input_files=False, + for index, output in outputs.items(): + file_descriptor, temporary_name = tempfile.mkstemp( + prefix=f".{output.name}.reindex-", suffix=".cottas", dir=str(output.parent), ) - if not temporary_path.is_file() or temporary_path.stat().st_size == 0: - raise RuntimeError("pycottas did not create a non-empty reindexed file") - os.replace(temporary_path, cottas_path) - temporary_path = None + os.close(file_descriptor) + temporary_paths[index] = Path(temporary_name) + temporary_paths[index].unlink() + with cottas_scratch_workspace(): + reindex_cottas(cottas_path, temporary_paths) + for temporary_path in temporary_paths.values(): + if not temporary_path.is_file() or temporary_path.stat().st_size == 0: + raise RuntimeError("pycottas did not create a non-empty reindexed file") + # Publish the primary last so failed builds preserve the original. + for index in reversed(args.index): + os.replace(temporary_paths[index], outputs[index]) finally: - if temporary_path is not None: + for temporary_path in temporary_paths.values(): temporary_path.unlink(missing_ok=True) return 0 diff --git a/src/partitioned_compression.py b/src/partitioned_compression.py index 175d8b0..c2f6d8a 100644 --- a/src/partitioned_compression.py +++ b/src/partitioned_compression.py @@ -26,6 +26,8 @@ resource = None from pathlib import Path +from vcf_rdfizer_cottas import cottas_index_paths, parse_cottas_indexes + HDT_METHODS = {"hdt", "hdt_gzip", "hdt_brotli"} COTTAS_METHODS = {"cottas", "cottas_gzip", "cottas_brotli"} @@ -108,6 +110,7 @@ def parse_args() -> argparse.Namespace: parser.add_argument("--output-dir", required=True) parser.add_argument("--output-name", required=True) parser.add_argument("--methods", required=True, help="comma-separated internal method names") + parser.add_argument("--cottas-indexes", default="spo", type=parse_cottas_indexes) parser.add_argument("--target-chunk-bytes", required=True, type=int) parser.add_argument("--min-chunk-bytes", required=True, type=int) parser.add_argument("--max-chunk-bytes", required=True, type=int) @@ -386,6 +389,7 @@ def cottas_merge_many_command( inputs: list[Path], merged: Path, *, + index: str = "spo", progress_path: Path | None = None, ) -> list[str]: """Build one bounded-memory, multi-input COTTAS merge command. @@ -404,7 +408,7 @@ def cottas_merge_many_command( "--output-cottas-file", str(merged), "--index", - "spo", + index, *(["--progress-path", str(progress_path)] if progress_path else []), ] @@ -414,6 +418,7 @@ def cottas_merge_command( left: Path, right: Path, merged: Path, + index: str = "spo", ) -> list[str]: """Build a compatible two-input bounded-memory COTTAS merge command. @@ -428,7 +433,7 @@ def cottas_merge_command( str(left), str(right), str(merged), - "spo", + index, ] @@ -832,6 +837,8 @@ def main() -> int: cottas_total = {"exit_code": 0, "wall_seconds": 0.0, "user_seconds": 0.0, "sys_seconds": 0.0, "max_rss_kb": 0, "has_user": False, "has_sys": False, "has_rss": False} output_hdt = output_dir / f"{args.output_name}.hdt" output_cottas = output_dir / f"{args.output_name}.cottas" + cottas_indexes = args.cottas_indexes + cottas_outputs = cottas_index_paths(output_cottas, cottas_indexes) cottas_failed = False cottas_warning = None @@ -913,7 +920,8 @@ def skipped_cottas_result(method: str) -> dict: def cleanup_cottas_intermediates() -> None: """Release COTTAS chunks/merge outputs after a failed attempt.""" for path in list(cottas_paths): - path.unlink(missing_ok=True) + for indexed_path in cottas_index_paths(path, cottas_indexes).values(): + indexed_path.unlink(missing_ok=True) cottas_paths.clear() for path in work_dir.glob("cottas-merge-*.cottas"): path.unlink(missing_ok=True) @@ -1046,7 +1054,7 @@ def validate_artifact( chunk_cottas = work_dir / f"chunk-{index:05d}.cottas" stage = runner.run( f"cottas-build-{index:05d}", - [cottas_python, "/opt/vcf-rdfizer/cottas_tool.py", "convert", str(chunk), str(chunk_cottas), "spo"], + [cottas_python, "/opt/vcf-rdfizer/cottas_tool.py", "convert", str(chunk), str(chunk_cottas), ",".join(cottas_indexes)], chunk_cottas, ) add_totals(cottas_total, stage) @@ -1065,7 +1073,8 @@ def validate_artifact( ), stage_result=stage, ) - chunk_cottas.unlink(missing_ok=True) + for path in cottas_index_paths(chunk_cottas, cottas_indexes).values(): + path.unlink(missing_ok=True) cleanup_cottas_intermediates() else: cottas_paths.append(chunk_cottas) @@ -1178,94 +1187,69 @@ def validate_artifact( } if cottas_paths and not cottas_failed: - # Each COTTAS chunk is already sorted by the `spo` index. The - # adapter k-way merges those ordered Parquet streams with one small - # batch per input, avoiding both pycottas.cat's global hash table - # and DuckDB external-sort spill files that exceeded 41 GiB for a - # 52-chunk condensed cohort. + # Merge one order at a time: memory does not multiply by index count. + indexed_results = {} cottas_stage = None - cottas_rounds = 0 - if len(cottas_paths) == 1: - final_cottas = cottas_paths[0] - else: - cottas_merged_path = work_dir / "cottas-merge-final.cottas" - cottas_stage = runner.run( - "cottas-merge-stream", - cottas_merge_many_command( - cottas_python, - cottas_paths, - cottas_merged_path, - progress_path=progress_path, - ), - cottas_merged_path, - ) - cottas_stage["cottas_merge_batch_rows"] = os.environ.get( - "COTTAS_MERGE_BATCH_ROWS", DEFAULT_COTTAS_MERGE_BATCH_ROWS - ) - runner.stages[-1].update(cottas_stage) - add_totals(cottas_total, cottas_stage) - cottas_rounds = 1 - final_cottas = ( - cottas_merged_path - if cottas_stage["exit_code"] == 0 and cottas_merged_path.is_file() - else None - ) - if final_cottas is None: - if not args.allow_index_failures: - raise RuntimeError( - failure_message(cottas_stage, "COTTAS merge/index creation failed") + try: + for order, output in cottas_outputs.items(): + paths = [cottas_index_paths(path, cottas_indexes)[order] for path in cottas_paths] + final_cottas = paths[0] + if len(paths) > 1: + final_cottas = work_dir / f"cottas-merge-{order}.cottas" + cottas_stage = runner.run( + "cottas-merge-stream" if order == cottas_indexes[0] else f"cottas-merge-{order}", + cottas_merge_many_command( + cottas_python, paths, final_cottas, + index=order, progress_path=progress_path, + ), + final_cottas, + ) + cottas_stage["cottas_merge_batch_rows"] = os.environ.get( + "COTTAS_MERGE_BATCH_ROWS", DEFAULT_COTTAS_MERGE_BATCH_ROWS + ) + runner.stages[-1].update(cottas_stage) + add_totals(cottas_total, cottas_stage) + if cottas_stage["exit_code"] != 0 or not final_cottas.is_file(): + raise RuntimeError(failure_message(cottas_stage, "COTTAS merge/index creation failed")) + shutil.copyfile(final_cottas, output) + final_cottas.unlink(missing_ok=True) + validation = validate_artifact( + name="cottas-validate" if order == cottas_indexes[0] else f"cottas-validate-{order}", + artifact=output, artifact_format="cottas", python_bin=cottas_python, ) + indexed_results[order] = { + "output_path": str(output), + "output_size_bytes": output.stat().st_size, + "validation": validation, + } + except RuntimeError as exc: + if not args.allow_index_failures: + raise cottas_failed = True cottas_total["exit_code"] = 0 - cleanup_cottas_intermediates() + for output in cottas_outputs.values(): + output.unlink(missing_ok=True) cottas_warning = record_index_warning( - "cottas", - "cottas-index", - output_cottas, - failure_message(cottas_stage, "COTTAS merge/index creation failed"), - stage_result=cottas_stage, + "cottas", "cottas-index", output_cottas, str(exc), stage_result=cottas_stage, ) - else: - shutil.copyfile(final_cottas, output_cottas) - final_cottas.unlink(missing_ok=True) + finally: cleanup_cottas_intermediates() - try: - cottas_validation = validate_artifact( - name="cottas-validate", - artifact=output_cottas, - artifact_format="cottas", - python_bin=cottas_python, - ) - except RuntimeError as exc: - if not args.allow_index_failures: - raise - cottas_failed = True - cottas_total["exit_code"] = 0 - output_cottas.unlink(missing_ok=True) - cottas_warning = record_index_warning( - "cottas", - "cottas-index", - output_cottas, - str(exc), - ) - if not cottas_failed: - results["cottas"] = { - **finalize_totals(cottas_total), - "output_path": str(output_cottas), - "output_size_bytes": output_cottas.stat().st_size, - "source": "partitioned_generated", - "details": { - **plan, - "merge_rounds": cottas_rounds, - "merge_strategy": "pyarrow_streaming_k_way", - "merge_batch_rows": os.environ.get( - "COTTAS_MERGE_BATCH_ROWS", - DEFAULT_COTTAS_MERGE_BATCH_ROWS, - ), - "index": "spo", - "validation": cottas_validation, - }, - } + if not cottas_failed: + results["cottas"] = { + **finalize_totals(cottas_total), + "output_path": str(output_cottas), + "output_size_bytes": output_cottas.stat().st_size, + "source": "partitioned_generated", + "details": { + **plan, + "merge_rounds": int(plan["chunk_count"] > 1), + "merge_strategy": "pyarrow_streaming_k_way", + "merge_batch_rows": os.environ.get("COTTAS_MERGE_BATCH_ROWS", DEFAULT_COTTAS_MERGE_BATCH_ROWS), + "index": cottas_indexes[0], + "indexes": indexed_results, + "validation": indexed_results[cottas_indexes[0]]["validation"], + }, + } if any(method in COTTAS_METHODS for method in methods) and cottas_failed: results["cottas"] = skipped_cottas_result("cottas") @@ -1277,18 +1261,27 @@ def validate_artifact( elif method == "hdt_brotli": artifact = output_dir / f"{args.output_name}.hdt.br" stage = runner.run("hdt-brotli", ["brotli", "-q", "7", "-c", str(output_hdt)], artifact, artifact) - elif method == "cottas_gzip": - if cottas_failed: - results[method] = skipped_cottas_result(method) - continue - artifact = output_dir / f"{args.output_name}.cottas.gz" - stage = runner.run("cottas-gzip", ["gzip", "-c", str(output_cottas)], artifact, artifact) - elif method == "cottas_brotli": + elif method in {"cottas_gzip", "cottas_brotli"}: if cottas_failed: results[method] = skipped_cottas_result(method) continue - artifact = output_dir / f"{args.output_name}.cottas.br" - stage = runner.run("cottas-brotli", ["brotli", "-q", "7", "-c", str(output_cottas)], artifact, artifact) + codec = "gzip" if method == "cottas_gzip" else "brotli" + suffix = ".gz" if codec == "gzip" else ".br" + packages = {} + total = {"exit_code": 0, "wall_seconds": 0.0, "user_seconds": 0.0, "sys_seconds": 0.0, "max_rss_kb": 0, "has_user": False, "has_sys": False, "has_rss": False} + for order, output in cottas_outputs.items(): + artifact = Path(str(output) + suffix) + command = ["gzip", "-c"] if codec == "gzip" else ["brotli", "-q", "7", "-c"] + stage = runner.run(f"cottas-{codec}-{order}", [*command, str(output)], artifact, artifact) + if stage["exit_code"] != 0: + raise RuntimeError(f"{method} packaging failed for {order}") + add_totals(total, stage) + packages[order] = stage + results[method] = { + **packages[cottas_indexes[0]], **finalize_totals(total), + "details": {"index": cottas_indexes[0], "indexes": packages}, + } + continue else: continue if stage["exit_code"] != 0: diff --git a/test/README.md b/test/README.md index 5c8e117..3465a59 100644 --- a/test/README.md +++ b/test/README.md @@ -186,3 +186,20 @@ Example (truncated): Ran 10 tests in 0.90s OK ``` + +## COTTAS index coverage + +Install the pinned container libraries to exercise real Parquet conversion, +all six orders, reindexing, deduplication, packaging dispatch and failure cleanup: + +```bash +python -m pip install pycottas==1.1.0 duckdb==1.5.5 pyarrow==22.0.0 coverage +coverage run --branch -m unittest test.test_cottas_indexes_unit test.test_cottas_tool +coverage report -m --include='*cottas_tool.py,*vcf_rdfizer_cottas.py' +``` + +For Docker end-to-end checks, generate a small fixture with `--cottas-indexes all` +and `--artifact-compression gzip,brotli`, then run `test/cottas_index_roundtrip.py` +inside that image with `--source --cottas +--indexes all --packages`. It checks exact RDF equality, ordering, metadata, +eight native triple-pattern shapes per order, and package contents. diff --git a/test/cottas_index_roundtrip.py b/test/cottas_index_roundtrip.py new file mode 100644 index 0000000..b162872 --- /dev/null +++ b/test/cottas_index_roundtrip.py @@ -0,0 +1,54 @@ +"""Verify a small COTTAS index family against RDF, including native queries. + +Run inside the built Docker image with its /opt/pycottas-venv/bin/python. +""" +import argparse +import gzip +import itertools +import json +import subprocess +import tempfile +from pathlib import Path + +import pyarrow.parquet as pq +import pycottas +from rdflib import Graph + +from vcf_rdfizer_cottas import cottas_index_paths, parse_cottas_indexes + + +def main(): + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument('--source', type=Path, required=True) + parser.add_argument('--cottas', type=Path, required=True) + parser.add_argument('--indexes', type=parse_cottas_indexes, default='all') + parser.add_argument('--packages', action='store_true') + args = parser.parse_args() + opener = gzip.open if args.source.name.endswith('.gz') else open + with opener(args.source, 'rt') as handle: + expected = set(Graph().parse(data=handle.read(), format='nt')) + report = {} + for order, path in cottas_index_paths(args.cottas, args.indexes).items(): + table = pq.read_table(path) + assert table.schema.metadata[b'index'].decode() == order, path + rows = list(zip(*(table[name].to_pylist() for name in 'spo'))) + assert rows == sorted(set(rows), key=lambda row: tuple(row['spo'.index(c)] for c in order)), path + with tempfile.TemporaryDirectory() as td: + decoded = Path(td) / 'decoded.nt' + pycottas.cottas2rdf(str(path), str(decoded)) + assert set(Graph().parse(decoded, format='nt')) == expected, path + # Exercise each possible bound/unbound triple-pattern shape. + probe = rows[0] + for bound in itertools.product((False, True), repeat=3): + pattern = ' '.join(probe[i] if bound[i] else '?'+name for i, name in enumerate('spo')) + wanted = {row for row in rows if all(not bound[i] or row[i] == probe[i] for i in range(3))} + assert set(pycottas.search(str(path), pattern)) == wanted, (path, pattern) + if args.packages: + assert gzip.decompress(Path(str(path)+'.gz').read_bytes()) == path.read_bytes(), path + assert subprocess.check_output(['brotli', '-d', '-c', str(path)+'.br']) == path.read_bytes(), path + report[order] = {'triples': len(rows), 'bytes': path.stat().st_size, 'native_query_shapes': 8} + print(json.dumps(report, indent=2)) + + +if __name__ == '__main__': + main() diff --git a/test/test_cottas_indexes_unit.py b/test/test_cottas_indexes_unit.py new file mode 100644 index 0000000..08bcdf8 --- /dev/null +++ b/test/test_cottas_indexes_unit.py @@ -0,0 +1,382 @@ +"""Index selection, real Parquet ordering, and the chunked pipeline contract.""" +import contextlib +import io +import json +import os +import sys +import tempfile +import unittest +from pathlib import Path +from unittest import mock + +import vcf_rdfizer as host +from vcf_rdfizer_cottas import COTTAS_INDEXES, cottas_index_paths, parse_cottas_indexes +from test.test_cottas_tool import load_cottas_tool +from test.test_partitioned_compression_unit import load_runner_module + +try: + import pyarrow as pa + import pyarrow.parquet as pq + import pycottas +except ImportError: + pa = None + + +class IndexSelectionTests(unittest.TestCase): + def test_selection_and_names(self): + self.assertEqual(parse_cottas_indexes(' SPO, pso,spo,POS '), ('spo', 'pso', 'pos')) + self.assertEqual(parse_cottas_indexes(' ALL '), COTTAS_INDEXES) + for value in ('', 'spo,', ',spo', 'sp', 'sspo', 'spog', 'all,spo', 'spp', 'sp o'): + with self.subTest(value=value), self.assertRaises(ValueError): + parse_cottas_indexes(value) + self.assertEqual(cottas_index_paths(Path('a.cottas'), ('pos', 'spo')), + {'pos': Path('a.cottas'), 'spo': Path('a.spo.cottas')}) + + def test_planned_outputs_include_every_index_and_package(self): + paths = host.planned_output_paths(out_dir=Path('/out'), output_name='a', rdf_name=None, + methods=['cottas', 'cottas_gzip', 'cottas_brotli'], partitioned=True, + cottas_indexes=('pos', 'pso')) + for name in ('a.cottas', 'a.pso.cottas'): + for suffix in ('', '.gz', '.br'): + self.assertIn(Path('/out/a') / (name + suffix), paths) + with tempfile.TemporaryDirectory() as td: + collision = Path(td) / 'a.pso.cottas' + collision.touch() + with self.assertRaises(ValueError): + host.validate_no_output_collisions({'indexes': {collision}}) + + def test_merge_commands_forward_orders(self): + runner = load_runner_module() + for order in COTTAS_INDEXES: + self.assertEqual(runner.cottas_merge_command('python', Path('a'), Path('b'), Path('c'), order)[-1], order) + self.assertEqual(runner.cottas_merge_many_command('python', [Path('a'), Path('b')], Path('c'), index=order)[-1], order) + + def test_cli_rejects_invalid_and_irrelevant_indexes(self): + with contextlib.redirect_stderr(io.StringIO()), contextlib.redirect_stdout(io.StringIO()): + with mock.patch.object(sys, 'argv', ['vcf-rdfizer', '--cottas-indexes', 'spp']), self.assertRaises(SystemExit): + host.main() + for mode in ('tsv', 'decompress', 'validation', 'index'): + with mock.patch.object(sys, 'argv', ['vcf-rdfizer', '--mode', mode, '--cottas-indexes', 'pos', '-o', '/tmp/unused']): + self.assertEqual(host.main(), 2) + + def test_cli_requires_selected_cottas_representation(self): + with tempfile.TemporaryDirectory() as td: + root = Path(td) + source = root / 'input.nt' + source.write_text(' .\n') + vcf = root / 'input.vcf' + vcf.write_text('##fileformat=VCFv4.3\n#CHROM\tPOS\tID\tREF\tALT\tQUAL\tFILTER\tINFO\n') + for mode, option, path in [('full','-i',vcf), ('compress','--rdf',source)]: + argv = ['vcf-rdfizer', '-m', mode, option, str(path), '-o', str(root/'out'), '--cottas-indexes', 'pos', '--representations', 'none'] + with self.subTest(mode=mode), mock.patch.object(sys, 'argv', argv), mock.patch.object(host, 'check_docker') as docker, contextlib.redirect_stderr(io.StringIO()), contextlib.redirect_stdout(io.StringIO()): + self.assertEqual(host.main(), 2) + docker.assert_not_called() + + +@unittest.skipIf(pa is None, 'install pycottas and pyarrow for real Parquet tests') +class CottasIndexTests(unittest.TestCase): + def setUp(self): + self.temporary = tempfile.TemporaryDirectory() + self.addCleanup(self.temporary.cleanup) + self.root = Path(self.temporary.name).resolve() + self.tool = load_cottas_tool() + self.env = mock.patch.dict(os.environ, {'COTTAS_SCRATCH_DIR': str(self.root / 'scratch'), 'COTTAS_MERGE_BATCH_ROWS': '2'}) + self.env.start() + self.addCleanup(self.env.stop) + self.rows = [('', '', '"é"'), ('', '', ''), + ('', '', '"escaped \\"quote\\""'), ('', '', '"1"^^')] + + def write(self, path, rows, order='spo', metadata=True): + fields = {name: pa.array([row[i] for row in rows], type=pa.string()) for i, name in enumerate('spo')} + table = pa.table(fields) + if metadata: + table = table.replace_schema_metadata({b'index': order.encode()}) + pq.write_table(table, path) + + def sorted(self, rows, order): + return sorted(rows, key=lambda row: tuple(row['spo'.index(c)] for c in order)) + + def assert_graph(self, path, order, rows=None): + table = pq.read_table(path) + self.assertEqual(table.schema.metadata[b'index'], order.encode()) + actual = list(zip(*(table[name].to_pylist() for name in 'spo'))) + self.assertEqual(actual, self.sorted(set(self.rows if rows is None else rows), order)) + + def main(self, *args): + with mock.patch.object(sys, 'argv', ['cottas_tool.py', *map(str, args)]): + return self.tool.main() + + def test_all_conversion_indexes_parse_rdf_once(self): + source = self.root / 'input.nt' + source.write_text('\n'.join(' '.join(row) + ' .' for row in self.rows + self.rows[:1]) + '\n') + output = self.root / 'out.cottas' + with mock.patch.object(pycottas, 'rdf2cottas', wraps=pycottas.rdf2cottas) as convert: + self.assertEqual(self.main('convert', source, output, 'ALL'), 0) + self.assertEqual(convert.call_count, 1) + for order, path in cottas_index_paths(output, COTTAS_INDEXES).items(): + self.assert_graph(path, order) + self.assertEqual(list((self.root / 'scratch').iterdir()), []) + + def test_each_single_conversion_order(self): + source = self.root / 'input.nt' + source.write_text(' "value" .\n .\n') + rows = [('', '', '"value"'), ('', '', '')] + for order in COTTAS_INDEXES: + output = self.root / f'{order}.cottas' + self.assertEqual(self.main('convert', source, output, order.upper()), 0) + self.assert_graph(output, order, rows) + + def test_merge_all_orders_batches_duplicates_and_progress(self): + self.tool.COTTAS_OUTPUT_BATCH_ROWS = 2 + self.tool.COTTAS_MERGE_PROGRESS_ROWS = 2 + for order in COTTAS_INDEXES: + inputs = [self.root / f'{order}-{n}.cottas' for n in range(3)] + for path, rows in zip(inputs, (self.rows[:3], self.rows[1:], [])): + self.write(path, self.sorted(rows + rows[:1], order), order) + output = self.root / f'{order}.cottas' + progress = self.root / f'{order}.jsonl' + self.tool.streaming_cottas_merge(list(map(str, inputs)), str(output), index=order.upper(), remove_input_files=True, progress_path=progress) + self.assert_graph(output, order) + self.assertFalse(any(path.exists() for path in inputs)) + events = [json.loads(line) for line in progress.read_text().splitlines()] + self.assertEqual(events[-1]['phase'], 'complete') + self.assertIn('merging', [event['phase'] for event in events]) + + def test_merge_rejects_bad_inputs_and_preserves_output(self): + good = self.root / 'good.cottas' + bad = self.root / 'bad.cottas' + output = self.root / 'output.cottas' + self.write(good, self.sorted(self.rows, 'spo')) + for defect in ('missing_index', 'wrong_index', 'unsorted', 'null', 'missing_column', 'type'): + with self.subTest(defect=defect): + if defect == 'missing_column': + pq.write_table(pa.table({'s': ['a'], 'p': ['b']}).replace_schema_metadata({b'index': b'spo'}), bad) + elif defect == 'type': + pq.write_table(pa.table({'s': [1], 'p': [2], 'o': [3]}).replace_schema_metadata({b'index': b'spo'}), bad) + else: + rows = [(None, '', '')] if defect == 'null' else self.sorted(self.rows, 'spo') + self.write(bad, rows[::-1] if defect == 'unsorted' else rows, + 'pos' if defect == 'wrong_index' else 'spo', metadata=defect != 'missing_index') + output.write_bytes(b'original') + with self.assertRaises(RuntimeError): + self.tool.streaming_cottas_merge([str(good), str(bad)], str(output), index='spo', remove_input_files=True) + self.assertEqual(output.read_bytes(), b'original') + self.assertTrue(good.exists() and bad.exists()) + self.assertFalse(list(self.root.glob('.output.cottas.merge-*'))) + + def test_merge_flushes_final_partial_batch(self): + self.tool.COTTAS_OUTPUT_BATCH_ROWS = 3 + source = self.root / 'source.cottas' + self.write(source, self.sorted(self.rows, 'spo')) + output = self.root / 'output.cottas' + self.tool.streaming_cottas_merge([str(source)], str(output), index='spo', remove_input_files=False) + self.assert_graph(output, 'spo') + self.assertEqual(pq.ParquetFile(output).metadata.num_row_groups, 2) + + def test_empty_batches_and_metadata_variants(self): + from types import SimpleNamespace + parquet = SimpleNamespace(metadata=SimpleNamespace(metadata={'index': 'POS'})) + self.assertEqual(self.tool.cottas_file_index(parquet), 'pos') + schema = pa.schema([(name, pa.string()) for name in 'spo']) + empty = pa.RecordBatch.from_arrays([pa.array([], type=pa.string()) for _ in 'spo'], schema=schema) + row = pa.RecordBatch.from_arrays([pa.array([term]) for term in self.rows[0]], schema=schema) + parquet = SimpleNamespace(schema_arrow=schema, metadata=SimpleNamespace(metadata={b'index': b'spo'}, num_rows=1), iter_batches=lambda **kwargs: iter([empty, row])) + stream = self.tool.CottasTripleStream(SimpleNamespace(ParquetFile=lambda path: parquet), Path('empty-batch'), 'spo', 2) + self.assertEqual(stream.current, self.rows[0]) + stream.close() # compatible with readers that do not expose close() + + def test_empty_merge_and_invalid_configuration(self): + empty = self.root / 'empty.cottas' + self.write(empty, []) + output = self.root / 'out.cottas' + self.tool.streaming_cottas_merge([str(empty)], str(output), index='spo', remove_input_files=False) + self.assert_graph(output, 'spo', []) + for inputs, order in (([], 'spo'), ([str(empty)], ''), ([str(empty)], 'spp')): + with self.assertRaises(ValueError): + self.tool.streaming_cottas_merge(inputs, str(output), index=order, remove_input_files=False) + for value in ('0', '-2', 'abc'): + with mock.patch.dict(os.environ, {'COTTAS_MERGE_BATCH_ROWS': value}), self.assertRaises(ValueError): + self.tool.cottas_merge_batch_rows() + self.tool.emit_merge_progress(self.root / 'absent' / 'progress', 'test', completed=0, total=0, detail='') + with self.tool.cottas_scratch_workspace(), self.assertRaises(ValueError): + self.tool.sort_cottas_chunk(empty, {'bad': output}) + + def test_reindex_all_orders_and_same_order_rewrite(self): + self.tool.COTTAS_OUTPUT_BATCH_ROWS = 2 + self.tool.COTTAS_REINDEX_FAN_IN = 2 + source = self.root / 'source.cottas' + self.write(source, self.rows + self.rows, 'pos') # deliberately unsorted: reindex repairs it + self.assertEqual(self.main('reindex', source, 'all'), 0) + for order, path in cottas_index_paths(source, COTTAS_INDEXES).items(): + self.assert_graph(path, order) + self.assertEqual(self.main('reindex', source, 'spo'), 0) + self.assert_graph(source, 'spo') + self.assertEqual(list((self.root / 'scratch').iterdir()), []) + + def test_reindex_empty_without_metadata(self): + source = self.root / 'empty.cottas' + self.write(source, [], metadata=False) + self.assertEqual(self.main('reindex', source, 'pos,ops'), 0) + for order, path in cottas_index_paths(source, ('pos', 'ops')).items(): + self.assert_graph(path, order, []) + + def test_reindex_failure_and_collision_preserve_original(self): + source = self.root / 'source.cottas' + self.write(source, self.sorted(self.rows, 'spo')) + original = source.read_bytes() + with mock.patch.object(self.tool, 'reindex_cottas', side_effect=RuntimeError('failure')): + with self.assertRaisesRegex(RuntimeError, 'failure'): + self.main('reindex', source, 'pos,pso') + with mock.patch.object(self.tool, 'reindex_cottas'): + with self.assertRaisesRegex(RuntimeError, 'non-empty'): + self.main('reindex', source, 'pos') + sibling = self.root / 'source.pso.cottas' + sibling.write_bytes(b'keep') + with self.assertRaises(FileExistsError): + self.main('reindex', source, 'pos,pso') + self.assertEqual(source.read_bytes(), original) + self.assertEqual(sibling.read_bytes(), b'keep') + self.assertFalse(list(self.root.glob('.*.reindex-*'))) + with contextlib.redirect_stderr(io.StringIO()): + self.assertEqual(self.main('reindex', self.root / 'absent'), 2) + self.assertEqual(self.main('merge-many', '--input-cottas-files', source, '--output-cottas-file', sibling), 2) + + def test_missing_dependencies(self): + with mock.patch.dict(sys.modules, {'pycottas': None}), contextlib.redirect_stderr(io.StringIO()): + self.assertEqual(self.main('convert', 'a', 'b'), 127) + with mock.patch.dict(sys.modules, {'pyarrow': None}), self.assertRaisesRegex(RuntimeError, 'PyArrow'): + self.tool.streaming_cottas_merge(['a'], 'b', index='spo', remove_input_files=False) + + +class PartitionedIndexTests(unittest.TestCase): + def run_pipeline(self, root, *, indexes=('pos', 'spo', 'ops'), chunk_bytes=55, failure=None, allow=False, methods='cottas,cottas_gzip,cottas_brotli'): + runner = load_runner_module() + work, out = root / 'work', root / 'out' + work.mkdir(); out.mkdir() + source = root / 'source.nt' + source.write_text(''.join(f' .\n' for i in range(7))) + result_path = out / 'result.json' + calls = [] + class StageRunner(runner.StageRunner): + def run(self, name, command, output_path=None, *args): + calls.append((name, command)) + failed = failure and failure in name + if 'validate' in name: + output_path.write_text(json.dumps({'valid': not failed, 'count_match': not failed})) + elif not failed and output_path: + output_path.write_bytes(b'artifact') + if 'build' in name: + for path in cottas_index_paths(output_path, parse_cottas_indexes(command[-1])).values(): + path.write_bytes(b'chunk') + stage = {'name': name, 'exit_code': int(bool(failed)), 'wall_seconds': 1, 'output_path': str(output_path), 'output_size_bytes': 8} + self.stages.append(stage) + return stage + argv = ['runner', '--source', str(source), '--output-dir', str(out), '--output-name', 'sample', '--methods', methods, + '--cottas-indexes', ','.join(indexes), '--target-chunk-bytes', str(chunk_bytes), + '--min-chunk-bytes', '1', '--max-chunk-bytes', str(chunk_bytes*2), '--result-path', str(result_path)] + if allow: argv.append('--allow-index-failures') + with mock.patch.object(runner, 'StageRunner', StageRunner), mock.patch.object(runner, 'Path', wraps=Path, side_effect=lambda p: work if p == '/work' else Path(p)), mock.patch.object(sys, 'argv', argv), contextlib.redirect_stderr(io.StringIO()): + code = runner.main() + return code, json.loads(result_path.read_text()), calls, work, out + + def test_multiple_indexes_merge_validate_and_package(self): + for chunk_bytes in (55, 10000): + with self.subTest(chunk_bytes=chunk_bytes), tempfile.TemporaryDirectory() as td: + code, report, calls, work, out = self.run_pipeline(Path(td), chunk_bytes=chunk_bytes) + self.assertEqual(code, 0, report) + details = report['methods']['cottas']['details'] + self.assertEqual(list(details['indexes']), ['pos', 'spo', 'ops']) + self.assertEqual(details['index'], 'pos') + self.assertEqual(details['merge_rounds'], int(chunk_bytes == 55)) + self.assertFalse(list(work.glob('*.cottas'))) + for order, path in cottas_index_paths(out / 'sample.cottas', ('pos','spo','ops')).items(): + for suffix in ('', '.gz', '.br'): + self.assertTrue(Path(str(path)+suffix).exists()) + merges = [command for name, command in calls if name.startswith('cottas-merge')] + self.assertEqual([command[command.index('--index')+1] for command in merges], ['pos','spo','ops'] if chunk_bytes==55 else []) + builds = [command for name, command in calls if name.startswith('cottas-build')] + self.assertEqual(len(builds), details['chunk_count']) + self.assertTrue(all(command[-1] == 'pos,spo,ops' for command in builds)) + + def test_failures_skip_indexes_and_packages_or_fail_strictly(self): + for failure in ('cottas-build', 'cottas-merge-ops', 'cottas-validate-ops'): + for allow in (True, False): + with self.subTest(failure=failure, allow=allow), tempfile.TemporaryDirectory() as td: + code, report, calls, work, out = self.run_pipeline(Path(td), failure=failure, allow=allow) + self.assertEqual(code, 0 if allow else 1, report) + if allow: + self.assertEqual(report['methods']['cottas']['details']['index_status'], 'failed') + self.assertFalse(list(out.glob('*.cottas'))) + self.assertFalse(list(work.glob('*.cottas'))) + self.assertEqual(report['methods']['cottas_gzip']['source'], 'index_unavailable') + self.assertFalse(any(name.startswith('cottas-gzip') for name, _ in calls)) + with tempfile.TemporaryDirectory() as td: + code, report, *_ = self.run_pipeline(Path(td), failure='cottas-gzip-spo') + self.assertEqual(code, 1) + self.assertIn('packaging failed for spo', report['error']) + +class HostIndexRoutingTests(unittest.TestCase): + def test_host_forwards_indexes_and_translates_every_artifact(self): + with tempfile.TemporaryDirectory() as td: + root = Path(td) + source = root / 'source.nt' + source.write_text(' .\n') + out = root / 'out' + def run(command, **kwargs): + if host.PARTITIONED_COMPRESSION_RUNNER_CONTAINER in command: + self.assertEqual(command[command.index('--cottas-indexes')+1], 'pos,pso') + indexes = {} + for order, name in [('pos','sample.cottas'), ('pso','sample.pso.cottas')]: + (out / name).write_bytes(b'cottas') + indexes[order] = {'output_path': '/data/out/'+name, 'output_size_bytes': 6} + payload = {'exit_code': 0, 'methods': {'cottas': {'details': {'index': 'pos', 'indexes': indexes}}}} + (out / '.sample.partitioned-results.json').write_text(json.dumps(payload)) + return 0 + with mock.patch.object(host, 'run', side_effect=run): + ok, results = host.run_partitioned_representation_methods_for_rdf_files( + rdf_paths=[source], out_dir=out, image_ref='test', methods=['cottas'], + wrapper_log_path=root/'log', output_name='sample', target_chunk_bytes=100, + min_chunk_bytes=1, max_chunk_bytes=200, cottas_indexes=('pos','pso')) + self.assertTrue(ok) + for item in results['cottas']['details']['indexes'].values(): + self.assertTrue(Path(item['output_path']).is_file()) + self.assertEqual(Path(item['output_path']).parent, out) + + def test_index_mode_records_requested_indexes(self): + with tempfile.TemporaryDirectory() as td: + root = Path(td) + source = root / 'source.cottas' + source.write_bytes(b'old') + def run(command, **kwargs): + self.assertIn('pos,pso', command[-1]) + (root / 'source.pso.cottas').write_bytes(b'new') + return 0 + with mock.patch.object(host, 'run', side_effect=run), contextlib.redirect_stdout(io.StringIO()): + self.assertEqual(host.run_index_mode(index_path=source, index_format='cottas', + metrics_dir=root/'metrics', image_ref='test', wrapper_log_path=root/'log', + cottas_indexes=('pos','pso')), 0) + payload = json.loads((root/'metrics/stages/index/cottas-source.cottas.json').read_text()) + self.assertEqual(payload['indexes'], {'pos': str(source), 'pso': str(root/'source.pso.cottas')}) + + def test_failed_reindex_reports_failure_even_when_original_survives(self): + with tempfile.TemporaryDirectory() as td: + root = Path(td) + source = root / 'source.cottas' + source.write_bytes(b'original') + for code in (0, 1): + with self.subTest(code=code), mock.patch.object(host, 'run', return_value=code), contextlib.redirect_stdout(io.StringIO()), contextlib.redirect_stderr(io.StringIO()): + self.assertEqual(host.run_index_mode(index_path=source, index_format='cottas', + metrics_dir=root/'metrics', image_ref='test', wrapper_log_path=root/'log', + cottas_indexes=('pos','pso')), 1) + payload = json.loads((root/'metrics/stages/index/cottas-source.cottas.json').read_text()) + self.assertEqual(payload['index_status'], 'failed') + self.assertEqual(source.read_bytes(), b'original') + + +class AdapterEntryPointTests(unittest.TestCase): + def test_script_entry_point_returns_dependency_exit_code(self): + import runpy + from test.test_cottas_tool import COTTAS_TOOL_PATH + with mock.patch.dict(sys.modules, {'pycottas': None}), mock.patch.object(sys, 'argv', ['cottas_tool.py', 'convert', 'a', 'b']), contextlib.redirect_stderr(io.StringIO()), self.assertRaises(SystemExit) as raised: + runpy.run_path(str(COTTAS_TOOL_PATH), run_name='__main__') + self.assertEqual(raised.exception.code, 127) diff --git a/test/test_cottas_tool.py b/test/test_cottas_tool.py index 81f3c82..e4b73d9 100644 --- a/test/test_cottas_tool.py +++ b/test/test_cottas_tool.py @@ -201,8 +201,9 @@ def test_reindex_rewrites_atomically_without_removing_input(self): module = load_cottas_tool() calls = [] - def fake_streaming_merge(paths, cottas_path, *, index, remove_input_files, progress_path=None): - calls.append((paths, cottas_path, index, remove_input_files)) + def fake_reindex(source_path, outputs): + index, cottas_path = next(iter(outputs.items())) + calls.append(([str(source_path)], str(cottas_path), index, False)) self.assertNotEqual(Path(cottas_path), source) self.assertTrue(Path(cottas_path).name.startswith(f".{source.name}.reindex-")) Path(cottas_path).write_text("reindexed COTTAS\n") @@ -219,7 +220,7 @@ def fake_streaming_merge(paths, cottas_path, *, index, remove_input_files, progr ), mock.patch.dict( os.environ, {"COTTAS_SCRATCH_DIR": str(scratch_root)}, clear=False ), mock.patch.object( - module, "streaming_cottas_merge", side_effect=fake_streaming_merge + module, "reindex_cottas", side_effect=fake_reindex ), mock.patch.object( sys, "argv", @@ -253,7 +254,7 @@ def failing_streaming_merge(*args, **kwargs): ), mock.patch.dict( os.environ, {"COTTAS_SCRATCH_DIR": str(scratch_root)}, clear=False ), mock.patch.object( - module, "streaming_cottas_merge", side_effect=failing_streaming_merge + module, "reindex_cottas", side_effect=failing_streaming_merge ), mock.patch.object( sys, "argv", diff --git a/vcf_rdfizer.py b/vcf_rdfizer.py index cb45f3d..ac85c69 100644 --- a/vcf_rdfizer.py +++ b/vcf_rdfizer.py @@ -74,6 +74,8 @@ except ImportError: # pragma: no cover - shipped alongside this module vcf_rdfizer_gzip = None +from vcf_rdfizer_cottas import cottas_index_paths, parse_cottas_indexes + import vcf_rdfizer_vocab as vocab from vcf_rdfizer_vocab import ( CONTIG_ATTRIBUTES, @@ -1599,6 +1601,7 @@ def planned_output_paths( rdf_name: str | None, methods: list[str], partitioned: bool, + cottas_indexes: tuple[str, ...] = ("spo",), ) -> set[Path]: """List final output paths that a compression plan would create.""" target_dir = out_dir / output_name @@ -1616,7 +1619,12 @@ def planned_output_paths( # The pinned Java-free indexer generates this canonical sidecar. planned.add(target_dir / f"{output_name}.hdt.index.v1-1") if any(method in COTTAS_COMPRESSION_METHODS for method in methods): - planned.add(target_dir / f"{output_name}.cottas") + for path in cottas_index_paths(target_dir / f"{output_name}.cottas", cottas_indexes).values(): + planned.add(path) + if "cottas_gzip" in methods: + planned.add(Path(str(path) + ".gz")) + if "cottas_brotli" in methods: + planned.add(Path(str(path) + ".br")) if partitioned: planned.add(target_dir / f".{safe_metrics_name(output_name)}.partitioned-results.json") return planned @@ -5815,6 +5823,8 @@ def timing_payload(result: dict): }, "cottas_conversion": { "output_cottas_path": cottas_result.get("output_path", ""), + "index": cottas_result.get("details", {}).get("index", "spo"), + "indexes": cottas_result.get("details", {}).get("indexes", {}), "output_cottas_size_bytes": int(cottas_result.get("output_size_bytes") or 0), "exit_code": int(cottas_result.get("exit_code") or 0), "timing": timing_payload(cottas_result), @@ -5843,12 +5853,14 @@ def timing_payload(result: dict): }, "gzip_on_cottas": { "output_cottas_gz_path": cottas_gzip_result.get("output_path", ""), + "indexes": cottas_gzip_result.get("details", {}).get("indexes", {}), "output_cottas_gz_size_bytes": int(cottas_gzip_result.get("output_size_bytes") or 0), "exit_code": int(cottas_gzip_result.get("exit_code") or 0), "timing": timing_payload(cottas_gzip_result), }, "brotli_on_cottas": { "output_cottas_br_path": cottas_brotli_result.get("output_path", ""), + "indexes": cottas_brotli_result.get("details", {}).get("indexes", {}), "output_cottas_br_size_bytes": int(cottas_brotli_result.get("output_size_bytes") or 0), "exit_code": int(cottas_brotli_result.get("exit_code") or 0), "timing": timing_payload(cottas_brotli_result), @@ -6924,6 +6936,7 @@ def run_containerized_partitioned_representation_methods( max_chunk_bytes: int, expected_triples: int | None = None, index_warnings: list[dict] | None = None, + cottas_indexes: tuple[str, ...] = ("spo",), ): """Run partitioned compression in an ephemeral Docker-managed volume. @@ -7013,6 +7026,8 @@ def run_containerized_partitioned_representation_methods( output_name, "--methods", ",".join(methods), + "--cottas-indexes", + ",".join(cottas_indexes), "--target-chunk-bytes", str(target_chunk_bytes), "--min-chunk-bytes", @@ -7065,6 +7080,8 @@ def run_containerized_partitioned_representation_methods( result["output_path"] = str(artifact_path) result["output_size_bytes"] = int(file_size_bytes(artifact_path) or 0) details = result.setdefault("details", {}) + for indexed in details.get("indexes", {}).values(): + indexed["output_path"] = str(out_dir / Path(indexed["output_path"]).name) if method == "hdt": index_path = find_hdt_index_sidecar(artifact_path) details["index_path"] = str(index_path) if index_path else "" @@ -7152,6 +7169,7 @@ def run_partitioned_representation_methods_for_rdf_files( max_chunk_bytes: int, expected_triples: int | None = None, index_warnings: list[dict] | None = None, + cottas_indexes: tuple[str, ...] = ("spo",), ): """Dispatch aggregate RDF to the ephemeral container pipeline. @@ -7170,6 +7188,7 @@ def run_partitioned_representation_methods_for_rdf_files( return False, {} source_rdf_path = rdf_paths[0] return run_containerized_partitioned_representation_methods( + cottas_indexes=cottas_indexes, source_rdf_path=source_rdf_path, out_dir=out_dir, image_ref=image_ref, @@ -7260,6 +7279,7 @@ def run_full_mode( linking_manifests: list | None = None, linking_options: dict | None = None, allow_cohort_expansion: bool = False, + cottas_indexes: tuple[str, ...] = ("spo",), ): """Execute full pipeline: per-input TSV -> RDF -> compression -> validation.""" linking_options = dict(linking_options or {}) @@ -7841,6 +7861,7 @@ def fail_current(stage: str, message: str): if not input_failed and use_partitioned_compression and partitioned_methods: ok, partitioned_representation_results = run_partitioned_representation_methods_for_rdf_files( + cottas_indexes=cottas_indexes, rdf_paths=[], source_rdf_path=raw_rdf_files[0], out_dir=out_dir / output_name, @@ -8284,6 +8305,7 @@ def run_compress_mode( chunk_min_bytes: int, chunk_max_bytes: int, wrapper_log_path: Path, + cottas_indexes: tuple[str, ...] = ("spo",), ): """Execute compression-only mode for a designated RDF file.""" print("Step 3/3: Compressing RDF input") @@ -8343,6 +8365,7 @@ def run_compress_mode( return 1 method_results.update(raw_method_results) ok, partitioned_results = run_partitioned_representation_methods_for_rdf_files( + cottas_indexes=cottas_indexes, rdf_paths=[], source_rdf_path=rdf_path, out_dir=out_dir / input_stem, @@ -8523,6 +8546,7 @@ def run_index_mode( wrapper_log_path: Path, run_id: str | None = None, timestamp: str | None = None, + cottas_indexes: tuple[str, ...] = ("spo",), ): """Generate or regenerate the query index for one existing artifact. @@ -8551,7 +8575,7 @@ def run_index_mode( 'if [[ -z "$PYTHON_BIN" || ! -x "$PYTHON_BIN" ]]; then ' 'echo "Missing pycottas Python executable in container" >&2; exit 127; fi; ' f'"$PYTHON_BIN" {shlex.quote(COTTAS_TOOL_CONTAINER)} reindex ' - f"{shlex.quote(source_container)} spo" + f"{shlex.quote(source_container)} {shlex.quote(','.join(cottas_indexes))}" ) safe_input = safe_metrics_name(index_path.name) timing_host = metrics_dir / "timings" / "index" / f"{safe_input}.txt" @@ -8591,6 +8615,10 @@ def run_index_mode( else (index_path if file_size_bytes(index_path) else None) ) index_ready = index_path_after is not None + if index_format == "cottas": + index_ready = index_ready and all( + file_size_bytes(path) for path in cottas_index_paths(index_path, cottas_indexes).values() + ) final_code = int(exit_code) if int(exit_code) != 0 else (0 if index_ready else 1) index_was_present = existing_index_path is not None or index_format == "cottas" payload = { @@ -8611,7 +8639,7 @@ def run_index_mode( }, "index_status": ( "regenerated" if index_was_present else "generated" - ) if index_ready else "failed", + ) if final_code == 0 else "failed", "index_size_bytes": file_size_bytes(index_path_after) if index_path_after else 0, } if index_format == "hdt": @@ -8619,6 +8647,9 @@ def run_index_mode( payload["hdt_path"] = str(index_path) else: payload["cottas_path"] = str(index_path) + payload["indexes"] = { + order: str(path) for order, path in cottas_index_paths(index_path, cottas_indexes).items() + } metrics_path = metrics_dir / "stages" / "index" / f"{index_format}-{safe_input}.json" metrics_path.parent.mkdir(parents=True, exist_ok=True) metrics_path.write_text( @@ -9294,6 +9325,10 @@ def main(): f"(default: {DEFAULT_REPRESENTATIONS})" ), ) + parser.add_argument( + "--cottas-indexes", type=parse_cottas_indexes, default="spo", + help="COTTAS orders: spo,sop,pso,pos,osp,ops (comma-separated), or all; first is primary (default: spo)", + ) parser.add_argument( "--artifact-compression", default=DEFAULT_ARTIFACT_COMPRESSION, @@ -9532,6 +9567,12 @@ def main(): from vcf_rdfizer_link import add_link_arguments, run_options, run_posthoc, selected_linkers add_link_arguments(parser) args = parser.parse_args() + if args.cottas_indexes != ("spo",) and ( + args.mode not in {"full", "compress", "index"} + or (args.mode == "index" and not args.cottas) + ): + eprint("Error: --cottas-indexes requires COTTAS conversion or --mode index --cottas") + return 2 if args.mode not in {"full", "link"} and (args.link or args.linker_path or args.offline or args.links_cache_only or args.assembly or args.links_contact_email): eprint("Error: linking options require --mode full or --mode link") return 2 @@ -9724,6 +9765,8 @@ def main(): representations=args.representations, artifact_compression=args.artifact_compression, ) + if args.cottas_indexes != ("spo",) and not any(method in COTTAS_COMPRESSION_METHODS for method in full_methods): + raise ValueError("--cottas-indexes requires --representations cottas") hdt_strategy_error = hdt_strategy_rejection( hdt_strategy=args.hdt_strategy, methods=full_methods, @@ -9760,6 +9803,7 @@ def main(): rdf_name=rdf_name, methods=full_methods, partitioned=full_uses_partitioning, + cottas_indexes=args.cottas_indexes, ) if linking_manifests: output_plans[f"input {index} ({prefix})"].add( @@ -9836,6 +9880,8 @@ def main(): representations=args.representations, artifact_compression=args.artifact_compression, ) + if args.cottas_indexes != ("spo",) and not any(method in COTTAS_COMPRESSION_METHODS for method in methods): + raise ValueError("--cottas-indexes requires --representations cottas") hdt_strategy_error = hdt_strategy_rejection( hdt_strategy=args.hdt_strategy, methods=methods, @@ -9865,6 +9911,7 @@ def main(): rdf_name=None, methods=methods, partitioned=compression_uses_partitioning_for_input, + cottas_indexes=args.cottas_indexes, ) } ) @@ -9889,6 +9936,10 @@ def main(): raise ValueError( f"{index_format.upper()} input file not found: {index_path}" ) + if index_format == "cottas": + validate_no_output_collisions({ + "COTTAS indexes": set(list(cottas_index_paths(index_path, args.cottas_indexes).values())[1:]) + }) validate_mode_dirs([out_root, out_dir, metrics_root]) else: if args.spark_partitions is not None: @@ -10172,6 +10223,7 @@ def execute_mode(): if mode == "full": # Full-mode orchestrates conversion + compression pipeline. return run_full_mode( + cottas_indexes=args.cottas_indexes, input_mount_dir=input_mount_dir, container_inputs=container_inputs, expected_prefixes=expected_prefixes, @@ -10228,6 +10280,7 @@ def execute_mode(): if mode == "compress": # Compression-only mode. return run_compress_mode( + cottas_indexes=args.cottas_indexes, rdf_path=rdf_path, out_dir=out_dir, metrics_dir=metrics_dir, @@ -10266,6 +10319,7 @@ def execute_mode(): ) if mode == "index": return run_index_mode( + cottas_indexes=args.cottas_indexes, index_path=index_path, index_format=index_format, metrics_dir=metrics_dir, diff --git a/vcf_rdfizer_cottas.py b/vcf_rdfizer_cottas.py new file mode 100644 index 0000000..545abe0 --- /dev/null +++ b/vcf_rdfizer_cottas.py @@ -0,0 +1,24 @@ +"""COTTAS index selection and artifact naming shared by host and container.""" + +from pathlib import Path + + +COTTAS_INDEXES = ("spo", "sop", "pso", "pos", "osp", "ops") + + +def parse_cottas_indexes(value: str) -> tuple[str, ...]: + """Normalize an ordered selection, rejecting typos before conversion.""" + if value.strip().lower() == "all": + return COTTAS_INDEXES + indexes = tuple(dict.fromkeys(part.strip().lower() for part in value.split(","))) + if any(index not in COTTAS_INDEXES for index in indexes): + raise ValueError("COTTAS indexes must be spo,sop,pso,pos,osp,ops or all") + return indexes + + +def cottas_index_paths(path: Path, indexes: tuple[str, ...]) -> dict[str, Path]: + """Keep the primary filename; suffix additional copies with their order.""" + return { + index: path if number == 0 else path.with_suffix(f".{index}.cottas") + for number, index in enumerate(indexes) + } From 0dc836a46814fde6b2c7908da2d853270a8b3668 Mon Sep 17 00:00:00 2001 From: ecrum19 Date: Fri, 25 Sep 2026 10:05:29 +0200 Subject: [PATCH 2/2] Support graph-aware COTTAS dataset indexes --- README.md | 7 +- docs/representations.md | 35 +++- src/cottas_tool.py | 195 ++++++++++++++++---- src/partitioned_compression.py | 12 +- src/validate_compression.py | 18 +- test/README.md | 8 +- test/cottas_index_roundtrip.py | 52 ++++-- test/fixtures/cottas-dataset.nq | 11 ++ test/test_cottas_indexes_unit.py | 16 +- test/test_cottas_quads_unit.py | 238 +++++++++++++++++++++++++ test/test_cottas_tool.py | 2 +- test/test_validate_compression_unit.py | 2 +- vcf_rdfizer.py | 40 +++-- vcf_rdfizer_cottas.py | 15 +- 14 files changed, 551 insertions(+), 100 deletions(-) create mode 100644 test/fixtures/cottas-dataset.nq create mode 100644 test/test_cottas_quads_unit.py diff --git a/README.md b/README.md index e5d0366..0151555 100644 --- a/README.md +++ b/README.md @@ -269,7 +269,7 @@ host filesystem. - `space-optimized`: gzip each part into one `.nt.gz` aggregate and delete the source part immediately - `--rdf-compression {gzip,brotli,none}` raw RDF artifacts to retain - `--representations {hdt,cottas,none}` queryable primary representations -- `--cottas-indexes` one or more of `spo,sop,pso,pos,osp,ops`, or `all` (default `spo`) +- `--cottas-indexes` permutations of `spo` or `spog`, comma-separated; `all` selects six triple orders, `all-quads` selects 24 dataset orders (default `spo`) - `--artifact-compression {gzip,brotli,none}` optional packaging applied to each selected representation - `--hdt-strategy {auto,partitioned,single}` - `auto`: in full mode, build smaller HDT chunks and merge them with native `hdtc` @@ -647,6 +647,11 @@ when memory is available. Temporary files live in the container's `/work` area and are removed after the attempt. Choose COTTAS orders with `--cottas-indexes spo,pso,pos` or `--cottas-indexes all`. +Datasets (`--mode compress --rdf dataset.nq` or `.nq.gz`) use graph-aware orders, +for example `--representations cottas --cottas-indexes spog,gspo`, or +`--cottas-indexes all-quads` for all 24 permutations. Named and default graphs +are preserved; triple-only orders and HDT are rejected for dataset inputs. +Graph-aware orders also work on VCF/triple input, using the default graph. The first order keeps `sample.cottas`; additional whole-graph copies are named `sample.pso.cottas`, `sample.pos.cottas`, etc. Query one copy at a time. Each RDF chunk is parsed once; extra orders sort its deduplicated Parquet data. Each order diff --git a/docs/representations.md b/docs/representations.md index 734f2b5..3134c49 100644 --- a/docs/representations.md +++ b/docs/representations.md @@ -42,7 +42,7 @@ queries must run without a decompression step. ## 3. Record-safe chunking When HDT or COTTAS is selected, the aggregate is read sequentially and split -into chunks on **complete N-Triples line boundaries**. Only one uncompressed +into chunks on **complete N-Triples or N-Quads line boundaries**. Only one uncompressed chunk exists at a time: it is consumed by both converters and removed before the next is read. That property is what makes a `space-optimized` `.nt.gz` aggregate usable without ever expanding a second full raw copy. @@ -122,7 +122,9 @@ there is no sidecar. Chunk conversion uses `pycottas.rdf2cottas(..., disk=True)` with a fresh container-local DuckDB workspace per operation. `--cottas-indexes` accepts `spo` (default), `sop`, `pso`, `pos`, `osp`, `ops`, a comma-separated selection, or `all`; case is -ignored and repeated orders are built once. The first order uses `name.cottas`; +ignored and repeated orders are built once. Dataset orders accept all 24 +permutations of `spog` (such as `spog,gspo,pgos`), or `all-quads` for all 24. +The first order uses `name.cottas`; additional orders use `name..cottas`. Each file contains the whole graph in one order. Query one copy at a time; loading them together repeats the graph. @@ -133,12 +135,24 @@ multiplying merge memory. Gzip/Brotli packaging and round-trip checks cover each copy. Per-index paths, sizes and validation are in JSON `details.indexes`; the existing CSV size columns describe the primary copy, with timings for all orders. +`--mode compress --rdf dataset.nq` (also `.nq.gz`) preserves named and default +graphs. Every selected order must include `g`; HDT and triple-only orders are +rejected for datasets. N-Quads uses a streaming RDFLib parser with bounded +Parquet batches because the pinned pycottas parser renames blank nodes per +chunk. This preserves shared blank nodes, literal lexical forms and graph +identity across chunks. Default graphs are stored as NULL, sorted last. +Graph-aware indexes on VCF/N-Triples input add the default graph. Reindexing +also accepts dataset orders and normalizes legacy pycottas `DEFAULT` values. +When using pycottas 1.1.0 directly, pass a four-term RDFLib tuple to `search` +for quad patterns; its string-pattern parser ignores the fourth term. + The final merge deliberately does **not** call `pycottas.cat`. In version 1.1.0 that runs a global `DISTINCT` plus `ORDER BY` through an unbounded in-memory DuckDB connection, which gets OOM-killed on large condensed graphs. VCF-RDFizer instead performs a **k-way PyArrow merge** of the already index-sorted Parquet chunks: it holds at most `COTTAS_MERGE_BATCH_ROWS` (default 2048) rows per -input, drops adjacent duplicate triples, and writes the result incrementally. +input, drops adjacent duplicate triples/quads, and writes the result incrementally. +The same triple in different graphs remains a separate quad in each graph. There is no graph-wide hash table and no external-sort spill directory; memory is a function of the batch size and the chunk count, not of the graph size. @@ -179,10 +193,13 @@ The same operation runs automatically after each partitioned HDT merge. ## 8. Decompression -`--mode decompress` decodes `.nt.gz`, `.nt.br`, `.hdt`, `.cottas`, `.cottas.gz` -and `.cottas.br` back to N-Triples. A packaged COTTAS is unwrapped **inside the -container** before `pycottas` writes the decoded output, so the intermediate -unwrapped file never appears on the host. +`--mode decompress` decodes raw RDF archives and HDT/COTTAS artifacts. For a +COTTAS dataset, specify `--decompress-out /dataset.nq` to preserve named graphs; +the default `.nt` output is accepted only for triple/default-graph data. +The streaming quad decoder omits the graph term for default-graph statements. +A packaged COTTAS is unwrapped **inside the container**, so the intermediate +unwrapped file never appears on the host. Raw `.nq.gz`/`.nq.br` archives keep +their `.nq` extension when decompressed. ## 9. Choosing @@ -211,8 +228,8 @@ unwrapped file never appears on the host. - **Packaged representations are not queryable**, which is easy to forget when `--artifact-compression` is set and `--remove-rdf-storage-output` has removed the alternative. -- **Neither HDT nor COTTAS carries named graphs**, matching the conversion's - triples-only output. +- **HDT cannot carry named graphs.** COTTAS preserves them with dataset indexes; + VCF conversion itself still produces triples in the default graph. - **No incremental update.** Adding variants means reconverting and rebuilding the representation from scratch. diff --git a/src/cottas_tool.py b/src/cottas_tool.py index 49b9168..0ab8991 100644 --- a/src/cottas_tool.py +++ b/src/cottas_tool.py @@ -7,10 +7,12 @@ import os import sys import tempfile -from contextlib import contextmanager +from contextlib import contextmanager, nullcontext from pathlib import Path -from vcf_rdfizer_cottas import COTTAS_INDEXES, cottas_index_paths, parse_cottas_indexes +from vcf_rdfizer_cottas import ( + COTTAS_ALL_INDEXES, cottas_index_paths, parse_cottas_indexes, require_dataset_indexes, +) DEFAULT_COTTAS_MERGE_BATCH_ROWS = 2048 @@ -19,21 +21,118 @@ COTTAS_REINDEX_FAN_IN = 128 -def sort_cottas_chunk(source: Path, outputs: dict[str, Path]) -> None: +def parse_nquads_chunk(source: str, output: Path) -> None: + """Stream N-Quads without pycottas 1.1.0's per-chunk blank-node renaming.""" + import pyarrow as pa + import pyarrow.parquet as pq + import pyoxigraph + import rdflib + from rdflib.store import Store + + schema = pa.schema([(name, pa.string()) for name in "spog"]) + rows = [] + + class BlankNodeLabels(dict): + def get(self, key, default=None): + return key + + class QuadSink(Store): + context_aware = graph_aware = True + + def add_graph(self, graph): + pass + + def remove_graph(self, graph): + pass + + def add(self, triple, context, quoted=False): + terms = [] + for term in triple: + if isinstance(term, rdflib.Literal): + datatype = pyoxigraph.NamedNode(str(term.datatype)) if term.datatype else None + terms.append(str(pyoxigraph.Literal(str(term), language=term.language, datatype=datatype))) + else: + terms.append(term.n3()) + graph = context.identifier + rows.append((*terms, None if graph == dataset.default_context.identifier else graph.n3())) + if len(rows) >= COTTAS_OUTPUT_BATCH_ROWS: + flush() + + def flush(): + writer.write_table(pa.table({name: [row[i] for row in rows] for i, name in enumerate("spog")}, schema=schema)) + rows.clear() + + dataset = rdflib.Dataset(store=QuadSink()) + dataset.default_context = rdflib.Graph(store=dataset.store, identifier=rdflib.BNode()) + normalize = rdflib.NORMALIZE_LITERALS + try: + # Preserve lexical forms such as "01"^^xsd:integer. + rdflib.NORMALIZE_LITERALS = False + with pq.ParquetWriter(output, schema, compression="zstd") as writer: + dataset.parse(source, format="nquads", bnode_context=BlankNodeLabels()) + if rows: + flush() + finally: + rdflib.NORMALIZE_LITERALS = normalize + + +def sort_cottas_chunk(source: Path, outputs: dict[str, Path], *, deduplicate: bool = False) -> None: """Build extra orders from one parsed chunk, with disk-backed sorting.""" import duckdb + import pyarrow.parquet as pq + + has_graph = "g" in pq.read_schema(source).names + if has_graph: + require_dataset_indexes(tuple(outputs)) with duckdb.connect("index-sort.duckdb") as connection: connection.execute("SET preserve_insertion_order = false") for index, output in outputs.items(): - if index not in COTTAS_INDEXES: - raise ValueError("COTTAS index must be a permutation of spo") + if index not in COTTAS_ALL_INDEXES: + raise ValueError("COTTAS index must be a permutation of spo or spog") + columns = "s, p, o" + if "g" in index: + # pycottas 1.1.0 emits DEFAULT for the N-Quads default graph. + columns += ", NULLIF(g, 'DEFAULT') AS g" if has_graph else ", NULL::VARCHAR AS g" connection.execute( - f"COPY (SELECT s, p, o FROM read_parquet($source) ORDER BY {', '.join(index)}) " + f"COPY (SELECT {'DISTINCT ' if deduplicate else ''}{columns} FROM read_parquet($source) ORDER BY {', '.join(column + ' NULLS LAST' for column in index)}) " f"TO $target (FORMAT PARQUET, COMPRESSION ZSTD, COMPRESSION_LEVEL 22, " f"PARQUET_VERSION v2, KV_METADATA {{index: '{index}'}})", {"source": str(source), "target": str(output)}, ) + if deduplicate: + source = output + deduplicate = False + + +def normalize_graph_column(table): + """Represent default graphs consistently as NULL, including legacy files.""" + import pyarrow as pa + import pyarrow.compute as pc + + if "g" not in table.column_names: + return table.append_column("g", pa.nulls(table.num_rows, type=pa.string())) + graph = table["g"] + return table.set_column(table.column_names.index("g"), "g", pc.if_else(pc.equal(graph, "DEFAULT"), None, graph)) + + +def decompress_cottas(source: str, output: str) -> None: + """Stream triples or quads, omitting the graph term for the default graph.""" + import pyarrow.parquet as pq + import pycottas + + with pq.ParquetFile(source) as parquet: + if "g" not in parquet.schema_arrow.names: + pycottas.cottas2rdf(source, output) + return + with (nullcontext(sys.stdout) if output == "/dev/stdout" else open(output, "w", encoding="utf-8")) as handle: + for batch in parquet.iter_batches(batch_size=COTTAS_OUTPUT_BATCH_ROWS, columns=list("spog")): + for s, p, o, g in zip(*(column.to_pylist() for column in batch.columns)): + if any(term is None for term in (s, p, o)): + raise RuntimeError("COTTAS input contains a null RDF term") + if g not in (None, "DEFAULT") and Path(output).suffix == ".nt": + raise ValueError("Named graphs require N-Quads output; choose --decompress-out ending in .nq") + handle.write(f"{s} {p} {o}" + (f" {g}" if g not in (None, "DEFAULT") else "") + " .\n") def reindex_cottas(source: Path, outputs: dict[str, Path]) -> None: @@ -42,18 +141,22 @@ def reindex_cottas(source: Path, outputs: dict[str, Path]) -> None: import pyarrow.parquet as pq with pq.ParquetFile(source) as parquet: - if len(outputs) == 1 and cottas_file_index(parquet) == next(iter(outputs)): + if "g" in parquet.schema_arrow.names: + require_dataset_indexes(tuple(outputs)) + if len(outputs) == 1 and "g" not in next(iter(outputs)) and cottas_file_index(parquet) == next(iter(outputs)): index, output = next(iter(outputs.items())) streaming_cottas_merge([str(source)], str(output), index=index, remove_input_files=False) return # Sorting batches keeps reindexing bounded even for a cohort-sized file. with tempfile.TemporaryDirectory(prefix="reindex-runs-", dir=Path.cwd()) as directory: runs = {index: [] for index in outputs} - for number, batch in enumerate(parquet.iter_batches(batch_size=COTTAS_OUTPUT_BATCH_ROWS, columns=["s", "p", "o"])): + columns = list("spog" if "g" in parquet.schema_arrow.names else "spo") + for number, batch in enumerate(parquet.iter_batches(batch_size=COTTAS_OUTPUT_BATCH_ROWS, columns=columns)): table = pa.Table.from_batches([batch]) for index in outputs: path = Path(directory) / f"{number}.{index}.cottas" - ordered = table.sort_by([(column, "ascending") for column in index]) + indexed = normalize_graph_column(table) if "g" in index else table + ordered = indexed.sort_by([(column, "ascending") for column in index]) pq.write_table(ordered.replace_schema_metadata({b"index": index.encode()}), path, compression="zstd") runs[index].append(str(path)) for index, output in outputs.items(): @@ -71,7 +174,10 @@ def reindex_cottas(source: Path, outputs: dict[str, Path]) -> None: streaming_cottas_merge(paths, str(output), index=index, remove_input_files=True) else: schema = pa.schema([parquet.schema_arrow.field(name) for name in ("s", "p", "o")], metadata={b"index": index.encode()}) - pq.write_table(pa.Table.from_batches([], schema=schema), output, compression="zstd") + table = pa.Table.from_batches([], schema=schema) + if "g" in index: + table = normalize_graph_column(table) + pq.write_table(table, output, compression="zstd") @contextmanager @@ -166,15 +272,22 @@ def __init__(self, parquet_module, path: Path, index: str, batch_rows: int): f"COTTAS input {path.name} is indexed as {found!r}, not {index!r}; " "a streaming merge requires every input to use the requested index" ) - self.fields = tuple(self.schema.field(name) for name in ("s", "p", "o")) - positions = {"s": 0, "p": 1, "o": 2} + import pyarrow as pa + if "g" in self.schema.names: + require_dataset_indexes((index,)) + self.columns = "spog" if "g" in index else "spo" + self.fields = tuple( + self.schema.field(name) if name in self.schema.names else pa.field(name, pa.string()) + for name in self.columns + ) + positions = {name: position for position, name in enumerate(self.columns)} self._sort_positions = tuple(positions[column] for column in index) self._batches = self.parquet_file.iter_batches( batch_size=batch_rows, - columns=["s", "p", "o"], + columns=[name for name in self.columns if name in self.schema.names], use_threads=False, ) - self._values: tuple[list, list, list] | None = None + self._values: tuple[list, ...] | None = None self._row = 0 self.current: tuple | None = None self.sort_key: tuple | None = None @@ -202,14 +315,21 @@ def advance(self) -> None: values = tuple(column.to_pylist() for column in batch.columns) if not values[0]: continue + if len(self.columns) == 4 and len(values) == 3: + values += ([None] * len(values[0]),) self._values = values self._row = 0 - triple = (self._values[0][self._row], self._values[1][self._row], self._values[2][self._row]) + triple = tuple(values[self._row] for values in self._values) + if len(triple) == 4 and triple[3] == "DEFAULT": + triple = triple[:3] + (None,) self._row += 1 - if any(value is None for value in triple): + if any(value is None for value in triple[:3]): raise RuntimeError(f"COTTAS input {self.path.name} contains a null RDF term") - sort_key = tuple(triple[position] for position in self._sort_positions) + sort_key = tuple( + (triple[position] is None, triple[position] or "") if position == 3 else triple[position] + for position in self._sort_positions + ) if self._previous_sort_key is not None and sort_key < self._previous_sort_key: raise RuntimeError( f"COTTAS input {self.path.name} is not sorted by its declared {self.index!r} index" @@ -244,8 +364,8 @@ def streaming_cottas_merge( """ if not input_paths: raise ValueError("at least one COTTAS input is required for a merge") - if not index or set(index.lower()) != {"s", "p", "o"} or len(index) != 3: - raise ValueError("COTTAS merge index must be a permutation of spo") + if index.lower() not in COTTAS_ALL_INDEXES: + raise ValueError("COTTAS merge index must be a permutation of spo or spog") try: import pyarrow as pa @@ -323,7 +443,7 @@ def streaming_cottas_merge( if len(output_rows) >= COTTAS_OUTPUT_BATCH_ROWS: arrays = [ pa.array([row[column] for row in output_rows], type=fields[column].type) - for column in range(3) + for column in range(len(fields)) ] writer.write_batch(pa.RecordBatch.from_arrays(arrays, schema=output_schema)) written_rows += len(output_rows) @@ -346,7 +466,7 @@ def streaming_cottas_merge( if output_rows: arrays = [ pa.array([row[column] for row in output_rows], type=fields[column].type) - for column in range(3) + for column in range(len(fields)) ] writer.write_batch(pa.RecordBatch.from_arrays(arrays, schema=output_schema)) written_rows += len(output_rows) @@ -398,7 +518,7 @@ def main() -> int: merge.add_argument("left_path") merge.add_argument("right_path") merge.add_argument("cottas_path") - merge.add_argument("index", nargs="?", default="spo", type=str.lower, choices=COTTAS_INDEXES) + merge.add_argument("index", nargs="?", default="spo", type=str.lower, choices=COTTAS_ALL_INDEXES) merge_many = subparsers.add_parser( "merge-many", @@ -411,7 +531,7 @@ def main() -> int: help="COTTAS inputs to merge", ) merge_many.add_argument("--output-cottas-file", required=True) - merge_many.add_argument("--index", default="spo", type=str.lower, choices=COTTAS_INDEXES) + merge_many.add_argument("--index", default="spo", type=str.lower, choices=COTTAS_ALL_INDEXES) merge_many.add_argument( "--progress-path", help="optional JSONL sidecar for bounded streaming-merge progress", @@ -438,28 +558,33 @@ def main() -> int: if args.command == "convert": rdf_path = str(Path(args.rdf_path).resolve()) cottas_path = str(Path(args.cottas_path).resolve()) - # disk=True keeps parser/index construction from requiring the whole - # RDF chunk in Python memory. + if Path(rdf_path).suffix.lower() in {".nq", ".trig"}: + require_dataset_indexes(args.index) + # Keep one RDF parse/DISTINCT per chunk. Quad orders are sorted after + # normalizing pycottas's DEFAULT sentinel, so NULL ordering is uniform. with cottas_scratch_workspace(): - pycottas.rdf2cottas( - rdf_path, - cottas_path, - index=args.index[0], - disk=True, - ) + primary_is_quad = "g" in args.index[0] + parsed = Path.cwd() / "parsed.cottas" if primary_is_quad else Path(cottas_path) + if Path(rdf_path).suffix.lower() == ".nq": + parse_nquads_chunk(rdf_path, parsed) + else: + pycottas.rdf2cottas( + rdf_path, str(parsed), index="" if primary_is_quad else args.index[0], disk=True, + ) outputs = cottas_index_paths(Path(cottas_path), args.index) - outputs.pop(args.index[0]) + if not primary_is_quad: + outputs.pop(args.index[0]) if outputs: - sort_cottas_chunk(Path(cottas_path), outputs) + sort_cottas_chunk(parsed, outputs, deduplicate=Path(rdf_path).suffix.lower() == ".nq") return 0 if args.command == "decompress": cottas_path = str(Path(args.cottas_path).resolve()) - rdf_path = str(Path(args.rdf_path).resolve()) + rdf_path = str(Path(args.rdf_path).absolute()) # Keep DuckDB scratch state in the container-local workspace while # pycottas writes the decoded RDF directly to the mounted output. with cottas_scratch_workspace(): - pycottas.cottas2rdf(cottas_path, rdf_path) + decompress_cottas(cottas_path, rdf_path) return 0 if args.command == "reindex": diff --git a/src/partitioned_compression.py b/src/partitioned_compression.py index c2f6d8a..ca4bbe4 100644 --- a/src/partitioned_compression.py +++ b/src/partitioned_compression.py @@ -26,7 +26,7 @@ resource = None from pathlib import Path -from vcf_rdfizer_cottas import cottas_index_paths, parse_cottas_indexes +from vcf_rdfizer_cottas import cottas_index_paths, parse_cottas_indexes, require_dataset_indexes HDT_METHODS = {"hdt", "hdt_gzip", "hdt_brotli"} @@ -106,7 +106,7 @@ def is_triple_line(line: bytes) -> bool: def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser(description="VCF-RDFizer partitioned compression runner") - parser.add_argument("--source", required=True, help="plain or gzip-compressed N-Triples input") + parser.add_argument("--source", required=True, help="plain or gzip-compressed N-Triples/N-Quads input") parser.add_argument("--output-dir", required=True) parser.add_argument("--output-name", required=True) parser.add_argument("--methods", required=True, help="comma-separated internal method names") @@ -161,6 +161,7 @@ def stream_chunks( chunk_dir.mkdir(parents=True, exist_ok=True) prepare_progress_path(progress_path) + chunk_suffix = ".nq" if source.name.endswith((".nq", ".nq.gz")) else ".nt" plan = { "source_file_count": 1, "source_paths": [str(source)], @@ -196,7 +197,7 @@ def generate(): def open_chunk(): nonlocal handle, chunk_path, chunk_size, chunk_start_offset, chunk_start_record, chunk_index, chunk_started chunk_started = time.perf_counter() - chunk_path = chunk_dir / f"chunk-{chunk_index:05d}.nt" + chunk_path = chunk_dir / f"chunk-{chunk_index:05d}{chunk_suffix}" chunk_index += 1 handle = chunk_path.open("wb") chunk_size = 0 @@ -914,6 +915,11 @@ def skipped_cottas_result(method: str) -> dict: if not methods or not all(method in HDT_METHODS | COTTAS_METHODS for method in methods): raise ValueError(f"Unsupported partitioned method list: {methods}") + if source.name.endswith((".nq", ".nq.gz")): + if any(method in HDT_METHODS for method in methods): + raise ValueError("HDT does not preserve named graphs; select --representations cottas") + require_dataset_indexes(cottas_indexes) + hdt_paths: list[Path] = [] cottas_paths: list[Path] = [] diff --git a/src/validate_compression.py b/src/validate_compression.py index ca833e9..9692eac 100644 --- a/src/validate_compression.py +++ b/src/validate_compression.py @@ -111,20 +111,12 @@ def validate(args: argparse.Namespace) -> dict: except ImportError as exc: raise RuntimeError(f"COTTAS dependency is unavailable: {exc}") from exc - # pycottas exposes cottas2rdf; its output path can be /dev/stdout, so - # the COTTAS reader is exercised without materializing decoded RDF. + # The adapter also preserves default/named graphs when decoding quads. decoded_triples = count_decoded( - [ - sys.executable, - "-c", - ( - "import pycottas, sys; " - "pycottas.cottas2rdf(sys.argv[1], '/dev/stdout')" - ), - str(artifact), - ] + [sys.executable, str(Path(__file__).with_name("cottas_tool.py")), + "decompress", str(artifact), "/dev/stdout"] ) - validator = "pycottas.cottas2rdf" + validator = "cottas_tool.decompress" return { "valid": True, @@ -143,7 +135,7 @@ def validate(args: argparse.Namespace) -> dict: 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") + 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) diff --git a/test/README.md b/test/README.md index 3465a59..fcfdc39 100644 --- a/test/README.md +++ b/test/README.md @@ -190,11 +190,12 @@ OK ## COTTAS index coverage Install the pinned container libraries to exercise real Parquet conversion, -all six orders, reindexing, deduplication, packaging dispatch and failure cleanup: +all six triple and 24 dataset orders, reindexing, deduplication, packaging dispatch +and failure cleanup: ```bash python -m pip install pycottas==1.1.0 duckdb==1.5.5 pyarrow==22.0.0 coverage -coverage run --branch -m unittest test.test_cottas_indexes_unit test.test_cottas_tool +coverage run --branch -m unittest test.test_cottas_indexes_unit test.test_cottas_quads_unit test.test_cottas_tool coverage report -m --include='*cottas_tool.py,*vcf_rdfizer_cottas.py' ``` @@ -203,3 +204,6 @@ and `--artifact-compression gzip,brotli`, then run `test/cottas_index_roundtrip. inside that image with `--source --cottas --indexes all --packages`. It checks exact RDF equality, ordering, metadata, eight native triple-pattern shapes per order, and package contents. +For datasets, use `test/fixtures/cottas-dataset.nq` and `--indexes all-quads`. +This adds 16 quad-pattern shapes per order and checks default/named graphs, +shared blank nodes and lexical literal values across chunk boundaries. diff --git a/test/cottas_index_roundtrip.py b/test/cottas_index_roundtrip.py index b162872..1e76f70 100644 --- a/test/cottas_index_roundtrip.py +++ b/test/cottas_index_roundtrip.py @@ -8,15 +8,32 @@ import json import subprocess import tempfile +import sys from pathlib import Path import pyarrow.parquet as pq import pycottas -from rdflib import Graph +import rdflib from vcf_rdfizer_cottas import cottas_index_paths, parse_cottas_indexes +class BlankNodeLabels(dict): + def get(self, key, default=None): + return key + + +def read_rdf(text, dataset): + rdflib.NORMALIZE_LITERALS = False + graph = rdflib.Dataset() if dataset else rdflib.Graph() + if dataset: + graph.default_context = rdflib.Graph(store=graph.store, identifier=rdflib.BNode()) + graph.parse(data=text, format='nquads' if dataset else 'nt', bnode_context=BlankNodeLabels()) + if dataset: + return {(s, p, o, None if g == graph.default_context.identifier else g) for s, p, o, g in graph.quads()} + return set(graph) + + def main(): parser = argparse.ArgumentParser(description=__doc__) parser.add_argument('--source', type=Path, required=True) @@ -24,29 +41,38 @@ def main(): parser.add_argument('--indexes', type=parse_cottas_indexes, default='all') parser.add_argument('--packages', action='store_true') args = parser.parse_args() + dataset = args.source.name.endswith(('.nq', '.nq.gz')) opener = gzip.open if args.source.name.endswith('.gz') else open with opener(args.source, 'rt') as handle: - expected = set(Graph().parse(data=handle.read(), format='nt')) + expected = read_rdf(handle.read(), dataset) report = {} for order, path in cottas_index_paths(args.cottas, args.indexes).items(): table = pq.read_table(path) assert table.schema.metadata[b'index'].decode() == order, path - rows = list(zip(*(table[name].to_pylist() for name in 'spo'))) - assert rows == sorted(set(rows), key=lambda row: tuple(row['spo'.index(c)] for c in order)), path + columns = 'spog' if 'g' in order else 'spo' + rows = list(zip(*(table[name].to_pylist() for name in columns))) + assert rows == sorted(set(rows), key=lambda row: tuple((row[columns.index(c)] is None, row[columns.index(c)] or '') for c in order)), path with tempfile.TemporaryDirectory() as td: - decoded = Path(td) / 'decoded.nt' - pycottas.cottas2rdf(str(path), str(decoded)) - assert set(Graph().parse(decoded, format='nt')) == expected, path - # Exercise each possible bound/unbound triple-pattern shape. - probe = rows[0] - for bound in itertools.product((False, True), repeat=3): - pattern = ' '.join(probe[i] if bound[i] else '?'+name for i, name in enumerate('spo')) - wanted = {row for row in rows if all(not bound[i] or row[i] == probe[i] for i in range(3))} + decoded = Path(td) / ('decoded.nq' if 'g' in order else 'decoded.nt') + subprocess.run([sys.executable, '/opt/vcf-rdfizer/cottas_tool.py', 'decompress', str(path), str(decoded)], check=True) + wanted = expected if dataset or 'g' not in order else {(*triple, None) for triple in expected} + assert read_rdf(decoded.read_text(), 'g' in order) == wanted, path + # Exercise every bound/unbound shape with a named-graph probe when present. + probe = next((row for row in rows if all(term is not None and not term.startswith('_:') for term in row)), rows[0]) + shapes = 0 + for bound in itertools.product((False, True), repeat=len(columns)): + if any(bound[i] and probe[i] is None for i in range(len(columns))): + continue # None denotes a wildcard in pycottas, not a graph constant. + # pycottas 1.1.0 truncates string patterns to three terms; tuples + # are its documented API for a four-term quad pattern. + pattern = tuple(rdflib.util.from_n3(probe[i]) if bound[i] else None for i in range(len(columns))) + wanted = {row for row in rows if all(not bound[i] or row[i] == probe[i] for i in range(len(columns)))} assert set(pycottas.search(str(path), pattern)) == wanted, (path, pattern) + shapes += 1 if args.packages: assert gzip.decompress(Path(str(path)+'.gz').read_bytes()) == path.read_bytes(), path assert subprocess.check_output(['brotli', '-d', '-c', str(path)+'.br']) == path.read_bytes(), path - report[order] = {'triples': len(rows), 'bytes': path.stat().st_size, 'native_query_shapes': 8} + report[order] = {'quads' if 'g' in order else 'triples': len(rows), 'bytes': path.stat().st_size, 'native_query_shapes': shapes} print(json.dumps(report, indent=2)) diff --git a/test/fixtures/cottas-dataset.nq b/test/fixtures/cottas-dataset.nq new file mode 100644 index 0000000..3dd3b33 --- /dev/null +++ b/test/fixtures/cottas-dataset.nq @@ -0,0 +1,11 @@ +# Shared blank nodes, named graphs and the default graph span chunk boundaries. + "same" . + "same" . + "same" . +_:shared _:object _:graph . + "é"@fr . + "01"^^ . + "line\nquote\"slash\\" _:graph . +_:shared _:object . + . + "" _:graph . diff --git a/test/test_cottas_indexes_unit.py b/test/test_cottas_indexes_unit.py index 08bcdf8..3cc0ee4 100644 --- a/test/test_cottas_indexes_unit.py +++ b/test/test_cottas_indexes_unit.py @@ -26,7 +26,7 @@ class IndexSelectionTests(unittest.TestCase): def test_selection_and_names(self): self.assertEqual(parse_cottas_indexes(' SPO, pso,spo,POS '), ('spo', 'pso', 'pos')) self.assertEqual(parse_cottas_indexes(' ALL '), COTTAS_INDEXES) - for value in ('', 'spo,', ',spo', 'sp', 'sspo', 'spog', 'all,spo', 'spp', 'sp o'): + for value in ('', 'spo,', ',spo', 'sp', 'sspo', 'spogg', 'all,spo', 'spp', 'sp o'): with self.subTest(value=value), self.assertRaises(ValueError): parse_cottas_indexes(value) self.assertEqual(cottas_index_paths(Path('a.cottas'), ('pos', 'spo')), @@ -249,11 +249,11 @@ def test_missing_dependencies(self): class PartitionedIndexTests(unittest.TestCase): - def run_pipeline(self, root, *, indexes=('pos', 'spo', 'ops'), chunk_bytes=55, failure=None, allow=False, methods='cottas,cottas_gzip,cottas_brotli'): + def run_pipeline(self, root, *, indexes=('pos', 'spo', 'ops'), chunk_bytes=55, failure=None, allow=False, methods='cottas,cottas_gzip,cottas_brotli', dataset=False): runner = load_runner_module() work, out = root / 'work', root / 'out' work.mkdir(); out.mkdir() - source = root / 'source.nt' + source = root / ('source.nq' if dataset else 'source.nt') source.write_text(''.join(f' .\n' for i in range(7))) result_path = out / 'result.json' calls = [] @@ -315,6 +315,16 @@ def test_failures_skip_indexes_and_packages_or_fail_strictly(self): self.assertEqual(code, 1) self.assertIn('packaging failed for spo', report['error']) + def test_dataset_pipeline_and_graph_loss_guards(self): + for methods, indexes, expected in [('cottas', ('gspo', 'spog'), 0), ('cottas', ('spo',), 1), ('hdt', ('gspo',), 1)]: + with tempfile.TemporaryDirectory() as td, self.subTest(methods=methods, indexes=indexes): + code, report, calls, *_ = self.run_pipeline(Path(td), dataset=True, indexes=indexes, methods=methods) + self.assertEqual(code, expected, report) + if expected: + self.assertFalse(calls) + else: + self.assertTrue(all(command[3].endswith('.nq') for name, command in calls if name.startswith('cottas-build'))) + class HostIndexRoutingTests(unittest.TestCase): def test_host_forwards_indexes_and_translates_every_artifact(self): with tempfile.TemporaryDirectory() as td: diff --git a/test/test_cottas_quads_unit.py b/test/test_cottas_quads_unit.py new file mode 100644 index 0000000..47edc3d --- /dev/null +++ b/test/test_cottas_quads_unit.py @@ -0,0 +1,238 @@ +"""Dataset ordering and identity through chunk conversion, merging and export.""" +import contextlib +import gzip +import io +import itertools +import os +import sys +import subprocess +import tempfile +import unittest +from pathlib import Path +from unittest import mock + +import vcf_rdfizer as host +from vcf_rdfizer_cottas import COTTAS_INDEXES, COTTAS_QUAD_INDEXES, cottas_index_paths, parse_cottas_indexes +from test.test_cottas_tool import load_cottas_tool +from test.test_partitioned_compression_unit import load_runner_module + +try: + import pyarrow as pa + import pyarrow.parquet as pq + import pycottas +except ImportError: + pa = None + +FIXTURE = Path(__file__).with_name('fixtures') / 'cottas-dataset.nq' + + +class DatasetSelectionTests(unittest.TestCase): + def test_all_permutations_and_names(self): + self.assertEqual(len(set(COTTAS_QUAD_INDEXES)), 24) + self.assertTrue(all(set(index) == set('spog') for index in COTTAS_QUAD_INDEXES)) + self.assertEqual(parse_cottas_indexes(' ALL-QUADS '), COTTAS_QUAD_INDEXES) + self.assertEqual(parse_cottas_indexes(' GsPo,SPoG,gspo '), ('gspo', 'spog')) + self.assertEqual(parse_cottas_indexes('all'), COTTAS_INDEXES) + for selection in ('g', 'spg', 'spogg', 'all-quads,spo'): + with self.assertRaises(ValueError): + parse_cottas_indexes(selection) + for name in ('sample.nq', 'sample.nq.gz'): + source = Path(name) + self.assertEqual(host.rdf_output_basename(source), 'sample') + self.assertEqual(host.compression_artifact_name_for_method(source, 'gzip'), 'sample.nq.gz') + self.assertEqual(host.compression_artifact_name_for_method(source, 'brotli'), 'sample.nq.br') + self.assertEqual(host.compression_method_label_for_path(source, 'gzip'), 'gzip (.nq.gz)') + planned = host.planned_output_paths(out_dir=Path('/out'), output_name='sample', rdf_name=None, + source_rdf_name=name, methods=['gzip', 'brotli', 'cottas'], partitioned=True, cottas_indexes=('gspo', 'spog')) + self.assertIn(Path('/out/sample/sample.nq.gz'), planned) + self.assertIn(Path('/out/sample/sample.spog.cottas'), planned) + self.assertNotIn(Path('/out/sample/sample.nq'), planned) + for codec in ('gzip', 'brotli'): + suffix = 'gz' if codec == 'gzip' else 'br' + self.assertEqual(host.default_decompressed_name(Path('sample.nq.' + suffix), codec), 'sample.nq') + + def test_dataset_preflight_prevents_graph_loss(self): + for representations, indexes in [('hdt', 'spo'), ('cottas', 'spo'), ('cottas', 'gspo,spo')]: + with tempfile.TemporaryDirectory() as td, self.subTest(representations=representations, indexes=indexes): + argv = ['vcf-rdfizer', '-m', 'compress', '--rdf', str(FIXTURE), '-o', td, + '--representations', representations, '--cottas-indexes', indexes] + with mock.patch.object(sys, 'argv', argv), mock.patch.object(host, 'check_docker') as docker, contextlib.redirect_stderr(io.StringIO()), contextlib.redirect_stdout(io.StringIO()): + self.assertEqual(host.main(), 2) + docker.assert_not_called() + + def test_gzip_dataset_chunks_keep_format_and_contents(self): + runner = load_runner_module() + with tempfile.TemporaryDirectory() as td: + root = Path(td) + source = root / 'dataset.nq.gz' + source.write_bytes(gzip.compress(FIXTURE.read_bytes())) + chunks = root / 'chunks'; chunks.mkdir() + iterator, plan = runner.stream_chunks(source, chunks, target_bytes=100, min_bytes=1, max_bytes=200) + paths = [path for path, metadata in iterator] + self.assertGreater(len(paths), 1) + self.assertTrue(all(path.suffix == '.nq' for path in paths)) + self.assertEqual(b''.join(path.read_bytes() for path in paths), FIXTURE.read_bytes()) + self.assertEqual(plan['record_count'], 10) + + +@unittest.skipIf(pa is None, 'install pycottas and pyarrow for real Parquet tests') +class QuadIndexTests(unittest.TestCase): + def setUp(self): + temporary = tempfile.TemporaryDirectory() + self.addCleanup(temporary.cleanup) + self.root = Path(temporary.name).resolve() + self.tool = load_cottas_tool() + patch = mock.patch.dict(os.environ, {'COTTAS_SCRATCH_DIR': str(self.root / 'scratch'), 'COTTAS_MERGE_BATCH_ROWS': '2'}) + patch.start(); self.addCleanup(patch.stop) + self.rows = [('', '', '"same"', graph) for graph in ('', '', None, '_:graph')] + self.rows += [('_:shared', '', '_:object', '_:graph')] + + def main(self, *args): + with mock.patch.object(sys, 'argv', ['cottas_tool', *map(str, args)]): + return self.tool.main() + + def ordered(self, rows, order): + return sorted(set(rows), key=lambda row: tuple((row['spog'.index(c)] is None, row['spog'.index(c)] or '') for c in order)) + + def write(self, path, rows, order='spog', columns='spog'): + table = pa.table({name: pa.array([row[i] for row in rows], type=pa.string()) for i, name in enumerate(columns)}) + pq.write_table(table.replace_schema_metadata({b'index': order.encode()}), path) + + def assert_rows(self, path, order, rows): + table = pq.read_table(path) + self.assertEqual(table.schema.metadata[b'index'], order.encode()) + self.assertEqual(list(zip(*(table[name].to_pylist() for name in 'spog'))), self.ordered(rows, order)) + + def test_all_quad_orders_parse_once_and_preserve_lexical_terms(self): + self.tool.COTTAS_OUTPUT_BATCH_ROWS = 3 + output = self.root / 'all.cottas' + with mock.patch.object(self.tool, 'parse_nquads_chunk', wraps=self.tool.parse_nquads_chunk) as parse: + self.assertEqual(self.main('convert', FIXTURE, output, 'all-quads'), 0) + self.assertEqual(parse.call_count, 1) + rows = list(zip(*(pq.read_table(output)[name].to_pylist() for name in 'spog'))) + self.assertEqual(len(rows), 10) + self.assertIn(('_:shared', '', '_:object', '_:graph'), rows) + self.assertIn(('', '', '"01"^^', ''), rows) + self.assertIn(('', '', '"é"@fr', ''), rows) + for order, path in cottas_index_paths(output, COTTAS_QUAD_INDEXES).items(): + self.assert_rows(path, order, rows) + import rdflib + probe = self.rows[0] + for bound in itertools.product((False, True), repeat=4): + pattern = tuple(rdflib.util.from_n3(probe[i]) if bound[i] else None for i in range(4)) + expected = {row for row in rows if all(not bound[i] or row[i] == probe[i] for i in range(4))} + self.assertEqual(set(pycottas.search(str(path), pattern)), expected) + decoded = self.root / 'decoded.nq' + self.assertEqual(self.main('decompress', output, decoded), 0) + expected = {line for line in FIXTURE.read_text().splitlines() if not line.startswith('#')} + self.assertEqual(set(decoded.read_text().splitlines()), expected) + self.assertEqual(list((self.root / 'scratch').iterdir()), []) + + def test_every_single_quad_order(self): + source = self.root / 'input.nq' + source.write_text(' "same" .\n' * 2) + for order in COTTAS_QUAD_INDEXES: + path = self.root / f'{order}.cottas' + self.main('convert', source, path, order) + self.assert_rows(path, order, [self.rows[2]]) + + def test_named_graph_can_use_rdflib_default_identifier(self): + source = self.root / 'default.nq' + source.write_text(' "same" .\n "same" .\n') + output = self.root / 'default.cottas' + self.main('convert', source, output, 'gspo') + self.assert_rows(output, 'gspo', [self.rows[2], self.rows[2][:3] + ('',)]) + + def test_chunk_merge_all_orders_retains_graph_identity(self): + self.tool.COTTAS_OUTPUT_BATCH_ROWS = 2 + for order in COTTAS_QUAD_INDEXES: + paths = [self.root / f'{order}-{n}.cottas' for n in range(3)] + for path, rows in zip(paths, (self.rows[:4], self.rows[1:] + self.rows[1:], [])): + self.write(path, self.ordered(rows, order), order) + output = self.root / f'{order}.cottas' + self.tool.streaming_cottas_merge(list(map(str, paths)), str(output), index=order, remove_input_files=True) + self.assert_rows(output, order, self.rows) + + def test_reindex_legacy_default_and_unsorted_quads_all_orders(self): + self.tool.COTTAS_OUTPUT_BATCH_ROWS = 2 + self.tool.COTTAS_REINDEX_FAN_IN = 2 + source = self.root / 'legacy.cottas' + rows = [row[:3] + ('DEFAULT' if row[3] is None else row[3],) for row in self.rows] + self.write(source, rows * 2, 'gspo') + self.main('reindex', source, 'all-quads') + for order, path in cottas_index_paths(source, COTTAS_QUAD_INDEXES).items(): + self.assert_rows(path, order, self.rows) + + def test_triples_can_gain_default_graph_and_mix_index_families(self): + source = self.root / 'input.nt' + source.write_text(' "same" .\n') + for selection in ('spo,gspo', 'gspo,spo'): + output = self.root / (selection.replace(',', '-') + '.cottas') + self.main('convert', source, output, selection) + for order, path in cottas_index_paths(output, parse_cottas_indexes(selection)).items(): + self.assertEqual(pq.read_table(path).num_rows, 1) + if 'g' in order: + self.assert_rows(path, order, [self.rows[2]]) + triple = self.root / 'triple.cottas' + self.write(triple, [self.rows[2][:3]], 'spo', columns='spo') + self.main('decompress', triple, self.root / 'triple.nt') + self.assertEqual((self.root / 'triple.nt').read_text(), source.read_text()) + self.main('reindex', triple, 'gspo,spog') + self.assert_rows(triple, 'gspo', [self.rows[2]]) + decoded = self.root / 'default.nt' + self.main('decompress', triple, decoded) + self.assertEqual(decoded.read_text(), source.read_text()) + + def test_empty_dataset_conversion_and_reindex(self): + source = self.root / 'empty.nq'; source.touch() + output = self.root / 'empty.cottas' + self.main('convert', source, output, 'gspo') + self.assert_rows(output, 'gspo', []) + self.main('reindex', output, 'spog,gspo') + self.assert_rows(output, 'spog', []) + + def test_legacy_default_only_chunk_without_graph_column(self): + source = self.root / 'triple.cottas' + self.write(source, [self.rows[2][:3]], 'gspo', columns='spo') + output = self.root / 'out.cottas' + self.tool.streaming_cottas_merge([str(source)], str(output), index='gspo', remove_input_files=False) + self.assert_rows(output, 'gspo', [self.rows[2]]) + self.write(source, [self.rows[2][:3] + ('DEFAULT',)], 'gspo') + self.tool.streaming_cottas_merge([str(source)], str(output), index='gspo', remove_input_files=False) + self.assert_rows(output, 'gspo', [self.rows[2]]) + self.main('decompress', source, self.root / 'legacy.nq') + self.assertNotIn('DEFAULT', (self.root / 'legacy.nq').read_text()) + + def test_quad_safety_guards(self): + source = self.root / 'quads.cottas' + output = self.root / 'out.cottas' + self.write(source, self.rows, 'spo') + for action in (lambda: self.main('convert', FIXTURE, output, 'spo'), + lambda: self.main('reindex', source, 'spo'), + lambda: self.tool.sort_cottas_chunk(source, {'spo': output})): + with self.assertRaisesRegex(ValueError, 'containing g'): + action() + with self.assertRaisesRegex(RuntimeError, 'containing g'): + self.tool.streaming_cottas_merge([str(source)], str(output), index='spo', remove_input_files=False) + with self.assertRaisesRegex(ValueError, 'N-Quads'): + self.main('decompress', source, self.root / 'bad.nt') + self.write(source, [(None, '', '', None)]) + with self.assertRaisesRegex(RuntimeError, 'null RDF term'): + self.main('decompress', source, self.root / 'bad.nq') + + def test_decoder_stdout_pipe_for_validation(self): + source = self.root / 'quads.cottas' + self.write(source, self.rows) + environment = dict(os.environ, PYTHONPATH=str(Path(__file__).resolve().parents[1])) + result = subprocess.run([sys.executable, self.tool.__file__, 'decompress', str(source), '/dev/stdout'], + env=environment, check=True, capture_output=True, text=True) + self.assertEqual(len(result.stdout.splitlines()), len(self.rows)) + self.assertIn(' "same" .', result.stdout) + + def test_parse_failure_restores_rdflib_setting(self): + import rdflib + source = self.root / 'bad.nq'; source.write_text('not RDF') + before = rdflib.NORMALIZE_LITERALS + with self.assertRaises(Exception): + self.main('convert', source, self.root / 'bad.cottas', 'gspo') + self.assertEqual(rdflib.NORMALIZE_LITERALS, before) diff --git a/test/test_cottas_tool.py b/test/test_cottas_tool.py index e4b73d9..3ab4010 100644 --- a/test/test_cottas_tool.py +++ b/test/test_cottas_tool.py @@ -182,7 +182,7 @@ def fake_cottas2rdf(cottas_path, rdf_path): with mock.patch.dict( sys.modules, {"pycottas": types.SimpleNamespace(cottas2rdf=fake_cottas2rdf)}, - ), mock.patch.dict( + ), mock.patch.object(module, "decompress_cottas", side_effect=fake_cottas2rdf), mock.patch.dict( os.environ, {"COTTAS_SCRATCH_DIR": str(scratch_root)}, clear=False ), mock.patch.object( sys, diff --git a/test/test_validate_compression_unit.py b/test/test_validate_compression_unit.py index 8f14fac..d219526 100644 --- a/test/test_validate_compression_unit.py +++ b/test/test_validate_compression_unit.py @@ -252,7 +252,7 @@ def test_the_cottas_path_names_its_own_validator(self): with mock.patch.dict(sys.modules, {"pycottas": mock.Mock()}), \ mock.patch.object(M, "count_decoded", return_value=2): result = M.validate(self.args(format="cottas")) - self.assertEqual(result["validator"], "pycottas.cottas2rdf") + self.assertEqual(result["validator"], "cottas_tool.decompress") self.assertTrue(result["count_match"]) diff --git a/vcf_rdfizer.py b/vcf_rdfizer.py index ac85c69..b8c076a 100644 --- a/vcf_rdfizer.py +++ b/vcf_rdfizer.py @@ -74,7 +74,7 @@ except ImportError: # pragma: no cover - shipped alongside this module vcf_rdfizer_gzip = None -from vcf_rdfizer_cottas import cottas_index_paths, parse_cottas_indexes +from vcf_rdfizer_cottas import cottas_index_paths, parse_cottas_indexes, require_dataset_indexes import vcf_rdfizer_vocab as vocab from vcf_rdfizer_vocab import ( @@ -1560,18 +1560,17 @@ def rdf_label_for_path(path: Path) -> str: def rdf_output_basename(path: Path) -> str: """Return the common output basename for ``.nt`` and ``.nt.gz`` RDF.""" - if path.name.endswith(".nt.gz"): - return path.name[: -len(".nt.gz")] - if path.name.endswith(".nt"): - return path.name[: -len(".nt")] + for suffix in (".nt.gz", ".nq.gz", ".nt", ".nq"): + if path.name.endswith(suffix): + return path.name[:-len(suffix)] return path.stem def compression_artifact_name_for_method(path: Path, method: str) -> str: """Compute expected compressed artifact filename for a method.""" - if path.name.endswith(".nt.gz"): + if path.name.endswith((".nt.gz", ".nq.gz")): stem = rdf_output_basename(path) - ext = "nt" + ext = Path(path.stem).suffix.lstrip(".") else: stem = rdf_output_basename(path) ext = path.suffix.lstrip(".") or "nt" @@ -1602,6 +1601,7 @@ def planned_output_paths( methods: list[str], partitioned: bool, cottas_indexes: tuple[str, ...] = ("spo",), + source_rdf_name: str | None = None, ) -> set[Path]: """List final output paths that a compression plan would create.""" target_dir = out_dir / output_name @@ -1609,7 +1609,7 @@ def planned_output_paths( if rdf_name is not None: planned.add(target_dir / rdf_name) - rdf_path = Path(rdf_name or f"{output_name}.nt") + rdf_path = Path(source_rdf_name or rdf_name or f"{output_name}.nt") planned.update( target_dir / compression_artifact_name_for_method(rdf_path, method) for method in methods @@ -1661,8 +1661,8 @@ def validate_no_output_collisions(plans: dict[str, set[Path]]): def compression_method_label_for_path(path: Path, method: str) -> str: """Return human-readable compression method label for a path.""" ext = path.suffix.lstrip(".") or "nt" - if path.name.endswith(".nt.gz"): - ext = "nt" + if path.name.endswith((".nt.gz", ".nq.gz")): + ext = Path(path.stem).suffix.lstrip(".") labels = { "gzip": f"gzip (.{ext}.gz)", "brotli": f"brotli (.{ext}.br)", @@ -6355,7 +6355,7 @@ def run_compression_methods_for_rdf( in_dir = rdf_path.parent input_container = f"/data/in/{rdf_path.name}" input_stem = rdf_output_basename(rdf_path) - input_ext = "nt" if rdf_path.name.endswith(".nt.gz") else rdf_path.suffix.lstrip(".") or "nt" + input_ext = Path(rdf_path.stem).suffix.lstrip(".") if rdf_path.name.endswith(".gz") else rdf_path.suffix.lstrip(".") or "nt" if target_out_dir is None: target_out_dir = out_dir / input_stem ensure_dir(target_out_dir) @@ -8523,11 +8523,11 @@ def detect_compressed_format(path: Path): def default_decompressed_name(path: Path, fmt: str): """Compute default output filename for decompression mode.""" if fmt == "gzip": - if path.name.endswith(".nt.gz"): + if path.name.endswith((".nt.gz", ".nq.gz")): return path.name[: -len(".gz")] return f"{path.stem}.nt" if fmt == "brotli": - if path.name.endswith(".nt.br"): + if path.name.endswith((".nt.br", ".nq.br")): return path.name[: -len(".br")] return f"{path.stem}.nt" if fmt == "cottas": @@ -9160,7 +9160,7 @@ def main(): parser.add_argument( "--rdf", default=None, - help="Input RDF file (.nt or .nt.gz) for --mode compress; .nt.gz required for --mode validation", + help="Input RDF file (.nt/.nt.gz or .nq/.nq.gz) for --mode compress; RDF artifact for --mode validation", ) parser.add_argument( "-C", @@ -9327,7 +9327,7 @@ def main(): ) parser.add_argument( "--cottas-indexes", type=parse_cottas_indexes, default="spo", - help="COTTAS orders: spo,sop,pso,pos,osp,ops (comma-separated), or all; first is primary (default: spo)", + help="COTTAS orders: permutations of spo or spog (comma-separated); all selects six triple orders, all-quads selects 24 dataset orders (default: spo)", ) parser.add_argument( "--artifact-compression", @@ -9861,8 +9861,8 @@ def main(): rdf_path = Path(args.rdf).expanduser().resolve() if not rdf_path.exists() or not rdf_path.is_file(): raise ValueError(f"RDF input file not found: {rdf_path}") - if rdf_path.suffix != ".nt" and not rdf_path.name.endswith(".nt.gz"): - raise ValueError("Compression input must be a .nt or .nt.gz file") + if not rdf_path.name.endswith((".nt", ".nt.gz", ".nq", ".nq.gz")): + raise ValueError("Compression input must be a .nt, .nt.gz, .nq or .nq.gz file") if args.legacy_compression is not None: if ( args.rdf_compression != DEFAULT_RDF_COMPRESSION @@ -9882,6 +9882,11 @@ def main(): ) if args.cottas_indexes != ("spo",) and not any(method in COTTAS_COMPRESSION_METHODS for method in methods): raise ValueError("--cottas-indexes requires --representations cottas") + if rdf_path.name.endswith((".nq", ".nq.gz")): + if any(method in HDT_COMPRESSION_METHODS for method in methods): + raise ValueError("HDT does not preserve named graphs; select --representations cottas") + if any(method in COTTAS_COMPRESSION_METHODS for method in methods): + require_dataset_indexes(args.cottas_indexes) hdt_strategy_error = hdt_strategy_rejection( hdt_strategy=args.hdt_strategy, methods=methods, @@ -9912,6 +9917,7 @@ def main(): methods=methods, partitioned=compression_uses_partitioning_for_input, cottas_indexes=args.cottas_indexes, + source_rdf_name=rdf_path.name, ) } ) diff --git a/vcf_rdfizer_cottas.py b/vcf_rdfizer_cottas.py index 545abe0..f680fea 100644 --- a/vcf_rdfizer_cottas.py +++ b/vcf_rdfizer_cottas.py @@ -1,21 +1,32 @@ """COTTAS index selection and artifact naming shared by host and container.""" from pathlib import Path +from itertools import permutations COTTAS_INDEXES = ("spo", "sop", "pso", "pos", "osp", "ops") +COTTAS_QUAD_INDEXES = tuple("".join(order) for order in permutations("spog")) +COTTAS_ALL_INDEXES = COTTAS_INDEXES + COTTAS_QUAD_INDEXES def parse_cottas_indexes(value: str) -> tuple[str, ...]: """Normalize an ordered selection, rejecting typos before conversion.""" if value.strip().lower() == "all": return COTTAS_INDEXES + if value.strip().lower() == "all-quads": + return COTTAS_QUAD_INDEXES indexes = tuple(dict.fromkeys(part.strip().lower() for part in value.split(","))) - if any(index not in COTTAS_INDEXES for index in indexes): - raise ValueError("COTTAS indexes must be spo,sop,pso,pos,osp,ops or all") + if any(index not in COTTAS_ALL_INDEXES for index in indexes): + raise ValueError("COTTAS indexes must be permutations of spo or spog, all, or all-quads") return indexes +def require_dataset_indexes(indexes: tuple[str, ...]) -> None: + """Prevent triple-only orders from discarding a dataset's graph column.""" + if any("g" not in index for index in indexes): + raise ValueError("RDF datasets require indexes containing g (e.g. spog,gspo or all-quads)") + + def cottas_index_paths(path: Path, indexes: tuple[str, ...]) -> dict[str, Path]: """Keep the primary filename; suffix additional copies with their order.""" return {