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 5704dbd..8451ef4 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 @@ -22,13 +22,14 @@ import platform import struct import sys +import threading from collections.abc import Sequence from pathlib import Path import numpy as np +_DEFLATE_CODEC_CACHE = threading.local() _REPACK_DEFLATE_KERNEL = None -_DEFLATE_CODEC = None _CPP_EXT = None _CPP_EXT_PROBED = False _CPP_EXT_IMPORT_ERROR: str | None = None @@ -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(): @@ -651,12 +673,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 @@ -749,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, () @@ -813,6 +844,7 @@ def gpu_gzip_decompress_batch( uncompressed_sizes: Sequence[int], *, gzip_wrapped: bool = True, + gzip_decoder: str = "auto", use_cpp_helper: str | bool = "auto", use_native_pool: bool = False, keepalive: list | None = None, @@ -825,10 +857,16 @@ 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 - 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. + 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 @@ -842,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 ------- @@ -850,7 +889,19 @@ 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'") + 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) @@ -908,24 +959,36 @@ 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: 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) + 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 - d_concat, deflate_offsets, aligned_owners = _align_deflate_inputs( - d_concat, deflate_offsets, deflate_lengths - ) - if keepalive is not None: - keepalive.extend(aligned_owners) + 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 +1007,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 @@ -956,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/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/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp b/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp index 218a652..528c759 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,34 @@ 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. + // Pass actual-size and status buffers together: nvCOMP 5.2 + // rejects a mixed null/non-null pair on the CUDA backend. 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), + gzip_wrapped ? "batched gzip decompression" : "batched deflate decompression"); } if (use_native_pool) { @@ -294,7 +354,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 +384,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/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", }, diff --git a/tests/xdr/test_gzip_gpu.py b/tests/xdr/test_gzip_gpu.py new file mode 100644 index 0000000..e98b1b1 --- /dev/null +++ b/tests/xdr/test_gzip_gpu.py @@ -0,0 +1,208 @@ +# 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 weakref +import zlib +from types import SimpleNamespace + +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("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, decoder): + 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, + gzip_decoder=decoder, + 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("