diff --git a/docs/components/xdr.md b/docs/components/xdr.md index 3d172e4e..0cadf871 100644 --- a/docs/components/xdr.md +++ b/docs/components/xdr.md @@ -12,6 +12,35 @@ Current scope: - pipelined loading with `batch_to_device_stream` - explicit `NotImplementedError` for unsupported compression formats + +## Runtime choices + +The default `postprocess="auto"` uses fused unshuffle, byte-order conversion, +and scatter for supported GZIP FITS tiles. `"fused"` selects that path explicitly; +`"separate"` runs those steps with individual kernels for comparison. +Both preserve FITS values, including integer masks and floating-point bits. +Both leave the decoded input buffer unchanged. The separate path allocates +another buffer the size of the decoded tiles and copies GZIP_1 input before +byte-order conversion. Its timings therefore include that copy and are not +an exact baseline for the former in-place GZIP_1 implementation. + +```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..c0c919fb 100644 --- a/src/cuphoton/core/cli/fits.py +++ b/src/cuphoton/core/cli/fits.py @@ -2,12 +2,53 @@ # # 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, + 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 " + "input choices or uses xDR defaults; 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/__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/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 ee5ef438..549c424a 100644 --- a/src/cuphoton/xdr/kernels.py +++ b/src/cuphoton/xdr/kernels.py @@ -91,6 +91,116 @@ def byteswap_inplace(d_buf, itemsize: int) -> None: def scatter_tiles_2d( + d_tiles_concat, + d_out, + tile_byte_offsets, + tile_byte_lengths, + tile_origins_row, + tile_origins_col, + tile_heights, + tile_widths, + tile_src_off_row, + tile_src_off_col, + tile_full_widths, + itemsize: int, + *, + shuffled: bool, +) -> None: + """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)`` + 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. ``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 + 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) + threads_per_block = 256 + kern( + (n_tiles,), + (threads_per_block,), + ( + d_tiles_concat, + d_out, + tile_byte_offsets, + tile_byte_lengths, + 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_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 {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, + 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 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]; + 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); + 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) {{ + 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, name) + + +def scatter_native_tiles_2d( d_tiles_concat, d_out, tile_byte_offsets, @@ -115,10 +225,11 @@ def scatter_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_tiles_2d_kernel(itemsize) + kern = _scatter_native_tiles_2d_kernel(itemsize) threads_per_block = 256 kern( (n_tiles,), @@ -140,7 +251,7 @@ def scatter_tiles_2d( @functools.cache -def _scatter_tiles_2d_kernel(itemsize: int): +def _scatter_native_tiles_2d_kernel(itemsize: int): import cupy as cp src = rf""" 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 5e6c89e6..a591ba95 100644 --- a/src/cuphoton/xdr/reader.py +++ b/src/cuphoton/xdr/reader.py @@ -15,17 +15,19 @@ 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, - _native_device_empty_uint8, gpu_gzip_decompress_batch, ) @@ -62,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) @@ -168,8 +175,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"}) @@ -462,8 +469,16 @@ 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 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 use_native_pool = keepalive is not None @@ -483,31 +498,26 @@ 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) - ) + 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_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"]) + 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 - # 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"]) @@ -527,7 +537,7 @@ def postprocess_decoded_tiles( d_tile_full_w, ] ) - scatter_tiles_2d( + scatter_args = ( d_pixels, d_scatter_target, d_tile_offsets, @@ -540,8 +550,17 @@ def postprocess_decoded_tiles( d_tile_full_w, scatter_itemsize, ) + 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, + ) - # 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"]) @@ -569,6 +588,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. @@ -577,6 +597,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 = ( @@ -598,9 +619,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 @@ -619,6 +649,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) @@ -687,7 +718,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) @@ -729,6 +765,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..7c7c7bdd 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,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") + 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 @@ -182,6 +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.get("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 new file mode 100644 index 00000000..f1a68087 --- /dev/null +++ b/tests/xdr/test_postprocess_gpu.py @@ -0,0 +1,198 @@ +# 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_native_tiles_2d, 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("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, 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) + 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, + postprocess=postprocess, + ) + 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("postprocess", ["auto", "separate"]) +@pytest.mark.parametrize("compression", ["GZIP_1", "GZIP_2"]) +def test_scatter_preserves_nondithered_dequantization( + cp, compression, postprocess +): + # 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="f4", + zbitpix=-32, + 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, postprocess=postprocess + ) + np.testing.assert_array_equal( + 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_[:, :]) + 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 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..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) @@ -146,15 +149,64 @@ 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": xdr_options, + }, ) worker = executor.XFitWorker(options) assert calls[0]["device"] is True assert calls[0]["reader"] == "xdr" + 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")], +) +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