From f181d0bd5119ee9a6e451b45cdd51380fddbb850 Mon Sep 17 00:00:00 2001 From: Trent Nelson Date: Tue, 29 Sep 2026 15:09:13 -0700 Subject: [PATCH 1/4] Decode FITS gzip tiles without stripping their framing Signed-off-by: Trent Nelson --- src/cuphoton/xdr/nvcomp_batch.py | 53 ++++++--- src/cuphoton/xdr/src/nvcomp_batch_ext.cpp | 126 +++++++++++++++++----- tests/xdr/test_gzip_gpu.py | 103 ++++++++++++++++++ 3 files changed, 242 insertions(+), 40 deletions(-) create mode 100644 tests/xdr/test_gzip_gpu.py diff --git a/src/cuphoton/xdr/nvcomp_batch.py b/src/cuphoton/xdr/nvcomp_batch.py index 5704dbd..8e0f6e3 100644 --- a/src/cuphoton/xdr/nvcomp_batch.py +++ b/src/cuphoton/xdr/nvcomp_batch.py @@ -2,7 +2,7 @@ # # SPDX-License-Identifier: Apache-2.0 -"""Device-in / device-out nvCOMP batched DEFLATE decompression. +"""Device-in / device-out nvCOMP batched gzip/DEFLATE decompression. Accepts a single contiguous device buffer holding N concatenated gzip-compressed tile payloads plus per-tile (offset, length) tables, and @@ -651,12 +651,19 @@ def bounded_groups(probe_lengths): for group in bounded_groups(fixed_lengths): group_indices = all_indices[group] peeks = gather_prefixes(group_indices, fixed_lengths[group]) - for cursor, index in enumerate(group_indices): - head = peeks[cursor * peek_len : (cursor + 1) * peek_len] - try: - header_sizes[index] = _parse_gzip_header_len(head) - except _TruncatedGzipHeader: - pending.append(int(index)) + fixed = np.frombuffer(peeks, dtype=np.uint8).reshape(-1, peek_len) + invalid = ( + (fixed[:, 0] != 0x1F) + | (fixed[:, 1] != 0x8B) + | (fixed[:, 2] != 0x08) + | ((fixed[:, 3] & 0xE0) != 0) + ) + if np.any(invalid): + # Preserve the parser's specific error for the first bad prefix. + first = int(np.flatnonzero(invalid)[0]) + _parse_gzip_header_len(fixed[first].tobytes()) + header_sizes[group_indices] = peek_len + pending.extend(group_indices[(fixed[:, 3] & 0x1E) != 0].tolist()) # Neither the trailer nor a valid two-byte DEFLATE stream can contain # header bytes. Re-gather still-truncated tiles in bounded groups per @@ -828,7 +835,8 @@ def gpu_gzip_decompress_batch( Per-tile metadata arrays. gzip_wrapped : bool If True (default), each tile is an RFC 1952 gzip stream; the - variable header + 8-byte trailer are stripped before nvcomp. + native gzip decoder consumes its header and trailer directly. Older + extensions and the Python fallback strip the framing for DEFLATE. use_cpp_helper : "auto" | True | False "auto" (default) → use the C++ pybind11 helper when importable, else fall back to the Python loop. True forces the C++ path (raises if @@ -921,11 +929,15 @@ def gpu_gzip_decompress_batch( if n == 0: return cp.empty(0, dtype=cp.uint8), out_offsets - d_concat, deflate_offsets, aligned_owners = _align_deflate_inputs( - d_concat, deflate_offsets, deflate_lengths - ) - if keepalive is not None: - keepalive.extend(aligned_owners) + gzip_fn = getattr(ext, "batch_gzip_decompress", None) + use_gzip = gzip_wrapped and gzip_fn is not None + aligned_owners = () + if not use_gzip: + d_concat, deflate_offsets, aligned_owners = _align_deflate_inputs( + d_concat, deflate_offsets, deflate_lengths + ) + if keepalive is not None: + keepalive.extend(aligned_owners) try: # Allocate one concatenated output buffer and slice it per tile. @@ -944,6 +956,21 @@ def gpu_gzip_decompress_batch( stream_ptr = int(cp.cuda.get_current_stream().ptr) d_concat_ptr = int(d_concat.data.ptr) d_out_ptr = int(d_out.data.ptr) + if use_gzip: + assert gzip_fn is not None + scratch_owner = gzip_fn( + d_concat_ptr, + offset_array, + length_array, + d_out_ptr, + out_offsets, + size_array, + stream_ptr, + use_native_pool and keepalive is not None, + ) + if keepalive is not None and scratch_owner is not None: + keepalive.append(scratch_owner) + return d_out, out_offsets pooled_fn = getattr(ext, "batch_deflate_decompress_pooled", None) if ( use_native_pool diff --git a/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp b/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp index 218a652..0785ad8 100644 --- a/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp +++ b/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp @@ -3,11 +3,11 @@ * SPDX-License-Identifier: Apache-2.0 */ -// Batched DEFLATE decompression — pybind11 helper that bypasses the +// Batched Gzip/DEFLATE decompression — pybind11 helper that bypasses the // per-tile Python `nvcomp.as_array` loop. // // Takes device pointers + offset/length arrays and calls -// `nvcompBatchedDeflateDecompressAsync` directly. The whole call releases the +// the corresponding nvCOMP batched API directly. The decode launch releases the // GIL so a Python prefetch thread can run during decode. #include "nvcomp_batch_ext.h" @@ -17,6 +17,7 @@ #include #include +#include #include #include @@ -36,7 +37,7 @@ inline void check_nvcomp(nvcompStatus_t s, const char* ctx) { } } -py::object batch_deflate_decompress_impl( +py::object batch_decompress_impl( std::uintptr_t d_concat_ptr, py::array_t rel_offsets, py::array_t lengths, @@ -44,7 +45,8 @@ py::object batch_deflate_decompress_impl( py::array_t out_offsets, py::array_t out_sizes, std::uintptr_t stream_ptr, - bool use_native_pool); + bool use_native_pool, + bool gzip_wrapped); // Batched DEFLATE decompress across N tiles packed inside one device buffer. // @@ -74,8 +76,16 @@ void batch_deflate_decompress( py::array_t out_offsets, py::array_t out_sizes, std::uintptr_t stream_ptr) { - batch_deflate_decompress_impl( - d_concat_ptr, rel_offsets, lengths, d_out_ptr, out_offsets, out_sizes, stream_ptr, false); + batch_decompress_impl( + d_concat_ptr, + rel_offsets, + lengths, + d_out_ptr, + out_offsets, + out_sizes, + stream_ptr, + false, + false); } py::object batch_deflate_decompress_pooled( @@ -86,11 +96,19 @@ py::object batch_deflate_decompress_pooled( py::array_t out_offsets, py::array_t out_sizes, std::uintptr_t stream_ptr) { - return batch_deflate_decompress_impl( - d_concat_ptr, rel_offsets, lengths, d_out_ptr, out_offsets, out_sizes, stream_ptr, true); + return batch_decompress_impl( + d_concat_ptr, + rel_offsets, + lengths, + d_out_ptr, + out_offsets, + out_sizes, + stream_ptr, + true, + false); } -py::object batch_deflate_decompress_impl( +py::object batch_gzip_decompress( std::uintptr_t d_concat_ptr, py::array_t rel_offsets, py::array_t lengths, @@ -99,6 +117,28 @@ py::object batch_deflate_decompress_impl( py::array_t out_sizes, std::uintptr_t stream_ptr, bool use_native_pool) { + return batch_decompress_impl( + d_concat_ptr, + rel_offsets, + lengths, + d_out_ptr, + out_offsets, + out_sizes, + stream_ptr, + use_native_pool, + true); +} + +py::object batch_decompress_impl( + std::uintptr_t d_concat_ptr, + py::array_t rel_offsets, + py::array_t lengths, + std::uintptr_t d_out_ptr, + py::array_t out_offsets, + py::array_t out_sizes, + std::uintptr_t stream_ptr, + bool use_native_pool, + bool gzip_wrapped) { const std::size_t n = static_cast(rel_offsets.size()); if (lengths.size() != static_cast(n) || out_offsets.size() != static_cast(n) @@ -130,7 +170,7 @@ py::object batch_deflate_decompress_impl( } h_comp_ptrs[i] = reinterpret_cast(d_concat_ptr + static_cast(ro(i))); - if (reinterpret_cast(h_comp_ptrs[i]) % 4 != 0) { + if (!gzip_wrapped && reinterpret_cast(h_comp_ptrs[i]) % 4 != 0) { throw std::invalid_argument("raw DEFLATE inputs must be 4-byte aligned"); } h_out_ptrs[i] = reinterpret_cast(d_out_ptr + static_cast(oo(i))); @@ -153,10 +193,17 @@ py::object batch_deflate_decompress_impl( // The hardware backend requires non-stream-ordered scratch allocations. // Keep this path on CUDA until that ownership contract is supported. opts.backend = NVCOMP_DECOMPRESS_BACKEND_CUDA; + auto gzip_opts = nvcompBatchedGzipDecompressDefaultOpts; + gzip_opts.backend = NVCOMP_DECOMPRESS_BACKEND_CUDA; + // NAIVE accepts byte-aligned FITS gzip tile starts. LOOKAHEAD requires + // additional input alignment and is intended for much larger chunks. + gzip_opts.algorithm = NVCOMP_GZIP_DECOMPRESS_ALGORITHM_NAIVE; std::size_t temp_bytes = 0; check_nvcomp( - nvcompBatchedDeflateDecompressGetTempSizeAsync( - n, max_uncomp, opts, &temp_bytes, total_uncomp), + gzip_wrapped ? nvcompBatchedGzipDecompressGetTempSizeAsync( + n, max_uncomp, gzip_opts, &temp_bytes, total_uncomp) + : nvcompBatchedDeflateDecompressGetTempSizeAsync( + n, max_uncomp, opts, &temp_bytes, total_uncomp), "get temp size"); std::vector> pooled_keepalive; @@ -250,21 +297,33 @@ py::object batch_deflate_decompress_impl( // can progress while the async kernel launches + returns. { py::gil_scoped_release release; - // Supply both output buffers for nvCOMP 5.2 compatibility. + // nvCOMP 5.2 requires a non-null status buffer for both codecs. check_nvcomp( - nvcompBatchedDeflateDecompressAsync( - static_cast(d_comp_ptrs), - static_cast(d_comp_sizes), - static_cast(d_out_sizes), - static_cast(d_actual_sizes), - n, - d_temp, - temp_bytes, - static_cast(d_out_ptrs), - opts, - static_cast(d_statuses), - stream), - "nvcompBatchedDeflateDecompressAsync"); + gzip_wrapped ? nvcompBatchedGzipDecompressAsync( + static_cast(d_comp_ptrs), + static_cast(d_comp_sizes), + static_cast(d_out_sizes), + static_cast(d_actual_sizes), + n, + d_temp, + temp_bytes, + static_cast(d_out_ptrs), + gzip_opts, + static_cast(d_statuses), + stream) + : nvcompBatchedDeflateDecompressAsync( + static_cast(d_comp_ptrs), + static_cast(d_comp_sizes), + static_cast(d_out_sizes), + static_cast(d_actual_sizes), + n, + d_temp, + temp_bytes, + static_cast(d_out_ptrs), + opts, + static_cast(d_statuses), + stream), + "batched gzip/deflate decompression"); } if (use_native_pool) { @@ -294,7 +353,7 @@ py::object batch_deflate_decompress_impl( } // namespace xdr_gpu PYBIND11_MODULE(_nvcomp_batch_ext, m) { - m.doc() = "Batched DEFLATE decompression — device-pointer interface to nvcomp."; + m.doc() = "Batched Gzip/DEFLATE decompression — device-pointer interface to nvcomp."; xdr_gpu::bind_io(m); xdr_gpu::bind_memory_manager(m); @@ -324,4 +383,17 @@ PYBIND11_MODULE(_nvcomp_batch_ext, m) { py::arg("stream_ptr") = 0, "Batched DEFLATE decompress using native pooled scratch buffers.\n" "Returns an owner capsule that must stay alive until stream work completes."); + m.def( + "batch_gzip_decompress", + &xdr_gpu::batch_gzip_decompress, + py::arg("d_concat_ptr"), + py::arg("rel_offsets"), + py::arg("lengths"), + py::arg("d_out_ptr"), + py::arg("out_offsets"), + py::arg("out_sizes"), + py::arg("stream_ptr") = 0, + py::arg("use_native_pool") = false, + "Batched Gzip decompress including RFC-1952 headers and trailers.\n" + "Pooled calls return an owner capsule retained until stream completion."); } diff --git a/tests/xdr/test_gzip_gpu.py b/tests/xdr/test_gzip_gpu.py new file mode 100644 index 0000000..ba157af --- /dev/null +++ b/tests/xdr/test_gzip_gpu.py @@ -0,0 +1,103 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# +# SPDX-License-Identifier: Apache-2.0 + +"""Native gzip tile decoding with standard and optional RFC 1952 headers.""" + +from __future__ import annotations + +import gzip +import struct +import zlib + +import numpy as np +import pytest + +from cuphoton.xdr import nvcomp_batch + + +@pytest.fixture +def cp(): + module = pytest.importorskip("cupy") + try: + if module.cuda.runtime.getDeviceCount() == 0: + pytest.skip("CUDA device unavailable") + except module.cuda.runtime.CUDARuntimeError: + pytest.skip("CUDA runtime unavailable") + ext = nvcomp_batch._try_get_cpp_ext() + if ext is None or not hasattr(ext, "batch_gzip_decompress"): + pytest.skip("native gzip extension unavailable") + return module + + +@pytest.mark.parametrize("pooled", [False, True]) +@pytest.mark.parametrize("preparsed", [False, True]) +def test_native_gzip_decodes_optional_headers(cp, pooled, preparsed): + payloads, originals, header_sizes = [], [], [] + for flags in range(32): + original = bytes(range(256)) * 9 + bytes([flags]) * 11 + payload = _gzip_with_flags(original, flags) + assert gzip.decompress(payload) == original + payloads.append(payload) + originals.append(original) + header_sizes.append(nvcomp_batch._parse_gzip_header_len(payload)) + lengths = np.array([len(payload) for payload in payloads], dtype="i8") + offsets = np.cumsum(lengths) - lengths + sizes = np.array([len(value) for value in originals], dtype="i8") + # Dense tile starts deliberately do not obey raw DEFLATE's alignment. + assert np.any(offsets % 4 != 0) + stream = cp.cuda.Stream(non_blocking=True) + owners = [] + with stream: + data = cp.asarray(np.frombuffer(b"".join(payloads), dtype="u1")) + result, out_offsets = nvcomp_batch.gpu_gzip_decompress_batch( + data, + offsets, + lengths, + sizes, + use_cpp_helper=True, + use_native_pool=pooled, + keepalive=owners, + header_sizes=header_sizes if preparsed else None, + ) + stream.synchronize() + np.testing.assert_array_equal(out_offsets, np.cumsum(sizes) - sizes) + np.testing.assert_array_equal( + cp.asnumpy(result), np.frombuffer(b"".join(originals), dtype="u1") + ) + + +@pytest.mark.parametrize( + "header,error", + [ + (b"\x00\x8b\x08\x00" + bytes(6), "wrong magic"), + (b"\x1f\x8b\x09\x00" + bytes(6), "DEFLATE compression"), + (b"\x1f\x8b\x08\x80" + bytes(6), "reserved header flags"), + (b"\x1f\x8b\x08\x08" + bytes(6) + b"name", "truncated gzip"), + ], +) +def test_native_gzip_rejects_invalid_headers_before_decode(cp, header, error): + # Ten final bytes reserve the two-byte payload and eight-byte trailer. + data = cp.asarray(np.frombuffer(header + bytes(10), dtype="u1")) + with pytest.raises(ValueError, match=error): + nvcomp_batch.gpu_gzip_decompress_batch( + data, [0], [data.size], [16], use_cpp_helper=True + ) + + +def _gzip_with_flags(data, flags): + header = b"\x1f\x8b\x08" + bytes([flags]) + bytes(6) + if flags & 4: + extra = b"extra\x00payload" + header += struct.pack(" Date: Tue, 29 Sep 2026 18:40:10 -0700 Subject: [PATCH 2/4] Allow explicit selection of FITS gzip decoders Align raw DEFLATE payloads on the device before decoding and retain those buffers through asynchronous work. Keep automatic selection compatible with older native extensions. Signed-off-by: Trent Nelson --- src/cuphoton/xdr/nvcomp_batch.py | 27 ++++++++++-- tests/xdr/test_gzip_gpu.py | 73 +++++++++++++++++++++++++++++++- tests/xdr/test_nvcomp_batch.py | 71 +++++++++++++++++++++++++++++++ 3 files changed, 166 insertions(+), 5 deletions(-) diff --git a/src/cuphoton/xdr/nvcomp_batch.py b/src/cuphoton/xdr/nvcomp_batch.py index 8e0f6e3..a28e957 100644 --- a/src/cuphoton/xdr/nvcomp_batch.py +++ b/src/cuphoton/xdr/nvcomp_batch.py @@ -24,11 +24,12 @@ import sys from collections.abc import Sequence from pathlib import Path +from typing import Literal import numpy as np -_REPACK_DEFLATE_KERNEL = None _DEFLATE_CODEC = None +_REPACK_DEFLATE_KERNEL = None _CPP_EXT = None _CPP_EXT_PROBED = False _CPP_EXT_IMPORT_ERROR: str | None = None @@ -820,6 +821,7 @@ def gpu_gzip_decompress_batch( uncompressed_sizes: Sequence[int], *, gzip_wrapped: bool = True, + gzip_decoder: Literal["auto", "gzip", "deflate"] = "auto", use_cpp_helper: str | bool = "auto", use_native_pool: bool = False, keepalive: list | None = None, @@ -837,6 +839,10 @@ def gpu_gzip_decompress_batch( If True (default), each tile is an RFC 1952 gzip stream; the native gzip decoder consumes its header and trailer directly. Older extensions and the Python fallback strip the framing for DEFLATE. + gzip_decoder : "auto" | "gzip" | "deflate" + "auto" uses native Gzip when available, otherwise raw DEFLATE. + "gzip" requires the native Gzip helper and gzip-wrapped inputs. + "deflate" strips gzip framing and aligns payloads before decoding. use_cpp_helper : "auto" | True | False "auto" (default) → use the C++ pybind11 helper when importable, else fall back to the Python loop. True forces the C++ path (raises if @@ -859,6 +865,11 @@ def gpu_gzip_decompress_batch( out_offsets : numpy.ndarray (int64) Start offset of each decompressed tile inside ``d_out``. """ + if gzip_decoder not in ("auto", "gzip", "deflate"): + raise ValueError("gzip_decoder must be 'auto', 'gzip', or 'deflate'") + if gzip_decoder == "gzip" and not gzip_wrapped: + raise ValueError("gzip_decoder='gzip' requires gzip_wrapped=True") + import cupy as cp concat_size = int(d_concat.size) @@ -916,7 +927,7 @@ def gpu_gzip_decompress_batch( "`bash src/cuphoton/xdr/src/build.sh` from source. " f"Import error: {_CPP_EXT_IMPORT_ERROR}" ) - elif n == 0: + elif n == 0 and gzip_decoder != "gzip": # No backend work is needed, so an automatic fallback is not useful. ext = None elif use_cpp_helper is False: @@ -926,11 +937,19 @@ def gpu_gzip_decompress_batch( if ext is None: _warn_python_fallback_once() + gzip_fn = getattr(ext, "batch_gzip_decompress", None) + if gzip_decoder == "gzip" and gzip_fn is None: + raise RuntimeError( + "gzip_decoder='gzip' requires a native extension with Gzip " + "support and use_cpp_helper enabled" + ) + use_gzip = ( + gzip_wrapped and gzip_decoder != "deflate" and gzip_fn is not None + ) + if n == 0: return cp.empty(0, dtype=cp.uint8), out_offsets - gzip_fn = getattr(ext, "batch_gzip_decompress", None) - use_gzip = gzip_wrapped and gzip_fn is not None aligned_owners = () if not use_gzip: d_concat, deflate_offsets, aligned_owners = _align_deflate_inputs( diff --git a/tests/xdr/test_gzip_gpu.py b/tests/xdr/test_gzip_gpu.py index ba157af..aa57637 100644 --- a/tests/xdr/test_gzip_gpu.py +++ b/tests/xdr/test_gzip_gpu.py @@ -9,6 +9,7 @@ import gzip import struct import zlib +from types import SimpleNamespace import numpy as np import pytest @@ -30,9 +31,10 @@ def cp(): return module +@pytest.mark.parametrize("decoder", ["auto", "gzip", "deflate"]) @pytest.mark.parametrize("pooled", [False, True]) @pytest.mark.parametrize("preparsed", [False, True]) -def test_native_gzip_decodes_optional_headers(cp, pooled, preparsed): +def test_native_gzip_decodes_optional_headers(cp, pooled, preparsed, decoder): payloads, originals, header_sizes = [], [], [] for flags in range(32): original = bytes(range(256)) * 9 + bytes([flags]) * 11 @@ -56,6 +58,7 @@ def test_native_gzip_decodes_optional_headers(cp, pooled, preparsed): lengths, sizes, use_cpp_helper=True, + gzip_decoder=decoder, use_native_pool=pooled, keepalive=owners, header_sizes=header_sizes if preparsed else None, @@ -101,3 +104,71 @@ def _gzip_with_flags(data, flags): + zlib.compress(data, wbits=-15) + struct.pack(" Date: Tue, 29 Sep 2026 18:53:05 -0700 Subject: [PATCH 3/4] Expose gzip decoder choice through FITS APIs and CLI Signed-off-by: Trent Nelson --- docs/components/xdr.md | 8 ++++++++ src/cuphoton/core/cli/fits.py | 11 +++++++++++ src/cuphoton/core/fits_options.py | 1 + src/cuphoton/xdr/benchmark_fits.py | 14 ++++++++++++-- src/cuphoton/xdr/convenience.py | 5 ++++- src/cuphoton/xdr/nvcomp_batch.py | 3 +-- src/cuphoton/xdr/prefetch.py | 25 +++++++++++++++++++++++-- src/cuphoton/xdr/reader.py | 16 ++++++++++++++-- tests/core/test_cli_contract.py | 12 ++++++------ tests/core/test_fits_io.py | 16 +++++++++++++--- tests/xdr/test_benchmark_fits.py | 4 ++++ tests/xdr/test_frontend.py | 3 +++ 12 files changed, 100 insertions(+), 18 deletions(-) diff --git a/docs/components/xdr.md b/docs/components/xdr.md index 0cadf87..78f35da 100644 --- a/docs/components/xdr.md +++ b/docs/components/xdr.md @@ -37,6 +37,14 @@ CLI flag is `--xdr-postprocess {auto,fused,separate}`, including on `cuphoton xdr benchmark-fits`. Explicit CLI values override matching descriptor fields; omitted flags preserve the descriptor's choices. +Gzip decoding defaults to `gzip_decoder="auto"`: use the native Gzip helper +when available, or aligned raw DEFLATE with an older extension or the Python +fallback. Select `gzip_decoder="gzip"` to require native Gzip support, or +`gzip_decoder="deflate"` to strip the wrapper and decode aligned raw payloads. +The same choice is available as `xdr_options={"gzip_decoder": "deflate"}` in +workflow APIs and descriptors, or `--xdr-gzip-decoder {auto,gzip,deflate}` on +the CLI. Both decoders preserve the FITS pixel values. + These controls apply when the existing reader policy selects xDR. Reader selection and CPU fallback remain controlled by `--fits-reader`. Invalid options fail during validation; decode, I/O, and CUDA failures propagate to the caller. diff --git a/src/cuphoton/core/cli/fits.py b/src/cuphoton/core/cli/fits.py index c0c919f..e5f6e38 100644 --- a/src/cuphoton/core/cli/fits.py +++ b/src/cuphoton/core/cli/fits.py @@ -27,6 +27,17 @@ class XdrPostprocessArg(SetInvariant): _set = XDR_OPTION_CHOICES["postprocess"] _default = None + xdr_gzip_decoder = None + + class XdrGzipDecoderArg(SetInvariant): + _arg = "--xdr-gzip-decoder" + _help = ( + "xDR gzip decoder: auto, gzip, or deflate. Omitted preserves " + "manifest choices; applies when the FITS reader uses xDR." + ) + _set = XDR_OPTION_CHOICES["gzip_decoder"] + _default = None + def xdr_options_from_cli(command) -> dict[str, str]: """Collect explicit flags without replacing omitted manifest choices.""" diff --git a/src/cuphoton/core/fits_options.py b/src/cuphoton/core/fits_options.py index c005a2f..c098a09 100644 --- a/src/cuphoton/core/fits_options.py +++ b/src/cuphoton/core/fits_options.py @@ -8,6 +8,7 @@ XDR_OPTION_CHOICES = { "postprocess": frozenset({"auto", "fused", "separate"}), + "gzip_decoder": frozenset({"auto", "gzip", "deflate"}), } diff --git a/src/cuphoton/xdr/benchmark_fits.py b/src/cuphoton/xdr/benchmark_fits.py index 8f58957..0b49e07 100644 --- a/src/cuphoton/xdr/benchmark_fits.py +++ b/src/cuphoton/xdr/benchmark_fits.py @@ -163,6 +163,7 @@ class _BatchOptions(TypedDict): native_plan_threads: int native_batcher: str | bool postprocess: str + gzip_decoder: str def parse_hdu_indices(value: str) -> tuple[int, ...]: @@ -420,6 +421,7 @@ def bench_batch_load( use_stream: bool, data_mb: float, postprocess: str = "auto", + gzip_decoder: str = "auto", ) -> PhaseResult: fn = batch_to_device_stream if use_stream else batch_to_device phase = "batch_to_device_stream" if use_stream else "batch_to_device" @@ -432,12 +434,14 @@ def bench_batch_load( native_plan_threads=native_plan_threads, native_batcher=native_batcher, postprocess=postprocess, + gzip_decoder=gzip_decoder, ) times: list[float] = [] note = ( f"{len(paths)} files x {len(tuple(hdu_indices))} HDUs, " - f"native_batcher={native_batcher}, postprocess={postprocess}" + f"native_batcher={native_batcher}, postprocess={postprocess}, " + f"gzip_decoder={gzip_decoder}" ) with nvtx_range(f"xdr.{phase}"): try: @@ -559,6 +563,7 @@ def run_benchmark( native_plan_threads: int = max(1, os.cpu_count() or 1), native_batcher: str = "auto", postprocess: str = "auto", + gzip_decoder: str = "auto", mock_storage_kind: str | None = None, skip_gds_read: bool = False, output_json: Path | None = None, @@ -571,7 +576,9 @@ def run_benchmark( helpers; the return value remains a list of phases. """ - normalize_xdr_options({"postprocess": postprocess}) + normalize_xdr_options( + {"postprocess": postprocess, "gzip_decoder": gzip_decoder} + ) paths = resolve_paths( fits_files, scan_dir=scan_dir, @@ -648,6 +655,7 @@ def run_benchmark( native_plan_threads=native_plan_threads, native_batcher=resolved_native_batcher, postprocess=postprocess, + gzip_decoder=gzip_decoder, use_stream=False, data_mb=total_data_mb, ) @@ -664,6 +672,7 @@ def run_benchmark( native_plan_threads=native_plan_threads, native_batcher=resolved_native_batcher, postprocess=postprocess, + gzip_decoder=gzip_decoder, use_stream=True, data_mb=total_data_mb, ) @@ -707,6 +716,7 @@ def run_benchmark( "native_plan_threads": native_plan_threads, "native_batcher": native_batcher, "postprocess": postprocess, + "gzip_decoder": gzip_decoder, "native_batcher_enabled": batcher_enabled, "native_batcher_error": batcher_error, "skip_gds_read": skip_gds_read, diff --git a/src/cuphoton/xdr/convenience.py b/src/cuphoton/xdr/convenience.py index b601bd8..da46cf8 100644 --- a/src/cuphoton/xdr/convenience.py +++ b/src/cuphoton/xdr/convenience.py @@ -43,6 +43,7 @@ def batch_to_device( native_plan_threads: int | None = None, native_batcher: str | bool = "auto", postprocess: str = "auto", + gzip_decoder: str = "auto", ): """Load image HDUs from FITS files into stacked device arrays. @@ -64,7 +65,8 @@ def batch_to_device( When True, use the streaming batch reader. Set False for a depth-1 path that preserves the existing API without requiring Astropy HDU objects. prefetch_depth, decode_batch_files, batch_queue_depth, - native_read_threads, native_plan_threads, native_batcher, postprocess + native_read_threads, native_plan_threads, native_batcher, postprocess, + gzip_decoder Passed through to `batch_to_device_stream` when ``parallel=True``. Returns @@ -89,6 +91,7 @@ def batch_to_device( native_plan_threads=native_plan_threads, native_batcher=native_batcher, postprocess=postprocess, + gzip_decoder=gzip_decoder, section=section, stream=stream, ) diff --git a/src/cuphoton/xdr/nvcomp_batch.py b/src/cuphoton/xdr/nvcomp_batch.py index a28e957..b9807f6 100644 --- a/src/cuphoton/xdr/nvcomp_batch.py +++ b/src/cuphoton/xdr/nvcomp_batch.py @@ -24,7 +24,6 @@ import sys from collections.abc import Sequence from pathlib import Path -from typing import Literal import numpy as np @@ -821,7 +820,7 @@ def gpu_gzip_decompress_batch( uncompressed_sizes: Sequence[int], *, gzip_wrapped: bool = True, - gzip_decoder: Literal["auto", "gzip", "deflate"] = "auto", + gzip_decoder: str = "auto", use_cpp_helper: str | bool = "auto", use_native_pool: bool = False, keepalive: list | None = None, diff --git a/src/cuphoton/xdr/prefetch.py b/src/cuphoton/xdr/prefetch.py index 2bfb7a8..35ee932 100644 --- a/src/cuphoton/xdr/prefetch.py +++ b/src/cuphoton/xdr/prefetch.py @@ -602,6 +602,7 @@ def _consume_comp_batch( keepalive=None, *, postprocess="auto", + gzip_decoder: str = "auto", ): """Decode many compressed HDUs with one batched nvCOMP call. @@ -717,6 +718,7 @@ def _consume_comp_batch( lengths=lengths, out_bytes=out_bytes, stream=stream, + gzip_decoder=gzip_decoder, keepalive=keepalive, header_sizes=header_sizes, ) @@ -801,6 +803,7 @@ def _consume_prefetched_group( keepalive=None, *, postprocess="auto", + gzip_decoder: str = "auto", ): """Consume prefetched files, batching compressed HDUs together.""" comp_entries: list[_CompBatchEntry] = [] @@ -833,7 +836,11 @@ def _consume_prefetched_group( raise AssertionError(f"unknown kind {plan_item.kind!r}") _consume_comp_batch( - comp_entries, stream, keepalive=keepalive, postprocess=postprocess + comp_entries, + stream, + keepalive=keepalive, + postprocess=postprocess, + gzip_decoder=gzip_decoder, ) @@ -1066,6 +1073,7 @@ def _submit_prefetched_group( owner, in_flight: deque[_GpuBatchHandle], postprocess="auto", + gzip_decoder: str = "auto", ) -> None: """Queue GPU work and register its lifetime handle.""" import cupy as cp @@ -1091,6 +1099,7 @@ def _submit_prefetched_group( stream, keepalive=keepalive, postprocess=postprocess, + gzip_decoder=gzip_decoder, ) candidate_event = cp.cuda.Event() candidate_event.record(use_stream) @@ -1417,6 +1426,7 @@ def _consume_native_batches( NativeBatchBuilder, native_plan_files, postprocess="auto", + gzip_decoder: str = "auto", ): """Consume device batches built by the C++ KvikIO worker pool.""" import cupy as cp @@ -1493,6 +1503,7 @@ def _consume_native_batches( owner=native_batch, in_flight=in_flight, postprocess=postprocess, + gzip_decoder=gzip_decoder, ) planner.join() if planner.error is not None: @@ -1533,6 +1544,7 @@ def _consume_python_batches( section, stream, postprocess="auto", + gzip_decoder: str = "auto", ) -> None: """Consume pinned-host batches with bounded event-owned lifetimes.""" n_files = len(paths) @@ -1554,6 +1566,7 @@ def submit(group: list[PrefetchedFile]) -> None: owner=group, in_flight=in_flight, postprocess=postprocess, + gzip_decoder=gzip_decoder, ) prefetcher.start() @@ -1622,6 +1635,7 @@ def batch_to_device_stream( native_plan_threads: int | None = None, native_batcher: str | bool = "auto", postprocess: str = "auto", + gzip_decoder: str = "auto", section=None, stream=None, ): @@ -1668,6 +1682,9 @@ def batch_to_device_stream( postprocess "auto" (default) and "fused" restore FITS pixel order in one kernel. "separate" uses individual unshuffle, byteswap and scatter kernels. + gzip_decoder + "auto" uses native Gzip when available, otherwise aligned DEFLATE. + "gzip" requires native Gzip support; "deflate" selects raw DEFLATE. section Optional 2D ROI applied uniformly to CompImageHDUs. stream @@ -1688,7 +1705,9 @@ def batch_to_device_stream( length ``len(paths)``. """ - normalize_xdr_options({"postprocess": postprocess}) + normalize_xdr_options( + {"postprocess": postprocess, "gzip_decoder": gzip_decoder} + ) hdu_indices = tuple(int(i) for i in hdu_indices) resolved_paths = [Path(p) for p in paths] n_files = len(resolved_paths) @@ -1739,6 +1758,7 @@ def batch_to_device_stream( NativeBatchBuilder=NativeBatchBuilder, native_plan_files=native_plan_files, postprocess=postprocess, + gzip_decoder=gzip_decoder, ) return tuple(outs) @@ -1752,6 +1772,7 @@ def batch_to_device_stream( section=section, stream=stream, postprocess=postprocess, + gzip_decoder=gzip_decoder, ) return tuple(outs) diff --git a/src/cuphoton/xdr/reader.py b/src/cuphoton/xdr/reader.py index a591ba9..c23d278 100644 --- a/src/cuphoton/xdr/reader.py +++ b/src/cuphoton/xdr/reader.py @@ -437,6 +437,7 @@ def inflate_device_heap( stream=None, keepalive=None, header_sizes=None, + gzip_decoder: str = "auto", ): """Run batched nvCOMP inflate for gzip tiles resident on the device. @@ -446,6 +447,7 @@ def inflate_device_heap( byte counts. Returns concatenated decompressed bytes plus per-tile offsets into that output buffer. """ + normalize_xdr_options({"gzip_decoder": gzip_decoder}) import cupy as cp with stream or cp.cuda.Stream.null: @@ -455,6 +457,7 @@ def inflate_device_heap( lengths, out_bytes, gzip_wrapped=True, + gzip_decoder=gzip_decoder, use_native_pool=keepalive is not None, keepalive=keepalive, header_sizes=header_sizes, @@ -589,6 +592,7 @@ def decode_from_device_heap( stream=None, keepalive=None, postprocess: str = "auto", + gzip_decoder: str = "auto", ): """Run the decode pipeline on a pre-loaded compressed heap buffer. @@ -597,7 +601,9 @@ def decode_from_device_heap( start offset of tile i inside `d_concat`. This is what the prefetch consumer calls after staging its pinned host heap to device. """ - normalize_xdr_options({"postprocess": postprocess}) + normalize_xdr_options( + {"postprocess": postprocess, "gzip_decoder": gzip_decoder} + ) if keepalive is not None: keepalive.append(d_concat) d_pixels, tile_byte_offsets_np = ( @@ -607,6 +613,7 @@ def decode_from_device_heap( lengths=plan["sel_lengths"], out_bytes=plan["out_bytes"], stream=stream, + gzip_decoder=gzip_decoder, keepalive=keepalive, ) ) @@ -630,6 +637,7 @@ def read( section=None, loader=None, postprocess: str = "auto", + gzip_decoder: str = "auto", ): """Read and decode this HDU into a `cupy.ndarray`. @@ -649,7 +657,9 @@ def read( Failure before an I/O handle is returned leaves completion unknown and requires a process restart before further GPU submissions. """ - normalize_xdr_options({"postprocess": postprocess}) + normalize_xdr_options( + {"postprocess": postprocess, "gzip_decoder": gzip_decoder} + ) import cupy as cp plan = self.prepare_plan(section=section) @@ -724,6 +734,7 @@ def abandon(error): out=out, stream=None, postprocess=postprocess, + gzip_decoder=gzip_decoder, ) except BaseException as error: abandon(error) @@ -766,6 +777,7 @@ def abandon(error): stream=stream, keepalive=keepalive, postprocess=postprocess, + gzip_decoder=gzip_decoder, ) except (KeyboardInterrupt, SystemExit) as error: abandon(error) diff --git a/tests/core/test_cli_contract.py b/tests/core/test_cli_contract.py index a97dc7a..939af0a 100644 --- a/tests/core/test_cli_contract.py +++ b/tests/core/test_cli_contract.py @@ -70,18 +70,18 @@ def test_public_command_surface_counts_are_exact() -> None: ) assert per_group == [ - ("xdr", 1, 0, 15, 0), - ("xfit", 3, 3, 30, 1), - ("xpois", 7, 7, 154, 1), - ("xscan", 44, 44, 214, 1), - ("xrep", 6, 6, 113, 1), + ("xdr", 1, 0, 16, 0), + ("xfit", 3, 3, 33, 1), + ("xpois", 7, 7, 158, 1), + ("xscan", 44, 44, 219, 1), + ("xrep", 6, 6, 118, 1), ("xray", 33, 31, 392, 1), ] assert len(per_group) == 6 assert sum(item[1] for item in per_group) == 94 assert sum(item[1] + item[4] for item in per_group) == 99 assert sum(item[2] for item in per_group) == 91 - assert sum(item[3] for item in per_group) == 918 + assert sum(item[3] for item in per_group) == 936 def test_public_registry_order_and_component_derivations() -> None: diff --git a/tests/core/test_fits_io.py b/tests/core/test_fits_io.py index 438f43e..f8e171d 100644 --- a/tests/core/test_fits_io.py +++ b/tests/core/test_fits_io.py @@ -16,7 +16,10 @@ @pytest.mark.parametrize("postprocess", ["fused", "separate"]) -def test_xdr_options_reach_batch_dispatch(tmp_path, monkeypatch, postprocess): +@pytest.mark.parametrize("gzip_decoder", ["auto", "gzip", "deflate"]) +def test_xdr_options_reach_batch_dispatch( + tmp_path, monkeypatch, postprocess, gzip_decoder +): import cuphoton.xdr as xdr data = np.arange(12, dtype=np.float32).reshape(3, 4) @@ -34,10 +37,17 @@ def batch(paths, hdus, **options): [1], reader="xdr", device=True, - xdr_options={"postprocess": postprocess}, + xdr_options={ + "postprocess": postprocess, + "gzip_decoder": gzip_decoder, + }, ) assert observed["postprocess"] == postprocess - assert result.metadata()["xdr_options"] == {"postprocess": postprocess} + assert observed["gzip_decoder"] == gzip_decoder + assert result.metadata()["xdr_options"] == { + "postprocess": postprocess, + "gzip_decoder": gzip_decoder, + } np.testing.assert_array_equal(result.arrays[0], data) diff --git a/tests/xdr/test_benchmark_fits.py b/tests/xdr/test_benchmark_fits.py index 487ceee..47dbb27 100644 --- a/tests/xdr/test_benchmark_fits.py +++ b/tests/xdr/test_benchmark_fits.py @@ -165,6 +165,8 @@ def fake_run_benchmark(fits_files, **kwargs): "off", "--xdr-postprocess", "separate", + "--xdr-gzip-decoder", + "deflate", "--mock-storage", "host", "--skip-gds-read", @@ -186,6 +188,7 @@ def fake_run_benchmark(fits_files, **kwargs): assert received["hdu_indices"] == [2, 1, 2] assert received["native_batcher"] == "off" assert received["postprocess"] == "separate" + assert received["gzip_decoder"] == "deflate" assert received["mock_storage_kind"] == "host" assert received["skip_gds_read"] is True assert received["output_json"] == Path("report.json") @@ -400,6 +403,7 @@ def test_json_report_matches_returned_and_printed_phases( "native_plan_threads": 6, "native_batcher": "auto", "postprocess": "auto", + "gzip_decoder": "auto", "native_batcher_enabled": mode == "real", "native_batcher_error": None, "skip_gds_read": mode != "real", diff --git a/tests/xdr/test_frontend.py b/tests/xdr/test_frontend.py index 23bf49d..21af613 100644 --- a/tests/xdr/test_frontend.py +++ b/tests/xdr/test_frontend.py @@ -46,6 +46,7 @@ def test_batch_api_keeps_existing_keyword_surface(): "native_plan_threads", "native_batcher", "postprocess", + "gzip_decoder", ) actual = tuple(inspect.signature(xdr.batch_to_device).parameters) @@ -77,6 +78,7 @@ def fake_stream(paths, **kwargs): native_plan_threads=7, native_batcher=True, postprocess="separate", + gzip_decoder="deflate", ) assert result == ("ok",) @@ -93,6 +95,7 @@ def fake_stream(paths, **kwargs): "native_plan_threads": 7, "native_batcher": True, "postprocess": "separate", + "gzip_decoder": "deflate", "section": (slice(0, 1), slice(0, 1)), "stream": "stream", }, From 735b135081f15958477d95eddd9a2d846b62a172 Mon Sep 17 00:00:00 2001 From: Trent Nelson Date: Tue, 29 Sep 2026 19:47:14 -0700 Subject: [PATCH 4/4] Bind xDR fallback decoding to the caller stream Bind Python DEFLATE decoding to the current device and CUDA stream. Retain that stream until its cached codec is destroyed, ordering decode after compressed-input uploads and alignment copies. Keep Gzip dispatch and alignment covered by CPU checks and a delayed producer GPU regression. Signed-off-by: Trent Nelson --- src/cuphoton/xdr/nvcomp_batch.py | 57 +++++-- src/cuphoton/xdr/src/nvcomp_batch_ext.cpp | 5 +- tests/xdr/test_gzip_gpu.py | 48 +++++- tests/xdr/test_nvcomp_batch.py | 193 +++++++++++++++++++++- 4 files changed, 281 insertions(+), 22 deletions(-) diff --git a/src/cuphoton/xdr/nvcomp_batch.py b/src/cuphoton/xdr/nvcomp_batch.py index b9807f6..8451ef4 100644 --- a/src/cuphoton/xdr/nvcomp_batch.py +++ b/src/cuphoton/xdr/nvcomp_batch.py @@ -22,12 +22,13 @@ import platform import struct import sys +import threading from collections.abc import Sequence from pathlib import Path import numpy as np -_DEFLATE_CODEC = None +_DEFLATE_CODEC_CACHE = threading.local() _REPACK_DEFLATE_KERNEL = None _CPP_EXT = None _CPP_EXT_PROBED = False @@ -225,18 +226,39 @@ def _preload_cpp_ext_libraries() -> None: _CPP_EXT_LIBS_PRELOADED = True -def _get_codec(): - """Return a cached nvcomp deflate+RAW codec. Creating a codec allocates a - CUDA stream (~3-4 ms); cache it.""" - global _DEFLATE_CODEC - if _DEFLATE_CODEC is None: +class _StreamBoundCodec: + """Destroy the codec before releasing its CuPy stream owner.""" + + def __init__(self, stream, device_id): import nvidia.nvcomp as nvcomp - _DEFLATE_CODEC = nvcomp.Codec( + self.stream = stream + self.device_id = device_id + self.codec = nvcomp.Codec( algorithm="deflate", bitstream_kind=nvcomp.BitstreamKind.RAW, + cuda_stream=int(stream.ptr), ) - return _DEFLATE_CODEC + + def __del__(self): + self.codec = None + + +def _get_codec(): + """Reuse one codec per thread, bound to its current device and stream.""" + import cupy as cp + + stream = cp.cuda.get_current_stream() + device_id = cp.cuda.runtime.getDevice() + cached = getattr(_DEFLATE_CODEC_CACHE, "entry", None) + if ( + cached is None + or cached.device_id != device_id + or cached.stream.ptr != stream.ptr + ): + cached = _StreamBoundCodec(stream, device_id) + _DEFLATE_CODEC_CACHE.entry = cached + return cached.codec def _try_get_cpp_ext(): @@ -756,6 +778,8 @@ def _checked_output_offsets( def _align_deflate_inputs(d_concat, offsets, lengths): """Pack raw DEFLATE payloads at nvCOMP's required four-byte alignment.""" + offsets = np.asarray(offsets, dtype=np.int64) + lengths = np.asarray(lengths, dtype=np.int64) if not np.any((int(d_concat.data.ptr) + offsets) % 4): return d_concat, offsets, () @@ -833,7 +857,8 @@ def gpu_gzip_decompress_batch( d_concat : cupy.ndarray (uint8) Contiguous device buffer with concatenated compressed tile bytes. rel_offsets, lengths, uncompressed_sizes - Per-tile metadata arrays. + Per-tile metadata arrays. Decoded sizes must match the stream + contents. gzip_wrapped : bool If True (default), each tile is an RFC 1952 gzip stream; the native gzip decoder consumes its header and trailer directly. Older @@ -855,7 +880,8 @@ def gpu_gzip_decompress_batch( header_sizes : array-like of int64 or None Optional precomputed per-tile gzip header lengths. When provided, the blocking device-to-host header probe is skipped. Only valid when - ``gzip_wrapped=True``. + ``gzip_wrapped=True``. These lengths locate raw DEFLATE payloads; + the native Gzip decoder parses the intact headers itself. Returns ------- @@ -863,6 +889,13 @@ def gpu_gzip_decompress_batch( Contiguous device buffer with decompressed tile bytes concatenated. out_offsets : numpy.ndarray (int64) Start offset of each decompressed tile inside ``d_out``. + + Notes + ----- + Inputs must contain valid compressed tiles and matching size metadata. + Header checks do not validate the compressed body or its checksum. + The Python fallback uses the current CuPy stream. Callers using an + externally owned stream must keep its underlying CUDA stream alive. """ if gzip_decoder not in ("auto", "gzip", "deflate"): raise ValueError("gzip_decoder must be 'auto', 'gzip', or 'deflate'") @@ -933,7 +966,7 @@ def gpu_gzip_decompress_batch( ext = None else: ext = _try_get_cpp_ext() - if ext is None: + if ext is None and gzip_decoder != "gzip": _warn_python_fallback_once() gzip_fn = getattr(ext, "batch_gzip_decompress", None) @@ -1001,7 +1034,7 @@ def gpu_gzip_decompress_batch( deflate_lengths, d_out_ptr, out_offsets, - uncompressed_sizes, + size_array, stream_ptr, ) keepalive.append(scratch_owner) diff --git a/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp b/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp index 0785ad8..528c759 100644 --- a/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp +++ b/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp @@ -297,7 +297,8 @@ py::object batch_decompress_impl( // can progress while the async kernel launches + returns. { py::gil_scoped_release release; - // nvCOMP 5.2 requires a non-null status buffer for both codecs. + // Pass actual-size and status buffers together: nvCOMP 5.2 + // rejects a mixed null/non-null pair on the CUDA backend. check_nvcomp( gzip_wrapped ? nvcompBatchedGzipDecompressAsync( static_cast(d_comp_ptrs), @@ -323,7 +324,7 @@ py::object batch_decompress_impl( opts, static_cast(d_statuses), stream), - "batched gzip/deflate decompression"); + gzip_wrapped ? "batched gzip decompression" : "batched deflate decompression"); } if (use_native_pool) { diff --git a/tests/xdr/test_gzip_gpu.py b/tests/xdr/test_gzip_gpu.py index aa57637..e98b1b1 100644 --- a/tests/xdr/test_gzip_gpu.py +++ b/tests/xdr/test_gzip_gpu.py @@ -8,6 +8,7 @@ import gzip import struct +import weakref import zlib from types import SimpleNamespace @@ -131,14 +132,23 @@ def test_deflate_alignment_without_keepalive(cp, wrapped): gzip_decoder="deflate", use_cpp_helper=True, ) + input_owner = weakref.ref(allocation) + del data, allocation + assert input_owner() is not None + cp.get_default_memory_pool().free_all_blocks() stream.synchronize() np.testing.assert_array_equal( cp.asnumpy(result), np.frombuffer(b"".join(originals), dtype="u1") ) + del result + assert input_owner() is None @pytest.mark.parametrize("native", [False, True]) -def test_auto_decoder_supports_legacy_deflate(cp, monkeypatch, native): +@pytest.mark.parametrize("non_blocking", [False, True]) +def test_auto_decoder_orders_legacy_deflate( + cp, monkeypatch, native, non_blocking +): extension = nvcomp_batch._try_get_cpp_ext() monkeypatch.setattr( nvcomp_batch, @@ -152,13 +162,37 @@ def test_auto_decoder_supports_legacy_deflate(cp, monkeypatch, native): ), ) original = b"legacy decoder tile" * 123 - payload = gzip.compress(original) - if not native: - # Warm the codec so stream ordering cannot rely on setup latency. - nvcomp_batch._get_codec() - stream = cp.cuda.Stream(non_blocking=True) + payload = gzip.compress(original, compresslevel=0) + stale = gzip.compress(bytes(len(original)), compresslevel=0) + assert len(payload) == len(stale) + delayed_copy = cp.RawKernel( + r""" + extern "C" __global__ void delayed_copy( + const unsigned char* src, unsigned char* dst, long long size) { + unsigned long long start, now; + asm volatile("mov.u64 %0, %%globaltimer;" : "=l"(start)); + do { + asm volatile("mov.u64 %0, %%globaltimer;" : "=l"(now)); + } while (now - start < 150000000ULL); + for (long long i = threadIdx.x; i < size; i += blockDim.x) { + dst[i] = src[i]; + } + } + """, + "delayed_copy", + ) + delayed_copy.compile() + staged = cp.asarray(np.frombuffer(payload, dtype="u1")) + # Keep the stripped payload aligned and valid even before the producer + # runs, so a stream-ordering regression returns stale pixels safely. + allocation = cp.asarray(np.frombuffer(b"xx" + stale, dtype="u1")) + data = allocation[2:] + cp.cuda.get_current_stream().synchronize() + stream = cp.cuda.Stream(non_blocking=non_blocking) with stream: - data = cp.asarray(np.frombuffer(payload, dtype="u1")) + if not native: + nvcomp_batch._get_codec() + delayed_copy((1,), (256,), (staged, data, np.int64(len(payload)))) result, _ = nvcomp_batch.gpu_gzip_decompress_batch( data, [0], diff --git a/tests/xdr/test_nvcomp_batch.py b/tests/xdr/test_nvcomp_batch.py index 386c589..07224bd 100644 --- a/tests/xdr/test_nvcomp_batch.py +++ b/tests/xdr/test_nvcomp_batch.py @@ -5,6 +5,8 @@ from __future__ import annotations import sys +import threading +import weakref from types import SimpleNamespace import numpy as np @@ -13,6 +15,115 @@ from cuphoton.xdr import nvcomp_batch +def test_python_codec_cache_tracks_stream_device_and_thread(monkeypatch): + events, created = [], [] + state = SimpleNamespace(device=0) + + class Stream: + def __init__(self, ptr): + self.ptr = ptr + + def __del__(self): + events.append(("stream", self.ptr)) + + class Codec: + def __init__(self, *, algorithm, bitstream_kind, cuda_stream): + assert algorithm == "deflate" and bitstream_kind == "raw" + self.stream = weakref.ref(state.stream) + self.ptr = cuda_stream + created.append((cuda_stream, state.device)) + + def __del__(self): + events.append(("codec", self.ptr, self.stream() is not None)) + + fake_nvcomp = SimpleNamespace( + Codec=Codec, BitstreamKind=SimpleNamespace(RAW="raw") + ) + monkeypatch.setitem( + sys.modules, "nvidia", SimpleNamespace(nvcomp=fake_nvcomp) + ) + monkeypatch.setitem(sys.modules, "nvidia.nvcomp", fake_nvcomp) + monkeypatch.setitem( + sys.modules, + "cupy", + SimpleNamespace( + cuda=SimpleNamespace( + get_current_stream=lambda: state.stream, + runtime=SimpleNamespace(getDevice=lambda: state.device), + ) + ), + ) + monkeypatch.setattr( + nvcomp_batch, "_DEFLATE_CODEC_CACHE", threading.local() + ) + state.stream = Stream(11) + assert nvcomp_batch._get_codec() is nvcomp_batch._get_codec() + old_stream = weakref.ref(state.stream) + state.stream = Stream(22) + assert old_stream() is not None + nvcomp_batch._get_codec() + assert events[:2] == [("codec", 11, True), ("stream", 11)] + state.device = 1 + nvcomp_batch._get_codec() + thread = threading.Thread(target=nvcomp_batch._get_codec) + thread.start() + thread.join() + assert created == [(11, 0), (22, 0), (22, 1), (22, 1)] + assert all(event[2] for event in events if event[0] == "codec") + + +@pytest.mark.parametrize("aligned", [False, True]) +def test_deflate_repack_uses_int64_tables_and_retains_inputs( + monkeypatch, aligned +): + data = SimpleNamespace(data=SimpleNamespace(ptr=100)) + packed = SimpleNamespace() + calls = [] + offsets = np.array( + [0, 8, 8, 24] if aligned else [1, 7, 7, 30], + dtype=np.int64 if aligned else np.int32, + ) + lengths = np.array([5, 0, 13, 3], dtype=np.int32) + + def allocate(size, dtype): + assert size == 28 and dtype == np.uint8 + return packed + + monkeypatch.setattr(nvcomp_batch, "_REPACK_DEFLATE_KERNEL", None) + monkeypatch.setitem( + sys.modules, + "cupy", + SimpleNamespace( + RawKernel=lambda *a: lambda *args: calls.append(args), + empty=allocate, + asarray=np.asarray, + uint8=np.uint8, + ), + ) + result, out_offsets, owners = nvcomp_batch._align_deflate_inputs( + data, offsets, lengths + ) + if aligned: + assert result is data and out_offsets is offsets and owners == () + assert not calls + else: + assert result is packed + np.testing.assert_array_equal(out_offsets, [0, 8, 8, 24]) + grid, block, args = calls[0] + assert grid == (4,) and block == (256,) + assert args[0] is data and args[-1] is packed + for actual, expected in zip( + args[1:-1], (offsets, lengths, out_offsets), strict=True + ): + assert actual.dtype == np.int64 + np.testing.assert_array_equal(actual, expected) + assert len(owners) == 5 + assert owners[0] is data and owners[1] is packed + assert all( + a is b for a, b in zip(owners[2:], args[1:-1], strict=True) + ) + + @pytest.mark.parametrize("flag", [0x08, 0x10]) def test_parse_gzip_header_rejects_unterminated_optional_string(flag): head = bytearray(16) @@ -544,6 +655,84 @@ def batch_deflate_decompress(self, *args): assert keepalive == [output] +@pytest.mark.parametrize( + "decoder,gzip_available,native_pool,retain,expected_codec", + [ + ("auto", True, False, True, "gzip"), + ("gzip", True, True, True, "gzip"), + ("auto", True, True, False, "gzip"), + ("deflate", True, True, True, "deflate"), + ("auto", False, False, True, "deflate"), + ], +) +def test_native_decoder_receives_framing_and_retains_scratch( + monkeypatch, decoder, gzip_available, native_pool, retain, expected_codec +): + output = SimpleNamespace(data=SimpleNamespace(ptr=200)) + scratch = object() + keepalive = [] if retain else None + calls = [] + + def record(codec, args, pooled): + if keepalive is not None: + assert keepalive == [output] + calls.append((codec, args, pooled)) + return scratch if pooled else None + + extension = SimpleNamespace( + batch_deflate_decompress=lambda *args: record("deflate", args, False), + batch_deflate_decompress_pooled=lambda *args: record( + "deflate", args, True + ), + ) + if gzip_available: + extension.batch_gzip_decompress = lambda *args: record( + "gzip", args[:7], args[7] + ) + monkeypatch.setattr(nvcomp_batch, "_try_get_cpp_ext", lambda: extension) + monkeypatch.setattr( + nvcomp_batch, "_native_device_empty_uint8", lambda *a, **kw: output + ) + monkeypatch.setitem( + sys.modules, + "cupy", + SimpleNamespace( + empty=lambda *a, **kw: output, + uint8=np.uint8, + cuda=SimpleNamespace( + get_current_stream=lambda: SimpleNamespace(ptr=300) + ), + ), + ) + + result, offsets = nvcomp_batch.gpu_gzip_decompress_batch( + SimpleNamespace(size=60, data=SimpleNamespace(ptr=102)), + [0, 20, 40], + [20, 20, 20], + [1, 2, 3], + header_sizes=[10, 10, 10], + gzip_decoder=decoder, + use_native_pool=native_pool, + keepalive=keepalive, + ) + + assert result is output + np.testing.assert_array_equal(offsets, [0, 1, 3]) + assert len(calls) == 1 + codec, args, pooled = calls[0] + assert codec == expected_codec + assert pooled is (native_pool and retain) + assert (args[0], args[3], args[6]) == (102, 200, 300) + expected_offsets = [0, 20, 40] if codec == "gzip" else [10, 30, 50] + expected_lengths = [20, 20, 20] if codec == "gzip" else [2, 2, 2] + np.testing.assert_array_equal(args[1], expected_offsets) + np.testing.assert_array_equal(args[2], expected_lengths) + np.testing.assert_array_equal(args[4], offsets) + np.testing.assert_array_equal(args[5], [1, 2, 3]) + if keepalive is not None: + assert keepalive == ([output, scratch] if pooled else [output]) + + def test_precomputed_header_sizes_skip_device_probe(monkeypatch): class StopAfterHeaderValidation(Exception): pass @@ -695,7 +884,9 @@ def test_forced_gzip_rejects_missing_native_capability( monkeypatch.setitem(sys.modules, "cupy", SimpleNamespace()) monkeypatch.setattr(nvcomp_batch, "_try_get_cpp_ext", lambda: extension) monkeypatch.setattr( - nvcomp_batch, "_warn_python_fallback_once", lambda: None + nvcomp_batch, + "_warn_python_fallback_once", + lambda: pytest.fail("a forced decoder must not warn about fallback"), ) with pytest.raises( RuntimeError, match="requires a native extension with Gzip"