Files
Jungfraujoch/image_analysis/image_preprocessing/ImagePreprocessorGPU.cu
T
leonarski_fandClaude Opus 5 bec7e2e922 image_preprocessing: fuse the bitshuffle inverse with preprocessing, and verify the decode
The device decoder was byte-exact on every valid input - 994 production-compressed
images, 927 hand-built LZ4 blocks covering engineered (offset, matchlen) pairs across
the overlap branch boundary, 18000 repeat decodes, sanitizer-clean - and an audit
against LZ4_decompress_generic could not construct a valid block it mis-decodes. What
it did not do was notice when the input was NOT valid, and that mattered more than it
looks: the decode buffers are reused frame to frame, so a block that stopped early left
the PREVIOUS image in place, and in the bitshuffled layout the untouched tail is the
most significant byte-plane. A corrupt chunk therefore did not look like a missing
corner. It looked like thousands of real pixels several powers of two too bright, fed
to spot finding with no diagnostic, where the host decoder had raised an error.

So the kernel now flags a block that fails to reach its declared length while consuming
exactly its payload, and the host turns that into an exception once the caller has
synchronised. Reads are clamped against the end of the payload as well as the output,
both length chains are bounded exactly as read_variable_length bounds them, the two
offset bytes are bounded, and LZ4's parsing restrictions are enforced. On the host side
a block size that is not a multiple of 8 elements is rejected (it made the un-transpose
read uninitialised shared memory), the block count is bounded by what the chunk could
hold before it becomes an allocation (twelve header bytes could demand hundreds of MB
of pinned memory, permanently, per worker), trailing bytes are rejected, and the stream
is synchronised before any throw that happens after work is queued. An image of fewer
than 8 elements is all verbatim tail and now decodes rather than throwing. When the
device route fails for any reason the host decoder gets its turn, so it costs speed
rather than the acquisition.

The lanes cooperate on the copies and a later match can read bytes another lane wrote,
which since Volta needs an explicit __syncwarp(); it worked only because ptxas happened
to reconverge at the post-dominator. The prototype's offset == 1 and power-of-two fast
paths are also restored - the shipped kernel ran a runtime modulo, an emulated 32-bit
division per output byte, on the path its own comment calls the common case.

The un-transpose is now fused with preprocessing. One thread owns one group of 8
elements across every byte-plane, so once it has transposed its 8 bytes out of each
plane it holds 8 complete elements and emits 8 finished int32 pixels with the mask, the
error marker, the saturation cap and the statistics applied. The decompressed image is
never materialised: 0.623 -> 0.411 ms/frame at 18 Mpx, 0.523 -> 0.340 with 8 concurrent
workers. Staging nothing in shared memory also drops the 48 kB ceiling, which had made
any file whose bitshuffle blocks exceed it a hard failure; 64 kB blocks now decode.
gpu_compressed is sized from the chunk with grow-on-demand instead of from the
uncompressed size - it was reserving ~73 MB per worker to hold ~4 MB. Measured on a
1630x1553 uint32 rotation set at -N 32, peak GPU memory falls 3756 -> 3084 MiB; the
same model gives ~144 MB per worker on an 18 Mpx frame.

Decoding on the device also stopped reporting a decompression time, which blanked the
broker's compression plot trace and filled /entry/profiling/compressionTime with NaN.
The decoder brackets the decode with CUDA events and reports it again.

Tests: a differential fuzz suite against the CPU decoder - incompressible and highly
compressible data, engineered offsets, a size sweep hitting every rem%8 value twice,
all six element sizes, an 18 Mpx frame, decoder reuse, concurrency, hand-built LZ4
blocks across the overlap boundary, 26 foreign bitshuffle block sizes from 128 B to
64 kB, corrupt payloads and malformed containers, with a coverage report that proves
which LZ4 paths were reached rather than assuming it. Plus the fused path held byte for
byte against ImagePreprocessorCPU, statistics included, and against the host-upload
path on the same frame.

Battery: 37 crystals, every merged number identical to the host-decode run.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-03 14:13:55 +02:00

400 lines
17 KiB
Plaintext

