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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 29 additions & 0 deletions docs/components/xdr.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 7 additions & 0 deletions docs/components/xfit.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions docs/components/xpois.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
6 changes: 6 additions & 0 deletions docs/components/xrep.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
12 changes: 12 additions & 0 deletions docs/components/xscan.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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 -- \
Expand Down
45 changes: 43 additions & 2 deletions src/cuphoton/core/cli/fits.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
33 changes: 29 additions & 4 deletions src/cuphoton/core/fits_io.py
Original file line number Diff line number Diff line change
Expand Up @@ -13,14 +13,16 @@

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

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"}


Expand All @@ -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."""
Expand All @@ -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 {}
),
}


Expand Down Expand Up @@ -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(
Expand All @@ -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,
)


Expand All @@ -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.

Expand All @@ -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")
Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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
)
41 changes: 41 additions & 0 deletions src/cuphoton/core/fits_options.py
Original file line number Diff line number Diff line change
@@ -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)}
5 changes: 3 additions & 2 deletions src/cuphoton/xdr/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
11 changes: 10 additions & 1 deletion src/cuphoton/xdr/benchmark_fits.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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, ...]:
Expand Down Expand Up @@ -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"
Expand All @@ -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:
Expand Down Expand Up @@ -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,
Expand All @@ -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,
Expand Down Expand Up @@ -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,
)
Expand All @@ -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,
)
Expand Down Expand Up @@ -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,
Expand Down
Loading
Loading