diff --git a/CHANGELOG.md b/CHANGELOG.md index 9325fc3..5c908a8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -60,6 +60,10 @@ See [Getting started](docs/getting-started.md) and Benchmark rounds preserve separate warmup and measurement artifacts. - xDR retains buffers until GPU work completes, permits independent reads during GDS waits, and exports benchmark results as JSON. +- xDR defaults to fused FITS pixel restoration and prefers native Gzip decoding, with + automatic selection of compatible decompression hardware. CLI flags and + Python options select separate restoration kernels, aligned raw DEFLATE, + or CUDA decompression explicitly; see [runtime choices](docs/components/xdr.md#runtime-choices). See the [component guides](docs/README.md#workflow-guides), [data contracts](docs/data-artifacts.md), and diff --git a/README.md b/README.md index 184dd2e..7b609ae 100644 --- a/README.md +++ b/README.md @@ -149,6 +149,9 @@ path and xDR supplies eligible GPU reads. Reader selection is separate from the numerical backend. Read receipts record the requested and actual reader, selected HDUs, and any fallback reason; see [FITS reading in a workflow](docs/components/xdr.md#read-fits-images-in-a-workflow). +The [xDR runtime controls](docs/components/xdr.md#runtime-choices) select +postprocessing, Gzip decoding and decompression backends through the CLI, +Python APIs or input manifests. ## Installation profiles diff --git a/docs/README.md b/docs/README.md index b18a702..19cdcb4 100644 --- a/docs/README.md +++ b/docs/README.md @@ -32,7 +32,8 @@ requirements. - [Core](components/core.md): shared CLI and configuration behavior. - [xDataReader](components/xdr.md): GPU FITS loading and the shared - Astropy/xDR reader interface; native GDS depends on the storage setup. + Astropy/xDR reader interface, with [runtime compression controls](components/xdr.md#runtime-choices) + for CLI, Python and manifests; native GDS depends on the storage setup. - [xFit](components/xfit.md): batched nonlinear least-squares dipole fitting. - [xPois](components/xpois.md): kernel fitting and image subtraction. - [xScan](components/xscan.md): transient datasets, classification, diff --git a/docs/cli.md b/docs/cli.md index e3c65b3..2854bfd 100644 --- a/docs/cli.md +++ b/docs/cli.md @@ -27,6 +27,13 @@ Check command help for the applicable default and the [FITS reader guide](components/xdr.md#read-fits-images-in-a-workflow) for scaling, compression, and section-read rules. +Applicable FITS commands and `xdr benchmark-fits` also expose +`--xdr-postprocess`, `--xdr-gzip-decoder` and `--xdr-decompression-backend`. +Omitted flags preserve manifest choices; otherwise the xDR runtime defaults +are `auto`. These settings apply when xDR is the selected FITS reader. See +[runtime choices](components/xdr.md#runtime-choices) for values, Python +keywords and capability fallback. + ## xDataReader: `cuphoton xdr` `benchmark-fits` runs the GPU-native FITS loading benchmark for individual diff --git a/docs/components/xdr.md b/docs/components/xdr.md index 78f35da..49a5959 100644 --- a/docs/components/xdr.md +++ b/docs/components/xdr.md @@ -12,42 +12,151 @@ Current scope: - pipelined loading with `batch_to_device_stream` - explicit `NotImplementedError` for unsupported compression formats - ## Runtime choices -The default `postprocess="auto"` uses fused unshuffle, byte-order conversion, -and scatter for supported GZIP FITS tiles. `"fused"` selects that path explicitly; -`"separate"` runs those steps with individual kernels for comparison. -Both preserve FITS values, including integer masks and floating-point bits. -Both leave the decoded input buffer unchanged. The separate path allocates -another buffer the size of the decoded tiles and copies GZIP_1 input before -byte-order conversion. Its timings therefore include that copy and are not -an exact baseline for the former in-place GZIP_1 implementation. +Three independent controls select how xDR restores supported compressed FITS +pixels. Each defaults to `auto`. They preserve image values, mask bits and +floating-point bit patterns; choose explicit modes for comparisons or debugging. + +| Direct Python keyword / `xdr_options` key | CLI flag | Values | Meaning of `auto` | +| --- | --- | --- | --- | +| `postprocess` | `--xdr-postprocess` | `auto`, `fused`, `separate` | Fuse unshuffle, byte-order conversion and scatter. | +| `gzip_decoder` | `--xdr-gzip-decoder` | `auto`, `gzip`, `deflate` | Use native Gzip when available; otherwise decode aligned raw DEFLATE. | +| `decompression_backend` | `--xdr-decompression-backend` | `auto`, `cuda` | Let nvCOMP choose a compatible decompression engine for native Gzip, with CUDA fallback. | + +`postprocess="fused"` selects the same kernel as `auto`; `"separate"` runs +individual restoration kernels. `gzip_decoder="gzip"` requires native Gzip +support. `"deflate"` strips the Gzip framing and aligns payloads for raw +DEFLATE. `decompression_backend="cuda"` requires a native extension with +backend selection support and selects CUDA kernels explicitly. The native +raw DEFLATE decoder uses CUDA in either backend mode. + +Both postprocessing modes leave the decoded input buffer unchanged. The +separate path allocates another buffer the size of the decoded tiles and +copies GZIP_1 input before byte-order conversion. Its timings therefore +include that copy and are not an exact baseline for the former in-place +GZIP_1 implementation. + +### Direct Python APIs + +Both batch APIs and `GpuCompImageReader.read` accept these named keywords: ```python from cuphoton.xdr import batch_to_device -(images,) = batch_to_device(paths, postprocess="auto") +(images,) = batch_to_device( + ["exposure-1.fits", "exposure-2.fits"], + hdu_indices=[1], + postprocess="auto", + gzip_decoder="auto", + decompression_backend="auto", +) +``` + +The same controls apply with `batch_to_device_stream` or `parallel=False`. +Python prefetching (`native_batcher=False`, or CLI `--native-batcher off`) +still uses the selected GPU decoder. It does not select CPU decompression. + +### Workflow APIs and manifests + +Workflow APIs accept the three keys in an `xdr_options` mapping. For example, +to require xDR and compare CUDA decompression with separate pixel restoration: + +```python +from cuphoton.core.fits_io import read_fits_images + +result = read_fits_images( + "exposure.fits", [1], reader="xdr", device=True, + xdr_options={ + "postprocess": "separate", + "gzip_decoder": "gzip", + "decompression_backend": "cuda", + }, +) +image = result.arrays[0] ``` -Workflow FITS APIs accept `xdr_options={"postprocess": "separate"}` alongside -`fits_reader="xdr"` (or `reader="xdr"` for `read_fits_images`). FITS input -descriptors can persist the same `xdr_options` mapping. The corresponding -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. +The mapping is also accepted by the FITS-consuming APIs in +[xFit](xfit.md), [xPois](xpois.md), [xRep](xrep.md), and [xScan](xscan.md). +Manifest-based workflows preserve their input choices through worker +serialization and execution. See the component guide for the mapping's +placement in its manifest or FITS descriptor. + +An explicit API override or CLI flag replaces only the corresponding manifest +key. Omitted flags preserve manifest values; keys absent from both use `auto`. +Passing `--xdr-gzip-decoder auto` explicitly therefore replaces a manifest's +`gzip_decoder` value. Reader receipts include explicit `xdr_options`; they +record requested settings rather than which nvCOMP engine actually ran. + +### Reader selection and fallback + +`--fits-reader` (Python `reader` or `fits_reader`) selects Astropy versus xDR. +The three controls above apply after that selection and do not request a GPU +reader by themselves. See [reader policies below](#read-fits-images-in-a-workflow) +for automatic CPU fallback and FITS semantic restrictions. + +Within xDR, native Gzip with `decompression_backend="auto"` may use CUDA when +a hardware engine or suitable buffers are unavailable. An older extension +without native Gzip support can use aligned raw DEFLATE. The low-level decoder's +Python fallback also runs on the GPU and retains nvCOMP's existing backend +defaults. Full batch FITS loading +still requires the native FITS planner; codec fallback does not remove that +requirement. + +When xDR decodes compressed HDUs, explicit `gzip` or `cuda` requests fail clearly +if the installed extension cannot honor them. Rebuild an older source extension +using the [native build instructions](#native-extension-availability). The `auto` +choices select available capabilities before decoding. I/O or decoder exceptions +propagate; they do not trigger another reader, codec or postprocessing attempt. +The supported FITS compression formats remain `GZIP_1` and `GZIP_2`. + +### Hardware eligibility and diagnostics + +nvCOMP selects an engine for each native Gzip call. Its +[Decompression Engine FAQ](https://docs.nvidia.com/cuda/nvcomp/decompression_engine_faq.html) +lists B200, B300, GB200 and GB300 support. GB10 and RTX PRO 6000 Blackwell +use CUDA kernels; the Blackwell name alone does not imply engine support. +The compressed data, output and decoded-size buffers must all +use compatible allocations. B200 has a 4 MiB hardware chunk limit; the limit +on a device is available through `CU_DEVICE_ATTRIBUTE_MEM_DECOMPRESS_MAXIMUM_LENGTH`. +Use `decompression_backend="cuda"` when comparing execution paths or reading +tiles beyond the hardware limit. + +Batch readers retain native pooled scratch buffers. The low-level non-pooled +path, including `GpuCompImageReader.read()` without an explicit stream, uses +`cudaMallocAsync` scratch, which is not hardware-decompression capable. That +path selects CUDA directly, including with `decompression_backend="auto"`, +to avoid a failed hardware attempt on each call. Custom CuPy allocators can +also affect eligibility in pooled calls. +An `auto` receipt and a compatible GPU therefore do not establish hardware use. +See [nvCOMP logging](../troubleshooting.md#confirm-the-decompression-engine) +to inspect the actual decoder calls. + +The benchmark's `--output-json` report includes +`capabilities.hardware_decompression`: the current device's algorithm mask, +`supports_deflate`, and `max_chunk_bytes`. Unavailable queries leave those +values `null` and record an `error`. These device limits are collected after +the timed phases and do not identify the engine used by an individual call. + +The native decoder expects valid compressed payloads and correct output sizes. +Header validation and a successful launch do not verify payload integrity; +nvCOMP's [C API](https://docs.nvidia.com/cuda/nvcomp/c_api.html) +does not guarantee safe decoding of corrupt Gzip or DEFLATE streams. + +### Tile size and batch size + +Gzip tiles provide independent work for the decoder. A batch with only a few +large tiles can leave much of the GPU idle, including on GPUs without a +hardware decompression engine. Increasing `decode_batch_files` can supply +more tiles per decode call when GPU memory allows it. + +For newly written FITS files, start with the usual row-sized tiles or tiles +of a few tens of KiB, then measure the complete read with representative +data. In one 128 MiB workload, 8–64 KiB tiles gave similar read times; +multi-MiB tiles were much slower. That result is a starting point, not a +universal optimum. Tiles beyond the device's hardware chunk limit use CUDA +in automatic mode; changing the backend alone does not restore the missing +tile parallelism. ## Read FITS images in a workflow @@ -221,6 +330,24 @@ cuphoton xdr benchmark-fits \ /path/to/file1.fits /path/to/file2.fits ``` +All three [runtime controls](#runtime-choices) are available on this command. +For a comparison against separate restoration and the native CUDA DEFLATE path, +repeat the same workload with a distinct report: + +```bash +cuphoton xdr benchmark-fits \ + --hdu-indices 1,2,3 --output-json separate-deflate-cuda.json \ + --xdr-postprocess separate --xdr-gzip-decoder deflate \ + --xdr-decompression-backend cuda \ + /path/to/file1.fits /path/to/file2.fits +``` + +Use the same files, HDUs, batching settings and storage mode for both runs. +These controls affect the `batch_to_device` and `batch_to_device_stream` +phases; planner and raw-read timings do not measure decompression. To isolate +one change, vary only that control. The Python `run_benchmark` API accepts +`postprocess`, `gzip_decoder` and `decompression_backend` directly. + Use `--dir` and `--max-files` to scan directories of FITS files. The benchmark defaults to `--native-read-threads=4`; the loading APIs default to the available CPU core count when `native_read_threads` is omitted. @@ -243,8 +370,11 @@ The version 1 JSON report contains: measured read used GDS, including when storage is mocked. - `storage`: the effective `real`, `host`, or `device` mode, including an ambient mock-storage context or environment setting. -- `options`: HDUs, iteration count, thread/queue settings, and native-batcher - selection. `native_batcher_enabled=null` records an invalid forced selection +- `options`: HDUs, iteration count, thread/queue settings, native-batcher + selection, and requested `postprocess`, `gzip_decoder` and + `decompression_backend`. An `auto` value does not identify which nvCOMP + engine ran or prove hardware decompression. `native_batcher_enabled=null` + records an invalid forced selection with the reason in `native_batcher_error`. - `workload`: ordered input files, counts, planned raw bytes, and decoded MiB. Failed planning can leave these planned sizes at zero. diff --git a/docs/components/xfit.md b/docs/components/xfit.md index 5bdd720..71e5757 100644 --- a/docs/components/xfit.md +++ b/docs/components/xfit.md @@ -182,12 +182,20 @@ Selecting xDR does not assert native GPUDirect Storage use. Run artifacts record the manifest and source-file hashes, selected reader and any automatic fallback. Existing NPZ loading is unaffected by this option. -`--xdr-postprocess auto|fused|separate` selects xDR postprocessing for FITS -reads. Python `load_xfit_dataset` accepts the equivalent -`xdr_options={"postprocess": "separate"}`. FITS manifests may specify -`xdr_options` at the top level and on individual image or `stamp_basis` -descriptors. Per-image choices override the manifest default; explicit -Python/CLI choices override only supplied keys and survive worker dispatch. +FITS commands also accept `--xdr-postprocess`, `--xdr-gzip-decoder`, and +`--xdr-decompression-backend`. These control xDR's pixel conversion, Gzip +decoder, and decompression backend when xDR is the selected reader. See +[xDR runtime choices](xdr.md#runtime-choices) for values, automatic defaults, +and capability fallback. + +Python `load_xfit_dataset` accepts the corresponding `xdr_options` keys +`postprocess`, `gzip_decoder`, and `decompression_backend`. For example, +`xdr_options={"gzip_decoder": "deflate", "decompression_backend": "cuda"}` +selects raw Deflate decoding on CUDA. In a FITS manifest, `xdr_options` may +appear at the top level and on each entry in `images` or the `stamp_basis` +descriptor. Descriptor choices override top-level defaults; explicit Python +options or CLI flags override only their matching keys. Omitted flags preserve +the manifest's choices, including when dispatching work to MPI or Dragon. Input archives contain candidate identifiers and exact image pixels. Fit artifacts contain identifiers, hashes, parameters, uncertainties, covariance, diff --git a/docs/components/xpois.md b/docs/components/xpois.md index 1f410e7..bcdfc09 100644 --- a/docs/components/xpois.md +++ b/docs/components/xpois.md @@ -48,10 +48,19 @@ to require xDR. With `--backend cpu`, automatic reading uses Astropy. The option also applies to batch, MPI and Dragon execution. NPY inputs keep their existing loading path. -`--xdr-postprocess auto|fused|separate` selects xDR's postprocessing path -when xDR is the active reader. Python workflows and `BatchFitOptions` accept -`xdr_options={"postprocess": "separate"}`; batch workers retain these choices. -Omitting the option uses xDR's default without changing FITS reader selection. +FITS commands accept `--xdr-postprocess`, `--xdr-gzip-decoder`, and +`--xdr-decompression-backend` when xDR is the active reader. Python workflows +and `BatchFitOptions` accept the corresponding `xdr_options` keys +`postprocess`, `gzip_decoder`, and `decompression_backend`. For a CUDA-only +decoding comparison, pass `--xdr-decompression-backend cuda` or +`xdr_options={"decompression_backend": "cuda"}`. See +[xDR runtime choices](xdr.md#runtime-choices) for all values and automatic +capability selection. + +Configure batch choices through the command flags or `BatchFitOptions`; +they apply to every pair and are retained by MPI and Dragon workers. +Omitted keys use xDR defaults. These controls leave the `fits_reader` policy +in charge of reader selection. Image, variance and mask HDU selection stays the same. `summary.json` records `fits_reader` and `fits_reads`, including the reader used and any fallback. diff --git a/docs/components/xrep.md b/docs/components/xrep.md index 0c1f254..fbc0a1c 100644 --- a/docs/components/xrep.md +++ b/docs/components/xrep.md @@ -32,11 +32,19 @@ decoded device arrays directly. Output images retain the usual host-array and FITS contracts. Summaries record the selected reader and any fallback under `fits_reads`. -FITS-consuming commands also accept `--xdr-postprocess auto|fused|separate` -for xDR's pixel conversion path. Python FITS APIs and workflow helpers accept -`xdr_options={"postprocess": "separate"}` (or `auto`/`fused`) and carry that -choice through image, mask, and stack reads. Omitted options use xDR defaults; -CPU Astropy reads retain their existing behavior. +FITS-consuming commands, including benchmarks and backend comparisons, +accept `--xdr-postprocess`, `--xdr-gzip-decoder`, and +`--xdr-decompression-backend`. Python FITS APIs and workflow helpers accept +the corresponding `xdr_options` keys `postprocess`, `gzip_decoder`, and +`decompression_backend`. For example, use `--xdr-decompression-backend cuda` +or `xdr_options={"decompression_backend": "cuda"}` to select CUDA +decompression when the reader uses xDR. See +[xDR runtime choices](xdr.md#runtime-choices) for values and automatic +capability selection. + +The same choices reach image, mask, and stack reads. Omitted keys use xDR +defaults; `fits_reader` continues to select the reader. Header-only WCS +inspection does not use these decompression controls. WCS mapping and target-grid setup read headers and dimensions without decompressing image pixels. Reading with xDR does not by itself imply native diff --git a/docs/components/xscan.md b/docs/components/xscan.md index b75b5a5..3ebbc06 100644 --- a/docs/components/xscan.md +++ b/docs/components/xscan.md @@ -55,10 +55,17 @@ cutouts remain bounded reads, and the prepared dataset format is unchanged. Automatic uncompressed cutouts use Astropy sections. Explicit `xdr` requires supported tile-compressed cutouts and rejects uncompressed section reads. -Builders also accept `--xdr-postprocess auto|fused|separate`, or the Python -keyword `xdr_options={"postprocess": "separate"}`. A manifest's -`xdr_options` mapping supplies defaults; explicit options override only -their matching keys. These choices apply when the selected reader uses xDR. +Raw builders also accept `--xdr-postprocess`, `--xdr-gzip-decoder`, and +`--xdr-decompression-backend`, or the Python keyword `xdr_options` with keys +`postprocess`, `gzip_decoder`, and `decompression_backend`. See +[xDR runtime choices](xdr.md#runtime-choices) for values and automatic +capability selection. These choices apply when the selected reader uses xDR. + +For raw builders, put defaults in the manifest's top-level `xdr_options` +mapping, for example `"xdr_options": {"gzip_decoder": "gzip"}`. Explicit +Python options or CLI flags override only their matching keys; omitted flags +preserve manifest choices. For example, `--xdr-decompression-backend cuda` +keeps that manifest's Gzip choice while requesting CUDA decompression. The raw builders (`data-build-autoscan-raw`, `data-build-nodiff-raw`, and `data-build-lsstcomcam-smoke`) also accept `--fits-reader auto|astropy|xdr`. @@ -454,12 +461,13 @@ from cuphoton.core.fits_io import inspect_fits_image from cuphoton.xscan.device_pipeline import FitsArrayDescriptor -def describe_fits(path, hdu, *, reader="auto"): +def describe_fits(path, hdu, *, reader="auto", xdr_options=None): path = Path(path).resolve() info = inspect_fits_image(path, hdu) return FitsArrayDescriptor( path=str(path), sha256=file_sha256(path), hdu=info.hdu, shape=info.shape, dtype=info.dtype.str, reader=reader, + xdr_options=xdr_options, ) ``` @@ -476,12 +484,19 @@ masks, before work is sent to MPI or Dragon workers or benchmark children. The effective policies are retained in workload identities and read receipts. Omitting the option preserves per-descriptor choices; NPY inputs are unchanged. -FITS descriptors can additionally carry `"xdr_options": {"postprocess": -"separate"}`. The same mapping is accepted by `predict_fits`, manifest -execution and benchmark input loading. `--xdr-postprocess` overrides that -key for all FITS descriptors. Effective choices travel with work items to -MPI/Dragon workers and benchmark subprocesses, and appear in read receipts. -HDUs grouped into one file read must use matching xDR choices. +In a pipeline manifest, put `xdr_options` on each FITS descriptor, including +variance and mask descriptors. It accepts all three +[xDR runtime choices](xdr.md#runtime-choices), for example +`"xdr_options": {"gzip_decoder": "gzip", "decompression_backend": "cuda"}`. +The Python `predict_fits` API, manifest execution, and benchmark input loading +accept the same mapping as an `xdr_options` keyword. + +`run-pipeline` and `benchmark-pipeline` accept all three `--xdr-*` flags. +An explicit flag overrides its matching key on every FITS descriptor; +omitted flags preserve per-descriptor choices. Requested choices travel with +work items to MPI/Dragon workers and benchmark subprocesses, and appear in +workload identities and read receipts. HDUs grouped into one file read must +use matching xDR choices. ```bash : "${CUDA_VISIBLE_DEVICES:?must enumerate the allocated GPUs}" diff --git a/docs/troubleshooting.md b/docs/troubleshooting.md index d4dc4a5..adc85a1 100644 --- a/docs/troubleshooting.md +++ b/docs/troubleshooting.md @@ -69,6 +69,39 @@ another reader. Check file accessibility, the selected HDU, and the native error. GPUDirect Storage also requires host and storage configuration; use `KVIKIO_COMPAT_MODE=ON` for ordinary file I/O when GDS is not configured. +### An explicit compression choice is unavailable + +`--xdr-gzip-decoder gzip` requires native Gzip support; +`--xdr-decompression-backend cuda` requires an extension with backend selection. +An importable extension may predate those entry points. After updating a source +checkout, rebuild it with the [native build instructions](components/xdr.md#native-extension-availability) +in the environment running your command. Use `auto` when capability fallback +is acceptable. A Python prefetcher still performs GPU decoding, and an `auto` +backend in a report does not prove that a hardware decompression engine ran. +See [runtime choices](components/xdr.md#runtime-choices) for the separate +reader, decoder and backend policies. + +### Confirm the decompression engine + +Enable [nvCOMP logging](https://docs.nvidia.com/cuda/nvcomp/getting_started.html#logging) +for a diagnostic run. Level 5 includes API calls as well as informational +messages; use a separate log file to keep them out of the command output: + +```bash +NVCOMP_LOG_LEVEL=5 NVCOMP_LOG_FILE=nvcomp.log \ + cuphoton xdr benchmark-fits --hdu-indices 1 \ + --xdr-gzip-decoder gzip --xdr-decompression-backend auto exposure.fits +``` + +In nvCOMP 5.2, `Trying HW decompression` records an attempt, while +`Launching HW decompression` and `HW Decompression successful` identify +hardware execution. `launching SM decompression` identifies the CUDA path. +Inspect the messages for the measured calls; a hardware attempt alone does +not prove success. Compare with `--xdr-decompression-backend cuda`, and disable +logging for timings. See [hardware eligibility](components/xdr.md#hardware-eligibility-and-diagnostics) +for buffer and tile-size requirements. FITS receipts retain requested choices, +including choices that did not apply when the selected reader was Astropy. + ## A command rejected the input Use the component inspection or validation command before the expensive step. diff --git a/src/cuphoton/core/cli/fits.py b/src/cuphoton/core/cli/fits.py index e5f6e38..45393c9 100644 --- a/src/cuphoton/core/cli/fits.py +++ b/src/cuphoton/core/cli/fits.py @@ -20,9 +20,10 @@ class XdrOptionsMixin: class XdrPostprocessArg(SetInvariant): _arg = "--xdr-postprocess" _help = ( - "xDR postprocessing: auto, fused, or separate. Omitted preserves " - "input choices or uses xDR defaults; applies when the FITS " - "reader uses xDR." + "xDR postprocessing: auto/fused combine restoration kernels; " + "separate runs individual kernels. Omitted preserves input " + "choices or uses xDR defaults; applies when the FITS reader " + "uses xDR." ) _set = XDR_OPTION_CHOICES["postprocess"] _default = None @@ -32,12 +33,27 @@ class XdrPostprocessArg(SetInvariant): 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." + "xDR gzip decoder: auto prefers native Gzip with raw DEFLATE " + "fallback; gzip or deflate selects explicitly. Omitted preserves " + "input choices or uses xDR defaults; applies when the FITS " + "reader uses xDR." ) _set = XDR_OPTION_CHOICES["gzip_decoder"] _default = None + xdr_decompression_backend = None + + class XdrDecompressionBackendArg(SetInvariant): + _arg = "--xdr-decompression-backend" + _help = ( + "xDR decompression backend: auto permits compatible Gzip " + "hardware with CUDA fallback; cuda requires explicit CUDA " + "selection support. Omitted preserves input choices or uses " + "xDR defaults; applies when the FITS reader uses xDR." + ) + _set = XDR_OPTION_CHOICES["decompression_backend"] + _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_io.py b/src/cuphoton/core/fits_io.py index e8365ce..035feee 100644 --- a/src/cuphoton/core/fits_io.py +++ b/src/cuphoton/core/fits_io.py @@ -251,9 +251,14 @@ def read_fits_images( rejects unsupported semantics or unavailable dependencies before reading pixels. All returned device arrays are ready on return. - ``xdr_options`` carries explicit xDR runtime choices, such as - ``{"postprocess": "separate"}``. Omitted choices use xDR defaults. - These choices do not change the reader or its Astropy fallback policy. + ``xdr_options`` accepts ``postprocess`` (auto/fused/separate), + ``gzip_decoder`` (auto/gzip/deflate), and ``decompression_backend`` + (auto/cuda). Omitted choices use xDR's ``auto`` defaults. For example, + ``xdr_options={"decompression_backend": "cuda"}`` requests CUDA kernels + when xDR decodes a compressed HDU; it requires a native helper with + backend selection. These choices do not change the reader or its Astropy + fallback policy. Metadata records explicit choices as requested settings, + including when Astropy is selected; it does not identify an nvCOMP engine. """ validate_fits_reader(reader) xdr_options = normalize_xdr_options(xdr_options) diff --git a/src/cuphoton/core/fits_options.py b/src/cuphoton/core/fits_options.py index c098a09..df7a554 100644 --- a/src/cuphoton/core/fits_options.py +++ b/src/cuphoton/core/fits_options.py @@ -9,6 +9,7 @@ XDR_OPTION_CHOICES = { "postprocess": frozenset({"auto", "fused", "separate"}), "gzip_decoder": frozenset({"auto", "gzip", "deflate"}), + "decompression_backend": frozenset({"auto", "cuda"}), } diff --git a/src/cuphoton/xdr/benchmark_fits.py b/src/cuphoton/xdr/benchmark_fits.py index 0b49e07..4bc3cc0 100644 --- a/src/cuphoton/xdr/benchmark_fits.py +++ b/src/cuphoton/xdr/benchmark_fits.py @@ -12,6 +12,7 @@ from __future__ import annotations import contextlib +import ctypes import json import os import platform @@ -164,6 +165,7 @@ class _BatchOptions(TypedDict): native_batcher: str | bool postprocess: str gzip_decoder: str + decompression_backend: str def parse_hdu_indices(value: str) -> tuple[int, ...]: @@ -422,6 +424,7 @@ def bench_batch_load( data_mb: float, postprocess: str = "auto", gzip_decoder: str = "auto", + decompression_backend: 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" @@ -435,13 +438,14 @@ def bench_batch_load( native_batcher=native_batcher, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, ) times: list[float] = [] note = ( f"{len(paths)} files x {len(tuple(hdu_indices))} HDUs, " f"native_batcher={native_batcher}, postprocess={postprocess}, " - f"gzip_decoder={gzip_decoder}" + f"gzip_decoder={gzip_decoder}, backend={decompression_backend}" ) with nvtx_range(f"xdr.{phase}"): try: @@ -549,6 +553,60 @@ def _nvcomp_version() -> str | None: return None +def _hardware_decompression_capabilities() -> dict[ + str, int | bool | str | None +]: + """Query device support, independently of the engine used by nvCOMP.""" + result: dict[str, int | bool | str | None] = { + "algorithm_mask": None, + "supports_deflate": None, + "max_chunk_bytes": None, + "error": None, + } + try: + ordinal = int(cp.cuda.runtime.getDevice()) + driver = ctypes.CDLL("libcuda.so.1") + driver.cuDeviceGet.argtypes = [ + ctypes.POINTER(ctypes.c_int), + ctypes.c_int, + ] + driver.cuDeviceGet.restype = ctypes.c_int + driver.cuDeviceGetAttribute.argtypes = [ + ctypes.POINTER(ctypes.c_int), + ctypes.c_int, + ctypes.c_int, + ] + driver.cuDeviceGetAttribute.restype = ctypes.c_int + device = ctypes.c_int() + status = driver.cuDeviceGet(ctypes.byref(device), ordinal) + if status != 0: + raise RuntimeError(f"cuDeviceGet returned CUDA error {status}") + + def attribute(identifier: int) -> int: + value = ctypes.c_int() + status = driver.cuDeviceGetAttribute( + ctypes.byref(value), identifier, device.value + ) + if status != 0: + raise RuntimeError( + f"cuDeviceGetAttribute({identifier}) returned " + f"CUDA error {status}" + ) + return value.value + + # CUDA Driver API: MEM_DECOMPRESS_ALGORITHM_MASK / MAXIMUM_LENGTH. + # CU_MEM_DECOMPRESS_ALGORITHM_DEFLATE is bit 0 of the mask. + mask, maximum = attribute(136), attribute(137) + result.update( + algorithm_mask=mask, + supports_deflate=bool(mask & 1), + max_chunk_bytes=maximum, + ) + except Exception as exc: + result["error"] = f"{type(exc).__name__}: {exc}" + return result + + def run_benchmark( fits_files: Sequence[str | Path], *, @@ -564,6 +622,7 @@ def run_benchmark( native_batcher: str = "auto", postprocess: str = "auto", gzip_decoder: str = "auto", + decompression_backend: str = "auto", mock_storage_kind: str | None = None, skip_gds_read: bool = False, output_json: Path | None = None, @@ -574,10 +633,21 @@ def run_benchmark( ``output_json`` requires an existing parent directory. Report metadata is collected after the timed phases so it cannot warm their native helpers; the return value remains a list of phases. + + ``postprocess``, ``gzip_decoder`` and ``decompression_backend`` use the + choices documented by :func:`cuphoton.xdr.batch_to_device_stream`. They + affect the two batch-load phases, independently of planner/raw-read + phases. JSON ``options`` records the requested values; ``auto`` does not + identify the actual nvCOMP engine. Keep inputs and scheduling fixed and + vary one control at a time when comparing implementations. """ normalize_xdr_options( - {"postprocess": postprocess, "gzip_decoder": gzip_decoder} + { + "postprocess": postprocess, + "gzip_decoder": gzip_decoder, + "decompression_backend": decompression_backend, + } ) paths = resolve_paths( fits_files, @@ -656,6 +726,7 @@ def run_benchmark( native_batcher=resolved_native_batcher, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, use_stream=False, data_mb=total_data_mb, ) @@ -673,6 +744,7 @@ def run_benchmark( native_batcher=resolved_native_batcher, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, use_stream=True, data_mb=total_data_mb, ) @@ -701,6 +773,9 @@ def run_benchmark( "cpp_helper_available": cpp_helper_available(), "gds_active": gds_active, "gds_probe": "cuphoton.xdr.is_gds_active", + "hardware_decompression": ( + _hardware_decompression_capabilities() + ), }, "storage": { "mode": storage_mode, @@ -717,6 +792,7 @@ def run_benchmark( "native_batcher": native_batcher, "postprocess": postprocess, "gzip_decoder": gzip_decoder, + "decompression_backend": decompression_backend, "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 da46cf8..aa591fe 100644 --- a/src/cuphoton/xdr/convenience.py +++ b/src/cuphoton/xdr/convenience.py @@ -44,6 +44,7 @@ def batch_to_device( native_batcher: str | bool = "auto", postprocess: str = "auto", gzip_decoder: str = "auto", + decompression_backend: str = "auto", ): """Load image HDUs from FITS files into stacked device arrays. @@ -62,12 +63,25 @@ def batch_to_device( section : tuple of slice or None Optional 2D ROI applied uniformly to the CompImageHDUs. parallel : bool - When True, use the streaming batch reader. Set False for a depth-1 path - that preserves the existing API without requiring Astropy HDU objects. + True uses the configured streaming batch reader. False sets queue + depths and ``decode_batch_files`` to one and selects Python + prefetching; decoding still runs on the GPU with the requested + runtime choices. prefetch_depth, decode_batch_files, batch_queue_depth, - native_read_threads, native_plan_threads, native_batcher, postprocess, - gzip_decoder - Passed through to `batch_to_device_stream` when ``parallel=True``. + native_read_threads, native_plan_threads, native_batcher + Scheduling controls for `batch_to_device_stream`. ``parallel=False`` + overrides the depths and native batcher as described above. + postprocess : {"auto", "fused", "separate"} + ``auto`` (default) uses fused FITS pixel restoration. ``separate`` + selects individual unshuffle, byteswap and scatter kernels. + gzip_decoder : {"auto", "gzip", "deflate"} + ``auto`` prefers native Gzip, with aligned raw DEFLATE fallback. + ``gzip`` requires native Gzip support; ``deflate`` selects raw decode. + decompression_backend : {"auto", "cuda"} + ``auto`` lets native Gzip select compatible hardware with CUDA + fallback. ``cuda`` requires a native helper with backend selection. + Native raw DEFLATE uses CUDA in either mode. All three runtime + controls apply with either value of ``parallel``. Returns ------- @@ -92,6 +106,7 @@ def batch_to_device( native_batcher=native_batcher, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, section=section, stream=stream, ) diff --git a/src/cuphoton/xdr/nvcomp_batch.py b/src/cuphoton/xdr/nvcomp_batch.py index 8451ef4..7e4e061 100644 --- a/src/cuphoton/xdr/nvcomp_batch.py +++ b/src/cuphoton/xdr/nvcomp_batch.py @@ -845,6 +845,7 @@ def gpu_gzip_decompress_batch( *, gzip_wrapped: bool = True, gzip_decoder: str = "auto", + decompression_backend: str = "auto", use_cpp_helper: str | bool = "auto", use_native_pool: bool = False, keepalive: list | None = None, @@ -867,6 +868,10 @@ def gpu_gzip_decompress_batch( "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. + decompression_backend : "auto" | "cuda" + "auto" lets native Gzip use compatible hardware with CUDA fallback. + Native raw DEFLATE uses CUDA in either mode. + "cuda" requires a rebuilt native helper and uses CUDA kernels. 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 @@ -897,6 +902,8 @@ def gpu_gzip_decompress_batch( The Python fallback uses the current CuPy stream. Callers using an externally owned stream must keep its underlying CUDA stream alive. """ + if decompression_backend not in ("auto", "cuda"): + raise ValueError("decompression_backend must be 'auto' or 'cuda'") 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: @@ -959,7 +966,9 @@ def gpu_gzip_decompress_batch( "`bash src/cuphoton/xdr/src/build.sh` from source. " f"Import error: {_CPP_EXT_IMPORT_ERROR}" ) - elif n == 0 and gzip_decoder != "gzip": + elif ( + n == 0 and gzip_decoder != "gzip" and decompression_backend == "auto" + ): # No backend work is needed, so an automatic fallback is not useful. ext = None elif use_cpp_helper is False: @@ -969,6 +978,17 @@ def gpu_gzip_decompress_batch( if ext is None and gzip_decoder != "gzip": _warn_python_fallback_once() + backend_selection = getattr(ext, "supports_decompression_backend", False) + if decompression_backend != "auto" and not backend_selection: + raise RuntimeError( + "decompression_backend='cuda' requires a rebuilt native " + "extension with backend selection and use_cpp_helper enabled" + ) + backend_kwargs = ( + {"decompression_backend": decompression_backend} + if backend_selection + else {} + ) gzip_fn = getattr(ext, "batch_gzip_decompress", None) if gzip_decoder == "gzip" and gzip_fn is None: raise RuntimeError( @@ -1018,6 +1038,7 @@ def gpu_gzip_decompress_batch( size_array, stream_ptr, use_native_pool and keepalive is not None, + **backend_kwargs, ) if keepalive is not None and scratch_owner is not None: keepalive.append(scratch_owner) @@ -1036,6 +1057,7 @@ def gpu_gzip_decompress_batch( out_offsets, size_array, stream_ptr, + **backend_kwargs, ) keepalive.append(scratch_owner) return d_out, out_offsets @@ -1048,6 +1070,7 @@ def gpu_gzip_decompress_batch( out_offsets, size_array, stream_ptr, + **backend_kwargs, ) return d_out, out_offsets diff --git a/src/cuphoton/xdr/prefetch.py b/src/cuphoton/xdr/prefetch.py index 35ee932..1a49d67 100644 --- a/src/cuphoton/xdr/prefetch.py +++ b/src/cuphoton/xdr/prefetch.py @@ -603,6 +603,7 @@ def _consume_comp_batch( *, postprocess="auto", gzip_decoder: str = "auto", + decompression_backend: str = "auto", ): """Decode many compressed HDUs with one batched nvCOMP call. @@ -719,6 +720,7 @@ def _consume_comp_batch( out_bytes=out_bytes, stream=stream, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, keepalive=keepalive, header_sizes=header_sizes, ) @@ -804,6 +806,7 @@ def _consume_prefetched_group( *, postprocess="auto", gzip_decoder: str = "auto", + decompression_backend: str = "auto", ): """Consume prefetched files, batching compressed HDUs together.""" comp_entries: list[_CompBatchEntry] = [] @@ -841,6 +844,7 @@ def _consume_prefetched_group( keepalive=keepalive, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, ) @@ -1074,6 +1078,7 @@ def _submit_prefetched_group( in_flight: deque[_GpuBatchHandle], postprocess="auto", gzip_decoder: str = "auto", + decompression_backend: str = "auto", ) -> None: """Queue GPU work and register its lifetime handle.""" import cupy as cp @@ -1100,6 +1105,7 @@ def _submit_prefetched_group( keepalive=keepalive, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, ) candidate_event = cp.cuda.Event() candidate_event.record(use_stream) @@ -1427,6 +1433,7 @@ def _consume_native_batches( native_plan_files, postprocess="auto", gzip_decoder: str = "auto", + decompression_backend: str = "auto", ): """Consume device batches built by the C++ KvikIO worker pool.""" import cupy as cp @@ -1504,6 +1511,7 @@ def _consume_native_batches( in_flight=in_flight, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, ) planner.join() if planner.error is not None: @@ -1545,6 +1553,7 @@ def _consume_python_batches( stream, postprocess="auto", gzip_decoder: str = "auto", + decompression_backend: str = "auto", ) -> None: """Consume pinned-host batches with bounded event-owned lifetimes.""" n_files = len(paths) @@ -1567,6 +1576,7 @@ def submit(group: list[PrefetchedFile]) -> None: in_flight=in_flight, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, ) prefetcher.start() @@ -1636,6 +1646,7 @@ def batch_to_device_stream( native_batcher: str | bool = "auto", postprocess: str = "auto", gzip_decoder: str = "auto", + decompression_backend: str = "auto", section=None, stream=None, ): @@ -1685,6 +1696,10 @@ def batch_to_device_stream( gzip_decoder "auto" uses native Gzip when available, otherwise aligned DEFLATE. "gzip" requires native Gzip support; "deflate" selects raw DEFLATE. + decompression_backend + "auto" lets Gzip select compatible hardware with CUDA fallback. "cuda" + requires a rebuilt native helper and uses CUDA kernels. Native raw + DEFLATE uses CUDA in either mode. section Optional 2D ROI applied uniformly to CompImageHDUs. stream @@ -1706,7 +1721,11 @@ def batch_to_device_stream( """ normalize_xdr_options( - {"postprocess": postprocess, "gzip_decoder": gzip_decoder} + { + "postprocess": postprocess, + "gzip_decoder": gzip_decoder, + "decompression_backend": decompression_backend, + } ) hdu_indices = tuple(int(i) for i in hdu_indices) resolved_paths = [Path(p) for p in paths] @@ -1759,6 +1778,7 @@ def batch_to_device_stream( native_plan_files=native_plan_files, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, ) return tuple(outs) @@ -1773,6 +1793,7 @@ def batch_to_device_stream( stream=stream, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, ) return tuple(outs) diff --git a/src/cuphoton/xdr/reader.py b/src/cuphoton/xdr/reader.py index c23d278..3d077c3 100644 --- a/src/cuphoton/xdr/reader.py +++ b/src/cuphoton/xdr/reader.py @@ -438,6 +438,7 @@ def inflate_device_heap( keepalive=None, header_sizes=None, gzip_decoder: str = "auto", + decompression_backend: str = "auto", ): """Run batched nvCOMP inflate for gzip tiles resident on the device. @@ -445,9 +446,16 @@ def inflate_device_heap( and `rel_offsets[i]` is the start offset of tile i inside `d_concat`. `lengths` and `out_bytes` are per-tile compressed and decompressed byte counts. Returns concatenated decompressed bytes plus per-tile - offsets into that output buffer. + offsets into that output buffer. ``gzip_decoder`` and + ``decompression_backend`` use the choices documented by + :func:`cuphoton.xdr.batch_to_device_stream`. """ - normalize_xdr_options({"gzip_decoder": gzip_decoder}) + normalize_xdr_options( + { + "gzip_decoder": gzip_decoder, + "decompression_backend": decompression_backend, + } + ) import cupy as cp with stream or cp.cuda.Stream.null: @@ -458,6 +466,7 @@ def inflate_device_heap( out_bytes, gzip_wrapped=True, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, use_native_pool=keepalive is not None, keepalive=keepalive, header_sizes=header_sizes, @@ -593,6 +602,7 @@ def decode_from_device_heap( keepalive=None, postprocess: str = "auto", gzip_decoder: str = "auto", + decompression_backend: str = "auto", ): """Run the decode pipeline on a pre-loaded compressed heap buffer. @@ -600,9 +610,15 @@ def decode_from_device_heap( of the tiles referenced in `plan["sel_*"]`; `rel_offsets[i]` is the start offset of tile i inside `d_concat`. This is what the prefetch consumer calls after staging its pinned host heap to device. + ``postprocess``, ``gzip_decoder`` and ``decompression_backend`` use + the choices documented by :func:`cuphoton.xdr.batch_to_device_stream`. """ normalize_xdr_options( - {"postprocess": postprocess, "gzip_decoder": gzip_decoder} + { + "postprocess": postprocess, + "gzip_decoder": gzip_decoder, + "decompression_backend": decompression_backend, + } ) if keepalive is not None: keepalive.append(d_concat) @@ -614,6 +630,7 @@ def decode_from_device_heap( out_bytes=plan["out_bytes"], stream=stream, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, keepalive=keepalive, ) ) @@ -638,6 +655,7 @@ def read( loader=None, postprocess: str = "auto", gzip_decoder: str = "auto", + decompression_backend: str = "auto", ): """Read and decode this HDU into a `cupy.ndarray`. @@ -648,7 +666,20 @@ def read( kvikio `CompatModeManager` setup per HDU). If `None`, a fresh loader is created for this call. Closure is deferred if failure leaves its I/O completion unknown. - + postprocess : {"auto", "fused", "separate"} + ``auto`` (default) and ``fused`` restore FITS pixels in one + kernel; ``separate`` uses individual restoration kernels. + gzip_decoder : {"auto", "gzip", "deflate"} + ``auto`` prefers native Gzip with aligned DEFLATE fallback. + ``gzip`` requires native Gzip support; ``deflate`` selects + aligned raw DEFLATE decoding. + decompression_backend : {"auto", "cuda"} + ``auto`` lets native Gzip select compatible hardware with CUDA + fallback. ``cuda`` requires a native helper with backend + selection. Native raw DEFLATE uses CUDA in either mode. + + Notes + ----- With an explicit stream, independent reads can submit while this call waits for GDS I/O. Reads without a stream retain the submission lock. @@ -658,7 +689,11 @@ def read( requires a process restart before further GPU submissions. """ normalize_xdr_options( - {"postprocess": postprocess, "gzip_decoder": gzip_decoder} + { + "postprocess": postprocess, + "gzip_decoder": gzip_decoder, + "decompression_backend": decompression_backend, + } ) import cupy as cp @@ -735,6 +770,7 @@ def abandon(error): stream=None, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, ) except BaseException as error: abandon(error) @@ -778,6 +814,7 @@ def abandon(error): keepalive=keepalive, postprocess=postprocess, gzip_decoder=gzip_decoder, + decompression_backend=decompression_backend, ) 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 528c759..de483dc 100644 --- a/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp +++ b/src/cuphoton/xdr/src/nvcomp_batch_ext.cpp @@ -46,7 +46,8 @@ py::object batch_decompress_impl( py::array_t out_sizes, std::uintptr_t stream_ptr, bool use_native_pool, - bool gzip_wrapped); + bool gzip_wrapped, + const std::string& decompression_backend); // Batched DEFLATE decompress across N tiles packed inside one device buffer. // @@ -66,8 +67,10 @@ py::object batch_decompress_impl( // Per-tile uncompressed size (must match actual decompressed size). // stream_ptr : uintptr_t // CUDA stream. 0 = default stream. +// decompression_backend : string +// Validated as "auto" or "cuda". Raw DEFLATE always uses CUDA. // -// Raises on nvcomp / CUDA failure. Returns None. +// Raises on nvCOMP / CUDA launch failure. Returns None; execution is asynchronous. void batch_deflate_decompress( std::uintptr_t d_concat_ptr, py::array_t rel_offsets, @@ -75,7 +78,8 @@ void batch_deflate_decompress( std::uintptr_t d_out_ptr, py::array_t out_offsets, py::array_t out_sizes, - std::uintptr_t stream_ptr) { + std::uintptr_t stream_ptr, + const std::string& decompression_backend) { batch_decompress_impl( d_concat_ptr, rel_offsets, @@ -85,7 +89,8 @@ void batch_deflate_decompress( out_sizes, stream_ptr, false, - false); + false, + decompression_backend); } py::object batch_deflate_decompress_pooled( @@ -95,7 +100,8 @@ py::object batch_deflate_decompress_pooled( std::uintptr_t d_out_ptr, py::array_t out_offsets, py::array_t out_sizes, - std::uintptr_t stream_ptr) { + std::uintptr_t stream_ptr, + const std::string& decompression_backend) { return batch_decompress_impl( d_concat_ptr, rel_offsets, @@ -105,9 +111,13 @@ py::object batch_deflate_decompress_pooled( out_sizes, stream_ptr, true, - false); + false, + decompression_backend); } +// Gzip uses the same pointer/size arrays, with offsets and lengths covering +// complete RFC-1952 streams. "auto" permits compatible hardware with CUDA +// fallback; "cuda" selects CUDA. A pooled owner must survive stream completion. py::object batch_gzip_decompress( std::uintptr_t d_concat_ptr, py::array_t rel_offsets, @@ -116,7 +126,8 @@ py::object batch_gzip_decompress( py::array_t out_offsets, py::array_t out_sizes, std::uintptr_t stream_ptr, - bool use_native_pool) { + bool use_native_pool, + const std::string& decompression_backend) { return batch_decompress_impl( d_concat_ptr, rel_offsets, @@ -126,7 +137,8 @@ py::object batch_gzip_decompress( out_sizes, stream_ptr, use_native_pool, - true); + true, + decompression_backend); } py::object batch_decompress_impl( @@ -138,7 +150,16 @@ py::object batch_decompress_impl( py::array_t out_sizes, std::uintptr_t stream_ptr, bool use_native_pool, - bool gzip_wrapped) { + bool gzip_wrapped, + const std::string& decompression_backend) { + nvcompDecompressBackend_t backend; + if (decompression_backend == "auto") { + backend = NVCOMP_DECOMPRESS_BACKEND_DEFAULT; + } else if (decompression_backend == "cuda") { + backend = NVCOMP_DECOMPRESS_BACKEND_CUDA; + } else { + throw std::invalid_argument("decompression_backend must be 'auto' or 'cuda'"); + } const std::size_t n = static_cast(rel_offsets.size()); if (lengths.size() != static_cast(n) || out_offsets.size() != static_cast(n) @@ -190,11 +211,14 @@ py::object batch_decompress_impl( void* d_actual_sizes = nullptr; auto opts = nvcompBatchedDeflateDecompressDefaultOpts; - // The hardware backend requires non-stream-ordered scratch allocations. - // Keep this path on CUDA until that ownership contract is supported. + // Raw DEFLATE retains its CUDA path in either mode. opts.backend = NVCOMP_DECOMPRESS_BACKEND_CUDA; auto gzip_opts = nvcompBatchedGzipDecompressDefaultOpts; - gzip_opts.backend = NVCOMP_DECOMPRESS_BACKEND_CUDA; + // Non-pooled scratch uses cudaMallocAsync, whose allocations are not + // hardware-decompression capable. Select CUDA directly to avoid a failed + // hardware launch followed by fallback on every call. Pooled calls retain + // nvCOMP's automatic selection when requested. + gzip_opts.backend = use_native_pool ? 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; @@ -356,6 +380,8 @@ py::object batch_decompress_impl( PYBIND11_MODULE(_nvcomp_batch_ext, m) { m.doc() = "Batched Gzip/DEFLATE decompression — device-pointer interface to nvcomp."; + m.attr("supports_decompression_backend") = true; + xdr_gpu::bind_io(m); xdr_gpu::bind_memory_manager(m); @@ -369,9 +395,11 @@ PYBIND11_MODULE(_nvcomp_batch_ext, m) { py::arg("out_offsets"), py::arg("out_sizes"), py::arg("stream_ptr") = 0, + py::arg("decompression_backend") = "auto", "Batched DEFLATE decompress across N tiles packed in one device buffer.\n" "Caller must have already stripped the RFC-1952 gzip wrapper from each\n" - "tile (advance rel_offsets past the header, shorten lengths by header+8)."); + "tile (advance rel_offsets past the header, shorten lengths by header+8).\n" + "decompression_backend accepts auto/cuda; raw DEFLATE always uses CUDA."); m.def( "batch_deflate_decompress_pooled", &xdr_gpu::batch_deflate_decompress_pooled, @@ -382,8 +410,10 @@ PYBIND11_MODULE(_nvcomp_batch_ext, m) { py::arg("out_offsets"), py::arg("out_sizes"), py::arg("stream_ptr") = 0, + py::arg("decompression_backend") = "auto", "Batched DEFLATE decompress using native pooled scratch buffers.\n" - "Returns an owner capsule that must stay alive until stream work completes."); + "Returns an owner capsule that must stay alive until stream work completes.\n" + "decompression_backend accepts auto/cuda; raw DEFLATE always uses CUDA."); m.def( "batch_gzip_decompress", &xdr_gpu::batch_gzip_decompress, @@ -395,6 +425,9 @@ PYBIND11_MODULE(_nvcomp_batch_ext, m) { py::arg("out_sizes"), py::arg("stream_ptr") = 0, py::arg("use_native_pool") = false, + py::arg("decompression_backend") = "auto", "Batched Gzip decompress including RFC-1952 headers and trailers.\n" - "Pooled calls return an owner capsule retained until stream completion."); + "Pooled calls return an owner capsule retained until stream completion.\n" + "decompression_backend=auto permits compatible hardware with CUDA fallback;\n" + "cuda selects CUDA explicitly."); } diff --git a/tests/core/test_cli_contract.py b/tests/core/test_cli_contract.py index 939af0a..45d41b0 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, 16, 0), - ("xfit", 3, 3, 33, 1), - ("xpois", 7, 7, 158, 1), - ("xscan", 44, 44, 219, 1), - ("xrep", 6, 6, 118, 1), + ("xdr", 1, 0, 17, 0), + ("xfit", 3, 3, 36, 1), + ("xpois", 7, 7, 162, 1), + ("xscan", 44, 44, 224, 1), + ("xrep", 6, 6, 123, 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) == 936 + assert sum(item[3] for item in per_group) == 954 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 f8e171d..b3fb6db 100644 --- a/tests/core/test_fits_io.py +++ b/tests/core/test_fits_io.py @@ -15,10 +15,11 @@ from cuphoton.core import fits_io +@pytest.mark.parametrize("decompression_backend", ["auto", "cuda"]) @pytest.mark.parametrize("postprocess", ["fused", "separate"]) @pytest.mark.parametrize("gzip_decoder", ["auto", "gzip", "deflate"]) def test_xdr_options_reach_batch_dispatch( - tmp_path, monkeypatch, postprocess, gzip_decoder + tmp_path, monkeypatch, postprocess, gzip_decoder, decompression_backend ): import cuphoton.xdr as xdr @@ -40,13 +41,16 @@ def batch(paths, hdus, **options): xdr_options={ "postprocess": postprocess, "gzip_decoder": gzip_decoder, + "decompression_backend": decompression_backend, }, ) assert observed["postprocess"] == postprocess assert observed["gzip_decoder"] == gzip_decoder + assert observed["decompression_backend"] == decompression_backend assert result.metadata()["xdr_options"] == { "postprocess": postprocess, "gzip_decoder": gzip_decoder, + "decompression_backend": decompression_backend, } 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 47dbb27..896dd2b 100644 --- a/tests/xdr/test_benchmark_fits.py +++ b/tests/xdr/test_benchmark_fits.py @@ -167,6 +167,8 @@ def fake_run_benchmark(fits_files, **kwargs): "separate", "--xdr-gzip-decoder", "deflate", + "--xdr-decompression-backend", + "cuda", "--mock-storage", "host", "--skip-gds-read", @@ -189,6 +191,7 @@ def fake_run_benchmark(fits_files, **kwargs): assert received["native_batcher"] == "off" assert received["postprocess"] == "separate" assert received["gzip_decoder"] == "deflate" + assert received["decompression_backend"] == "cuda" assert received["mock_storage_kind"] == "host" assert received["skip_gds_read"] is True assert received["output_json"] == Path("report.json") @@ -388,11 +391,9 @@ def test_json_report_matches_returned_and_printed_phases( assert report["ok"] is True assert report["phases"] == [asdict(phase) for phase in results] assert report["storage"]["mode"] == mode - assert report["capabilities"] == { - "cpp_helper_available": True, - "gds_active": True, - "gds_probe": "cuphoton.xdr.is_gds_active", - } + assert report["capabilities"]["cpp_helper_available"] is True + assert report["capabilities"]["gds_active"] is True + assert report["capabilities"]["gds_probe"] == "cuphoton.xdr.is_gds_active" assert report["options"] == { "hdu_indices": [2, 1], "iterations": 2, @@ -404,6 +405,7 @@ def test_json_report_matches_returned_and_printed_phases( "native_batcher": "auto", "postprocess": "auto", "gzip_decoder": "auto", + "decompression_backend": "auto", "native_batcher_enabled": mode == "real", "native_batcher_error": None, "skip_gds_read": mode != "real", @@ -428,6 +430,77 @@ def test_json_report_matches_returned_and_printed_phases( assert not list(tmp_path.glob(".benchmark.json.*")) +@pytest.mark.parametrize( + "mask,maximum,failed_attribute", + [ + (7, 4194304, None), + (0, 0, None), + (6, 4194304, None), + (0, 0, 136), + (7, 4194304, 137), + ], +) +def test_json_reports_device_decompression_capabilities( + json_benchmark, monkeypatch, tmp_path, mask, maximum, failed_attribute +): + import json + + bench, path, phases = json_benchmark + queries = [] + + def current_device(): + assert phases[-1].phase == "batch_to_device_stream" + return 3 + + def device_get(pointer, ordinal): + assert ordinal == 3 + pointer._obj.value = ordinal + return 0 + + def attribute_get(pointer, attribute, device): + assert device == 3 + queries.append(attribute) + if attribute == failed_attribute: + return 1 # CUDA_ERROR_INVALID_VALUE: attribute unsupported. + pointer._obj.value = {136: mask, 137: maximum}[attribute] + return 0 + + monkeypatch.setattr( + bench, + "cp", + SimpleNamespace( + __version__="test-cupy", + cuda=SimpleNamespace( + runtime=SimpleNamespace(getDevice=current_device) + ), + ), + ) + monkeypatch.setattr( + bench.ctypes, + "CDLL", + lambda _name: SimpleNamespace( + cuDeviceGet=device_get, cuDeviceGetAttribute=attribute_get + ), + ) + output = tmp_path / "capabilities.json" + bench.run_benchmark([path], output_json=output, out=lambda _: None) + report = json.loads(output.read_text()) + capability = report["capabilities"]["hardware_decompression"] + assert queries == ([136] if failed_attribute == 136 else [136, 137]) + if failed_attribute is None: + assert capability["algorithm_mask"] == mask + assert capability["supports_deflate"] is bool(mask & 1) + assert capability["max_chunk_bytes"] == maximum + assert capability["error"] is None + else: + assert capability["algorithm_mask"] is None + assert capability["supports_deflate"] is None + assert capability["max_chunk_bytes"] is None + assert f"({failed_attribute})" in capability["error"] + assert "CUDA error 1" in capability["error"] + assert report["ok"] is True + + @pytest.mark.parametrize("native_batcher", ["auto", "off", "on"]) def test_json_reports_missing_helper_and_failed_phases( json_benchmark, monkeypatch, tmp_path, native_batcher @@ -585,6 +658,9 @@ def unexpected_metadata(): pytest.fail("optional report probes must not run without output_json") monkeypatch.setattr(bench, "cpp_helper_available", unexpected_metadata) + monkeypatch.setattr( + bench, "_hardware_decompression_capabilities", unexpected_metadata + ) printed = [] assert bench.run_benchmark([path], out=printed.append) == phases assert any("batch_to_device_stream" in line for line in printed) diff --git a/tests/xdr/test_frontend.py b/tests/xdr/test_frontend.py index 21af613..c9a52bb 100644 --- a/tests/xdr/test_frontend.py +++ b/tests/xdr/test_frontend.py @@ -47,6 +47,7 @@ def test_batch_api_keeps_existing_keyword_surface(): "native_batcher", "postprocess", "gzip_decoder", + "decompression_backend", ) actual = tuple(inspect.signature(xdr.batch_to_device).parameters) @@ -79,6 +80,7 @@ def fake_stream(paths, **kwargs): native_batcher=True, postprocess="separate", gzip_decoder="deflate", + decompression_backend="cuda", ) assert result == ("ok",) @@ -96,6 +98,7 @@ def fake_stream(paths, **kwargs): "native_batcher": True, "postprocess": "separate", "gzip_decoder": "deflate", + "decompression_backend": "cuda", "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 index e98b1b1..fd0aa12 100644 --- a/tests/xdr/test_gzip_gpu.py +++ b/tests/xdr/test_gzip_gpu.py @@ -32,10 +32,13 @@ def cp(): return module +@pytest.mark.parametrize("backend", ["auto", "cuda"]) @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): +def test_native_gzip_decodes_optional_headers( + cp, pooled, preparsed, decoder, backend +): payloads, originals, header_sizes = [], [], [] for flags in range(32): original = bytes(range(256)) * 9 + bytes([flags]) * 11 @@ -60,6 +63,7 @@ def test_native_gzip_decodes_optional_headers(cp, pooled, preparsed, decoder): sizes, use_cpp_helper=True, gzip_decoder=decoder, + decompression_backend=backend, use_native_pool=pooled, keepalive=owners, header_sizes=header_sizes if preparsed else None, @@ -206,3 +210,20 @@ def test_auto_decoder_orders_legacy_deflate( np.testing.assert_array_equal( cp.asnumpy(result), np.frombuffer(original, dtype="u1") ) + + +@pytest.mark.parametrize( + "entrypoint", + [ + "batch_gzip_decompress", + "batch_deflate_decompress", + "batch_deflate_decompress_pooled", + ], +) +def test_native_decoder_rejects_invalid_backend(cp, entrypoint): + extension = nvcomp_batch._try_get_cpp_ext() + empty = np.empty(0, dtype="i8") + with pytest.raises(ValueError, match="decompression_backend must be"): + getattr(extension, entrypoint)( + 0, empty, empty, 0, empty, empty, decompression_backend="invalid" + ) diff --git a/tests/xdr/test_nvcomp_batch.py b/tests/xdr/test_nvcomp_batch.py index 07224bd..db3bf0d 100644 --- a/tests/xdr/test_nvcomp_batch.py +++ b/tests/xdr/test_nvcomp_batch.py @@ -933,3 +933,88 @@ def fail_decode(*args): [1], header_sizes=[10], ) + + +@pytest.mark.parametrize("backend", ["invalid", "hardware", "", True]) +def test_gpu_batch_rejects_invalid_backend_before_import(backend): + with pytest.raises(ValueError, match="decompression_backend must be"): + nvcomp_batch.gpu_gzip_decompress_batch( + None, [], [], [], decompression_backend=backend + ) + + +@pytest.mark.parametrize("extension", [None, SimpleNamespace()]) +def test_forced_cuda_rejects_missing_native_capability( + monkeypatch, extension +): + 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 + ) + with pytest.raises(RuntimeError, match="requires a rebuilt native"): + nvcomp_batch.gpu_gzip_decompress_batch( + SimpleNamespace(size=20), + [0], + [20], + [1], + header_sizes=[10], + decompression_backend="cuda", + ) + + +@pytest.mark.parametrize("backend", ["auto", "cuda"]) +@pytest.mark.parametrize("route", ["gzip", "deflate", "deflate_pooled"]) +def test_native_backend_forwarding(monkeypatch, backend, route): + calls = [] + output = SimpleNamespace(data=SimpleNamespace(ptr=200)) + owner = object() + pooled = route.endswith("_pooled") + + def record(*args, **kwargs): + calls.append((args, kwargs)) + return owner if pooled else None + + entrypoint = { + "gzip": "batch_gzip_decompress", + "deflate": "batch_deflate_decompress", + "deflate_pooled": "batch_deflate_decompress_pooled", + }[route] + extension = SimpleNamespace( + supports_decompression_backend=True, **{entrypoint: record} + ) + monkeypatch.setattr(nvcomp_batch, "_try_get_cpp_ext", lambda: extension) + monkeypatch.setattr( + nvcomp_batch, "_native_device_empty_uint8", lambda *a, **k: output + ) + monkeypatch.setitem( + sys.modules, + "cupy", + SimpleNamespace( + empty=lambda *a, **k: output, + uint8=np.uint8, + cuda=SimpleNamespace( + get_current_stream=lambda: SimpleNamespace(ptr=300) + ), + ), + ) + keepalive = [] if pooled else None + actual, _ = nvcomp_batch.gpu_gzip_decompress_batch( + SimpleNamespace(size=20, data=SimpleNamespace(ptr=100)), + [0], + [20], + [8], + gzip_wrapped=route == "gzip", + header_sizes=[10] if route == "gzip" else None, + use_native_pool=pooled, + keepalive=keepalive, + decompression_backend=backend, + ) + + assert actual is output + assert len(calls) == 1 + args, kwargs = calls[0] + assert args[6] == 300 + assert kwargs == {"decompression_backend": backend} + if pooled: + assert keepalive == [output, owner]