// SPDX-FileCopyrightText: 2026 Filip Leonarski, Paul Scherrer Institute <filip.leonarski@psi.ch>
// SPDX-License-Identifier: GPL-3.0-only
#include <type_traits>
#include "ImagePreprocessorGPU.h"
template<class T>
__global__ void preprocess_kernel(
const T *__restrict__ input,
const uint8_t *__restrict__ mask,
int32_t *__restrict__ output,
ImageStatistics *__restrict__ stats,
T saturation_limit,
T err_value,
int npixels) {
// Shared block accumulators
__shared__ unsigned long long s_masked;
__shared__ unsigned long long s_saturated;
__shared__ unsigned long long s_error;
__shared__ long long s_max;
__shared__ long long s_min;
if (threadIdx.x == 0) {
s_masked = 0;
s_saturated = 0;
s_error = 0;
s_max = INT64_MIN;
s_min = INT64_MAX;
}
__syncthreads();
// Thread-local accumulators
unsigned long long local_masked = 0;
unsigned long long local_saturated = 0;
unsigned long long local_error = 0;
long long local_max = INT64_MIN;
long long local_min = INT64_MAX;
for (int i = blockIdx.x * blockDim.x + threadIdx.x;
i < npixels;
i += blockDim.x * gridDim.x) {
T v = input[i];
bool is_masked = mask[i];
// Error/invalid marker = the pixel type's extreme value (0xFFFFFFFF for EIGER uint32); tested
// before saturation, since for unsigned types the marker also exceeds saturation_limit (which is
// clipped to the HDF5 saturation_value). Priority: masked > error > saturated.
bool is_err = (v == err_value);
bool is_sat = !is_err && (v >= saturation_limit);
bool valid = !(is_masked || is_sat || is_err);
// Output
output[i] =
is_masked ? INT32_MIN : is_err ? INT32_MIN : is_sat ? INT32_MAX : (int32_t) v;
// Counters
local_masked += is_masked;
local_error += (!is_masked && is_err);
local_saturated += (!is_masked && !is_err && is_sat);
// Min/max only for valid
if (valid) {
int64_t val = (int64_t) v;
if (val > local_max) local_max = val;
if (val < local_min) local_min = val;
}
}
// Reduce to shared memory
atomicAdd(&s_masked, local_masked);
atomicAdd(&s_saturated, local_saturated);
atomicAdd(&s_error, local_error);
if (local_min <= local_max) {
atomicMax((long long *) &s_max, (long long) local_max);
atomicMin((long long *) &s_min, (long long) local_min);
}
__syncthreads();
// One thread writes block result
if (threadIdx.x == 0) {
atomicAdd(&stats->masked_pixel_count, s_masked);
atomicAdd(&stats->saturated_pixel_count, s_saturated);
atomicAdd(&stats->error_pixel_count, s_error);
atomicMax((long long *) &stats->max_value, (long long) s_max);
atomicMin((long long *) &stats->min_value, (long long) s_min);
}
}
// The per-pixel decision preprocess_kernel makes, in a form the fused kernel can reuse so the two
// cannot drift apart. Priority: masked > error > saturated.
template<class T>
struct PreprocessAccum {
unsigned long long masked = 0, saturated = 0, error = 0;
long long max_v = INT64_MIN, min_v = INT64_MAX;
__device__ __forceinline__ int32_t Apply(T v, bool is_masked, T sat_value, T err_value) {
const bool is_err = (v == err_value);
const bool is_sat = !is_err && (v >= sat_value);
masked += is_masked;
error += (!is_masked && is_err);
saturated += (!is_masked && !is_err && is_sat);
if (!(is_masked || is_sat || is_err)) {
const int64_t val = (int64_t) v;
if (val > max_v) max_v = val;
if (val < min_v) min_v = val;
}
return is_masked ? INT32_MIN : is_err ? INT32_MIN : is_sat ? INT32_MAX : (int32_t) v;
}
};
// Reduce a block's thread-local accumulators into the image-wide statistics. Every accumulator is an
// integer, so the result does not depend on the order the blocks arrive in.
template<class T>
__device__ __forceinline__ void FlushStats(PreprocessAccum<T> &l, ImageStatistics *stats) {
__shared__ unsigned long long s_masked, s_saturated, s_error;
__shared__ long long s_max, s_min;
if (threadIdx.x == 0) {
s_masked = 0; s_saturated = 0; s_error = 0; s_max = INT64_MIN; s_min = INT64_MAX;
}
__syncthreads();
atomicAdd(&s_masked, l.masked);
atomicAdd(&s_saturated, l.saturated);
atomicAdd(&s_error, l.error);
if (l.min_v <= l.max_v) {
atomicMax(&s_max, l.max_v);
atomicMin(&s_min, l.min_v);
}
__syncthreads();
if (threadIdx.x == 0) {
atomicAdd(&stats->masked_pixel_count, s_masked);
atomicAdd(&stats->saturated_pixel_count, s_saturated);
atomicAdd(&stats->error_pixel_count, s_error);
atomicMax((long long *) &stats->max_value, s_max);
atomicMin((long long *) &stats->min_value, s_min);
}
}
__device__ __forceinline__ uint64_t transpose8_fused(uint64_t x) {
uint64_t t;
t = (x ^ (x >> 7)) & 0x00aa00aa00aa00aaULL; x = x ^ t ^ (t << 7);
t = (x ^ (x >> 14)) & 0x0000cccc0000ccccULL; x = x ^ t ^ (t << 14);
t = (x ^ (x >> 28)) & 0x00000000f0f0f0f0ULL; x = x ^ t ^ (t << 28);
return x;
}
// The bitshuffle inverse and the preprocessing in ONE pass. One thread owns one group of 8 elements
// across every byte-plane, so once it has transposed its 8 bytes out of each plane it holds 8
// complete elements and can emit 8 finished int32 pixels - the decompressed image never has to exist
// in device memory at all. That removes a full-frame buffer per worker and a full-frame write plus
// read from the pipeline.
//
// The last CUDA block (blockIdx.x == nblocks) finishes the handful of elements bitshuffle stores
// verbatim; they are already on the device inside the uploaded chunk.
template<class T, int ES>
__global__ __launch_bounds__(256) void untranspose_preprocess_kernel(
const uint8_t *__restrict__ shuffled,
const BSLZ4BlockDesc *__restrict__ desc,
const uint8_t *__restrict__ mask,
int32_t *__restrict__ out,
ImageStatistics *__restrict__ stats,
T sat_value, T err_value, int nblocks,
const uint8_t *__restrict__ tail_src, uint32_t tail_elems, uint32_t tail_elem0) {
PreprocessAccum<T> l;
// The bytes are assembled in the unsigned counterpart of T - shifting a byte into the top of a
// signed type overflows it - and converted back at the end, which C++20 defines as two's
// complement reinterpretation. That is exactly what the byte order in the file means.
using U = typename std::make_unsigned<T>::type;
if (blockIdx.x == nblocks) {
if (threadIdx.x < tail_elems) {
U uv = 0;
#pragma unroll
for (int p = 0; p < ES; p++)
uv |= (U)((U)tail_src[threadIdx.x * ES + p] << (8 * p));
const T v = (T) uv;
out[tail_elem0 + threadIdx.x] = l.Apply(v, mask[tail_elem0 + threadIdx.x] != 0, sat_value, err_value);
}
FlushStats<T>(l, stats);
return;
}
const int b = blockIdx.x;
const uint32_t size = desc[b].nelem; // bytes per plane
const uint8_t *in = shuffled + desc[b].out_off;
const uint32_t elem0 = desc[b].out_off / ES; // first pixel of this block
const uint32_t n = size / 8;
for (uint32_t i = threadIdx.x; i < n; i += blockDim.x) {
uint64_t x[ES];
#pragma unroll
for (int p = 0; p < ES; p++) {
const uint8_t *pin = in + p * size;
uint64_t a = 0;
#pragma unroll
for (int k = 0; k < 8; k++) a |= (uint64_t)pin[k * n + i] << (8 * k);
x[p] = transpose8_fused(a);
}
int32_t o[8];
#pragma unroll
for (int k = 0; k < 8; k++) {
U uv = 0;
#pragma unroll
for (int p = 0; p < ES; p++) uv |= (U)((U)((x[p] >> (8 * k)) & 0xff) << (8 * p));
o[k] = l.Apply((T) uv, mask[elem0 + i * 8 + k] != 0, sat_value, err_value);
}
// elem0 is a multiple of 8 (bitshuffle blocks are), so this is 32-byte aligned.
int4 *dst = reinterpret_cast<int4 *>(out + elem0 + i * 8);
dst[0] = make_int4(o[0], o[1], o[2], o[3]);
dst[1] = make_int4(o[4], o[5], o[6], o[7]);
}
FlushStats<T>(l, stats);
}
ImagePreprocessorGPU::ImagePreprocessorGPU(const DiffractionExperiment &experiment, const PixelMask &mask,
std::shared_ptr<CudaStream> stream, bool copy_image_to_host)
: ImagePreprocessor(experiment),
stream(stream),
copy_image_to_host(copy_image_to_host),
gpu_stats(1),
cpu_stats(1),
cpu_stats_reg(cpu_stats) {
// Setup mask. The same for every worker, so it is uploaded once per GPU and shared; keyed on the
// PixelMask's own vector, which the derived table is a pure function of.
std::vector<uint8_t> mask_vec(npixels);
for (int i = 0; i < npixels; i++)
mask_vec[i] = (mask.GetMask().at(i) != 0);
gpu_mask = SharedDeviceTable(mask.GetMask().data(), npixels, mask_vec.data(), *stream);
// Setup GPU settings. The current device, not device 0: workers are pinned round-robin across GPUs,
// so device 0's SM count can belong to a different card than the one these kernels launch on.
int device = 0;
cudaGetDevice(&device);
cudaDeviceProp prop{};
cudaGetDeviceProperties(&prop, device);
threads = 128;
blocks = 4 * prop.multiProcessorCount;
}
float ImagePreprocessorGPU::GetLastDecompressionTime_s() const {
return bslz4_decoder ? bslz4_decoder->GetDecodeTime_s() : 0.0f;
}
void ImagePreprocessorGPU::PinInputBuffer(std::vector<uint8_t> &buffer, size_t size) {
if (buffer.size() == size)
return;
// Unregister before the resize, which can move the buffer.
input_reg.unregister();
buffer.resize(size);
input_reg.rebind(buffer);
}
ImageStatistics ImagePreprocessorGPU::Analyze(ImagePreprocessorBuffer &processed_image, const uint8_t *image_ptr,
CompressedImageMode image_mode) {
switch (image_mode) {
case CompressedImageMode::Int8:
return Analyze<int8_t>(processed_image, image_ptr, INT8_MIN, INT8_MAX);
case CompressedImageMode::Int16:
return Analyze<int16_t>(processed_image, image_ptr, INT16_MIN, INT16_MAX);
case CompressedImageMode::Int32:
return Analyze<int32_t>(processed_image, image_ptr, INT32_MIN, INT32_MAX);
case CompressedImageMode::Uint8:
return Analyze<uint8_t>(processed_image, image_ptr, UINT8_MAX, UINT8_MAX);
case CompressedImageMode::Uint16:
return Analyze<uint16_t>(processed_image, image_ptr, UINT16_MAX, UINT16_MAX);
case CompressedImageMode::Uint32:
return Analyze<uint32_t>(processed_image, image_ptr, UINT32_MAX, UINT32_MAX);
default:
throw JFJochException(JFJochExceptionCategory::InputParameterInvalid, "RGB/float mode not supported");
}
}
bool ImagePreprocessorGPU::AnalyzeCompressed(ImagePreprocessorBuffer &processed_image,
const CompressedImage &image,
ImageStatistics &stats) {
if (!BSLZ4DecoderGPU::Supports(image))
return false; // caller decompresses on the host and uses Analyze()
if (image.GetUncompressedSize() != npixels * image.GetByteDepth())
return false;
if (!bslz4_decoder)
bslz4_decoder = std::make_unique<BSLZ4DecoderGPU>(npixels * sizeof(uint32_t), stream);
// LZ4 on the device, then ONE kernel that un-transposes the bitshuffle blocks and preprocesses
// them as it goes. The decompressed image is never materialised: the fused kernel reads the
// shuffled bytes and writes finished int32 pixels.
const BSLZ4ShuffledImage shuffled = bslz4_decoder->DecodeShuffled(image);
switch (image.GetMode()) {
case CompressedImageMode::Int8:
stats = UntransposeAndAnalyze<int8_t, 1>(processed_image, shuffled, INT8_MIN, INT8_MAX); return true;
case CompressedImageMode::Uint8:
stats = UntransposeAndAnalyze<uint8_t, 1>(processed_image, shuffled, UINT8_MAX, UINT8_MAX); return true;
case CompressedImageMode::Int16:
stats = UntransposeAndAnalyze<int16_t, 2>(processed_image, shuffled, INT16_MIN, INT16_MAX); return true;
case CompressedImageMode::Uint16:
stats = UntransposeAndAnalyze<uint16_t, 2>(processed_image, shuffled, UINT16_MAX, UINT16_MAX); return true;
case CompressedImageMode::Int32:
stats = UntransposeAndAnalyze<int32_t, 4>(processed_image, shuffled, INT32_MIN, INT32_MAX); return true;
case CompressedImageMode::Uint32:
stats = UntransposeAndAnalyze<uint32_t, 4>(processed_image, shuffled, UINT32_MAX, UINT32_MAX); return true;
default:
return false; // Supports() already excludes these; belt and braces
}
}
// The device-decode counterpart of AnalyzeOnDevice: same per-pixel decision, same statistics, but
// fed from the bitshuffled bytes rather than from a decompressed image.
template<class T, int ES>
ImageStatistics ImagePreprocessorGPU::UntransposeAndAnalyze(ImagePreprocessorBuffer &processed_image,
const BSLZ4ShuffledImage &shuffled,
T err_value, T sat_value) {
if (sat_value > saturation_limit)
sat_value = static_cast<T>(saturation_limit);
cpu_stats[0] = ImageStatistics{.max_value = INT64_MIN, .min_value = INT64_MAX};
cudaMemcpyAsync(gpu_stats, cpu_stats.data(), sizeof(ImageStatistics), cudaMemcpyHostToDevice, *stream);
// One CUDA block per bitshuffle block, plus one for the verbatim tail when there is one.
const int nb = shuffled.nblocks + (shuffled.tail_elems > 0 ? 1 : 0);
untranspose_preprocess_kernel<T, ES> <<< nb, 256, 0, *stream >>>(
shuffled.shuffled,
shuffled.desc,
gpu_mask->get(),
processed_image.getGPUBuffer(),
gpu_stats,
sat_value,
err_value,
shuffled.nblocks,
shuffled.tail_src,
shuffled.tail_elems,
shuffled.tail_elem0);
if (copy_image_to_host)
cudaMemcpyAsync(processed_image.data(), processed_image.getGPUBuffer(), npixels * sizeof(int32_t), cudaMemcpyDeviceToHost, *stream);
cudaMemcpyAsync(cpu_stats.data(), gpu_stats, sizeof(ImageStatistics), cudaMemcpyDeviceToHost, *stream);
cudaStreamSynchronize(*stream);
// Only now can the device tell us whether every block actually decoded.
bslz4_decoder->ThrowIfDecodeFailed();
return cpu_stats[0];
}
template<class T>
ImageStatistics ImagePreprocessorGPU::Analyze(ImagePreprocessorBuffer &processed_image,
const uint8_t *input,
T err_value,
T sat_value) {
// Allocated here rather than in the constructor: only the host-upload path needs it, and the
// device-decode path - which is what BSHUF_LZ4 images take - never touches it. Overshoot to
// 4 bytes per pixel so a 1- or 2-byte image fits the same buffer.
if (!gpu_decompressed_image.get())
gpu_decompressed_image = CudaDevicePtr<uint8_t>(npixels * sizeof(uint32_t));
// On this engine's own stream, not the NULL stream: a NULL-stream copy implicitly synchronises with
// every blocking stream in the process, which serialised all workers behind whichever one was
// uploading. The stream is synchronised at the end of this function, so the ordering is unchanged.
cudaMemcpyAsync(gpu_decompressed_image, input, npixels * sizeof(T), cudaMemcpyHostToDevice, *stream);
return AnalyzeOnDevice<T>(processed_image, err_value, sat_value);
}
// Everything after the image is on the device, shared by the host-upload and the device-decode
// entry points so the two cannot drift apart.
template<class T>
ImageStatistics ImagePreprocessorGPU::AnalyzeOnDevice(ImagePreprocessorBuffer &processed_image,
T err_value, T sat_value) {
if (sat_value > saturation_limit)
sat_value = static_cast<T>(saturation_limit);
cpu_stats[0] = ImageStatistics{.max_value = INT64_MIN, .min_value = INT64_MAX};
cudaMemcpyAsync(gpu_stats, cpu_stats.data(), sizeof(ImageStatistics), cudaMemcpyHostToDevice, *stream);
preprocess_kernel<T> <<< blocks, threads, 0, *stream >>>(
reinterpret_cast<const T *>(gpu_decompressed_image.get()),
gpu_mask->get(),
processed_image.getGPUBuffer(),
gpu_stats,
sat_value,
err_value,
npixels);
// The preprocessed image is 4 bytes per pixel - by far the largest transfer here - and every GPU
// engine reads it straight from the device buffer, so it only comes back when a CPU engine needs it.
if (copy_image_to_host)
cudaMemcpyAsync(processed_image.data(), processed_image.getGPUBuffer(), npixels * sizeof(int32_t), cudaMemcpyDeviceToHost, *stream);
cudaMemcpyAsync(cpu_stats.data(), gpu_stats, sizeof(ImageStatistics), cudaMemcpyDeviceToHost, *stream);
cudaStreamSynchronize(*stream);
return cpu_stats[0];
}