Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions docs/components/xdr.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
11 changes: 11 additions & 0 deletions src/cuphoton/core/cli/fits.py
Original file line number Diff line number Diff line change
Expand Up @@ -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."""
Expand Down
1 change: 1 addition & 0 deletions src/cuphoton/core/fits_options.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@

XDR_OPTION_CHOICES = {
"postprocess": frozenset({"auto", "fused", "separate"}),
"gzip_decoder": frozenset({"auto", "gzip", "deflate"}),
}


Expand Down
14 changes: 12 additions & 2 deletions src/cuphoton/xdr/benchmark_fits.py
Original file line number Diff line number Diff line change
Expand Up @@ -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, ...]:
Expand Down Expand Up @@ -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"
Expand All @@ -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:
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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,
)
Expand All @@ -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,
)
Expand Down Expand Up @@ -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,
Expand Down
5 changes: 4 additions & 1 deletion src/cuphoton/xdr/convenience.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand All @@ -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
Expand All @@ -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,
)
130 changes: 104 additions & 26 deletions src/cuphoton/xdr/nvcomp_batch.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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():
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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, ()

Expand Down Expand Up @@ -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,
Expand All @@ -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
Expand All @@ -842,15 +880,28 @@ 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
-------
d_out : cupy.ndarray (uint8)
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)
Expand Down Expand Up @@ -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.
Expand All @@ -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
Expand All @@ -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)
Expand Down
Loading
Loading