From 3d9230cad354a76c1d03e129103edf72e54d0b2a Mon Sep 17 00:00:00 2001 From: Trent Nelson Date: Tue, 29 Sep 2026 14:48:12 -0700 Subject: [PATCH 1/3] Fuse FITS tile byte conversion with GPU scatter Signed-off-by: Trent Nelson --- src/cuphoton/xdr/kernels.py | 89 ++++------------ src/cuphoton/xdr/reader.py | 36 ++----- tests/xdr/test_postprocess_gpu.py | 163 ++++++++++++++++++++++++++++++ 3 files changed, 190 insertions(+), 98 deletions(-) create mode 100644 tests/xdr/test_postprocess_gpu.py diff --git a/src/cuphoton/xdr/kernels.py b/src/cuphoton/xdr/kernels.py index ee5ef438..cbd805e6 100644 --- a/src/cuphoton/xdr/kernels.py +++ b/src/cuphoton/xdr/kernels.py @@ -94,6 +94,7 @@ def scatter_tiles_2d( d_tiles_concat, d_out, tile_byte_offsets, + tile_byte_lengths, tile_origins_row, tile_origins_col, tile_heights, @@ -102,8 +103,10 @@ def scatter_tiles_2d( tile_src_off_col, tile_full_widths, itemsize: int, + *, + shuffled: bool, ) -> None: - """Copy decompressed tile bytes into a 2D output image on the device. + """Unshuffle and byteswap decoded FITS tiles into a 2D output image. Supports partial-tile copies for section (ROI) reads: each tile copies the sub-region ``(src_off_row..src_off_row+h, src_off_col..src_off_col+w)`` @@ -111,14 +114,16 @@ def scatter_tiles_2d( row) into ``(orig_row..orig_row+h, orig_col..orig_col+w)`` in the output. For a full-image read pass ``src_off_*`` zeros and ``tile_full_widths`` - equal to the crop widths. + equal to the crop widths. ``shuffled`` selects the GZIP_2 byte-plane + layout; both layouts contain big-endian pixels on input. The decoded + bytes are read without modification. """ pix_stride_out = d_out.strides[0] # row stride in bytes - assert d_out.strides[1] == itemsize, "d_out must be C-contiguous" + assert d_out.strides[1] == itemsize, "output pixels must be contiguous" n_tiles = tile_byte_offsets.size - kern = _scatter_tiles_2d_kernel(itemsize) + kern = _scatter_tiles_2d_kernel(itemsize, shuffled) threads_per_block = 256 kern( (n_tiles,), @@ -127,6 +132,7 @@ def scatter_tiles_2d( d_tiles_concat, d_out, tile_byte_offsets, + tile_byte_lengths, tile_origins_row, tile_origins_col, tile_heights, @@ -140,15 +146,17 @@ def scatter_tiles_2d( @functools.cache -def _scatter_tiles_2d_kernel(itemsize: int): +def _scatter_tiles_2d_kernel(itemsize: int, shuffled: bool): import cupy as cp + name = f"scatter_fits_tiles_{itemsize}_{int(shuffled)}" src = rf""" extern "C" __global__ - void scatter_tiles_{itemsize}( + void {name}( const unsigned char* __restrict__ tiles, unsigned char* __restrict__ out, const long long* __restrict__ tile_byte_offsets, + const long long* __restrict__ tile_byte_lengths, const int* __restrict__ origins_row, const int* __restrict__ origins_col, const int* __restrict__ tile_h, @@ -161,6 +169,7 @@ def _scatter_tiles_2d_kernel(itemsize: int): const int t = blockIdx.x; const int ITEMSIZE = {itemsize}; const long long src_off = tile_byte_offsets[t]; + const long long n_pixels = tile_byte_lengths[t] / ITEMSIZE; const int r0 = origins_row[t]; const int c0 = origins_col[t]; const int h = tile_h[t]; @@ -173,77 +182,21 @@ def _scatter_tiles_2d_kernel(itemsize: int): int r = i / w; int c = i - r * w; int src_linear = (sr + r) * fw + (sc + c); - const unsigned char* src_pix = - tiles + src_off + (long long)src_linear * ITEMSIZE; unsigned char* dst_pix = out + (long long)(r0 + r) * out_row_stride_bytes + (long long)(c0 + c) * ITEMSIZE; #pragma unroll - for (int b = 0; b < ITEMSIZE; ++b) dst_pix[b] = src_pix[b]; - }} - }} - """ - return cp.RawKernel(src, f"scatter_tiles_{itemsize}") - - -@functools.cache -def _unshuffle_gzip2_kernel(itemsize: int): - import cupy as cp - - src = rf""" - extern "C" __global__ - void unshuffle_gzip2_{itemsize}( - const unsigned char* __restrict__ shuffled, - unsigned char* __restrict__ interleaved, - const long long* __restrict__ tile_byte_offsets, - const long long* __restrict__ tile_byte_lengths - ) {{ - const int t = blockIdx.x; - const int ITEMSIZE = {itemsize}; - const long long off = tile_byte_offsets[t]; - const long long total = tile_byte_lengths[t]; - const long long n_pixels = total / ITEMSIZE; - // Shuffled layout: - // [plane_0 (n_pixels bytes), plane_1, ..., plane_(ITEMSIZE-1)] - // Interleaved layout: [pix_0 (ITEMSIZE bytes), pix_1, ...] - for (long long p = threadIdx.x; p < n_pixels; p += blockDim.x) {{ for (int b = 0; b < ITEMSIZE; ++b) {{ - interleaved[off + p * ITEMSIZE + b] = - shuffled[off + (long long)b * n_pixels + p]; + const int source_byte = ITEMSIZE - 1 - b; + const long long byte_offset = {int(shuffled)} + ? (long long)source_byte * n_pixels + src_linear + : (long long)src_linear * ITEMSIZE + source_byte; + dst_pix[b] = tiles[src_off + byte_offset]; }} }} }} """ - return cp.RawKernel(src, f"unshuffle_gzip2_{itemsize}") - - -def unshuffle_gzip2_tiles( - d_shuffled, - d_interleaved, - tile_byte_offsets, - tile_byte_lengths, - itemsize: int, -): - """Reorder GZIP_2 tile bytes from plane-major to pixel-major. - - `d_shuffled` and `d_interleaved` must be the same size; both are uint8. - Each tile is processed independently (one CUDA block per tile). - """ - if itemsize == 1: - # Nothing to unshuffle; in-place no-op via alias. - return - kern = _unshuffle_gzip2_kernel(itemsize) - n_tiles = tile_byte_offsets.size - kern( - (n_tiles,), - (256,), - ( - d_shuffled, - d_interleaved, - tile_byte_offsets, - tile_byte_lengths, - ), - ) + return cp.RawKernel(src, name) @functools.cache diff --git a/src/cuphoton/xdr/reader.py b/src/cuphoton/xdr/reader.py index 5e6c89e6..8a9a4fba 100644 --- a/src/cuphoton/xdr/reader.py +++ b/src/cuphoton/xdr/reader.py @@ -21,11 +21,9 @@ byteswap_inplace, dequantize_int_to_float, scatter_tiles_2d, - unshuffle_gzip2_tiles, ) from .nvcomp_batch import ( _native_device_empty, - _native_device_empty_uint8, gpu_gzip_decompress_batch, ) @@ -168,8 +166,8 @@ def _validate_comp_geometry( class GpuCompImageReader: """Read a GZIP_1 (M3) / GZIP_2 (M4) CompImageHDU into a cupy.ndarray. - Pipeline: GDS heap read -> nvCOMP batched DEFLATE -> device-side byteswap - -> (GZIP_2 unshuffle, M4) -> scatter into output image. + Pipeline: GDS heap read -> nvCOMP batched DEFLATE -> fused device-side + unshuffle, byteswap and scatter into the output image. """ SUPPORTED_GZIP = frozenset({"GZIP_1", "GZIP_2"}) @@ -483,31 +481,7 @@ def postprocess_decoded_tiles( if keepalive is not None: keepalive.extend([d_tile_offsets, d_tile_byte_lengths]) - # 1. GZIP_2: byte-plane unshuffle per tile. - if plan["compression_type"] == "GZIP_2" and plan["itemsize"] > 1: - d_inter = ( - _native_device_empty_uint8(d_pixels.nbytes).reshape( - d_pixels.shape - ) - if use_native_pool - else cp.empty_like(d_pixels) - ) - if keepalive is not None: - keepalive.append(d_inter) - unshuffle_gzip2_tiles( - d_pixels, - d_inter, - d_tile_offsets, - d_tile_byte_lengths, - plan["itemsize"], - ) - d_pixels = d_inter - - # 2. Byteswap (FITS on-disk is big-endian, host is little-endian). - if plan["itemsize"] > 1: - byteswap_inplace(d_pixels, plan["itemsize"]) - - # 3. Scatter tiles into int/float output. + # Restore pixel order and native byte order while scattering. d_origins_r = cp.asarray(plan["origins_r"]) d_origins_c = cp.asarray(plan["origins_c"]) d_heights = cp.asarray(plan["heights"]) @@ -531,6 +505,7 @@ def postprocess_decoded_tiles( d_pixels, d_scatter_target, d_tile_offsets, + d_tile_byte_lengths, d_origins_r, d_origins_c, d_heights, @@ -539,9 +514,10 @@ def postprocess_decoded_tiles( d_src_off_c, d_tile_full_w, scatter_itemsize, + shuffled=plan["compression_type"] == "GZIP_2", ) - # 4. Dequantize int -> float when ZSCALE/ZZERO are present. + # Dequantize int -> float when ZSCALE/ZZERO are present. if plan["quantized"]: d_sel_zscale = cp.asarray(plan["sel_zscale"]) d_sel_zzero = cp.asarray(plan["sel_zzero"]) diff --git a/tests/xdr/test_postprocess_gpu.py b/tests/xdr/test_postprocess_gpu.py new file mode 100644 index 00000000..7634c798 --- /dev/null +++ b/tests/xdr/test_postprocess_gpu.py @@ -0,0 +1,163 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# +# SPDX-License-Identifier: Apache-2.0 + +"""Pixel-level checks of decoded FITS tiles, independent of the codec.""" + +from __future__ import annotations + +import numpy as np +import pytest + +from cuphoton.xdr.kernels import scatter_tiles_2d +from cuphoton.xdr.reader import GpuCompImageReader + + +@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") + return module + + +@pytest.mark.parametrize("dtype", ["u1", "i2", "i4", "i8", "f4", "f8"]) +@pytest.mark.parametrize("compression", ["GZIP_1", "GZIP_2"]) +@pytest.mark.parametrize("cropped", [False, True]) +def test_scatter_restores_fits_pixels(cp, dtype, compression, cropped): + rng = np.random.default_rng(42) + # Arbitrary bit patterns exercise signed values, NaNs and integer masks. + pixels = rng.bytes(17 * 19 * np.dtype(dtype).itemsize) + image = np.frombuffer(pixels, dtype=dtype).reshape(17, 19) + section = (slice(3, 15), slice(5, 18)) if cropped else np.s_[:, :] + raw, offsets, plan = _decoded_tiles(image, compression, section) + stream = cp.cuda.Stream(non_blocking=True) + with stream: + device_raw = cp.asarray(raw) + out = cp.empty(plan["out_shape"], dtype=dtype) + result = GpuCompImageReader.postprocess_decoded_tiles( + device_raw, offsets, plan, out=out, stream=stream + ) + stream.synchronize() + + assert result is out + np.testing.assert_array_equal( + cp.asnumpy(result).view("u1"), image[section].copy().view("u1") + ) + np.testing.assert_array_equal(cp.asnumpy(device_raw), raw) + + +@pytest.mark.parametrize("dtype", ["f4", "f8"]) +@pytest.mark.parametrize("compression", ["GZIP_1", "GZIP_2"]) +def test_scatter_preserves_nondithered_dequantization(cp, dtype, compression): + itemsize = np.dtype(dtype).itemsize + image = np.arange(-150, 173, dtype=f"i{itemsize}").reshape(17, 19) + section = np.s_[3:15, 5:18] + raw, offsets, plan = _decoded_tiles(image, compression, section) + plan.update( + quantized=True, + dtype_char=dtype, + zbitpix=-8 * itemsize, + sel_zscale=np.full(len(offsets), 0.5), + sel_zzero=np.full(len(offsets), 2.0), + ) + result = GpuCompImageReader.postprocess_decoded_tiles( + cp.asarray(raw), offsets, plan + ) + np.testing.assert_array_equal( + cp.asnumpy(result), image[section].astype(dtype) * 0.5 + 2 + ) + + +def test_scatter_respects_output_row_stride(cp): + image = np.arange(17 * 19, dtype="i4").reshape(17, 19) + raw, offsets, plan = _decoded_tiles(image, "GZIP_2", np.s_[:, :]) + padded = cp.full((17, 23), -1, dtype="i4") + out = padded[:, 2:21] + scatter_tiles_2d( + cp.asarray(raw), + out, + cp.asarray(offsets), + *( + cp.asarray(plan[name]) + for name in ( + "out_bytes", + "origins_r", + "origins_c", + "heights", + "widths", + "src_off_r", + "src_off_c", + "tile_full_w", + ) + ), + 4, + shuffled=True, + ) + actual = cp.asnumpy(padded) + np.testing.assert_array_equal(actual[:, 2:21], image) + assert np.all(actual[:, :2] == -1) + assert np.all(actual[:, 21:] == -1) + + +def _decoded_tiles(image, compression, section): + """Encode byte order and byte planes independently of FITS/GPU code.""" + y0, y1, _ = section[0].indices(image.shape[0]) + x0, x1, _ = section[1].indices(image.shape[1]) + plan = { + name: [] + for name in ( + "origins_r", + "origins_c", + "heights", + "widths", + "src_off_r", + "src_off_c", + "tile_full_w", + "out_bytes", + ) + } + payloads = [] + for r in range(0, image.shape[0], 6): + for c in range(0, image.shape[1], 8): + tile = image[r : r + 6, c : c + 8] + top, bottom = max(y0, r), min(y1, r + tile.shape[0]) + left, right = max(x0, c), min(x1, c + tile.shape[1]) + if top >= bottom or left >= right: + continue + big_endian = tile.astype(image.dtype.newbyteorder(">")) + data = np.frombuffer(big_endian.tobytes(), dtype="u1") + if compression == "GZIP_2": + data = data.reshape(-1, image.itemsize).T.ravel() + payloads.append(data) + for name, value in zip( + plan, + ( + top - y0, + left - x0, + bottom - top, + right - left, + top - r, + left - c, + tile.shape[1], + data.size, + ), + strict=True, + ): + plan[name].append(value) + plan = { + name: np.asarray(values, dtype="i8" if name == "out_bytes" else "i4") + for name, values in plan.items() + } + offsets = np.cumsum(plan["out_bytes"], dtype="i8") - plan["out_bytes"] + plan.update( + out_shape=(y1 - y0, x1 - x0), + itemsize=image.itemsize, + dtype_char=image.dtype.str[1:], + quantized=False, + compression_type=compression, + ) + return np.concatenate(payloads), offsets, plan From 9ddc04f600cdf64ef7c7858baab150820aa32c11 Mon Sep 17 00:00:00 2001 From: Trent Nelson Date: Tue, 29 Sep 2026 18:47:55 -0700 Subject: [PATCH 2/3] Expose FITS runtime choices through APIs and CLI Default to fused pixel restoration and retain a separate-kernel option. Carry explicit xDR options through workflow inputs, workers and reader receipts while preserving omitted manifest choices. Signed-off-by: Trent Nelson --- docs/components/xdr.md | 25 +++ docs/components/xfit.md | 7 + docs/components/xpois.md | 5 + docs/components/xrep.md | 6 + docs/components/xscan.md | 12 ++ src/cuphoton/core/cli/fits.py | 42 ++++- src/cuphoton/core/fits_io.py | 33 +++- src/cuphoton/core/fits_options.py | 41 +++++ src/cuphoton/xdr/benchmark_fits.py | 11 +- src/cuphoton/xdr/commands.py | 4 +- src/cuphoton/xdr/convenience.py | 4 +- src/cuphoton/xdr/kernels.py | 156 ++++++++++++++++++ src/cuphoton/xdr/prefetch.py | 38 ++++- src/cuphoton/xdr/reader.py | 67 +++++++- src/cuphoton/xfit/commands.py | 13 +- src/cuphoton/xfit/executor.py | 4 + src/cuphoton/xfit/fits_input.py | 29 +++- src/cuphoton/xfit/io.py | 8 +- src/cuphoton/xpois/batch.py | 6 + src/cuphoton/xpois/commands.py | 9 +- src/cuphoton/xpois/data.py | 39 ++++- src/cuphoton/xpois/workflows.py | 10 ++ src/cuphoton/xrep/commands.py | 8 +- src/cuphoton/xrep/io.py | 15 +- src/cuphoton/xrep/reproject.py | 12 ++ src/cuphoton/xrep/workflows.py | 22 ++- src/cuphoton/xscan/commands.py | 7 +- src/cuphoton/xscan/des.py | 24 ++- src/cuphoton/xscan/device_pipeline.py | 58 ++++++- src/cuphoton/xscan/dragon_pipeline.py | 5 + src/cuphoton/xscan/fits_inference.py | 15 +- src/cuphoton/xscan/lsstcomcam.py | 7 + .../xscan/pipeline_benchmark/runner.py | 15 +- src/cuphoton/xscan/pipeline_executor.py | 8 +- src/cuphoton/xscan/workflows.py | 6 + tests/core/test_cli_contract.py | 12 +- tests/core/test_fits_io.py | 36 ++++ tests/xdr/test_benchmark_fits.py | 4 + tests/xdr/test_frontend.py | 3 + tests/xdr/test_postprocess_gpu.py | 19 ++- tests/xdr/test_prefetch.py | 26 +-- tests/xdr/test_reader.py | 21 +++ tests/xfit/test_fits_input.py | 26 ++- tests/xpois/test_cli.py | 3 + tests/xpois/test_fits_ingestion.py | 1 + tests/xrep/test_cli.py | 41 +++++ tests/xrep/test_fits_ingestion.py | 29 +++- tests/xrep/test_workflows.py | 30 ++++ tests/xscan/test_cli.py | 2 + tests/xscan/test_fits_preparation.py | 2 + tests/xscan/test_pipeline_benchmark.py | 2 + tests/xscan/test_pipeline_benchmark_cli.py | 1 + tests/xscan/test_pipeline_executor.py | 3 + tests/xscan/test_pipeline_fits.py | 53 +++++- 54 files changed, 1000 insertions(+), 85 deletions(-) create mode 100644 src/cuphoton/core/fits_options.py diff --git a/docs/components/xdr.md b/docs/components/xdr.md index 3d172e4e..b3e8db04 100644 --- a/docs/components/xdr.md +++ b/docs/components/xdr.md @@ -12,6 +12,31 @@ 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"` restores those steps with individual kernels for comparison. +Both preserve FITS values, including integer masks and floating-point bits. + +```python +from cuphoton.xdr import batch_to_device + +(images,) = batch_to_device(paths, postprocess="auto") +``` + +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. + +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. + ## Read FITS images in a workflow The shared FITS reader selects explicit image HDUs and returns NumPy or CuPy diff --git a/docs/components/xfit.md b/docs/components/xfit.md index 5a917b74..5bdd7201 100644 --- a/docs/components/xfit.md +++ b/docs/components/xfit.md @@ -182,6 +182,13 @@ 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. + Input archives contain candidate identifiers and exact image pixels. Fit artifacts contain identifiers, hashes, parameters, uncertainties, covariance, and optional residuals. Confirm that the underlying data and metadata are cleared diff --git a/docs/components/xpois.md b/docs/components/xpois.md index ce4945b3..1f410e7e 100644 --- a/docs/components/xpois.md +++ b/docs/components/xpois.md @@ -48,6 +48,11 @@ 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. + Image, variance and mask HDU selection stays the same. `summary.json` records `fits_reader` and `fits_reads`, including the reader used and any fallback. The standalone fitting workflows retain host input arrays: xDR accelerates diff --git a/docs/components/xrep.md b/docs/components/xrep.md index 5002e825..0c1f2548 100644 --- a/docs/components/xrep.md +++ b/docs/components/xrep.md @@ -32,6 +32,12 @@ 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. + WCS mapping and target-grid setup read headers and dimensions without decompressing image pixels. Reading with xDR does not by itself imply native GPUDirect Storage; it can also decode images using KvikIO compatibility I/O. diff --git a/docs/components/xscan.md b/docs/components/xscan.md index bb206dc2..b75b5a54 100644 --- a/docs/components/xscan.md +++ b/docs/components/xscan.md @@ -55,6 +55,11 @@ 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. + The raw builders (`data-build-autoscan-raw`, `data-build-nodiff-raw`, and `data-build-lsstcomcam-smoke`) also accept `--fits-reader auto|astropy|xdr`. An explicit option overrides the manifest for that run; omitting it preserves @@ -471,6 +476,13 @@ 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. + ```bash : "${CUDA_VISIBLE_DEVICES:?must enumerate the allocated GPUs}" mpiexec -n 2 -x CUDA_VISIBLE_DEVICES cuphoton-openmpi-rank-exec -- \ diff --git a/src/cuphoton/core/cli/fits.py b/src/cuphoton/core/cli/fits.py index ead076bc..3c40481b 100644 --- a/src/cuphoton/core/cli/fits.py +++ b/src/cuphoton/core/cli/fits.py @@ -4,10 +4,50 @@ """Optional FITS reader overrides for manifest-driven commands.""" +from cuphoton.core.fits_options import ( + XDR_OPTION_CHOICES, + normalize_xdr_options, +) + from .invariants import SetInvariant -class FitsReaderOptions: +class XdrOptionsMixin: + """Optional xDR choices shared by FITS-consuming commands.""" + + xdr_postprocess = None + + class XdrPostprocessArg(SetInvariant): + _arg = "--xdr-postprocess" + _help = ( + "xDR postprocessing: auto, fused, or separate. Omitted preserves " + "manifest choices; applies when the FITS reader uses xDR." + ) + _set = XDR_OPTION_CHOICES["postprocess"] + _default = None + + +def xdr_options_from_cli(command) -> dict[str, str]: + """Collect explicit flags without replacing omitted manifest choices.""" + return normalize_xdr_options( + { + name: value + for name in XDR_OPTION_CHOICES + if (value := getattr(command, f"xdr_{name}", None)) is not None + } + ) + + +def xdr_option_cli_args(options) -> list[str]: + """Serialize explicit options for a component's benchmark subprocess.""" + return [ + argument + for name, value in normalize_xdr_options(options).items() + for argument in (f"--xdr-{name.replace('_', '-')}", value) + ] + + +class FitsReaderOptions(XdrOptionsMixin): """Preserve manifest policies unless a reader is explicitly selected.""" fits_reader = None diff --git a/src/cuphoton/core/fits_io.py b/src/cuphoton/core/fits_io.py index 1eb326d9..e8365ceb 100644 --- a/src/cuphoton/core/fits_io.py +++ b/src/cuphoton/core/fits_io.py @@ -13,7 +13,7 @@ from __future__ import annotations -from collections.abc import Sequence +from collections.abc import Mapping, Sequence from dataclasses import dataclass from pathlib import Path from typing import Any @@ -21,6 +21,8 @@ import numpy as np from astropy.io import fits +from .fits_options import normalize_xdr_options + _BITPIX_DTYPES = {8: "u1", 16: "i2", 32: "i4", 64: "i8", -32: "f4", -64: "f8"} @@ -45,6 +47,7 @@ class FitsReadResult: requested_reader: str fallback_reason: str | None device: bool + xdr_options: Mapping[str, str] | None = None def metadata(self) -> dict[str, Any]: """Return logical read provenance, without physical I/O claims.""" @@ -63,6 +66,11 @@ def metadata(self) -> dict[str, Any]: for info, array in zip(self.infos, self.arrays, strict=True) ], "decoded_bytes": sum(int(array.nbytes) for array in self.arrays), + **( + {"xdr_options": dict(self.xdr_options)} + if self.xdr_options + else {} + ), } @@ -195,7 +203,9 @@ def _xdr_available() -> bool: return gpu_available() and native_plan_files_available() -def _read_xdr(path: Path, hdus: tuple[int, ...], *, section, stream): +def _read_xdr( + path: Path, hdus: tuple[int, ...], *, section, stream, xdr_options +): from cuphoton.xdr import batch_to_device_stream return batch_to_device_stream( @@ -208,6 +218,7 @@ def _read_xdr(path: Path, hdus: tuple[int, ...], *, section, stream): batch_queue_depth=1, native_read_threads=1, native_plan_threads=1, + **xdr_options, ) @@ -230,6 +241,7 @@ def read_fits_images( device: bool = False, section=None, stream=None, + xdr_options: Mapping[str, str] | None = None, ) -> FitsReadResult: """Read selected planes to host or device, preserving FITS semantics. @@ -238,8 +250,13 @@ def read_fits_images( small stamp does not require a full-image GPU read. Explicit ``xdr`` 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. """ validate_fits_reader(reader) + xdr_options = normalize_xdr_options(xdr_options) selectors = tuple(hdus) if not selectors: raise ValueError("At least one FITS HDU must be selected") @@ -272,7 +289,13 @@ def read_fits_images( resolved = infos[0].path unique = tuple(dict.fromkeys(int(hdu) for hdu in selectors)) if actual == "xdr": - stacked = _read_xdr(resolved, unique, section=section, stream=stream) + stacked = _read_xdr( + resolved, + unique, + section=section, + stream=stream, + xdr_options=xdr_options, + ) if len(stacked) != len(unique): raise RuntimeError("xDR returned different HDU coverage") by_hdu = dict( @@ -325,4 +348,6 @@ def read_fits_images( raise RuntimeError( "FITS image shape or dtype changed during read" ) - return FitsReadResult(arrays, infos, actual, reader, fallback, device) + return FitsReadResult( + arrays, infos, actual, reader, fallback, device, xdr_options + ) diff --git a/src/cuphoton/core/fits_options.py b/src/cuphoton/core/fits_options.py new file mode 100644 index 00000000..c005a2f4 --- /dev/null +++ b/src/cuphoton/core/fits_options.py @@ -0,0 +1,41 @@ +# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +# +# SPDX-License-Identifier: Apache-2.0 + +"""CPU-only validation of optional xDR runtime choices.""" + +from collections.abc import Mapping + +XDR_OPTION_CHOICES = { + "postprocess": frozenset({"auto", "fused", "separate"}), +} + + +def normalize_xdr_options( + options: Mapping[str, str] | None = None, +) -> dict[str, str]: + """Copy explicit choices, leaving omitted choices to xDR defaults.""" + if options is None: + return {} + if not isinstance(options, Mapping): + raise ValueError("xdr_options must be a mapping") + result = {} + for name, value in options.items(): + if name not in XDR_OPTION_CHOICES: + raise ValueError(f"unsupported xDR option: {name!r}") + if ( + not isinstance(value, str) + or value not in XDR_OPTION_CHOICES[name] + ): + choices = ", ".join(sorted(XDR_OPTION_CHOICES[name])) + raise ValueError(f"xDR {name} must be one of: {choices}") + result[name] = value + return result + + +def merge_xdr_options( + base: Mapping[str, str] | None, + overrides: Mapping[str, str] | None, +) -> dict[str, str]: + """Override only supplied keys and preserve other manifest choices.""" + return {**normalize_xdr_options(base), **normalize_xdr_options(overrides)} diff --git a/src/cuphoton/xdr/benchmark_fits.py b/src/cuphoton/xdr/benchmark_fits.py index a000511c..8f58957b 100644 --- a/src/cuphoton/xdr/benchmark_fits.py +++ b/src/cuphoton/xdr/benchmark_fits.py @@ -25,6 +25,7 @@ from typing import TypedDict from cuphoton import __version__, xdr +from cuphoton.core.fits_options import normalize_xdr_options from cuphoton.xdr import ( batch_to_device, batch_to_device_stream, @@ -161,6 +162,7 @@ class _BatchOptions(TypedDict): native_read_threads: int native_plan_threads: int native_batcher: str | bool + postprocess: str def parse_hdu_indices(value: str) -> tuple[int, ...]: @@ -417,6 +419,7 @@ def bench_batch_load( native_batcher: str | bool, use_stream: bool, data_mb: float, + postprocess: 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" @@ -428,12 +431,13 @@ def bench_batch_load( native_read_threads=native_read_threads, native_plan_threads=native_plan_threads, native_batcher=native_batcher, + postprocess=postprocess, ) times: list[float] = [] note = ( f"{len(paths)} files x {len(tuple(hdu_indices))} HDUs, " - f"native_batcher={native_batcher}" + f"native_batcher={native_batcher}, postprocess={postprocess}" ) with nvtx_range(f"xdr.{phase}"): try: @@ -554,6 +558,7 @@ def run_benchmark( native_read_threads: int = 4, native_plan_threads: int = max(1, os.cpu_count() or 1), native_batcher: str = "auto", + postprocess: str = "auto", mock_storage_kind: str | None = None, skip_gds_read: bool = False, output_json: Path | None = None, @@ -566,6 +571,7 @@ def run_benchmark( helpers; the return value remains a list of phases. """ + normalize_xdr_options({"postprocess": postprocess}) paths = resolve_paths( fits_files, scan_dir=scan_dir, @@ -641,6 +647,7 @@ def run_benchmark( native_read_threads=native_read_threads, native_plan_threads=native_plan_threads, native_batcher=resolved_native_batcher, + postprocess=postprocess, use_stream=False, data_mb=total_data_mb, ) @@ -656,6 +663,7 @@ def run_benchmark( native_read_threads=native_read_threads, native_plan_threads=native_plan_threads, native_batcher=resolved_native_batcher, + postprocess=postprocess, use_stream=True, data_mb=total_data_mb, ) @@ -698,6 +706,7 @@ def run_benchmark( "native_read_threads": native_read_threads, "native_plan_threads": native_plan_threads, "native_batcher": native_batcher, + "postprocess": postprocess, "native_batcher_enabled": batcher_enabled, "native_batcher_error": batcher_error, "skip_gds_read": skip_gds_read, diff --git a/src/cuphoton/xdr/commands.py b/src/cuphoton/xdr/commands.py index c2518203..79a87901 100644 --- a/src/cuphoton/xdr/commands.py +++ b/src/cuphoton/xdr/commands.py @@ -18,6 +18,7 @@ StringInvariant, VariablePositionalInvariant, ) +from cuphoton.core.cli.fits import XdrOptionsMixin, xdr_options_from_cli from .benchmark_fits import parse_hdu_indices, run_benchmark @@ -41,7 +42,7 @@ class MockStorageInvariant(SetInvariant): _set = {"device", "host"} -class BenchmarkFitsCommand(InvariantAwareCommand): +class BenchmarkFitsCommand(XdrOptionsMixin, InvariantAwareCommand): """Benchmark xdr GPU FITS loading. It uses xDataReader's native CFITSIO planning path instead of Astropy @@ -165,6 +166,7 @@ def run(self) -> None: native_read_threads=self.native_read_threads, native_plan_threads=self.native_plan_threads, native_batcher=self.native_batcher, + **xdr_options_from_cli(self), mock_storage_kind=self.mock_storage, skip_gds_read=self.skip_gds_read, output_json=( diff --git a/src/cuphoton/xdr/convenience.py b/src/cuphoton/xdr/convenience.py index 5d5695cd..b601bd8f 100644 --- a/src/cuphoton/xdr/convenience.py +++ b/src/cuphoton/xdr/convenience.py @@ -42,6 +42,7 @@ def batch_to_device( native_read_threads: int | None = None, native_plan_threads: int | None = None, native_batcher: str | bool = "auto", + postprocess: str = "auto", ): """Load image HDUs from FITS files into stacked device arrays. @@ -63,7 +64,7 @@ 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 + native_read_threads, native_plan_threads, native_batcher, postprocess Passed through to `batch_to_device_stream` when ``parallel=True``. Returns @@ -87,6 +88,7 @@ def batch_to_device( native_read_threads=native_read_threads, native_plan_threads=native_plan_threads, native_batcher=native_batcher, + postprocess=postprocess, section=section, stream=stream, ) diff --git a/src/cuphoton/xdr/kernels.py b/src/cuphoton/xdr/kernels.py index cbd805e6..53dbf611 100644 --- a/src/cuphoton/xdr/kernels.py +++ b/src/cuphoton/xdr/kernels.py @@ -199,6 +199,162 @@ def _scatter_tiles_2d_kernel(itemsize: int, shuffled: bool): return cp.RawKernel(src, name) +def scatter_native_tiles_2d( + d_tiles_concat, + d_out, + tile_byte_offsets, + tile_origins_row, + tile_origins_col, + tile_heights, + tile_widths, + tile_src_off_row, + tile_src_off_col, + tile_full_widths, + itemsize: int, +) -> None: + """Copy decompressed tile bytes into a 2D output image on the device. + + Supports partial-tile copies for section (ROI) reads: each tile copies the + sub-region ``(src_off_row..src_off_row+h, src_off_col..src_off_col+w)`` + from its decompressed block (laid out as ``tile_full_width`` pixels per + row) into ``(orig_row..orig_row+h, orig_col..orig_col+w)`` in the output. + + For a full-image read pass ``src_off_*`` zeros and ``tile_full_widths`` + equal to the crop widths. + """ + + pix_stride_out = d_out.strides[0] # row stride in bytes + assert d_out.strides[1] == itemsize, "d_out must be C-contiguous" + n_tiles = tile_byte_offsets.size + + kern = _scatter_native_tiles_2d_kernel(itemsize) + threads_per_block = 256 + kern( + (n_tiles,), + (threads_per_block,), + ( + d_tiles_concat, + d_out, + tile_byte_offsets, + tile_origins_row, + tile_origins_col, + tile_heights, + tile_widths, + tile_src_off_row, + tile_src_off_col, + tile_full_widths, + np.int64(pix_stride_out), + ), + ) + + +@functools.cache +def _scatter_native_tiles_2d_kernel(itemsize: int): + import cupy as cp + + src = rf""" + extern "C" __global__ + void scatter_tiles_{itemsize}( + const unsigned char* __restrict__ tiles, + unsigned char* __restrict__ out, + const long long* __restrict__ tile_byte_offsets, + const int* __restrict__ origins_row, + const int* __restrict__ origins_col, + const int* __restrict__ tile_h, + const int* __restrict__ tile_w, + const int* __restrict__ src_off_row, + const int* __restrict__ src_off_col, + const int* __restrict__ tile_full_w, + long long out_row_stride_bytes + ) {{ + const int t = blockIdx.x; + const int ITEMSIZE = {itemsize}; + const long long src_off = tile_byte_offsets[t]; + const int r0 = origins_row[t]; + const int c0 = origins_col[t]; + const int h = tile_h[t]; + const int w = tile_w[t]; + const int sr = src_off_row[t]; + const int sc = src_off_col[t]; + const int fw = tile_full_w[t]; + const int total = h * w; + for (int i = threadIdx.x; i < total; i += blockDim.x) {{ + int r = i / w; + int c = i - r * w; + int src_linear = (sr + r) * fw + (sc + c); + const unsigned char* src_pix = + tiles + src_off + (long long)src_linear * ITEMSIZE; + unsigned char* dst_pix = + out + (long long)(r0 + r) * out_row_stride_bytes + + (long long)(c0 + c) * ITEMSIZE; + #pragma unroll + for (int b = 0; b < ITEMSIZE; ++b) dst_pix[b] = src_pix[b]; + }} + }} + """ + return cp.RawKernel(src, f"scatter_tiles_{itemsize}") + + +@functools.cache +def _unshuffle_gzip2_kernel(itemsize: int): + import cupy as cp + + src = rf""" + extern "C" __global__ + void unshuffle_gzip2_{itemsize}( + const unsigned char* __restrict__ shuffled, + unsigned char* __restrict__ interleaved, + const long long* __restrict__ tile_byte_offsets, + const long long* __restrict__ tile_byte_lengths + ) {{ + const int t = blockIdx.x; + const int ITEMSIZE = {itemsize}; + const long long off = tile_byte_offsets[t]; + const long long total = tile_byte_lengths[t]; + const long long n_pixels = total / ITEMSIZE; + // Shuffled layout: + // [plane_0 (n_pixels bytes), plane_1, ..., plane_(ITEMSIZE-1)] + // Interleaved layout: [pix_0 (ITEMSIZE bytes), pix_1, ...] + for (long long p = threadIdx.x; p < n_pixels; p += blockDim.x) {{ + for (int b = 0; b < ITEMSIZE; ++b) {{ + interleaved[off + p * ITEMSIZE + b] = + shuffled[off + (long long)b * n_pixels + p]; + }} + }} + }} + """ + return cp.RawKernel(src, f"unshuffle_gzip2_{itemsize}") + + +def unshuffle_gzip2_tiles( + d_shuffled, + d_interleaved, + tile_byte_offsets, + tile_byte_lengths, + itemsize: int, +): + """Reorder GZIP_2 tile bytes from plane-major to pixel-major. + + `d_shuffled` and `d_interleaved` must be the same size; both are uint8. + Each tile is processed independently (one CUDA block per tile). + """ + if itemsize == 1: + # Nothing to unshuffle; in-place no-op via alias. + return + kern = _unshuffle_gzip2_kernel(itemsize) + n_tiles = tile_byte_offsets.size + kern( + (n_tiles,), + (256,), + ( + d_shuffled, + d_interleaved, + tile_byte_offsets, + tile_byte_lengths, + ), + ) + + @functools.cache def _dequantize_int_to_float_kernel(int_itemsize: int, float_dtype_char: str): """Per-tile `out[p] = (float)int_buf[p] * zscale[t] + zzero[t]`. diff --git a/src/cuphoton/xdr/prefetch.py b/src/cuphoton/xdr/prefetch.py index a15d2ac3..2bfb7a86 100644 --- a/src/cuphoton/xdr/prefetch.py +++ b/src/cuphoton/xdr/prefetch.py @@ -29,6 +29,8 @@ import numpy as np +from cuphoton.core.fits_options import normalize_xdr_options + from .gds import available_cpu_cores, configure_kvikio_parallelism # Sentinel value pushed on the queue when the prefetcher has finished so the @@ -595,7 +597,11 @@ def run(self) -> None: def _consume_comp_batch( - entries: list[_CompBatchEntry], stream, keepalive=None + entries: list[_CompBatchEntry], + stream, + keepalive=None, + *, + postprocess="auto", ): """Decode many compressed HDUs with one batched nvCOMP call. @@ -731,6 +737,7 @@ def _consume_comp_batch( d_pixels[decoded_start:decoded_end], local_tile_offsets, entry.plan_item.plan, + postprocess=postprocess, out=entry.out, stream=stream, keepalive=keepalive, @@ -788,7 +795,12 @@ def _consume_image( def _consume_prefetched_group( - items: list[PrefetchedFile], outs: list, stream, keepalive=None + items: list[PrefetchedFile], + outs: list, + stream, + keepalive=None, + *, + postprocess="auto", ): """Consume prefetched files, batching compressed HDUs together.""" comp_entries: list[_CompBatchEntry] = [] @@ -820,7 +832,9 @@ def _consume_prefetched_group( else: raise AssertionError(f"unknown kind {plan_item.kind!r}") - _consume_comp_batch(comp_entries, stream, keepalive=keepalive) + _consume_comp_batch( + comp_entries, stream, keepalive=keepalive, postprocess=postprocess + ) def _native_file_plans( @@ -1051,6 +1065,7 @@ def _submit_prefetched_group( stream, owner, in_flight: deque[_GpuBatchHandle], + postprocess="auto", ) -> None: """Queue GPU work and register its lifetime handle.""" import cupy as cp @@ -1071,7 +1086,11 @@ def _submit_prefetched_group( try: with use_stream: _consume_prefetched_group( - items, outs, stream, keepalive=keepalive + items, + outs, + stream, + keepalive=keepalive, + postprocess=postprocess, ) candidate_event = cp.cuda.Event() candidate_event.record(use_stream) @@ -1397,6 +1416,7 @@ def _consume_native_batches( stream, NativeBatchBuilder, native_plan_files, + postprocess="auto", ): """Consume device batches built by the C++ KvikIO worker pool.""" import cupy as cp @@ -1472,6 +1492,7 @@ def _consume_native_batches( stream=stream, owner=native_batch, in_flight=in_flight, + postprocess=postprocess, ) planner.join() if planner.error is not None: @@ -1511,6 +1532,7 @@ def _consume_python_batches( batch_queue_depth: int, section, stream, + postprocess="auto", ) -> None: """Consume pinned-host batches with bounded event-owned lifetimes.""" n_files = len(paths) @@ -1531,6 +1553,7 @@ def submit(group: list[PrefetchedFile]) -> None: stream=stream, owner=group, in_flight=in_flight, + postprocess=postprocess, ) prefetcher.start() @@ -1598,6 +1621,7 @@ def batch_to_device_stream( native_read_threads: int | None = None, native_plan_threads: int | None = None, native_batcher: str | bool = "auto", + postprocess: str = "auto", section=None, stream=None, ): @@ -1641,6 +1665,9 @@ def batch_to_device_stream( "auto" uses the C++ KvikIO batch builder when available and falls back to the Python prefetcher otherwise. True requires the native builder. False always uses the Python prefetcher. + postprocess + "auto" (default) and "fused" restore FITS pixel order in one kernel. + "separate" uses individual unshuffle, byteswap and scatter kernels. section Optional 2D ROI applied uniformly to CompImageHDUs. stream @@ -1661,6 +1688,7 @@ def batch_to_device_stream( length ``len(paths)``. """ + normalize_xdr_options({"postprocess": postprocess}) hdu_indices = tuple(int(i) for i in hdu_indices) resolved_paths = [Path(p) for p in paths] n_files = len(resolved_paths) @@ -1710,6 +1738,7 @@ def batch_to_device_stream( stream=stream, NativeBatchBuilder=NativeBatchBuilder, native_plan_files=native_plan_files, + postprocess=postprocess, ) return tuple(outs) @@ -1722,6 +1751,7 @@ def batch_to_device_stream( batch_queue_depth=batch_queue_depth, section=section, stream=stream, + postprocess=postprocess, ) return tuple(outs) diff --git a/src/cuphoton/xdr/reader.py b/src/cuphoton/xdr/reader.py index 8a9a4fba..a29fda18 100644 --- a/src/cuphoton/xdr/reader.py +++ b/src/cuphoton/xdr/reader.py @@ -15,12 +15,16 @@ import numpy as np +from cuphoton.core.fits_options import normalize_xdr_options + from .gds import GdsHeapLoader, pread_to_device from .kernels import ( apply_bzero_bscale, byteswap_inplace, dequantize_int_to_float, + scatter_native_tiles_2d, scatter_tiles_2d, + unshuffle_gzip2_tiles, ) from .nvcomp_batch import ( _native_device_empty, @@ -60,6 +64,11 @@ def done(self): return False return bool(self.stream.done) + def synchronize(self) -> None: + if not self.io_complete: + raise RuntimeError("Cannot synchronize before GDS I/O completes") + self.stream.synchronize() + def _hdu_path(hdu) -> str: f = getattr(hdu, "_file", None) @@ -460,8 +469,10 @@ def postprocess_decoded_tiles( out=None, stream=None, keepalive=None, + postprocess: str = "auto", ): - """Scatter already-inflated tile bytes into the final image output.""" + """Restore FITS pixels; auto uses the fused kernel.""" + normalize_xdr_options({"postprocess": postprocess}) import cupy as cp use_native_pool = keepalive is not None @@ -481,6 +492,25 @@ def postprocess_decoded_tiles( if keepalive is not None: keepalive.extend([d_tile_offsets, d_tile_byte_lengths]) + shuffled = plan["compression_type"] == "GZIP_2" + if postprocess == "separate": + # Keep the caller's decoded bytes intact in both modes. + d_separate = cp.empty_like(d_pixels) + if keepalive is not None: + keepalive.append(d_separate) + if shuffled and scatter_itemsize > 1: + unshuffle_gzip2_tiles( + d_pixels, + d_separate, + d_tile_offsets, + d_tile_byte_lengths, + scatter_itemsize, + ) + else: + d_separate[...] = d_pixels + byteswap_inplace(d_separate, scatter_itemsize) + d_pixels = d_separate + # Restore pixel order and native byte order while scattering. d_origins_r = cp.asarray(plan["origins_r"]) d_origins_c = cp.asarray(plan["origins_c"]) @@ -501,11 +531,10 @@ def postprocess_decoded_tiles( d_tile_full_w, ] ) - scatter_tiles_2d( + scatter_args = ( d_pixels, d_scatter_target, d_tile_offsets, - d_tile_byte_lengths, d_origins_r, d_origins_c, d_heights, @@ -514,8 +543,16 @@ def postprocess_decoded_tiles( d_src_off_c, d_tile_full_w, scatter_itemsize, - shuffled=plan["compression_type"] == "GZIP_2", ) + if postprocess == "separate": + scatter_native_tiles_2d(*scatter_args) + else: + scatter_tiles_2d( + *scatter_args[:3], + d_tile_byte_lengths, + *scatter_args[3:], + shuffled=shuffled, + ) # Dequantize int -> float when ZSCALE/ZZERO are present. if plan["quantized"]: @@ -545,6 +582,7 @@ def decode_from_device_heap( out=None, stream=None, keepalive=None, + postprocess: str = "auto", ): """Run the decode pipeline on a pre-loaded compressed heap buffer. @@ -553,6 +591,7 @@ 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}) if keepalive is not None: keepalive.append(d_concat) d_pixels, tile_byte_offsets_np = ( @@ -574,9 +613,18 @@ def decode_from_device_heap( out=out, stream=stream, keepalive=keepalive, + postprocess=postprocess, ) - def read(self, *, out=None, stream=None, section=None, loader=None): + def read( + self, + *, + out=None, + stream=None, + section=None, + loader=None, + postprocess: str = "auto", + ): """Read and decode this HDU into a `cupy.ndarray`. Parameters @@ -595,6 +643,7 @@ def read(self, *, out=None, stream=None, section=None, loader=None): 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}) import cupy as cp plan = self.prepare_plan(section=section) @@ -663,7 +712,12 @@ def abandon(error): io_completion.io_complete = True keepalive.append(d_concat) return GpuCompImageReader.decode_from_device_heap( - d_concat, rel_offsets, plan, out=out, stream=None + d_concat, + rel_offsets, + plan, + out=out, + stream=None, + postprocess=postprocess, ) except BaseException as error: abandon(error) @@ -705,6 +759,7 @@ def abandon(error): out=out, stream=stream, keepalive=keepalive, + postprocess=postprocess, ) except (KeyboardInterrupt, SystemExit) as error: abandon(error) diff --git a/src/cuphoton/xfit/commands.py b/src/cuphoton/xfit/commands.py index 0cfa12f1..cb7f4506 100644 --- a/src/cuphoton/xfit/commands.py +++ b/src/cuphoton/xfit/commands.py @@ -25,6 +25,7 @@ StringInvariant, ) from cuphoton.core.cli.executor import ExecutorOptions +from cuphoton.core.cli.fits import XdrOptionsMixin, xdr_options_from_cli from ._types import ( BACKEND_REQUESTS, @@ -138,7 +139,7 @@ def _emit_json(self, payload: object) -> None: ) -class DataInspectCommand(XFitCommand): +class DataInspectCommand(XdrOptionsMixin, XFitCommand): """Inspect an NPZ candidate batch or FITS candidate manifest.""" # This CLI path intentionally replaces the base command input stream. @@ -151,11 +152,15 @@ class InputArg(ExistingNPZInvariant): def run(self) -> None: assert self.input is not None - dataset = self._call(load_xfit_dataset, self.input) + dataset = self._call( + load_xfit_dataset, + self.input, + xdr_options=xdr_options_from_cli(self), + ) self._emit_json(inspect_xfit_dataset(dataset)) -class _ValidatedDatasetCommand(XFitCommand): +class _ValidatedDatasetCommand(XdrOptionsMixin, XFitCommand): # This CLI path intentionally replaces the base command input stream. input: str | None = None # type: ignore[assignment] model: ModelName | None = None @@ -199,6 +204,7 @@ def _load_dataset(self, *, device: bool = False) -> XFitDataset: mode=self.mode, reader=reader, device=device, + xdr_options=xdr_options_from_cli(self), ) @@ -358,6 +364,7 @@ def run(self) -> None: "fits_reader", ) } + fit_options["xdr_options"] = xdr_options_from_cli(self) execution_result = self._call( run_workload, executor=self.executor, diff --git a/src/cuphoton/xfit/executor.py b/src/cuphoton/xfit/executor.py index 0f015341..05f74c39 100644 --- a/src/cuphoton/xfit/executor.py +++ b/src/cuphoton/xfit/executor.py @@ -24,6 +24,7 @@ json_mapping, read_json_mapping, ) +from cuphoton.core.fits_options import normalize_xdr_options from .api import DipoleFitResult, _floating_dtype, fit_dipoles from .io import ( @@ -62,6 +63,7 @@ "max_evaluations": None, "use_finite_difference": False, "fits_reader": "auto", + "xdr_options": None, } @@ -103,6 +105,7 @@ def _plan_xfit_chunks( raise ValueError("unsupported xFit compute dtype") if settings["fits_reader"] not in {"auto", "astropy", "xdr"}: raise ValueError("unsupported xFit FITS reader") + settings["xdr_options"] = normalize_xdr_options(settings["xdr_options"]) fits_plan = None if Path(input_path).suffix.lower() == ".json": from .fits_input import plan_fits_input @@ -182,6 +185,7 @@ def _load_input(options: Mapping[str, Any], *, device=False) -> XFitDataset: reader_options = { "reader": settings["fits_reader"] if device else "astropy", "device": device, + "xdr_options": settings["xdr_options"], } dataset = load_xfit_dataset( options["input_path"], diff --git a/src/cuphoton/xfit/fits_input.py b/src/cuphoton/xfit/fits_input.py index 90eb48a4..a891a8ea 100644 --- a/src/cuphoton/xfit/fits_input.py +++ b/src/cuphoton/xfit/fits_input.py @@ -15,6 +15,10 @@ from cuphoton.core.artifacts import file_sha256 from cuphoton.core.fits_io import inspect_fits_image, read_fits_images +from cuphoton.core.fits_options import ( + merge_xdr_options, + normalize_xdr_options, +) from ._types import FIT_MODES @@ -33,6 +37,7 @@ class FitsInputPlan: sources: tuple[dict[str, Any], ...] initial: np.ndarray | None stamp_basis: dict[str, Any] | None + xdr_options: dict[str, str] @property def batch_size(self) -> int: @@ -58,7 +63,7 @@ def _integer(value, name, *, minimum=0): def _plane(value, root, *, name, auxiliary=True): - keys = {"path", "hdu"} + keys = {"path", "hdu", "xdr_options"} if auxiliary: keys |= {"mask_hdu", "variance_hdu", "bad_mask_bits"} _mapping(value, keys, required=("path", "hdu"), name=name) @@ -67,6 +72,8 @@ def _plane(value, root, *, name, auxiliary=True): path = Path(value["path"]).expanduser() path = (root / path).resolve() result = {**value, "path": str(path)} + if "xdr_options" in result: + result["xdr_options"] = normalize_xdr_options(result["xdr_options"]) image = inspect_fits_image(path, _integer(value["hdu"], "hdu")) for key in ("mask_hdu", "variance_hdu"): if key not in value: @@ -103,6 +110,7 @@ def plan_fits_input(path, *, mode=None, model=None) -> FitsInputPlan: "candidates", "initial", "stamp_basis", + "xdr_options", }, required=("schema", "mode", "stamp_shape", "images", "candidates"), name="FITS manifest", @@ -235,6 +243,7 @@ def plan_fits_input(path, *, mode=None, model=None) -> FitsInputPlan: tuple(sources), initial, basis, + normalize_xdr_options(data.get("xdr_options")), ) @@ -293,12 +302,26 @@ def copy_stamps(): def load_fits_input( - path, *, mode=None, model=None, reader="auto", device=False + path, + *, + mode=None, + model=None, + reader="auto", + device=False, + xdr_options=None, ): """Read the candidate bounding region and keep GPU stamps on device.""" from .io import XFitDataset, _validate_dataset_contract plan = plan_fits_input(path, mode=mode, model=model) + xdr_options = normalize_xdr_options(xdr_options) + + def plane_options(plane): + return merge_xdr_options( + merge_xdr_options(plan.xdr_options, plane.get("xdr_options")), + xdr_options, + ) + ap = np if device: import cupy as ap @@ -333,6 +356,7 @@ def load_fits_input( reader=reader, device=device, section=None if full else section, + xdr_options=plane_options(plane), ) arrays = _candidate_stamps( result, @@ -372,6 +396,7 @@ def combine(values): [plan.stamp_basis["hdu"]], reader=reader, device=False, + xdr_options=plane_options(plan.stamp_basis), ) basis = result.arrays[0] metadata.append( diff --git a/src/cuphoton/xfit/io.py b/src/cuphoton/xfit/io.py index 71902890..fdef792c 100644 --- a/src/cuphoton/xfit/io.py +++ b/src/cuphoton/xfit/io.py @@ -307,6 +307,7 @@ def load_xfit_dataset( model: ModelName | None = None, reader: str = "auto", device: bool = False, + xdr_options=None, ) -> XFitDataset: """Load a safe NPZ batch or an explicit FITS candidate manifest.""" @@ -314,7 +315,12 @@ def load_xfit_dataset( from .fits_input import load_fits_input return load_fits_input( - path, mode=mode, model=model, reader=reader, device=device + path, + mode=mode, + model=model, + reader=reader, + device=device, + xdr_options=xdr_options, ) resolved = _require_npz_path(path) diff --git a/src/cuphoton/xpois/batch.py b/src/cuphoton/xpois/batch.py index d87bf3eb..4cb527fa 100644 --- a/src/cuphoton/xpois/batch.py +++ b/src/cuphoton/xpois/batch.py @@ -20,6 +20,7 @@ from astropy.io import fits from cuphoton.core.bulk import WorkItem, validate_identifier +from cuphoton.core.fits_options import normalize_xdr_options from .data import ( ERROR_EXTENSION_NAMES, @@ -442,6 +443,7 @@ class BatchFitOptions: flux_conserve: bool = False backend: str = "cupy" fits_reader: str = "auto" + xdr_options: Mapping[str, str] | None = None solver: str = "constant" spatial_degree: int | None = None als_iterations: int | None = None @@ -501,6 +503,9 @@ def __post_init__(self) -> None: raise ValueError(f"unsupported mask policy: {self.mask_policy}") if self.fits_reader not in {"auto", "astropy", "xdr"}: raise ValueError(f"unsupported FITS reader: {self.fits_reader}") + object.__setattr__( + self, "xdr_options", normalize_xdr_options(self.xdr_options) + ) if self.backend not in _SUPPORTED_BACKENDS: raise ValueError(f"unsupported fit backend: {self.backend}") if self.solver not in _SUPPORTED_SOLVERS: @@ -861,6 +866,7 @@ def run_image_pair_item( flux_conserve=options.flux_conserve, backend=options.backend, fits_reader=options.fits_reader, + xdr_options=options.xdr_options, solver=options.solver, spatial_degree=options.spatial_degree, als_iterations=options.als_iterations, diff --git a/src/cuphoton/xpois/commands.py b/src/cuphoton/xpois/commands.py index a158f20e..4fd78e3f 100644 --- a/src/cuphoton/xpois/commands.py +++ b/src/cuphoton/xpois/commands.py @@ -24,6 +24,7 @@ SetInvariant, StringInvariant, ) +from cuphoton.core.cli.fits import XdrOptionsMixin, xdr_options_from_cli from .batch import BatchFitOptions from .data import inspect_hsc_data_tree @@ -317,7 +318,9 @@ def _output_root(self) -> Path: return self.context.runs_dir -class _SpatialSolverOptionsCommand(_KernelSolveOptionsCommand): +class _SpatialSolverOptionsCommand( + XdrOptionsMixin, _KernelSolveOptionsCommand +): fits_reader = None class FitsReaderArg(SetInvariant): @@ -563,6 +566,7 @@ def run(self) -> None: als_regularization=self.als_regularization, workflow_name="fit_kernel", run_prefix="fit-kernel", + xdr_options=xdr_options_from_cli(self), ) self._warn_if_not_converged(result.summary) self._emit_json(result.summary) @@ -618,6 +622,7 @@ def run(self) -> None: als_regularization=self.als_regularization, workflow_name="subtract", run_prefix="subtract", + xdr_options=xdr_options_from_cli(self), ) self._warn_if_not_converged(result.summary) self._emit_json(result.summary) @@ -768,6 +773,7 @@ def run(self) -> None: als_iterations=self.als_iterations, als_tolerance=self.als_tolerance, als_regularization=self.als_regularization, + xdr_options=xdr_options_from_cli(self), ) common = { "manifest_path": Path(self.manifest).expanduser(), @@ -992,6 +998,7 @@ def run(self) -> None: als_iterations=self.als_iterations, als_tolerance=self.als_tolerance, als_regularization=self.als_regularization, + xdr_options=xdr_options_from_cli(self), ) self._emit_json(result.summary) diff --git a/src/cuphoton/xpois/data.py b/src/cuphoton/xpois/data.py index f2a57f7a..a7a771ae 100644 --- a/src/cuphoton/xpois/data.py +++ b/src/cuphoton/xpois/data.py @@ -141,7 +141,11 @@ def inspect_hsc_data_tree(base: Path | None = None) -> dict[str, Any]: def load_image_array( - path: Path, hdu: int | None = None, *, fits_reader: str = "astropy" + path: Path, + hdu: int | None = None, + *, + fits_reader: str = "astropy", + xdr_options=None, ) -> np.ndarray: """Load a two-dimensional FITS or NumPy image. @@ -158,28 +162,38 @@ def load_image_array( Floating-point image data. """ - array, _, _ = load_image_with_wcs(path, hdu=hdu, fits_reader=fits_reader) + array, _, _ = load_image_with_wcs( + path, hdu=hdu, fits_reader=fits_reader, xdr_options=xdr_options + ) return array def load_variance_array( - path: Path, hdu: int | None = None, *, fits_reader: str = "astropy" + path: Path, + hdu: int | None = None, + *, + fits_reader: str = "astropy", + xdr_options=None, ) -> np.ndarray: """Load a two-dimensional variance image from FITS or NumPy.""" array, _, _ = load_variance_with_wcs( - path, hdu=hdu, fits_reader=fits_reader + path, hdu=hdu, fits_reader=fits_reader, xdr_options=xdr_options ) return array def load_mask_array( - path: Path, hdu: int | None = None, *, fits_reader: str = "astropy" + path: Path, + hdu: int | None = None, + *, + fits_reader: str = "astropy", + xdr_options=None, ) -> np.ndarray: """Load a two-dimensional integer mask from FITS or NumPy.""" array, _, _ = load_mask_with_planes( - path, hdu=hdu, fits_reader=fits_reader + path, hdu=hdu, fits_reader=fits_reader, xdr_options=xdr_options ) return array @@ -221,8 +235,10 @@ def load_fit_positions( return positions.astype(np.int64, copy=False) -def _load_plane(path, hdu, *, reader, dtype, read_metadata): - result = read_fits_images(path, [hdu], reader=reader) +def _load_plane(path, hdu, *, reader, dtype, read_metadata, xdr_options=None): + result = read_fits_images( + path, [hdu], reader=reader, xdr_options=xdr_options + ) if read_metadata is not None: read_metadata.append(result.metadata()) return np.asarray(result.arrays[0], dtype=dtype) @@ -234,6 +250,7 @@ def load_image_with_wcs( *, fits_reader: str = "astropy", read_metadata: list[dict[str, Any]] | None = None, + xdr_options=None, ) -> tuple[np.ndarray, WCS | None, int | None]: """Load a host image with optional GPU FITS decompression and its WCS.""" resolved = path.expanduser().resolve() @@ -259,6 +276,7 @@ def load_image_with_wcs( reader=fits_reader, dtype=np.float64, read_metadata=read_metadata, + xdr_options=xdr_options, ) return array, WCS(info.header), info.hdu @@ -269,6 +287,7 @@ def load_variance_with_wcs( *, fits_reader: str = "astropy", read_metadata: list[dict[str, Any]] | None = None, + xdr_options=None, ) -> tuple[np.ndarray, WCS | None, int | None]: """Load variance, preserving the unambiguous named-HDU policy.""" resolved = path.expanduser().resolve() @@ -278,6 +297,7 @@ def load_variance_with_wcs( hdu, fits_reader=fits_reader, read_metadata=read_metadata, + xdr_options=xdr_options, ) if resolved.suffix.lower() not in FITS_SUFFIXES: raise ValueError(f"Unsupported image format: {resolved}") @@ -314,6 +334,7 @@ def load_variance_with_wcs( selected.hdu, fits_reader=fits_reader, read_metadata=read_metadata, + xdr_options=xdr_options, ) @@ -323,6 +344,7 @@ def load_mask_with_planes( *, fits_reader: str = "astropy", read_metadata: list[dict[str, Any]] | None = None, + xdr_options=None, ) -> tuple[np.ndarray, int | None, dict[str, int] | None]: """Load an integer mask and preserve its FITS mask-plane mapping.""" resolved = path.expanduser().resolve() @@ -371,6 +393,7 @@ def load_mask_with_planes( reader=fits_reader, dtype=np.int64, read_metadata=read_metadata, + xdr_options=xdr_options, ) return array, info.hdu, plane_map or None diff --git a/src/cuphoton/xpois/workflows.py b/src/cuphoton/xpois/workflows.py index 3139f653..ba95815c 100644 --- a/src/cuphoton/xpois/workflows.py +++ b/src/cuphoton/xpois/workflows.py @@ -139,6 +139,7 @@ def _prepare_kernel_inputs( fit_positions_path: Path | None = None, review: bool = False, fits_reader: str = "astropy", + xdr_options=None, ) -> _PreparedKernelInputs: prepare_start = time.perf_counter() input_read_sec = 0.0 @@ -201,12 +202,14 @@ def _prepare_kernel_inputs( hdu=reference_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) target, _, used_target_hdu = load_image_with_wcs( target_path, hdu=target_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) variance = None used_variance_hdu = None @@ -216,6 +219,7 @@ def _prepare_kernel_inputs( hdu=variance_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) finally: input_read_sec += time.perf_counter() - input_read_start @@ -251,6 +255,7 @@ def _prepare_kernel_inputs( hdu=reference_mask_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) target_mask, used_target_mask_hdu, target_plane_map = ( load_mask_with_planes( @@ -258,6 +263,7 @@ def _prepare_kernel_inputs( hdu=target_mask_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) ) finally: @@ -434,6 +440,7 @@ def run_constant_kernel_fit( als_regularization: float | None = None, workflow_name: str = "fit_kernel", run_prefix: str = "fit-kernel", + xdr_options=None, ) -> WorkflowResult: """Fit a kernel model and persist a reproducible subtraction run. @@ -515,6 +522,7 @@ def run_constant_kernel_fit( kernel_shape=kernel_shape, review=review, fit_positions_path=fit_positions_path, + xdr_options=xdr_options, ) solve_start = time.perf_counter() result: ConstantKernelFitResult | SpatialALSFitResult @@ -893,6 +901,7 @@ def benchmark_constant_kernel_backends( als_iterations: int | None = None, als_tolerance: float | None = None, als_regularization: float | None = None, + xdr_options=None, ) -> WorkflowResult: """Benchmark kernel-solver backends with numerical parity checks. @@ -961,6 +970,7 @@ def benchmark_constant_kernel_backends( auto_stamp_count=auto_stamp_count, auto_peak_percentile=auto_peak_percentile, kernel_shape=kernel_shape, + xdr_options=xdr_options, ) warmup_timings_by_backend: dict[str, list[dict[str, float]]] = {} diff --git a/src/cuphoton/xrep/commands.py b/src/cuphoton/xrep/commands.py index 97fa80db..6cb6c4a3 100644 --- a/src/cuphoton/xrep/commands.py +++ b/src/cuphoton/xrep/commands.py @@ -22,6 +22,7 @@ SetInvariant, StringInvariant, ) +from cuphoton.core.cli.fits import XdrOptionsMixin, xdr_options_from_cli from .workflows import ( benchmark_backend_variants_reproject_image, @@ -283,7 +284,7 @@ def run(self) -> None: self._emit_json(payload) -class _FitsReprojectionCommand(_SharedReprojectionCommand): +class _FitsReprojectionCommand(XdrOptionsMixin, _SharedReprojectionCommand): fits_reader = None class FitsReaderArg(SetInvariant): @@ -332,6 +333,7 @@ def run(self) -> None: mask_hdu=self.mask_hdu, backend=self.backend, fits_reader=self.fits_reader, + xdr_options=xdr_options_from_cli(self), interpolation=self.interpolation, grid_crval_ra=self.grid_crval_ra, grid_crval_dec=self.grid_crval_dec, @@ -379,6 +381,7 @@ def run(self) -> None: target_hdu=self.target_hdu, backend=self.backend, fits_reader=self.fits_reader, + xdr_options=xdr_options_from_cli(self), interpolation=self.interpolation, grid_crval_ra=self.grid_crval_ra, grid_crval_dec=self.grid_crval_dec, @@ -425,6 +428,7 @@ def run(self) -> None: mask_hdu=self.mask_hdu, backend=self.backend, fits_reader=self.fits_reader, + xdr_options=xdr_options_from_cli(self), interpolation=self.interpolation, grid_crval_ra=self.grid_crval_ra, grid_crval_dec=self.grid_crval_dec, @@ -512,6 +516,7 @@ def run(self) -> None: mask_hdu=self.mask_hdu, variants=self._csv_backend_variants(self.variants), fits_reader=self.fits_reader, + xdr_options=xdr_options_from_cli(self), reference_variant=self.reference_variant, mask_cases=self._csv_mask_cases(self.mask_cases), interpolation=self.interpolation, @@ -595,6 +600,7 @@ def run(self) -> None: mask_hdu=self.mask_hdu, backends=self._csv_backends(self.backends), fits_reader=self.fits_reader, + xdr_options=xdr_options_from_cli(self), reference_backend=self.reference_backend, interpolation=self.interpolation, grid_crval_ra=self.grid_crval_ra, diff --git a/src/cuphoton/xrep/io.py b/src/cuphoton/xrep/io.py index 1ac126df..755c5418 100644 --- a/src/cuphoton/xrep/io.py +++ b/src/cuphoton/xrep/io.py @@ -6,6 +6,7 @@ from __future__ import annotations +from collections.abc import Mapping from pathlib import Path from typing import Any @@ -31,13 +32,18 @@ def load_fits_image_with_wcs( *, hdu: int | None = None, fits_reader: str = "astropy", + xdr_options: Mapping[str, str] | None = None, device: bool = False, read_metadata: list[dict[str, Any]] | None = None, ) -> tuple[Any, WCS, fits.Header, int]: """Load a FITS plane on the host or device, with its unchanged WCS.""" info = inspect_fits_image(path, hdu=hdu) result = read_fits_images( - path, [info.hdu], reader=fits_reader, device=device + path, + [info.hdu], + reader=fits_reader, + device=device, + xdr_options=xdr_options, ) if read_metadata is not None: read_metadata.append(result.metadata()) @@ -49,13 +55,18 @@ def load_fits_mask( *, hdu: int | None = None, fits_reader: str = "astropy", + xdr_options: Mapping[str, str] | None = None, device: bool = False, read_metadata: list[dict[str, Any]] | None = None, ) -> tuple[Any, fits.Header, int]: """Load a FITS mask without narrowing its integer bit representation.""" info = inspect_fits_image(path, hdu=hdu) result = read_fits_images( - path, [info.hdu], reader=fits_reader, device=device + path, + [info.hdu], + reader=fits_reader, + device=device, + xdr_options=xdr_options, ) if read_metadata is not None: read_metadata.append(result.metadata()) diff --git a/src/cuphoton/xrep/reproject.py b/src/cuphoton/xrep/reproject.py index fdb7cbb1..f98d6396 100644 --- a/src/cuphoton/xrep/reproject.py +++ b/src/cuphoton/xrep/reproject.py @@ -6,6 +6,7 @@ from __future__ import annotations +from collections.abc import Mapping from pathlib import Path from typing import Any @@ -285,6 +286,7 @@ def reproject_fits( interpolation: str = "lanczos3", backend: str | None = None, fits_reader: str = "auto", + xdr_options: Mapping[str, str] | None = None, mapping_grid_step: int = 100, area_scaling: bool = True, ) -> ReprojectionResult: @@ -310,6 +312,9 @@ def reproject_fits( Coarse WCS mapping interval in output pixels. area_scaling Apply relative pixel-area scaling. + xdr_options + Optional xDR runtime choices, such as ``{"postprocess": "separate"}``, + forwarded to image and mask reads when using the xDR reader. Returns ------- @@ -328,6 +333,7 @@ def reproject_fits( path, hdu=hdu, fits_reader=reader, + xdr_options=xdr_options, device=backend == "cupy", read_metadata=fits_reads, ) @@ -347,6 +353,7 @@ def reproject_fits( mask_path, hdu=mask_hdu, fits_reader=reader, + xdr_options=xdr_options, device=backend == "cupy", read_metadata=fits_reads, ) @@ -462,6 +469,7 @@ def build_stack_spec_from_fits( mapping_grid_step: int = 100, area_scaling: bool = True, fits_reader: str = "astropy", + xdr_options: Mapping[str, str] | None = None, device: bool = False, read_metadata: list[dict[str, Any]] | None = None, ) -> tuple[list[np.ndarray], StackReprojectionSpec]: @@ -483,6 +491,9 @@ def build_stack_spec_from_fits( Coarse WCS mapping interval in output pixels. area_scaling Apply relative pixel-area scaling. + xdr_options + Optional xDR runtime choices, such as ``{"postprocess": "separate"}``, + forwarded to every source image when using the xDR reader. Returns ------- @@ -501,6 +512,7 @@ def build_stack_spec_from_fits( path, hdu=hdu, fits_reader=fits_reader, + xdr_options=xdr_options, device=device, read_metadata=read_metadata, ) diff --git a/src/cuphoton/xrep/workflows.py b/src/cuphoton/xrep/workflows.py index b5c7681d..7acbb6ec 100644 --- a/src/cuphoton/xrep/workflows.py +++ b/src/cuphoton/xrep/workflows.py @@ -10,6 +10,7 @@ import os import time import warnings +from collections.abc import Mapping from contextlib import contextmanager from dataclasses import dataclass from datetime import UTC, datetime @@ -124,6 +125,7 @@ def run_reproject_image( *, input_path: Path, fits_reader: str = "auto", + xdr_options: Mapping[str, str] | None = None, output_root: Path | None, name: str | None, hdu: int | None, @@ -150,6 +152,7 @@ def run_reproject_image( result, grid, timings = _execute_single_reprojection( input_path=input_path, fits_reader=fits_reader, + xdr_options=xdr_options, hdu=hdu, grid_crval_ra=grid_crval_ra, grid_crval_dec=grid_crval_dec, @@ -217,6 +220,7 @@ def benchmark_reproject_image( *, input_path: Path, fits_reader: str = "auto", + xdr_options: Mapping[str, str] | None = None, output_root: Path | None, name: str | None, hdu: int | None, @@ -253,6 +257,7 @@ def benchmark_reproject_image( _execute_single_reprojection( input_path=input_path, fits_reader=fits_reader, + xdr_options=xdr_options, hdu=hdu, grid_crval_ra=grid_crval_ra, grid_crval_dec=grid_crval_dec, @@ -268,6 +273,7 @@ def benchmark_reproject_image( result, grid, timings = _execute_single_reprojection( input_path=input_path, fits_reader=fits_reader, + xdr_options=xdr_options, hdu=hdu, grid_crval_ra=grid_crval_ra, grid_crval_dec=grid_crval_dec, @@ -354,6 +360,7 @@ def compare_backends_reproject_image( *, input_path: Path, fits_reader: str = "auto", + xdr_options: Mapping[str, str] | None = None, output_root: Path | None, name: str | None, hdu: int | None, @@ -400,6 +407,7 @@ def compare_backends_reproject_image( _execute_single_reprojection( input_path=input_path, fits_reader=fits_reader, + xdr_options=xdr_options, hdu=hdu, grid_crval_ra=grid_crval_ra, grid_crval_dec=grid_crval_dec, @@ -417,6 +425,7 @@ def compare_backends_reproject_image( result, grid, timings = _execute_single_reprojection( input_path=input_path, fits_reader=fits_reader, + xdr_options=xdr_options, hdu=hdu, grid_crval_ra=grid_crval_ra, grid_crval_dec=grid_crval_dec, @@ -538,6 +547,7 @@ def benchmark_backend_variants_reproject_image( *, input_path: Path, fits_reader: str = "auto", + xdr_options: Mapping[str, str] | None = None, output_root: Path | None, name: str | None, hdu: int | None, @@ -585,7 +595,11 @@ def benchmark_backend_variants_reproject_image( ) load_start = time.perf_counter() image, wcs, _, _ = load_fits_image_with_wcs( - input_path, hdu=hdu, fits_reader=reader, read_metadata=fits_reads + input_path, + hdu=hdu, + fits_reader=reader, + xdr_options=xdr_options, + read_metadata=fits_reads, ) grid = _resolve_grid( grid_crval_ra=grid_crval_ra, @@ -600,6 +614,7 @@ def benchmark_backend_variants_reproject_image( mask_path, hdu=mask_hdu, fits_reader=reader, + xdr_options=xdr_options, read_metadata=fits_reads, ) load_seconds = float(time.perf_counter() - load_start) @@ -822,6 +837,7 @@ def run_reproject_stack( *, input_paths: list[Path], fits_reader: str = "auto", + xdr_options: Mapping[str, str] | None = None, output_root: Path | None, name: str | None, hdu: int | None, @@ -897,6 +913,7 @@ def run_reproject_stack( images, spec = build_stack_spec_from_fits( input_paths, fits_reader=reader, + xdr_options=xdr_options, device=backend == "cupy", read_metadata=fits_reads, grid=grid, @@ -1021,6 +1038,7 @@ def _execute_single_reprojection( *, input_path: Path, fits_reader: str = "auto", + xdr_options: Mapping[str, str] | None = None, hdu: int | None, grid_crval_ra: float | None, grid_crval_dec: float | None, @@ -1046,6 +1064,7 @@ def _execute_single_reprojection( input_path, hdu=hdu, fits_reader=reader, + xdr_options=xdr_options, device=backend == "cupy", read_metadata=fits_reads, ) @@ -1062,6 +1081,7 @@ def _execute_single_reprojection( mask_path, hdu=mask_hdu, fits_reader=reader, + xdr_options=xdr_options, device=backend == "cupy", read_metadata=fits_reads, ) diff --git a/src/cuphoton/xscan/commands.py b/src/cuphoton/xscan/commands.py index 611ef6c9..9d5b5f15 100644 --- a/src/cuphoton/xscan/commands.py +++ b/src/cuphoton/xscan/commands.py @@ -26,7 +26,7 @@ StringInvariant, ) from cuphoton.core.cli.executor import ExecutorOptions -from cuphoton.core.cli.fits import FitsReaderOptions +from cuphoton.core.cli.fits import FitsReaderOptions, xdr_options_from_cli T = TypeVar("T") _DEFAULT_REVIEW_DIR = os.environ.get("CUPHOTON_XSCAN_REVIEW_DIR") @@ -501,6 +501,7 @@ def run(self) -> None: manifest_path=Path(self.manifest).expanduser(), output_dir=Path(self.output_dir).expanduser(), fits_reader=self.fits_reader, + xdr_options=xdr_options_from_cli(self), ) self._emit_json(result.summary) @@ -533,6 +534,7 @@ def run(self) -> None: manifest_path=Path(self.manifest).expanduser(), output_dir=Path(self.output_dir).expanduser(), fits_reader=self.fits_reader, + xdr_options=xdr_options_from_cli(self), ) self._emit_json(result.summary) @@ -580,6 +582,7 @@ def run(self) -> None: manifest_path=Path(self.manifest).expanduser(), output_dir=Path(self.output_dir).expanduser(), fits_reader=self.fits_reader, + xdr_options=xdr_options_from_cli(self), ) self._emit_json(result.summary) @@ -981,6 +984,7 @@ def run(self) -> None: output_root=Path(self.output_dir), run_id=self.run_name, fits_reader=self.fits_reader, + xdr_options=xdr_options_from_cli(self), **options, ) if result is not None: @@ -2368,6 +2372,7 @@ def run(self) -> None: config=self._path(self.config), items=self._path(self.items), fits_reader=self.fits_reader, + xdr_options=xdr_options_from_cli(self), images=self.images, image_size=self.image_size, candidates=self.candidates, diff --git a/src/cuphoton/xscan/des.py b/src/cuphoton/xscan/des.py index 5bb0c59f..355df35f 100644 --- a/src/cuphoton/xscan/des.py +++ b/src/cuphoton/xscan/des.py @@ -20,6 +20,7 @@ read_fits_images, validate_fits_reader, ) +from cuphoton.core.fits_options import merge_xdr_options from .dataset import ( INDEX_TO_SPLIT, @@ -89,6 +90,7 @@ def build_autoscan_dataset_from_raw( manifest_path: Path, output_dir: Path, fits_reader: str | None = None, + xdr_options=None, ) -> DatasetBuildResult: manifest = load_manifest(manifest_path) fits_reader = validate_fits_reader( @@ -96,6 +98,7 @@ def build_autoscan_dataset_from_raw( if fits_reader is None else fits_reader ) + xdr_options = merge_xdr_options(manifest.get("xdr_options"), xdr_options) fits_reads: list[dict[str, Any]] = [] rows = load_table_rows(resolve_required_path(manifest, "records_path")) if not rows: @@ -137,6 +140,7 @@ def build_autoscan_dataset_from_raw( hdu=search_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) template = load_stamp_or_image_cutout( template_path, @@ -146,6 +150,7 @@ def build_autoscan_dataset_from_raw( hdu=template_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) if not difference_path: raise ValueError("autoScan raw rows must include difference_path") @@ -157,6 +162,7 @@ def build_autoscan_dataset_from_raw( hdu=difference_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) if search.shape != template.shape or search.shape != difference.shape: raise ValueError( @@ -251,6 +257,7 @@ def build_nodiff_dataset_from_raw( manifest_path: Path, output_dir: Path, fits_reader: str | None = None, + xdr_options=None, ) -> DatasetBuildResult: manifest = load_manifest(manifest_path) fits_reader = validate_fits_reader( @@ -258,6 +265,7 @@ def build_nodiff_dataset_from_raw( if fits_reader is None else fits_reader ) + xdr_options = merge_xdr_options(manifest.get("xdr_options"), xdr_options) fits_reads: list[dict[str, Any]] = [] exposures = load_table_rows( resolve_required_path(manifest, "exposures_path") @@ -304,6 +312,7 @@ def build_nodiff_dataset_from_raw( hdu=search_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) template = load_image_array( get_required_alias( @@ -314,6 +323,7 @@ def build_nodiff_dataset_from_raw( hdu=template_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) if search.shape != template.shape: raise ValueError("search and template images must share a shape") @@ -329,6 +339,7 @@ def build_nodiff_dataset_from_raw( hdu=difference_hdu, fits_reader=fits_reader, read_metadata=fits_reads, + xdr_options=xdr_options, ) if difference.shape != search.shape: raise ValueError( @@ -1109,6 +1120,7 @@ def load_stamp_or_image_cutout( hdu: str | int | None = None, fits_reader: str = "astropy", read_metadata: list[dict[str, Any]] | None = None, + xdr_options=None, ) -> np.ndarray: resolved = Path(path).expanduser().resolve() if resolved.suffix.lower() in {".fits", ".fit", ".fts"}: @@ -1133,7 +1145,11 @@ def load_stamp_or_image_cutout( raise ValueError("stamp exceeds image bounds") section = (slice(y0, y1), slice(x0, x1)) result = read_fits_images( - resolved, [info.hdu], reader=fits_reader, section=section + resolved, + [info.hdu], + reader=fits_reader, + section=section, + xdr_options=xdr_options, ) if read_metadata is not None: read_metadata.append(result.metadata()) @@ -1143,6 +1159,7 @@ def load_stamp_or_image_cutout( hdu=hdu, fits_reader=fits_reader, read_metadata=read_metadata, + xdr_options=xdr_options, ) if array.ndim != 2: raise ValueError(f"array at {path} must be 2D") @@ -1166,6 +1183,7 @@ def load_image_array( *, fits_reader: str = "astropy", read_metadata: list[dict[str, Any]] | None = None, + xdr_options=None, ) -> np.ndarray: """Load a host image, optionally decompressing FITS pixels with xDR.""" resolved = Path(path).expanduser().resolve() @@ -1176,7 +1194,9 @@ def load_image_array( info = inspect_fits_image( resolved, hdu=_resolve_fits_hdu(resolved, hdu) ) - result = read_fits_images(resolved, [info.hdu], reader=fits_reader) + result = read_fits_images( + resolved, [info.hdu], reader=fits_reader, xdr_options=xdr_options + ) if read_metadata is not None: read_metadata.append(result.metadata()) return np.asarray(result.arrays[0]) diff --git a/src/cuphoton/xscan/device_pipeline.py b/src/cuphoton/xscan/device_pipeline.py index f2b7b829..d6c2ee0d 100644 --- a/src/cuphoton/xscan/device_pipeline.py +++ b/src/cuphoton/xscan/device_pipeline.py @@ -31,6 +31,10 @@ from cuphoton.core.artifacts import file_sha256 from cuphoton.core.bulk import json_mapping, validate_identifier +from cuphoton.core.fits_options import ( + merge_xdr_options, + normalize_xdr_options, +) DEVICE_PIPELINE_CONFIG_SCHEMA: Final = ( "cuphoton.xscan.device-pipeline.config/v2" @@ -237,6 +241,7 @@ class FitsArrayDescriptor: dtype: str hdu: int reader: Literal["astropy", "auto", "xdr"] = "auto" + xdr_options: Mapping[str, str] | None = None def __post_init__(self) -> None: if not isinstance(self.path, str): @@ -255,6 +260,9 @@ def __post_init__(self) -> None: raise ValueError("FITS descriptor hdu must be non-negative") if self.reader not in {"astropy", "auto", "xdr"}: raise ValueError("FITS reader must be astropy, auto or xdr") + object.__setattr__( + self, "xdr_options", normalize_xdr_options(self.xdr_options) + ) if not isinstance(self.dtype, str): raise TypeError("FITS descriptor dtype must be a string") try: @@ -276,14 +284,26 @@ def to_payload(self) -> dict[str, Any]: "dtype": self.dtype, "hdu": self.hdu, "reader": self.reader, + **( + {"xdr_options": dict(self.xdr_options)} + if self.xdr_options + else {} + ), } @classmethod def from_payload(cls, payload: Mapping[str, Any]) -> FitsArrayDescriptor: """Restore one strict FITS descriptor.""" + if not isinstance(payload, Mapping): + raise TypeError("FITS descriptor must be a mapping") + xdr_options = normalize_xdr_options(payload.get("xdr_options")) values = _require_exact_fields( - payload, + { + name: value + for name, value in payload.items() + if name != "xdr_options" + }, expected=frozenset( { "format", @@ -299,7 +319,7 @@ def from_payload(cls, payload: Mapping[str, Any]) -> FitsArrayDescriptor: ) if values.pop("format") != "fits": raise ValueError("unsupported image descriptor format") - return cls(**values) + return cls(**values, xdr_options=xdr_options) ImageArrayDescriptor = NpyArrayDescriptor | FitsArrayDescriptor @@ -433,10 +453,15 @@ def to_payload(self) -> dict[str, Any]: @classmethod def from_payload( - cls, payload: Mapping[str, Any], *, fits_reader: str | None = None + cls, + payload: Mapping[str, Any], + *, + fits_reader: str | None = None, + xdr_options: Mapping[str, str] | None = None, ) -> DevicePipelineItem: """Restore an item, optionally overriding its FITS reader policies.""" + xdr_options = normalize_xdr_options(xdr_options) if fits_reader is not None: from cuphoton.core.fits_io import validate_fits_reader @@ -464,12 +489,18 @@ def from_payload( raise TypeError("device pipeline candidates must be a list") def descriptor(value: Any) -> ImageArrayDescriptor: - if ( - fits_reader is not None - and isinstance(value, Mapping) - and value.get("format") == "fits" - ): - value = {**value, "reader": fits_reader} + if isinstance(value, Mapping) and value.get("format") == "fits": + value = { + **value, + **( + {"reader": fits_reader} + if fits_reader is not None + else {} + ), + "xdr_options": merge_xdr_options( + value.get("xdr_options"), xdr_options + ), + } return _image_descriptor(value) def optional_descriptor(value: Any) -> ImageArrayDescriptor | None: @@ -2878,6 +2909,7 @@ def _fits_input_groups( if ( previous.sha256 != descriptor.sha256 or previous.reader != descriptor.reader + or previous.xdr_options != descriptor.xdr_options or ( previous.hdu == descriptor.hdu and previous != descriptor @@ -2911,6 +2943,11 @@ def _validate_fits_reads(item: DevicePipelineItem, reads: Any) -> None: "hdus", "decoded_bytes", } + | ( + {"xdr_options"} + if isinstance(record, Mapping) and "xdr_options" in record + else set() + ) ), field_name="FITS reader receipt", ) @@ -2921,6 +2958,8 @@ def _validate_fits_reads(item: DevicePipelineItem, reads: Any) -> None: != {role: value.hdu for role, value in roles.items()} or any(type(hdu) is not int for hdu in record["roles"].values()) or record["requested_reader"] != descriptor.reader + or normalize_xdr_options(record.get("xdr_options")) + != descriptor.xdr_options or record["location"] != "device" or record["reader"] not in {"astropy", "xdr"} or ( @@ -2996,6 +3035,7 @@ def _read_fits_device_inputs( descriptor.path, hdus, reader=descriptor.reader, + xdr_options=descriptor.xdr_options, device=True, stream=stream, ) diff --git a/src/cuphoton/xscan/dragon_pipeline.py b/src/cuphoton/xscan/dragon_pipeline.py index 5f901095..3eba9e99 100644 --- a/src/cuphoton/xscan/dragon_pipeline.py +++ b/src/cuphoton/xscan/dragon_pipeline.py @@ -648,6 +648,11 @@ def _preflight_items( "shape": list(descriptor.shape), "dtype": descriptor.dtype, "reader": descriptor.reader, + **( + {"xdr_options": dict(descriptor.xdr_options)} + if descriptor.xdr_options + else {} + ), "roles": [role], } ) diff --git a/src/cuphoton/xscan/fits_inference.py b/src/cuphoton/xscan/fits_inference.py index aed14a0c..c442623b 100644 --- a/src/cuphoton/xscan/fits_inference.py +++ b/src/cuphoton/xscan/fits_inference.py @@ -116,7 +116,7 @@ def _plan_inputs(planes, centers_yx, candidate_ids, stamp_shape): return ids, shape, tuple(slices) -def _read_planes(planes, *, reader, stream): +def _read_planes(planes, *, reader, stream, xdr_options=None): groups: dict[Path, list[int]] = {} for plane in planes: groups.setdefault(plane.path, []).append(plane.hdu) @@ -125,7 +125,12 @@ def _read_planes(planes, *, reader, stream): for path, hdus in groups.items(): selected = tuple(dict.fromkeys(hdus)) result = read_fits_images( - path, selected, reader=reader, device=True, stream=stream + path, + selected, + reader=reader, + device=True, + stream=stream, + xdr_options=xdr_options, ) arrays.update( ((path, hdu), array) @@ -178,6 +183,7 @@ def predict_fits( difference: FitsPlane | None = None, xfit_features: torch.Tensor | None = None, fits_reader: str = "auto", + xdr_options=None, ) -> FitsInferenceResult: """Crop aligned FITS planes on one GPU and run the existing tensor model. @@ -241,7 +247,10 @@ def predict_fits( "xfit_features must contain only finite values" ) arrays, receipts = _read_planes( - planes, reader=fits_reader, stream=producer + planes, + reader=fits_reader, + stream=producer, + xdr_options=xdr_options, ) owners.extend(arrays) stamps = _crop_planes(cp, arrays, slices, shape, owners) diff --git a/src/cuphoton/xscan/lsstcomcam.py b/src/cuphoton/xscan/lsstcomcam.py index b1f56509..b41300bb 100644 --- a/src/cuphoton/xscan/lsstcomcam.py +++ b/src/cuphoton/xscan/lsstcomcam.py @@ -19,6 +19,7 @@ import numpy as np from cuphoton.core.fits_io import read_fits_images, validate_fits_reader +from cuphoton.core.fits_options import merge_xdr_options from .butler import ( _is_missing, @@ -395,6 +396,7 @@ def build_lsstcomcam_smoke_dataset_from_manifest( manifest_path: Path, output_dir: Path, fits_reader: str | None = None, + xdr_options=None, ) -> DatasetBuildResult: payload = load_manifest(manifest_path) fits_reader = validate_fits_reader( @@ -402,6 +404,7 @@ def build_lsstcomcam_smoke_dataset_from_manifest( if fits_reader is None else fits_reader ) + xdr_options = merge_xdr_options(payload.get("xdr_options"), xdr_options) fits_reads: list[dict[str, Any]] = [] registry_path = _registry_path_from_manifest(payload) sample_count = int(payload.get("sample_count", payload.get("limit", 8))) @@ -544,6 +547,7 @@ def build_lsstcomcam_smoke_dataset_from_manifest( read_metadata=fits_reads, center_x=sample.center_x, center_y=sample.center_y, + xdr_options=xdr_options, ) difference_stamp = read_fits_stamp( Path(str(difference_row["path"])), @@ -553,6 +557,7 @@ def build_lsstcomcam_smoke_dataset_from_manifest( read_metadata=fits_reads, center_x=search_stamp.center_x, center_y=search_stamp.center_y, + xdr_options=xdr_options, ) if difference_stamp.image_shape != search_stamp.image_shape: raise ValueError( @@ -1626,6 +1631,7 @@ def read_fits_stamp( center_y: int | None = None, fits_reader: str = "astropy", read_metadata: list[dict[str, Any]] | None = None, + xdr_options=None, ) -> FitsStamp: """Read a small centered FITS stamp without materializing full images.""" try: @@ -1657,6 +1663,7 @@ def read_fits_stamp( [used_hdu], reader=fits_reader, section=(slice(y0, y1), slice(x0, x1)), + xdr_options=xdr_options, ) if read_metadata is not None: read_metadata.append(result.metadata()) diff --git a/src/cuphoton/xscan/pipeline_benchmark/runner.py b/src/cuphoton/xscan/pipeline_benchmark/runner.py index f02b964e..50eee9b4 100644 --- a/src/cuphoton/xscan/pipeline_benchmark/runner.py +++ b/src/cuphoton/xscan/pipeline_benchmark/runner.py @@ -21,6 +21,8 @@ import numpy as np +from cuphoton.core.cli.fits import xdr_option_cli_args + CLI = [sys.executable, "-m", "cuphoton", "xscan", "benchmark-pipeline"] @@ -33,6 +35,7 @@ def read_inputs( items_path: Path, *, fits_reader: str | None = None, + xdr_options=None, ) -> tuple[Any, list[Any]]: from cuphoton.xscan.device_pipeline import ( DevicePipelineConfig, @@ -43,7 +46,9 @@ def read_inputs( json.loads(config_path.read_text()) ) items = [ - DevicePipelineItem.from_payload(value, fits_reader=fits_reader) + DevicePipelineItem.from_payload( + value, fits_reader=fits_reader, xdr_options=xdr_options + ) for value in json.loads(items_path.read_text()) ] if not items or len({item.item_id for item in items}) != len(items): @@ -123,6 +128,7 @@ def pipeline_worker(args: SimpleNamespace) -> None: args.config, args.items, fits_reader=getattr(args, "fits_reader", None), + xdr_options=getattr(args, "xdr_options", None), ) ordinal = int(config.device.split(":")[1]) torch.set_num_threads(config.inference_policy["worker_cpu_threads"]) @@ -218,6 +224,7 @@ def measure_pipeline(args: SimpleNamespace) -> dict[str, Any]: "--repeat", str(args.repeat), *(["--fits-reader", reader] if reader is not None else []), + *xdr_option_cli_args(getattr(args, "xdr_options", None)), ], output / "process.log", timeout=args.timeout, @@ -262,6 +269,9 @@ def measure_stages(args: SimpleNamespace) -> dict[str, Any]: if reader is not None else [] ), + *xdr_option_cli_args( + getattr(args, "xdr_options", None) + ), ], root / f"{stage}.log", timeout=args.timeout, @@ -321,6 +331,7 @@ def audit(args: SimpleNamespace) -> dict[str, Any]: args.config, args.items, fits_reader=getattr(args, "fits_reader", None), + xdr_options=getattr(args, "xdr_options", None), ) checks = [] for index in range(args.warmup + args.repeat): @@ -474,6 +485,7 @@ def run_benchmark(args: SimpleNamespace) -> None: args.config, args.items, fits_reader=getattr(args, "fits_reader", None), + xdr_options=getattr(args, "xdr_options", None), ) run_stage(args.stage, config, items, args.output) return @@ -496,6 +508,7 @@ def run_benchmark(args: SimpleNamespace) -> None: args.config, args.items, fits_reader=getattr(args, "fits_reader", None), + xdr_options=getattr(args, "xdr_options", None), ) report: dict[str, Any] = { "schema": "cuphoton.pipeline-stage-benchmark/v1", diff --git a/src/cuphoton/xscan/pipeline_executor.py b/src/cuphoton/xscan/pipeline_executor.py index 4de20087..0a13bd61 100644 --- a/src/cuphoton/xscan/pipeline_executor.py +++ b/src/cuphoton/xscan/pipeline_executor.py @@ -43,6 +43,7 @@ def load_pipeline_manifest( path: Path, *, fits_reader: str | None = None, + xdr_options: Mapping[str, str] | None = None, ) -> tuple[DevicePipelineConfig, tuple[DevicePipelineItem, ...]]: """Read config/items and resolve NPY or FITS paths beside the manifest. @@ -94,7 +95,9 @@ def resolve(raw: Any) -> str: descriptor["path"] = resolve(descriptor["path"]) item[role] = descriptor items.append( - DevicePipelineItem.from_payload(item, fits_reader=fits_reader) + DevicePipelineItem.from_payload( + item, fits_reader=fits_reader, xdr_options=xdr_options + ) ) return config, tuple(items) @@ -209,6 +212,7 @@ def run_pipeline_manifest( output_root: Path, run_id: str | None = None, fits_reader: str | None = None, + xdr_options: Mapping[str, str] | None = None, **options: Any, ) -> ExecutionResult | None: """Run an xPois/xFit/xScan manifest through the selected executor.""" @@ -217,7 +221,7 @@ def run_pipeline_manifest( def prepare(rank: int): config, items = load_pipeline_manifest( - manifest_path, fits_reader=fits_reader + manifest_path, fits_reader=fits_reader, xdr_options=xdr_options ) return prepare_pipeline_workload(items, config, rank=rank) diff --git a/src/cuphoton/xscan/workflows.py b/src/cuphoton/xscan/workflows.py index 688c01a3..e9ad063a 100644 --- a/src/cuphoton/xscan/workflows.py +++ b/src/cuphoton/xscan/workflows.py @@ -749,11 +749,13 @@ def build_raw_autoscan_workflow( manifest_path: Path, output_dir: Path, fits_reader: str | None = None, + xdr_options=None, ) -> WorkflowResult: result = build_autoscan_dataset_from_raw( manifest_path=manifest_path, output_dir=output_dir, fits_reader=fits_reader, + xdr_options=xdr_options, ) return WorkflowResult(run_dir=result.output_dir, summary=result.summary) @@ -763,11 +765,13 @@ def build_raw_nodiff_workflow( manifest_path: Path, output_dir: Path, fits_reader: str | None = None, + xdr_options=None, ) -> WorkflowResult: result = build_nodiff_dataset_from_raw( manifest_path=manifest_path, output_dir=output_dir, fits_reader=fits_reader, + xdr_options=xdr_options, ) return WorkflowResult(run_dir=result.output_dir, summary=result.summary) @@ -801,11 +805,13 @@ def build_lsstcomcam_smoke_workflow( manifest_path: Path, output_dir: Path, fits_reader: str | None = None, + xdr_options=None, ) -> WorkflowResult: result = build_lsstcomcam_smoke_dataset_from_manifest( manifest_path=manifest_path, output_dir=output_dir, fits_reader=fits_reader, + xdr_options=xdr_options, ) return WorkflowResult(run_dir=result.output_dir, summary=result.summary) diff --git a/tests/core/test_cli_contract.py b/tests/core/test_cli_contract.py index be4af97c..a97dc7a9 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, 14, 0), - ("xfit", 3, 3, 27, 1), - ("xpois", 7, 7, 150, 1), - ("xscan", 44, 44, 209, 1), - ("xrep", 6, 6, 108, 1), + ("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), ("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) == 900 + assert sum(item[3] for item in per_group) == 918 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 2644d043..438f43ec 100644 --- a/tests/core/test_fits_io.py +++ b/tests/core/test_fits_io.py @@ -15,6 +15,42 @@ from cuphoton.core import fits_io +@pytest.mark.parametrize("postprocess", ["fused", "separate"]) +def test_xdr_options_reach_batch_dispatch(tmp_path, monkeypatch, postprocess): + import cuphoton.xdr as xdr + + data = np.arange(12, dtype=np.float32).reshape(3, 4) + path = _write(tmp_path, data, compression="GZIP_2") + observed = {} + + def batch(paths, hdus, **options): + observed.update(options) + return [data[None]] + + monkeypatch.setattr(fits_io, "_xdr_available", lambda: True) + monkeypatch.setattr(xdr, "batch_to_device_stream", batch) + result = fits_io.read_fits_images( + path, + [1], + reader="xdr", + device=True, + xdr_options={"postprocess": postprocess}, + ) + assert observed["postprocess"] == postprocess + assert result.metadata()["xdr_options"] == {"postprocess": postprocess} + np.testing.assert_array_equal(result.arrays[0], data) + + +@pytest.mark.parametrize( + "options", ["fused", {"unknown": "auto"}, {"postprocess": True}] +) +def test_invalid_xdr_options_fail_before_input_access(tmp_path, options): + with pytest.raises(ValueError, match="xDR|xdr_options"): + fits_io.read_fits_images( + tmp_path / "missing.fits", [1], xdr_options=options + ) + + def _write(tmp_path, array, *, compression=None, **kwargs): path = tmp_path / "images.fits" hdu = ( diff --git a/tests/xdr/test_benchmark_fits.py b/tests/xdr/test_benchmark_fits.py index db6607ec..487ceee7 100644 --- a/tests/xdr/test_benchmark_fits.py +++ b/tests/xdr/test_benchmark_fits.py @@ -163,6 +163,8 @@ def fake_run_benchmark(fits_files, **kwargs): "2,1,2", "--native-batcher", "off", + "--xdr-postprocess", + "separate", "--mock-storage", "host", "--skip-gds-read", @@ -183,6 +185,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["mock_storage_kind"] == "host" assert received["skip_gds_read"] is True assert received["output_json"] == Path("report.json") @@ -396,6 +399,7 @@ def test_json_report_matches_returned_and_printed_phases( "native_read_threads": 5, "native_plan_threads": 6, "native_batcher": "auto", + "postprocess": "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 198b70f7..23bf49dd 100644 --- a/tests/xdr/test_frontend.py +++ b/tests/xdr/test_frontend.py @@ -45,6 +45,7 @@ def test_batch_api_keeps_existing_keyword_surface(): "native_read_threads", "native_plan_threads", "native_batcher", + "postprocess", ) actual = tuple(inspect.signature(xdr.batch_to_device).parameters) @@ -75,6 +76,7 @@ def fake_stream(paths, **kwargs): native_read_threads=6, native_plan_threads=7, native_batcher=True, + postprocess="separate", ) assert result == ("ok",) @@ -90,6 +92,7 @@ def fake_stream(paths, **kwargs): "native_read_threads": 6, "native_plan_threads": 7, "native_batcher": True, + "postprocess": "separate", "section": (slice(0, 1), slice(0, 1)), "stream": "stream", }, diff --git a/tests/xdr/test_postprocess_gpu.py b/tests/xdr/test_postprocess_gpu.py index 7634c798..3768e4b2 100644 --- a/tests/xdr/test_postprocess_gpu.py +++ b/tests/xdr/test_postprocess_gpu.py @@ -24,10 +24,13 @@ def cp(): return module +@pytest.mark.parametrize("postprocess", ["auto", "fused", "separate"]) @pytest.mark.parametrize("dtype", ["u1", "i2", "i4", "i8", "f4", "f8"]) @pytest.mark.parametrize("compression", ["GZIP_1", "GZIP_2"]) @pytest.mark.parametrize("cropped", [False, True]) -def test_scatter_restores_fits_pixels(cp, dtype, compression, cropped): +def test_scatter_restores_fits_pixels( + cp, dtype, compression, cropped, postprocess +): rng = np.random.default_rng(42) # Arbitrary bit patterns exercise signed values, NaNs and integer masks. pixels = rng.bytes(17 * 19 * np.dtype(dtype).itemsize) @@ -39,7 +42,12 @@ def test_scatter_restores_fits_pixels(cp, dtype, compression, cropped): device_raw = cp.asarray(raw) out = cp.empty(plan["out_shape"], dtype=dtype) result = GpuCompImageReader.postprocess_decoded_tiles( - device_raw, offsets, plan, out=out, stream=stream + device_raw, + offsets, + plan, + out=out, + stream=stream, + postprocess=postprocess, ) stream.synchronize() @@ -50,9 +58,12 @@ def test_scatter_restores_fits_pixels(cp, dtype, compression, cropped): np.testing.assert_array_equal(cp.asnumpy(device_raw), raw) +@pytest.mark.parametrize("postprocess", ["auto", "separate"]) @pytest.mark.parametrize("dtype", ["f4", "f8"]) @pytest.mark.parametrize("compression", ["GZIP_1", "GZIP_2"]) -def test_scatter_preserves_nondithered_dequantization(cp, dtype, compression): +def test_scatter_preserves_nondithered_dequantization( + cp, dtype, compression, postprocess +): itemsize = np.dtype(dtype).itemsize image = np.arange(-150, 173, dtype=f"i{itemsize}").reshape(17, 19) section = np.s_[3:15, 5:18] @@ -65,7 +76,7 @@ def test_scatter_preserves_nondithered_dequantization(cp, dtype, compression): sel_zzero=np.full(len(offsets), 2.0), ) result = GpuCompImageReader.postprocess_decoded_tiles( - cp.asarray(raw), offsets, plan + cp.asarray(raw), offsets, plan, postprocess=postprocess ) np.testing.assert_array_equal( cp.asnumpy(result), image[section].astype(dtype) * 0.5 + 2 diff --git a/tests/xdr/test_prefetch.py b/tests/xdr/test_prefetch.py index 8894573c..0e27ed1c 100644 --- a/tests/xdr/test_prefetch.py +++ b/tests/xdr/test_prefetch.py @@ -548,7 +548,8 @@ def test_consume_group_carries_precomputed_header_sizes(monkeypatch): ) captured = [] - def consume(entries, stream, keepalive=None): + def consume(entries, stream, keepalive=None, **options): + assert options["postprocess"] == "separate" captured.extend(entries) monkeypatch.setattr(prefetch, "_consume_comp_batch", consume) @@ -557,6 +558,7 @@ def consume(entries, stream, keepalive=None): [item], [np.empty((1, 1), dtype=np.uint8)], stream=None, + postprocess="separate", ) assert captured[0].header_sizes is header_sizes @@ -843,7 +845,7 @@ class Owner: ), ) - def consume(items, outs, use_stream, keepalive): + def consume(items, outs, use_stream, keepalive, **options): assert use_stream is stream trace.append("consume") keepalive.append("device-owner") @@ -1929,7 +1931,7 @@ def synchronize(self): state["active"] -= 1 synchronized.append(self.index) - def submit(group, outs, *, stream, owner, in_flight): + def submit(group, outs, *, stream, owner, in_flight, **options): index = owner[0].file_index state["active"] += 1 state["maximum"] = max(state["maximum"], state["active"]) @@ -1972,7 +1974,7 @@ def synchronize(self): handle = prefetch._GpuBatchHandle(Event(), [items[0]], Stream()) - def submit(group, outs, *, stream, owner, in_flight): + def submit(group, outs, *, stream, owner, in_flight, **options): trace.append(f"submit-{owner[0].file_index}") in_flight.append(handle) @@ -2051,7 +2053,7 @@ def request_stop(self): def join(self): trace.append("planner-join") - def submit(group, outs, *, stream, owner, in_flight): + def submit(group, outs, *, stream, owner, in_flight, **options): handle = prefetch._GpuBatchHandle( Event(len(handles)), [owner, outs], object() ) @@ -2134,7 +2136,7 @@ def done(self): def synchronize(self): trace.append("event-synchronize") - def submit(group, outs, *, stream, owner, in_flight): + def submit(group, outs, *, stream, owner, in_flight, **options): trace.append("submit") in_flight.append( prefetch._GpuBatchHandle(Event(), [owner, outs], object()) @@ -2241,7 +2243,7 @@ def request_stop(self): ) monkeypatch.setattr(prefetch, "_NativeBatchPlanner", Planner) - def submit(items, outs, *, stream, owner, in_flight): + def submit(items, outs, *, stream, owner, in_flight, **options): index = owner[0].index trace.append(f"submit-{index}") state["active"] += 1 @@ -2360,7 +2362,7 @@ def request_stop(self): monkeypatch.setattr(prefetch, "_NativeBatchPlanner", Planner) submitted_handles = [] - def submit(items, outs, *, stream, owner, in_flight): + def submit(items, outs, *, stream, owner, in_flight, **options): index = owner[0].index trace.append(f"submit-{index}") if index == 1: @@ -2417,7 +2419,7 @@ def done(self): def synchronize(self): trace.append("wait-0") - def submit(group, outs, *, stream, owner, in_flight): + def submit(group, outs, *, stream, owner, in_flight, **options): index = owner[0].file_index trace.append(f"submit-{index}") if index == 1: @@ -2472,7 +2474,7 @@ def synchronize(self): trace.append("stream-sync") raise RuntimeError("sticky stream failure") - def submit(group, outs, *, stream, owner, in_flight): + def submit(group, outs, *, stream, owner, in_flight, **options): index = owner[0].file_index trace.append(f"submit-{index}") if index == 1: @@ -2524,7 +2526,7 @@ class Stream: def synchronize(self): trace.append("stream-sync") - def submit(group, outs, *, stream, owner, in_flight): + def submit(group, outs, *, stream, owner, in_flight, **options): trace.append("submit") handle = prefetch._GpuBatchHandle(Event(), [owner], Stream()) in_flight.append(handle) @@ -2647,7 +2649,7 @@ def join(self): output = np.zeros((2, 1), dtype=np.float32) consumed = [] - def consume(items, outs, stream): + def consume(items, outs, stream, **options): consumed.extend(items) outs[0][0] = 1.0 diff --git a/tests/xdr/test_reader.py b/tests/xdr/test_reader.py index c7e845ff..72b6f913 100644 --- a/tests/xdr/test_reader.py +++ b/tests/xdr/test_reader.py @@ -475,3 +475,24 @@ def test_comp_geometry_rejects_nonpositive_dims(data_shape, tile_shape): def test_comp_geometry_accepts_positive_dims_case(): _validate_comp_geometry((100, 100), (64, 64)) + + +@pytest.mark.parametrize("io_complete", [False, True]) +def test_read_completion_requires_io_before_cuda_wait(io_complete): + from cuphoton.xdr.reader import _ReadCompletion + + waits = [] + + class Stream: + def synchronize(self): + waits.append("cuda") + + completion = _ReadCompletion(Stream()) + completion.io_complete = io_complete + if io_complete: + completion.synchronize() + assert waits == ["cuda"] + else: + with pytest.raises(RuntimeError, match="GDS I/O completes"): + completion.synchronize() + assert waits == [] diff --git a/tests/xfit/test_fits_input.py b/tests/xfit/test_fits_input.py index b3b8aba6..b977d35f 100644 --- a/tests/xfit/test_fits_input.py +++ b/tests/xfit/test_fits_input.py @@ -146,15 +146,39 @@ def load(path, **kwargs): monkeypatch.setattr(executor, "load_xfit_dataset", load) items, options = executor.plan_xfit_chunks( - path, chunk_size=1, fit_options={"fits_reader": "xdr"} + path, + chunk_size=1, + fit_options={ + "fits_reader": "xdr", + "xdr_options": {"postprocess": "separate"}, + }, ) worker = executor.XFitWorker(options) assert calls[0]["device"] is True assert calls[0]["reader"] == "xdr" + assert calls[0]["xdr_options"] == {"postprocess": "separate"} assert len(items) == 2 worker.close() +@pytest.mark.parametrize( + "override,expected", + [(None, "separate"), ({"postprocess": "fused"}, "fused")], +) +def test_fits_manifest_options_preserved_and_explicitly_overridden( + tmp_path, override, expected +): + path, _, _, _, _ = _input(tmp_path) + manifest = json.loads(path.read_text()) + manifest["xdr_options"] = {"postprocess": "auto"} + manifest["images"][0]["xdr_options"] = {"postprocess": "separate"} + path.write_text(json.dumps(manifest)) + dataset = load_xfit_dataset(path, reader="astropy", xdr_options=override) + assert dataset.reader_metadata[0]["xdr_options"] == { + "postprocess": expected + } + + def test_source_hash_change_is_rejected_even_with_unchanged_stat(tmp_path): path, _, _, _, _ = _input(tmp_path) _, options = executor.plan_xfit_chunks(path, chunk_size=1, fit_options={}) diff --git a/tests/xpois/test_cli.py b/tests/xpois/test_cli.py index c96dc428..2af093dd 100644 --- a/tests/xpois/test_cli.py +++ b/tests/xpois/test_cli.py @@ -136,6 +136,8 @@ def fake_fit(**kwargs): "spatial-als", "--fits-reader", "xdr", + "--xdr-postprocess", + "separate", "--backend", "cupy", "--spatial-degree", @@ -155,6 +157,7 @@ def fake_fit(**kwargs): assert seen["solver"] == "spatial-als" assert seen["backend"] == "cupy" assert seen["fits_reader"] == "xdr" + assert seen["xdr_options"] == {"postprocess": "separate"} assert seen["spatial_degree"] == 3 assert seen["als_iterations"] == 17 assert seen["als_tolerance"] == pytest.approx(2e-7) diff --git a/tests/xpois/test_fits_ingestion.py b/tests/xpois/test_fits_ingestion.py index 677607b8..c3b985a3 100644 --- a/tests/xpois/test_fits_ingestion.py +++ b/tests/xpois/test_fits_ingestion.py @@ -91,5 +91,6 @@ def test_fits_reader_survives_distributed_options_round_trip(): basis_sigmas=(1,), basis_degrees=(0,), fits_reader="xdr", + xdr_options={"postprocess": "separate"}, ) assert BatchFitOptions.from_payload(options.to_payload()) == options diff --git a/tests/xrep/test_cli.py b/tests/xrep/test_cli.py index 3c44ccb6..7342d58b 100644 --- a/tests/xrep/test_cli.py +++ b/tests/xrep/test_cli.py @@ -5,8 +5,10 @@ from __future__ import annotations from pathlib import Path +from types import SimpleNamespace import numpy as np +import pytest from astropy.io import fits from cuphoton.core.cli import get_component, run_component @@ -107,6 +109,7 @@ def test_inspect_image_help_omits_fits_reader(capsys) -> None: assert rc == 0 assert "--fits-reader" not in captured.out + assert "--xdr-postprocess" not in captured.out def test_inspect_image_command_emits_json(tmp_path: Path, capsys) -> None: @@ -117,3 +120,41 @@ def test_inspect_image_command_emits_json(tmp_path: Path, capsys) -> None: assert rc == 0 assert str(fits_path) in captured.out + + +@pytest.mark.parametrize( + "command,workflow,input_flag", + [ + ("reproject-image", "run_reproject_image", "--input"), + ("reproject-stack", "run_reproject_stack", "--inputs"), + ("benchmark-reproject-image", "benchmark_reproject_image", "--input"), + ( + "benchmark-backend-variants", + "benchmark_backend_variants_reproject_image", + "--input", + ), + ("compare-backends", "compare_backends_reproject_image", "--input"), + ], +) +def test_fits_commands_forward_xdr_options( + tmp_path, monkeypatch, capsys, command, workflow, input_flag +): + from cuphoton.xrep import commands + + captured = {} + + def run(**kwargs): + captured.update(kwargs) + return SimpleNamespace(summary={}) + + monkeypatch.setattr(commands, workflow, run) + path = _write_test_fits(tmp_path / "image.fits") + assert ( + _run_cli( + [command, input_flag, str(path), "--xdr-postprocess", "separate"] + ) + == 0 + ) + assert captured["xdr_options"] == {"postprocess": "separate"} + assert captured["fits_reader"] == "auto" + capsys.readouterr() diff --git a/tests/xrep/test_fits_ingestion.py b/tests/xrep/test_fits_ingestion.py index 5e13266c..7d6ce523 100644 --- a/tests/xrep/test_fits_ingestion.py +++ b/tests/xrep/test_fits_ingestion.py @@ -5,6 +5,7 @@ """FITS geometry inspection must not decode image payloads.""" import numpy as np +import pytest from astropy.io import fits from cuphoton.xrep import Grid, io, make_north_up_wcs, mapping @@ -52,10 +53,12 @@ def record(*args, **kwargs): grid=Grid.from_wcs(wcs), backend="cpu", interpolation="bilinear", + xdr_options={"postprocess": "separate"}, ) assert len(calls) == 1 assert calls[0]["reader"] == "astropy" assert calls[0]["device"] is False + assert calls[0]["xdr_options"] == {"postprocess": "separate"} assert result.metadata["fits_reads"][0]["reader"] == "astropy" @@ -68,8 +71,11 @@ def test_mask_host_reader_preserves_unsigned_bits(tmp_path): np.testing.assert_array_equal(loaded, mask) -def test_device_image_loader_retains_reader_array_ownership( - tmp_path, monkeypatch +@pytest.mark.parametrize( + "loader", [io.load_fits_image_with_wcs, io.load_fits_mask] +) +def test_device_loader_retains_reader_array_ownership_and_options( + tmp_path, monkeypatch, loader ): from types import SimpleNamespace @@ -85,10 +91,23 @@ def device_read(path, hdus, **kwargs): monkeypatch.setattr(io, "read_fits_images", device_read) receipts = [] - array, _, _, hdu = io.load_fits_image_with_wcs( - path, fits_reader="xdr", device=True, read_metadata=receipts + array, *_, hdu = loader( + path, + fits_reader="xdr", + device=True, + read_metadata=receipts, + xdr_options={"postprocess": "separate"}, ) assert array is sentinel assert hdu == 0 - assert requests == [([0], {"reader": "xdr", "device": True})] + assert requests == [ + ( + [0], + { + "reader": "xdr", + "device": True, + "xdr_options": {"postprocess": "separate"}, + }, + ) + ] assert receipts == [{"reader": "xdr"}] diff --git a/tests/xrep/test_workflows.py b/tests/xrep/test_workflows.py index 9b458748..feb22a75 100644 --- a/tests/xrep/test_workflows.py +++ b/tests/xrep/test_workflows.py @@ -58,6 +58,7 @@ def test_run_reproject_image_persists_summary_and_arrays( ) result = run_reproject_image( + xdr_options={"postprocess": "separate"}, input_path=fits_path, output_root=tmp_path / "runs", name="reproject-image-test", @@ -74,6 +75,11 @@ def test_run_reproject_image_persists_summary_and_arrays( write_fits=True, ) + assert all( + receipt["xdr_options"] == {"postprocess": "separate"} + and receipt["reader"] == "astropy" + for receipt in result.summary["fits_reads"] + ) assert result.run_dir.name == "reproject-image-test" assert result.summary["requested_backend"] == "cpu" assert result.summary["backend"] == "cpu" @@ -132,6 +138,7 @@ def test_run_reproject_stack_persists_stack_outputs(tmp_path: Path) -> None: ) result = run_reproject_stack( + xdr_options={"postprocess": "separate"}, input_paths=[path_a, path_b], output_root=tmp_path / "runs", name="reproject-stack-test", @@ -146,6 +153,11 @@ def test_run_reproject_stack_persists_stack_outputs(tmp_path: Path) -> None: write_fits=True, ) + assert all( + receipt["xdr_options"] == {"postprocess": "separate"} + and receipt["reader"] == "astropy" + for receipt in result.summary["fits_reads"] + ) assert result.summary["stack_shape"][0] == 2 assert result.summary["runtime"]["backend"] == "cpu" assert result.summary["device"] == "cpu" @@ -165,6 +177,7 @@ def test_benchmark_reproject_image_reports_split_timings( ) result = benchmark_reproject_image( + xdr_options={"postprocess": "separate"}, input_path=fits_path, output_root=tmp_path / "runs", name="benchmark-reproject-image-test", @@ -183,6 +196,11 @@ def test_benchmark_reproject_image_reports_split_timings( warmup=0, ) + assert all( + receipt["xdr_options"] == {"postprocess": "separate"} + and receipt["reader"] == "astropy" + for receipt in result.summary["fits_reads"] + ) assert result.summary["workflow"] == "benchmark-reproject-image" timings = result.summary["timings"] assert set(timings) == { @@ -208,6 +226,7 @@ def test_compare_backends_reports_parity_and_timings( ) result = compare_backends_reproject_image( + xdr_options={"postprocess": "separate"}, input_path=fits_path, output_root=tmp_path / "runs", name="compare-backends-test", @@ -229,6 +248,11 @@ def test_compare_backends_reports_parity_and_timings( rtol=1e-12, ) + assert all( + receipt["xdr_options"] == {"postprocess": "separate"} + and receipt["reader"] == "astropy" + for receipt in result.summary["fits_reads"]["cpu"] + ) assert result.summary["workflow"] == "compare-backends" assert result.summary["backends"] == ["cpu"] assert result.summary["parity"]["ok"] is True @@ -249,6 +273,7 @@ def test_benchmark_backend_variants_reports_masks_and_parity( ) result = benchmark_backend_variants_reproject_image( + xdr_options={"postprocess": "separate"}, input_path=fits_path, output_root=tmp_path / "runs", name="benchmark-backend-variants-test", @@ -271,6 +296,11 @@ def test_benchmark_backend_variants_reports_masks_and_parity( rtol=1e-12, ) + assert all( + receipt["xdr_options"] == {"postprocess": "separate"} + and receipt["reader"] == "astropy" + for receipt in result.summary["fits_reads"] + ) assert result.summary["workflow"] == "benchmark-backend-variants" assert result.summary["mask_cases"] == ["none", "mask"] assert result.summary["parity"]["ok"] is True diff --git a/tests/xscan/test_cli.py b/tests/xscan/test_cli.py index 2da399a0..f458bcf7 100644 --- a/tests/xscan/test_cli.py +++ b/tests/xscan/test_cli.py @@ -77,8 +77,10 @@ def build(**kwargs): args = [command, "--manifest", "inputs.json", "--output-dir", "dataset"] if reader is not None: args.extend(["--fits-reader", reader]) + args.extend(["--xdr-postprocess", "fused"]) assert _run_cli(args) == 0 assert received[0]["fits_reader"] == reader + assert received[0]["xdr_options"] == {"postprocess": "fused"} received.clear() assert _run_cli([*args, "--fits-reader", "gds"]) != 0 assert not received diff --git a/tests/xscan/test_fits_preparation.py b/tests/xscan/test_fits_preparation.py index 6dc25e03..227c8333 100644 --- a/tests/xscan/test_fits_preparation.py +++ b/tests/xscan/test_fits_preparation.py @@ -176,6 +176,7 @@ def test_dataset_fits_reader_override_reaches_decoding_and_summary( expected, ): workflow, module, manifest, expected_reads = fits_dataset_builder + manifest["xdr_options"] = {"postprocess": "separate"} if manifest_reader is not None: manifest["fits_reader"] = manifest_reader path = tmp_path / "manifest.json" @@ -185,6 +186,7 @@ def test_dataset_fits_reader_override_reaches_decoding_and_summary( def read_on_cpu(*args, reader, **kwargs): readers.append(reader) + assert kwargs["xdr_options"] == {"postprocess": "separate"} result = original(*args, reader="astropy", **kwargs) return replace(result, requested_reader=reader) diff --git a/tests/xscan/test_pipeline_benchmark.py b/tests/xscan/test_pipeline_benchmark.py index 7b4a7bd5..fdd0d949 100644 --- a/tests/xscan/test_pipeline_benchmark.py +++ b/tests/xscan/test_pipeline_benchmark.py @@ -87,6 +87,7 @@ def test_reader_override_is_forwarded_to_benchmark_children( def child(command, log, **kwargs): commands.append(command) assert command[command.index("--fits-reader") + 1] == "astropy" + assert command[command.index("--xdr-postprocess") + 1] == "separate" stage = command[command.index("--stage") + 1] root = Path(command[command.index("--output") + 1]) if stage != "pipeline": @@ -103,6 +104,7 @@ def child(command, log, **kwargs): config=tmp_path / "config.json", items=tmp_path / "items.json", fits_reader="astropy", + xdr_options={"postprocess": "separate"}, warmup=1, repeat=1, timeout=10, diff --git a/tests/xscan/test_pipeline_benchmark_cli.py b/tests/xscan/test_pipeline_benchmark_cli.py index 2658f798..ed424daf 100644 --- a/tests/xscan/test_pipeline_benchmark_cli.py +++ b/tests/xscan/test_pipeline_benchmark_cli.py @@ -33,6 +33,7 @@ def test_benchmark_defaults_and_path_conversion(monkeypatch): "config": None, "items": None, "fits_reader": None, + "xdr_options": {}, "images": 4, "image_size": 256, "candidates": 9, diff --git a/tests/xscan/test_pipeline_executor.py b/tests/xscan/test_pipeline_executor.py index fd31bde6..99a2bc5c 100644 --- a/tests/xscan/test_pipeline_executor.py +++ b/tests/xscan/test_pipeline_executor.py @@ -338,6 +338,8 @@ def run(**kwargs): "3", "--fits-reader", "astropy", + "--xdr-postprocess", + "separate", ], ) == 0 @@ -345,6 +347,7 @@ def run(**kwargs): assert json.loads(capsys.readouterr().out) == {"executor": executor} assert seen["run_id"] == "qualified" assert seen["fits_reader"] == "astropy" + assert seen["xdr_options"] == {"postprocess": "separate"} assert seen["benchmark"].to_payload() == { "warmup_rounds": 1, "measure_rounds": 3, diff --git a/tests/xscan/test_pipeline_fits.py b/tests/xscan/test_pipeline_fits.py index d05c4528..9c04b650 100644 --- a/tests/xscan/test_pipeline_fits.py +++ b/tests/xscan/test_pipeline_fits.py @@ -262,14 +262,24 @@ def test_fits_worker_grouped_device_read_and_receipt(inputs, monkeypatch): from cuphoton.core import fits_io item, config = inputs + item = pipeline.DevicePipelineItem.from_payload( + item.to_payload(), xdr_options={"postprocess": "separate"} + ) calls = [] stream = object() cpu_read = fits_io.read_fits_images def device_read(path, hdus, **options): - assert options == {"reader": "auto", "device": True, "stream": stream} + assert options == { + "reader": "auto", + "device": True, + "stream": stream, + "xdr_options": {"postprocess": "separate"}, + } calls.append(tuple(hdus)) - result = cpu_read(path, hdus, reader="astropy") + result = cpu_read( + path, hdus, reader="astropy", xdr_options=options["xdr_options"] + ) return replace( result, requested_reader="auto", @@ -300,6 +310,7 @@ def device_read(path, hdus, **options): ("reader", "gds"), ("roles", {}), ("fallback_reason", None), + ("xdr_options", {"postprocess": "fused"}), ): changed = copy.deepcopy(receipt["fits_reads"]) changed[0][field] = replacement @@ -307,6 +318,44 @@ def device_read(path, hdus, **options): pipeline._validate_fits_reads(item, changed) +@pytest.mark.parametrize( + "override,expected", + [(None, "fused"), ({"postprocess": "separate"}, "separate")], +) +def test_xdr_options_survive_manifest_workers_and_identity( + inputs, tmp_path, override, expected +): + item, config = inputs + payload = item.to_payload() + for role in ("reference", "target", "variance", "fit_mask"): + payload[role]["xdr_options"] = {"postprocess": "fused"} + manifest = tmp_path / "runtime.json" + manifest.write_text( + json.dumps( + { + "schema": pipeline_executor.PIPELINE_MANIFEST_SCHEMA, + "configuration": config.to_payload(), + "items": [payload], + } + ) + ) + _, (loaded,) = pipeline_executor.load_pipeline_manifest( + manifest, xdr_options=override + ) + work, identities = dragon_pipeline._preflight_items((loaded,)) + restored = pipeline.DevicePipelineItem.from_payload(work[0].payload) + assert restored.reference.xdr_options == {"postprocess": expected} + assert identities[0]["files"][0]["images"][0]["xdr_options"] == { + "postprocess": expected + } + conflicting = replace( + loaded, + target=replace(loaded.target, xdr_options={"postprocess": "auto"}), + ) + with pytest.raises(ValueError, match="conflicting FITS"): + dragon_pipeline._preflight_items((conflicting,)) + + def test_fits_mask_rejects_bitplanes_and_changed_content(inputs, monkeypatch): from cuphoton.core import fits_io From 204c50e10b6fb16d7297c3511999bca83fabf038 Mon Sep 17 00:00:00 2001 From: Trent Nelson Date: Tue, 29 Sep 2026 19:38:57 -0700 Subject: [PATCH 3/3] Preserve omitted FITS options and validate scatter strides Keep xFit configuration identities stable when no xDR overrides are supplied. Reject unsupported output strides with runtime checks and exercise the valid int32 quantization layout. Clarify the separate postprocessing path and its input-preserving copy cost. Signed-off-by: Trent Nelson --- docs/components/xdr.md | 6 ++++- src/cuphoton/core/cli/fits.py | 5 ++-- src/cuphoton/xdr/__init__.py | 5 ++-- src/cuphoton/xdr/kernels.py | 6 +++-- src/cuphoton/xdr/reader.py | 8 ++++++- src/cuphoton/xfit/executor.py | 6 +++-- tests/xdr/test_postprocess_gpu.py | 40 ++++++++++++++++++++++++------- tests/xfit/test_fits_input.py | 34 +++++++++++++++++++++++--- 8 files changed, 89 insertions(+), 21 deletions(-) diff --git a/docs/components/xdr.md b/docs/components/xdr.md index b3e8db04..0cadf871 100644 --- a/docs/components/xdr.md +++ b/docs/components/xdr.md @@ -17,8 +17,12 @@ Current scope: The default `postprocess="auto"` uses fused unshuffle, byte-order conversion, and scatter for supported GZIP FITS tiles. `"fused"` selects that path explicitly; -`"separate"` restores those steps with individual kernels for comparison. +`"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. ```python from cuphoton.xdr import batch_to_device diff --git a/src/cuphoton/core/cli/fits.py b/src/cuphoton/core/cli/fits.py index 3c40481b..c0c919fb 100644 --- a/src/cuphoton/core/cli/fits.py +++ b/src/cuphoton/core/cli/fits.py @@ -2,7 +2,7 @@ # # SPDX-License-Identifier: Apache-2.0 -"""Optional FITS reader overrides for manifest-driven commands.""" +"""Optional FITS reader overrides shared by FITS-consuming commands.""" from cuphoton.core.fits_options import ( XDR_OPTION_CHOICES, @@ -21,7 +21,8 @@ class XdrPostprocessArg(SetInvariant): _arg = "--xdr-postprocess" _help = ( "xDR postprocessing: auto, fused, or separate. Omitted preserves " - "manifest choices; applies when the FITS reader uses xDR." + "input choices or uses xDR defaults; applies when the FITS " + "reader uses xDR." ) _set = XDR_OPTION_CHOICES["postprocess"] _default = None diff --git a/src/cuphoton/xdr/__init__.py b/src/cuphoton/xdr/__init__.py index a071f674..999ae2ec 100644 --- a/src/cuphoton/xdr/__init__.py +++ b/src/cuphoton/xdr/__init__.py @@ -5,8 +5,9 @@ """GPU-native FITS loading. Pipeline: NVMe -> GPU memory via kvikio (GPUDirect Storage) -> nvCOMP batched -decompression on device -> CuPy kernels for GZIP_2 unshuffle, dequantize, and -tile scatter -> cupy.ndarray. Headers remain parsed on the CPU. +decompression on device -> pixel restoration and tile scatter (fused by +default) -> dequantization when needed -> cupy.ndarray. Headers are parsed +on the CPU. Scope: GZIP_1 / GZIP_2 compressed images and uncompressed ImageHDU / PrimaryHDU pixel data. RICE_1 and HCOMPRESS_1 are explicitly out of scope diff --git a/src/cuphoton/xdr/kernels.py b/src/cuphoton/xdr/kernels.py index 53dbf611..549c424a 100644 --- a/src/cuphoton/xdr/kernels.py +++ b/src/cuphoton/xdr/kernels.py @@ -120,7 +120,8 @@ def scatter_tiles_2d( """ pix_stride_out = d_out.strides[0] # row stride in bytes - assert d_out.strides[1] == itemsize, "output pixels must be contiguous" + if d_out.strides[1] != itemsize: + raise ValueError("output pixels must be contiguous") n_tiles = tile_byte_offsets.size kern = _scatter_tiles_2d_kernel(itemsize, shuffled) @@ -224,7 +225,8 @@ def scatter_native_tiles_2d( """ pix_stride_out = d_out.strides[0] # row stride in bytes - assert d_out.strides[1] == itemsize, "d_out must be C-contiguous" + if d_out.strides[1] != itemsize: + raise ValueError("output pixels must be contiguous") n_tiles = tile_byte_offsets.size kern = _scatter_native_tiles_2d_kernel(itemsize) diff --git a/src/cuphoton/xdr/reader.py b/src/cuphoton/xdr/reader.py index a29fda18..a591ba95 100644 --- a/src/cuphoton/xdr/reader.py +++ b/src/cuphoton/xdr/reader.py @@ -471,7 +471,13 @@ def postprocess_decoded_tiles( keepalive=None, postprocess: str = "auto", ): - """Restore FITS pixels; auto uses the fused kernel.""" + """Restore FITS pixels without modifying the decoded input buffer. + + ``auto`` uses the fused kernel. ``separate`` uses individual kernels + and an extra decoded-buffer allocation; GZIP_1 also requires a copy. + This comparison mode preserves the input, unlike the former in-place + GZIP_1 path, so its timings are not an exact historical baseline. + """ normalize_xdr_options({"postprocess": postprocess}) import cupy as cp diff --git a/src/cuphoton/xfit/executor.py b/src/cuphoton/xfit/executor.py index 05f74c39..7c7c7bdd 100644 --- a/src/cuphoton/xfit/executor.py +++ b/src/cuphoton/xfit/executor.py @@ -105,7 +105,9 @@ def _plan_xfit_chunks( raise ValueError("unsupported xFit compute dtype") if settings["fits_reader"] not in {"auto", "astropy", "xdr"}: raise ValueError("unsupported xFit FITS reader") - settings["xdr_options"] = normalize_xdr_options(settings["xdr_options"]) + xdr_options = normalize_xdr_options(settings.pop("xdr_options")) + if xdr_options: + settings["xdr_options"] = xdr_options fits_plan = None if Path(input_path).suffix.lower() == ".json": from .fits_input import plan_fits_input @@ -185,7 +187,7 @@ def _load_input(options: Mapping[str, Any], *, device=False) -> XFitDataset: reader_options = { "reader": settings["fits_reader"] if device else "astropy", "device": device, - "xdr_options": settings["xdr_options"], + "xdr_options": settings.get("xdr_options"), } dataset = load_xfit_dataset( options["input_path"], diff --git a/tests/xdr/test_postprocess_gpu.py b/tests/xdr/test_postprocess_gpu.py index 3768e4b2..f1a68087 100644 --- a/tests/xdr/test_postprocess_gpu.py +++ b/tests/xdr/test_postprocess_gpu.py @@ -9,7 +9,7 @@ import numpy as np import pytest -from cuphoton.xdr.kernels import scatter_tiles_2d +from cuphoton.xdr.kernels import scatter_native_tiles_2d, scatter_tiles_2d from cuphoton.xdr.reader import GpuCompImageReader @@ -59,19 +59,19 @@ def test_scatter_restores_fits_pixels( @pytest.mark.parametrize("postprocess", ["auto", "separate"]) -@pytest.mark.parametrize("dtype", ["f4", "f8"]) @pytest.mark.parametrize("compression", ["GZIP_1", "GZIP_2"]) def test_scatter_preserves_nondithered_dequantization( - cp, dtype, compression, postprocess + cp, compression, postprocess ): - itemsize = np.dtype(dtype).itemsize - image = np.arange(-150, 173, dtype=f"i{itemsize}").reshape(17, 19) + # FITS quantization stores int32 tiles, including for float64 images. + # This checks the supported float32 layout, not synthetic int64 tiles. + image = np.arange(-150, 173, dtype="i4").reshape(17, 19) section = np.s_[3:15, 5:18] raw, offsets, plan = _decoded_tiles(image, compression, section) plan.update( quantized=True, - dtype_char=dtype, - zbitpix=-8 * itemsize, + dtype_char="f4", + zbitpix=-32, sel_zscale=np.full(len(offsets), 0.5), sel_zzero=np.full(len(offsets), 2.0), ) @@ -79,10 +79,34 @@ def test_scatter_preserves_nondithered_dequantization( cp.asarray(raw), offsets, plan, postprocess=postprocess ) np.testing.assert_array_equal( - cp.asnumpy(result), image[section].astype(dtype) * 0.5 + 2 + cp.asnumpy(result), image[section].astype("f4") * 0.5 + 2 ) +@pytest.mark.parametrize( + "scatter", [scatter_tiles_2d, scatter_native_tiles_2d] +) +def test_scatter_rejects_noncontiguous_output_pixels(scatter): + # Validation must not depend on GPU imports or Python assertions. + arguments = dict( + d_tiles_concat=None, + d_out=np.empty((2, 8), dtype="i4")[:, ::2], + tile_byte_offsets=None, + tile_origins_row=None, + tile_origins_col=None, + tile_heights=None, + tile_widths=None, + tile_src_off_row=None, + tile_src_off_col=None, + tile_full_widths=None, + itemsize=4, + ) + if scatter is scatter_tiles_2d: + arguments.update(tile_byte_lengths=None, shuffled=False) + with pytest.raises(ValueError, match="output pixels must be contiguous"): + scatter(**arguments) + + def test_scatter_respects_output_row_stride(cp): image = np.arange(17 * 19, dtype="i4").reshape(17, 19) raw, offsets, plan = _decoded_tiles(image, "GZIP_2", np.s_[:, :]) diff --git a/tests/xfit/test_fits_input.py b/tests/xfit/test_fits_input.py index b977d35f..21d48392 100644 --- a/tests/xfit/test_fits_input.py +++ b/tests/xfit/test_fits_input.py @@ -123,8 +123,11 @@ def test_fits_planning_does_not_decode_or_initialize_a_gpu( assert len(spec.options_payload["input_sources"]) == 1 +@pytest.mark.parametrize( + "xdr_options", [None, {}, {"postprocess": "separate"}] +) def test_worker_reads_fits_only_after_binding_and_requests_device( - tmp_path, monkeypatch + tmp_path, monkeypatch, xdr_options ): path, _, _, _, _ = _input(tmp_path) @@ -150,17 +153,42 @@ def load(path, **kwargs): chunk_size=1, fit_options={ "fits_reader": "xdr", - "xdr_options": {"postprocess": "separate"}, + "xdr_options": xdr_options, }, ) worker = executor.XFitWorker(options) assert calls[0]["device"] is True assert calls[0]["reader"] == "xdr" - assert calls[0]["xdr_options"] == {"postprocess": "separate"} + assert calls[0]["xdr_options"] == (xdr_options or None) assert len(items) == 2 worker.close() +def test_empty_xdr_options_preserve_configuration_and_work_items(tmp_path): + path, _, _, _, _ = _input(tmp_path) + items, options = executor.plan_xfit_chunks( + path, chunk_size=1, fit_options={} + ) + assert "xdr_options" not in options["fit_options"] + for empty in (None, {}): + same_items, same_options = executor.plan_xfit_chunks( + path, chunk_size=1, fit_options={"xdr_options": empty} + ) + assert same_options == options + assert same_items == items + + changed_items, changed_options = executor.plan_xfit_chunks( + path, + chunk_size=1, + fit_options={"xdr_options": {"postprocess": "separate"}}, + ) + assert ( + changed_options["configuration_sha256"] + != options["configuration_sha256"] + ) + assert changed_items != items + + @pytest.mark.parametrize( "override,expected", [(None, "separate"), ({"postprocess": "fused"}, "fused")],