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>
This commit is contained in:
2026-08-03 14:13:55 +02:00
co-authored by Claude Opus 5
parent 47277674fa
commit bec7e2e922
12 changed files with 2023 additions and 92 deletions
@@ -17,9 +17,20 @@ namespace {
// broadcast read, so no divergence - and the literal and match copies are split across the 32
// lanes so the stores coalesce. One thread per block instead has each thread streaming its own
// 8 kB region, which coalesces not at all and measured 13x slower.
//
// Because the lanes cooperate on the copies, a match can source bytes that OTHER lanes wrote in
// an earlier sequence. Since Volta that needs an explicit __syncwarp() - implicit reconvergence
// is not part of the programming model - so there is one after every copy loop. The full mask is
// correct: the early return and every break test warp-uniform values, so lanes never diverge
// permanently.
//
// Bounds: every read is clamped against iend and every write against oend, so a malformed or
// corrupt payload cannot walk off either buffer. It can still stop early, which leaves the block
// short; lane 0 flags that at the end and the host turns it into an exception.
__global__ void lz4_decode_blocks(const uint8_t *__restrict__ src,
const BSLZ4BlockDesc *__restrict__ desc,
uint8_t *__restrict__ dst,
uint32_t *__restrict__ status,
int nblocks, uint32_t elem_size) {
const int lane = threadIdx.x & 31;
const int b = (blockIdx.x * blockDim.x + threadIdx.x) >> 5;
@@ -27,34 +38,75 @@ namespace {
const uint8_t *ip = src + desc[b].in_off;
const uint8_t *const iend = ip + desc[b].in_len;
uint8_t *op = dst + desc[b].out_off;
uint8_t *const oend = op + desc[b].nelem * elem_size;
uint8_t *const obase = dst + desc[b].out_off;
uint8_t *op = obase;
uint8_t *const oend = obase + desc[b].nelem * elem_size;
bool malformed = false;
while (ip < iend) {
const uint32_t token = *ip++;
uint32_t litlen = token >> 4;
if (litlen == 15) {
// read_variable_length(&ip, iend - RUN_MASK, initial_check=1) in the reference: the
// chain may not start within, nor run into, the last RUN_MASK (15) input bytes. A
// valid stream never does - the literals it counts have to follow it - so a chain
// that reaches there is corruption, and this is the only place it shows up.
if ((size_t)(iend - ip) <= 15) { malformed = true; break; }
uint32_t s;
do { s = *ip++; litlen += s; } while (s == 255 && ip < iend);
do {
s = *ip++;
litlen += s;
if ((size_t)(iend - ip) < 15) { malformed = true; break; }
} while (s == 255);
if (malformed) break;
}
if (litlen) {
const uint32_t n = min(litlen, (uint32_t)(oend - op));
// Clamped by the INPUT as well as the output: a corrupt litlen must not read past the
// end of this block's payload or write past the end of the block. Clamping keeps the
// kernel in bounds; needing to clamp at all means the stream is not decodable, which
// is what the reference reports as an error, so record it.
if (litlen > (uint32_t)(oend - op) || litlen > (uint32_t)(iend - ip))
malformed = true;
const uint32_t n = min(min(litlen, (uint32_t)(oend - op)), (uint32_t)(iend - ip));
for (uint32_t i = lane; i < n; i += 32) op[i] = ip[i];
__syncwarp();
op += n; ip += litlen;
}
if (ip >= iend) break; // last sequence carries literals only
// LZ4's parsing restrictions: an encoder may not leave a match within MFLIMIT (12) bytes
// of the end of the block, nor fewer than 2+1+LASTLITERALS (8) input bytes after a
// literal run that is not the last one. So once either limit is reached this can ONLY be
// the final sequence, and the final sequence must consume the payload exactly. The
// reference applies this whether or not the run was empty, which is why the test sits
// outside the copy - a zero-length literal run near the end is just as illegal.
if ((size_t)(oend - op) < 12 || (size_t)(iend - ip) < 8) {
malformed = (ip != iend) || (op != oend);
break; // necessarily EOF
}
if (iend - ip < 2) break; // last sequence carries literals only
const uint32_t offset = (uint32_t)ip[0] | ((uint32_t)ip[1] << 8);
ip += 2;
uint32_t matchlen = token & 0x0F;
if (matchlen == 15) {
// read_variable_length(&ip, iend - LASTLITERALS + 1, initial_check=0): bounded by the
// last 4 input bytes rather than 15, and with no check before the first read.
uint32_t s;
do { s = *ip++; matchlen += s; } while (s == 255 && ip < iend);
do {
s = *ip++;
matchlen += s;
if ((size_t)(iend - ip) < 4) { malformed = true; break; }
} while (s == 255);
if (malformed) break;
}
matchlen += 4; // minmatch
if (offset == 0 || offset > (uint32_t)(op - (dst + desc[b].out_off))) break; // malformed
if (offset == 0 || offset > (uint32_t)(op - obase)) { malformed = true; break; }
const uint8_t *mp = op - offset;
// A match may reach the end of the block but never past it - the reference treats an
// overrun as an error rather than truncating, and so must this.
if (matchlen > (uint32_t)(oend - op))
malformed = true;
const uint32_t n = min(matchlen, (uint32_t)(oend - op));
if (offset >= matchlen) {
for (uint32_t i = lane; i < n; i += 32) op[i] = mp[i];
@@ -62,11 +114,28 @@ namespace {
// An overlapping match is a pattern of period `offset`. mp[0..offset-1] all lie
// before op and are already final, so each output byte can be sourced from them
// independently - which keeps this parallel rather than a serial byte loop. Long
// zero runs in sparse detector data arrive here with offset == 1.
for (uint32_t i = lane; i < n; i += 32) op[i] = mp[i % offset];
// zero runs in sparse detector data arrive here with offset == 1, and a runtime
// modulo is an emulated division, so the two cheap cases are peeled off first.
if (offset == 1) {
const uint8_t v = mp[0];
for (uint32_t i = lane; i < n; i += 32) op[i] = v;
} else if ((offset & (offset - 1)) == 0) {
const uint32_t m = offset - 1;
for (uint32_t i = lane; i < n; i += 32) op[i] = mp[i & m];
} else {
for (uint32_t i = lane; i < n; i += 32) op[i] = mp[i % offset];
}
}
__syncwarp();
op += n;
}
// A block must decode to exactly its declared length AND consume exactly its payload. Both
// are conditions LZ4_decompress_safe reports to the host path, and both are needed: a corrupt
// stream can land on the right output length while leaving input over, or run its input out
// early. Either way the bytes are not the ones that were compressed.
if (lane == 0 && (malformed || op != oend || ip != iend))
atomicExch(status, 1u);
}
__device__ __forceinline__ uint64_t transpose8(uint64_t x) {
@@ -77,14 +146,15 @@ namespace {
return x;
}
// The bitshuffle inverse, one CUDA block per bitshuffle block, mirroring bitshuf_decode_block:
// un-transpose the bits of each byte-plane, then interleave the planes back into elements. The
// planes are staged in shared memory so the final store is coalesced.
extern __shared__ uint8_t smem[];
// The bitshuffle inverse, mirroring bitshuf_decode_block. One thread owns one group of 8 elements
// across EVERY byte-plane, so after transposing its 8 bytes out of each plane it holds all bytes
// of 8 complete elements and can write them straight out. That needs no staging buffer, which is
// what keeps the kernel free of the 48 kB dynamic-shared-memory ceiling a block size taken from
// the file header would otherwise run into.
template<int ES>
__global__ void bitshuffle_untranspose(const uint8_t *__restrict__ shuffled,
const BSLZ4BlockDesc *__restrict__ desc,
uint8_t *__restrict__ out,
int nblocks, uint32_t elem_size) {
uint8_t *__restrict__ out, int nblocks) {
const int b = blockIdx.x;
if (b >= nblocks) return;
@@ -93,23 +163,22 @@ namespace {
uint8_t *dst = out + desc[b].out_off;
const uint32_t n = size / 8;
for (uint32_t p = 0; p < elem_size; p++) {
const uint8_t *pin = in + p * size;
uint8_t *pout = smem + p * size;
for (uint32_t i = threadIdx.x; i < n; i += blockDim.x) {
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);
const uint64_t x = transpose8(a);
#pragma unroll
for (int k = 0; k < 8; k++) pout[i * 8 + k] = (uint8_t)(x >> (8 * k));
x[p] = transpose8(a);
}
#pragma unroll
for (int k = 0; k < 8; k++)
#pragma unroll
for (int p = 0; p < ES; p++)
dst[(i * 8 + k) * ES + p] = (uint8_t)(x[p] >> (8 * k));
}
__syncthreads();
for (uint32_t i = threadIdx.x; i < size; i += blockDim.x)
for (uint32_t j = 0; j < elem_size; j++)
dst[i * elem_size + j] = smem[j * size + i];
}
uint64_t be64(const uint8_t *p) { uint64_t v = 0; for (int i = 0; i < 8; i++) v = (v << 8) | p[i]; return v; }
@@ -129,6 +198,11 @@ namespace {
default: return 0; // float and RGB modes are never bitshuffled by this pipeline
}
}
// The smallest a block can be on the wire: a 4-byte length plus at least one payload byte. Used
// to reject a header whose declared block size implies more blocks than the chunk could hold,
// before that count is turned into an allocation.
constexpr size_t MIN_BLOCK_BYTES_ON_WIRE = 5;
}
bool BSLZ4DecoderGPU::Supports(const CompressedImage &image) {
@@ -139,21 +213,35 @@ bool BSLZ4DecoderGPU::Supports(const CompressedImage &image) {
BSLZ4DecoderGPU::BSLZ4DecoderGPU(size_t in_max_uncompressed_bytes, std::shared_ptr<CudaStream> in_stream)
: stream(std::move(in_stream)),
max_uncompressed_bytes(in_max_uncompressed_bytes) {
// The compressed chunk is smaller than the image in every case worth having, but a pathological
// frame can expand slightly, so allow headroom rather than risk a per-frame reallocation.
max_compressed_bytes = max_uncompressed_bytes + max_uncompressed_bytes / 16 + 4096;
// Enough descriptors for the block size the compressor actually uses (8 kB); Decode() grows
// these if a stream turns up with smaller blocks. Sizing for the format's theoretical minimum
// block instead would allocate megabytes of pinned memory that no real file needs.
max_blocks = max_uncompressed_bytes / 8192 + 2;
gpu_compressed = CudaDevicePtr<uint8_t>(max_compressed_bytes);
gpu_shuffled = CudaDevicePtr<uint8_t>(max_uncompressed_bytes);
gpu_desc = CudaDevicePtr<BSLZ4BlockDesc>(max_blocks);
host_desc = CudaHostPtr<BSLZ4BlockDesc>(max_blocks);
gpu_status = CudaDevicePtr<uint32_t>(1);
host_status = CudaHostPtr<uint32_t>(1);
// The compressed buffer and the descriptors are grown to fit the first image instead of being
// sized for a worst case that no real frame reaches. A chunk is a few MB against an image of
// tens; sizing this from the UNCOMPRESSED size cost ~73 MB per worker to hold ~4 MB.
}
void BSLZ4DecoderGPU::Decode(const CompressedImage &image, uint8_t *gpu_out) {
void BSLZ4DecoderGPU::EnsureCompressedCapacity(size_t bytes) {
if (bytes <= compressed_capacity)
return;
// A little slack, so a frame that compresses slightly worse than the last does not reallocate.
const size_t want = bytes + bytes / 4;
cuda_err(cudaStreamSynchronize(*stream)); // nothing may still be reading the old buffer
gpu_compressed = CudaDevicePtr<uint8_t>(want);
compressed_capacity = want;
}
void BSLZ4DecoderGPU::EnsureBlockCapacity(size_t nblocks) {
if (nblocks <= max_blocks)
return;
const size_t want = nblocks + nblocks / 4 + 16;
cuda_err(cudaStreamSynchronize(*stream)); // the previous descriptor upload must have landed
gpu_desc = CudaDevicePtr<BSLZ4BlockDesc>(want);
host_desc = CudaHostPtr<BSLZ4BlockDesc>(want);
max_blocks = want;
}
BSLZ4ShuffledImage BSLZ4DecoderGPU::DecodeShuffled(const CompressedImage &image) {
const uint8_t *src = image.GetCompressed();
const size_t clen = image.GetCompressedSize();
const size_t elem_size = elem_size_of(image.GetMode());
@@ -161,7 +249,7 @@ void BSLZ4DecoderGPU::Decode(const CompressedImage &image, uint8_t *gpu_out) {
if (clen < 12)
throw JFJochException(JFJochExceptionCategory::Compression, "bslz4 chunk shorter than its header");
if (total_bytes > max_uncompressed_bytes || clen > max_compressed_bytes)
if (total_bytes > max_uncompressed_bytes)
throw JFJochException(JFJochExceptionCategory::Compression, "bslz4 image larger than the decoder was sized for");
if (be64(src) != total_bytes)
throw JFJochException(JFJochExceptionCategory::Compression, "bslz4 header size does not match the image");
@@ -171,6 +259,12 @@ void BSLZ4DecoderGPU::Decode(const CompressedImage &image, uint8_t *gpu_out) {
throw JFJochException(JFJochExceptionCategory::Compression, "bslz4 block size invalid");
const size_t block_elems = block_bytes / elem_size;
// bitshuffle transposes 8 elements at a time and refuses a block that is not a multiple of 8;
// the host decoder rejects this too (JFJochDecompress.h). Without the check the un-transpose
// would silently drop the last size % 8 elements of every block.
if (block_elems % BSHUF_BLOCKED_MULT != 0)
throw JFJochException(JFJochExceptionCategory::Compression, "bslz4 block size is not a multiple of 8 elements");
const size_t nelements = total_bytes / elem_size;
const size_t nfull = nelements / block_elems;
const size_t rem = nelements - nfull * block_elems;
@@ -179,49 +273,106 @@ void BSLZ4DecoderGPU::Decode(const CompressedImage &image, uint8_t *gpu_out) {
// Walk the container to locate the blocks. Lengths are only knowable in order, so this scan is
// inherent to the format rather than an implementation choice.
// An image of fewer than 8 elements has no bitshuffle block at all - it is entirely the verbatim
// tail. The host decoder handles that, so handle it here rather than declining: nblocks is simply
// zero and only the tail is copied.
const size_t nblocks_needed = nfull + (last > 0 ? 1 : 0);
if (nblocks_needed > max_blocks) { // a stream with smaller blocks than the compressor emits
max_blocks = nblocks_needed;
gpu_desc = CudaDevicePtr<BSLZ4BlockDesc>(max_blocks);
host_desc = CudaHostPtr<BSLZ4BlockDesc>(max_blocks);
}
// Bound the block count by what the chunk could actually hold BEFORE it becomes an allocation:
// a header declaring a one-element block size would otherwise ask for hundreds of MB of pinned
// memory, and only then fail on the first block header.
if (nblocks_needed > (clen - 12) / MIN_BLOCK_BYTES_ON_WIRE)
throw JFJochException(JFJochExceptionCategory::Compression, "bslz4 chunk too short for the blocks its header implies");
EnsureCompressedCapacity(clen);
EnsureBlockCapacity(nblocks_needed);
size_t nblk = 0, off = 12, out_off = 0;
for (size_t i = 0; i < nfull + (last > 0 ? 1 : 0); i++) {
for (size_t i = 0; i < nblocks_needed; i++) {
if (off + 4 > clen)
throw JFJochException(JFJochExceptionCategory::Compression, "truncated bslz4 block header");
const uint32_t block_clen = be32(src + off);
off += 4;
if (block_clen == 0 || off + block_clen > clen)
throw JFJochException(JFJochExceptionCategory::Compression, "bslz4 block extends past the chunk");
if (nblk >= max_blocks)
throw JFJochException(JFJochExceptionCategory::Compression, "bslz4 chunk has more blocks than expected");
const uint32_t ne = (i < nfull) ? static_cast<uint32_t>(block_elems) : static_cast<uint32_t>(last);
host_desc.get()[nblk++] = {static_cast<uint32_t>(off), block_clen, static_cast<uint32_t>(out_off), ne};
off += block_clen;
out_off += static_cast<size_t>(ne) * elem_size;
}
if (nblk == 0)
throw JFJochException(JFJochExceptionCategory::Compression, "bslz4 chunk contains no blocks");
// The tail that bitshuffle leaves uncompressed and copies verbatim, and then nothing else: the
// host decoder requires the chunk to be consumed exactly, so require it here too rather than
// ignoring trailing bytes that indicate the container is not what it claims to be.
if (off + leftover_bytes > clen)
throw JFJochException(JFJochExceptionCategory::Compression, "truncated bslz4 leftover bytes");
if (off + leftover_bytes != clen)
throw JFJochException(JFJochExceptionCategory::Compression, "bslz4 chunk has trailing bytes after the last block");
cuda_err(cudaMemcpyAsync(gpu_compressed.get(), src, clen, cudaMemcpyHostToDevice, *stream));
cuda_err(cudaMemcpyAsync(gpu_desc.get(), host_desc.get(), nblk * sizeof(BSLZ4BlockDesc),
// Everything that can be checked on the host has been checked; from here work is queued.
cuda_err(cudaEventRecord(decode_start, *stream));
host_status.get()[0] = 0;
cuda_err(cudaMemcpyAsync(gpu_status.get(), host_status.get(), sizeof(uint32_t),
cudaMemcpyHostToDevice, *stream));
cuda_err(cudaMemcpyAsync(gpu_compressed.get(), src, clen, cudaMemcpyHostToDevice, *stream));
const int nb = static_cast<int>(nblk);
lz4_decode_blocks<<<(nb * 32 + 255) / 256, 256, 0, *stream>>>(
gpu_compressed.get(), gpu_desc.get(),
gpu_shuffled.get(), nb, static_cast<uint32_t>(elem_size));
bitshuffle_untranspose<<<nb, 256, block_bytes, *stream>>>(
gpu_shuffled.get(), gpu_desc.get(),
gpu_out, nb, static_cast<uint32_t>(elem_size));
cuda_err(cudaGetLastError());
// The tail that bitshuffle leaves uncompressed and copies verbatim.
if (leftover_bytes > 0) {
if (off + leftover_bytes > clen)
throw JFJochException(JFJochExceptionCategory::Compression, "truncated bslz4 leftover bytes");
cuda_err(cudaMemcpyAsync(gpu_out + out_off, src + off, leftover_bytes,
if (nb > 0) {
cuda_err(cudaMemcpyAsync(gpu_desc.get(), host_desc.get(), nblk * sizeof(BSLZ4BlockDesc),
cudaMemcpyHostToDevice, *stream));
lz4_decode_blocks<<<(nb * 32 + 255) / 256, 256, 0, *stream>>>(
gpu_compressed.get(), gpu_desc.get(), gpu_shuffled.get(), gpu_status.get(),
nb, static_cast<uint32_t>(elem_size));
cuda_err(cudaGetLastError());
}
cuda_err(cudaMemcpyAsync(host_status.get(), gpu_status.get(), sizeof(uint32_t),
cudaMemcpyDeviceToHost, *stream));
// Stop the clock here rather than after the un-transpose: getting the chunk onto the device and
// LZ4-decoding it is the part that replaced the host decompression, and it is the same work on
// both the raw-bytes path and the fused one, where the un-transpose is inseparable from
// preprocessing and is reported with it.
cuda_err(cudaEventRecord(decode_stop, *stream));
decode_timed = true;
BSLZ4ShuffledImage ret;
ret.shuffled = gpu_shuffled.get();
ret.desc = gpu_desc.get();
ret.nblocks = nb;
ret.elem_size = static_cast<uint32_t>(elem_size);
ret.tail_elems = static_cast<uint32_t>(leftover_bytes / elem_size);
ret.tail_elem0 = static_cast<uint32_t>(out_off / elem_size);
ret.tail_src = leftover_bytes > 0 ? gpu_compressed.get() + off : nullptr;
return ret;
}
void BSLZ4DecoderGPU::Decode(const CompressedImage &image, uint8_t *gpu_out) {
const BSLZ4ShuffledImage s = DecodeShuffled(image);
if (s.nblocks > 0) { // an image of fewer than 8 elements is all tail and has no block
switch (s.elem_size) {
case 1: bitshuffle_untranspose<1><<<s.nblocks, 256, 0, *stream>>>(s.shuffled, s.desc, gpu_out, s.nblocks); break;
case 2: bitshuffle_untranspose<2><<<s.nblocks, 256, 0, *stream>>>(s.shuffled, s.desc, gpu_out, s.nblocks); break;
default: bitshuffle_untranspose<4><<<s.nblocks, 256, 0, *stream>>>(s.shuffled, s.desc, gpu_out, s.nblocks); break;
}
cuda_err(cudaGetLastError());
}
// The verbatim tail is already on the device inside the uploaded chunk.
if (s.tail_elems > 0)
cuda_err(cudaMemcpyAsync(gpu_out + static_cast<size_t>(s.tail_elem0) * s.elem_size, s.tail_src,
static_cast<size_t>(s.tail_elems) * s.elem_size,
cudaMemcpyDeviceToDevice, *stream));
}
void BSLZ4DecoderGPU::ThrowIfDecodeFailed() {
if (host_status.get()[0] != 0)
throw JFJochException(JFJochExceptionCategory::Compression,
"bslz4 block did not decode to its declared length - the compressed data is corrupt");
}
float BSLZ4DecoderGPU::GetDecodeTime_s() const {
if (!decode_timed)
return 0.0f;
float ms = 0.0f;
if (cudaEventElapsedTime(&ms, decode_start, decode_stop) != cudaSuccess)
return 0.0f;
return ms * 1e-3f;
}
@@ -16,6 +16,20 @@ struct BSLZ4BlockDesc {
uint32_t nelem; // elements in this block (the last one is usually shorter)
};
// What DecodeShuffled() leaves on the device: the LZ4 output, still bitshuffled, plus everything the
// un-transpose needs to finish the image. The tail is the handful of elements bitshuffle stores
// verbatim; it is already on the device inside the uploaded chunk, so it is handed over as a device
// pointer rather than copied again from the host.
struct BSLZ4ShuffledImage {
const uint8_t *shuffled = nullptr;
const BSLZ4BlockDesc *desc = nullptr;
int nblocks = 0;
uint32_t elem_size = 0;
const uint8_t *tail_src = nullptr;
uint32_t tail_elems = 0;
uint32_t tail_elem0 = 0; // index of the first tail element in the image
};
// Decompress a bitshuffle+LZ4 image ON THE DEVICE, so the compressed bytes are what crosses PCIe.
//
// The idea - upload the compressed chunk and decode it on the GPU rather than decompressing on the
@@ -30,6 +44,14 @@ struct BSLZ4BlockDesc {
//
// Only BSHUF_LZ4 is handled. The zstd variants have no device decoder, so Supports() returns false
// and the caller decompresses on the host exactly as before.
//
// The container comes off the network or off disk, so it is not trusted. Everything the host can
// check cheaply is checked before any work is queued and throws; what only the kernel can see - a
// block that does not decode to its declared length, which is what a corrupt LZ4 payload looks like
// - raises a device-side flag that ThrowIfDecodeFailed() reports once the caller has synchronised.
// The CPU decoder makes exactly the same checks (LZ4_decompress_safe's length check plus the
// consumed-input check in JFJochDecompress.h), so a chunk either decodes identically on both or
// fails on both. It is never silently wrong on one and right on the other.
class BSLZ4DecoderGPU {
std::shared_ptr<CudaStream> stream;
@@ -37,19 +59,42 @@ class BSLZ4DecoderGPU {
CudaDevicePtr<uint8_t> gpu_shuffled; // LZ4 output, still bitshuffled
CudaDevicePtr<BSLZ4BlockDesc> gpu_desc;
CudaHostPtr<BSLZ4BlockDesc> host_desc; // pinned, so the descriptor upload is truly async
CudaDevicePtr<uint32_t> gpu_status; // set by the kernel when a block decodes short
CudaHostPtr<uint32_t> host_status;
// Bracket the decode so the time it takes can still be reported as decompression, which is what
// it is. Both are recorded on the decoder's stream and read after the caller synchronises.
CudaEvent decode_start;
CudaEvent decode_stop;
bool decode_timed = false;
size_t max_compressed_bytes = 0;
size_t max_uncompressed_bytes = 0;
size_t compressed_capacity = 0;
size_t max_blocks = 0;
void EnsureCompressedCapacity(size_t bytes);
void EnsureBlockCapacity(size_t nblocks);
public:
BSLZ4DecoderGPU(size_t max_uncompressed_bytes, std::shared_ptr<CudaStream> stream);
// True when this image can be decoded on the device. Everything else must go the host route.
static bool Supports(const CompressedImage &image);
// Decode into gpu_out, which must hold image.GetUncompressedSize() bytes. Work is queued on the
// decoder's stream and the caller synchronises. Throws if the container is malformed - it comes
// off the network or off disk, so it is not trusted.
// Locate the blocks, upload the chunk, and run the LZ4 pass. The result is still bitshuffled -
// the caller finishes it, either with Decode()'s un-transpose or by fusing the un-transpose into
// its own kernel. Work is queued on the decoder's stream and the caller synchronises.
BSLZ4ShuffledImage DecodeShuffled(const CompressedImage &image);
// Decode into gpu_out, which must hold image.GetUncompressedSize() bytes. The plain raw-bytes
// path: DecodeShuffled() plus the un-transpose. Used by the tests and by any caller that wants
// the decompressed image rather than a preprocessed one.
void Decode(const CompressedImage &image, uint8_t *gpu_out);
// Report a block that did not decode to its declared length. MUST be called after the caller has
// synchronised the stream; until then the flag has not arrived. Throws on failure.
void ThrowIfDecodeFailed();
// Device time spent decoding the last image, in seconds. Valid after the caller synchronises.
[[nodiscard]] float GetDecodeTime_s() const;
};
@@ -39,6 +39,11 @@ public:
virtual bool AnalyzeCompressed(ImagePreprocessorBuffer &processed_image, const CompressedImage &image,
ImageStatistics &stats) { return false; }
// Device time the last AnalyzeCompressed() spent getting the chunk across and decompressing it,
// so the caller can still report a decompression cost once the host no longer does the work.
// Meaningless unless the previous call returned true.
[[nodiscard]] virtual float GetLastDecompressionTime_s() const { return 0.0f; }
// Resize the buffer an image will be decompressed into and page-lock it, so that the host->device
// copy of Analyze() is a real DMA. Without page-locking the driver stages the copy through its own
// pinned pool, which is a host-side copy on the calling thread: it does not overlap and it degrades
@@ -1,6 +1,8 @@
// 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>
@@ -89,12 +91,139 @@ __global__ void preprocess_kernel(
}
}
// 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_decompressed_image(npixels * sizeof(uint32_t)), // Overshoot - if input image is 1- or 2-byte, then it is still fine, while memory loss is minimal
gpu_stats(1),
cpu_stats(1),
cpu_stats_reg(cpu_stats) {
@@ -116,6 +245,10 @@ ImagePreprocessorGPU::ImagePreprocessorGPU(const DiffractionExperiment &experime
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;
@@ -156,33 +289,79 @@ bool ImagePreprocessorGPU::AnalyzeCompressed(ImagePreprocessorBuffer &processed_
if (!bslz4_decoder)
bslz4_decoder = std::make_unique<BSLZ4DecoderGPU>(npixels * sizeof(uint32_t), stream);
// Straight into the buffer Analyze() would have filled by copying the decompressed frame across
// the bus; from here the two paths are the same code.
bslz4_decoder->Decode(image, gpu_decompressed_image.get());
// 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 = AnalyzeOnDevice<int8_t>(processed_image, INT8_MIN, INT8_MAX); return true;
stats = UntransposeAndAnalyze<int8_t, 1>(processed_image, shuffled, INT8_MIN, INT8_MAX); return true;
case CompressedImageMode::Uint8:
stats = AnalyzeOnDevice<uint8_t>(processed_image, UINT8_MAX, UINT8_MAX); return true;
stats = UntransposeAndAnalyze<uint8_t, 1>(processed_image, shuffled, UINT8_MAX, UINT8_MAX); return true;
case CompressedImageMode::Int16:
stats = AnalyzeOnDevice<int16_t>(processed_image, INT16_MIN, INT16_MAX); return true;
stats = UntransposeAndAnalyze<int16_t, 2>(processed_image, shuffled, INT16_MIN, INT16_MAX); return true;
case CompressedImageMode::Uint16:
stats = AnalyzeOnDevice<uint16_t>(processed_image, UINT16_MAX, UINT16_MAX); return true;
stats = UntransposeAndAnalyze<uint16_t, 2>(processed_image, shuffled, UINT16_MAX, UINT16_MAX); return true;
case CompressedImageMode::Int32:
stats = AnalyzeOnDevice<int32_t>(processed_image, INT32_MIN, INT32_MAX); return true;
stats = UntransposeAndAnalyze<int32_t, 4>(processed_image, shuffled, INT32_MIN, INT32_MAX); return true;
case CompressedImageMode::Uint32:
stats = AnalyzeOnDevice<uint32_t>(processed_image, UINT32_MAX, UINT32_MAX); return true;
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.
@@ -17,6 +17,9 @@ class ImagePreprocessorGPU : public ImagePreprocessor {
int blocks;
// Geometry-only, so one copy per GPU shared with every other engine on it (CudaSharedTables.h).
std::shared_ptr<CudaDevicePtr<uint8_t>> gpu_mask;
// Landing buffer for the HOST-upload path only. The device-decode path un-transposes straight
// into the preprocessed image, so it never needs this - and at 4 bytes per pixel it is worth a
// frame per worker, so it is allocated on first use rather than always.
CudaDevicePtr<uint8_t> gpu_decompressed_image;
CudaDevicePtr<ImageStatistics> gpu_stats;
@@ -30,8 +33,13 @@ class ImagePreprocessorGPU : public ImagePreprocessor {
std::unique_ptr<BSLZ4DecoderGPU> bslz4_decoder;
template <class T> ImageStatistics Analyze(ImagePreprocessorBuffer &processed_image, const uint8_t *input, T err_value, T sat_value);
// Preprocess an image already sitting in gpu_decompressed_image, shared by both entry points.
// Preprocess an image already sitting in gpu_decompressed_image (the host-upload path).
template <class T> ImageStatistics AnalyzeOnDevice(ImagePreprocessorBuffer &processed_image, T err_value, T sat_value);
// Preprocess straight out of the bitshuffled bytes (the device-decode path). Same per-pixel
// decision and same statistics as AnalyzeOnDevice, with the un-transpose folded in.
template <class T, int ES> ImageStatistics UntransposeAndAnalyze(ImagePreprocessorBuffer &processed_image,
const BSLZ4ShuffledImage &shuffled,
T err_value, T sat_value);
public:
// copy_image_to_host copies the preprocessed image back after every frame. It is only needed when
// something on the CPU reads it - the GPU engines all work off the device buffer - and at 4 bytes
@@ -41,6 +49,7 @@ public:
ImageStatistics Analyze(ImagePreprocessorBuffer &processed_image, const uint8_t *decompressed_image, CompressedImageMode image_mode) override;
bool AnalyzeCompressed(ImagePreprocessorBuffer &processed_image, const CompressedImage &image,
ImageStatistics &stats) override;
[[nodiscard]] float GetLastDecompressionTime_s() const override;
void PinInputBuffer(std::vector<uint8_t> &buffer, size_t size) override;
};