diff --git a/programs/fileio.c b/programs/fileio.c index 3a6d0fcd6..b365a3d6d 100644 --- a/programs/fileio.c +++ b/programs/fileio.c @@ -1012,6 +1012,57 @@ typedef struct { FIO_rust_compress_status_display_fn display_status; } FIO_rust_compress_callbacks_t; +enum { + FIO_RUST_GZIP_OK = 0, + FIO_RUST_GZIP_INIT_ERROR = 1, + FIO_RUST_GZIP_DEFLATE_ERROR = 2, + FIO_RUST_GZIP_FINISH_ERROR = 3, + FIO_RUST_GZIP_END_ERROR = 4, + FIO_RUST_GZIP_INVALID_PROJECTION = 5 +}; + +typedef void (*FIO_rust_gzip_read_fill_fn)( + void* opaque, size_t requested, const unsigned char** buffer, size_t* loaded); +typedef void (*FIO_rust_gzip_read_consume_fn)(void* opaque, size_t consumed); +typedef void (*FIO_rust_gzip_write_acquire_fn)( + void* opaque, void** job, unsigned char** buffer, size_t* bufferSize); +typedef void (*FIO_rust_gzip_write_enqueue_fn)( + void* opaque, void** job, size_t usedBufferSize, + unsigned char** buffer, size_t* bufferSize); +typedef void (*FIO_rust_gzip_write_release_fn)(void* opaque, void* job); +typedef void (*FIO_rust_gzip_sparse_write_end_fn)(void* opaque); +typedef int (*FIO_rust_gzip_zlib_init_fn)(void* opaque, int compressionLevel); +typedef int (*FIO_rust_gzip_zlib_deflate_fn)( + void* opaque, const unsigned char* input, size_t inputSize, + unsigned char* output, size_t outputSize, int flush, + size_t* consumed, size_t* produced); +typedef int (*FIO_rust_gzip_zlib_end_fn)(void* opaque); +typedef void (*FIO_rust_gzip_progress_fn)( + void* opaque, U64 srcFileSize, U64 inFileSize, U64 outFileSize); + +typedef struct { + void* readOpaque; + void* writeOpaque; + void* zlibOpaque; + void* progressOpaque; + size_t readBufferSize; + FIO_rust_gzip_read_fill_fn readFill; + FIO_rust_gzip_read_consume_fn readConsume; + FIO_rust_gzip_write_acquire_fn writeAcquire; + FIO_rust_gzip_write_enqueue_fn writeEnqueue; + FIO_rust_gzip_write_release_fn writeRelease; + FIO_rust_gzip_sparse_write_end_fn sparseWriteEnd; + FIO_rust_gzip_zlib_init_fn zlibInit; + FIO_rust_gzip_zlib_deflate_fn zlibDeflate; + FIO_rust_gzip_zlib_end_fn zlibEnd; + FIO_rust_gzip_progress_fn progress; +} FIO_rust_gzip_compress_projection_t; + +int FIO_rust_compressGzipFrame( + const FIO_rust_gzip_compress_projection_t* projection, + const char* srcFileName, U64 srcFileSize, int compressionLevel, + U64* readsize, U64* compressedSize, int* zlibResult); + int FIO_rust_compressFilenameInternal( void* fCtx, FIO_prefs_t* prefs, void* ress, const char* dstFileName, const char* srcFileName, @@ -1158,96 +1209,106 @@ static void FIO_freeCResources(cRess_t* const ress) #ifdef ZSTD_GZCOMPRESS -static unsigned long long -FIO_compressGzFrame(const cRess_t* ress, /* buffers & handlers are used, but not changed */ - const char* srcFileName, U64 const srcFileSize, - int compressionLevel, U64* readsize) +static void FIO_rust_gzip_readFill(void* opaque, size_t requested, + const unsigned char** buffer, size_t* loaded) { - unsigned long long inFileSize = 0, outFileSize = 0; - z_stream strm; - IOJob_t *writeJob = NULL; + ReadPoolCtx_t* const readCtx = (ReadPoolCtx_t*)opaque; + AIO_ReadPool_fillBuffer(readCtx, requested); + *buffer = readCtx->srcBuffer; + *loaded = readCtx->srcBufferLoaded; +} - if (compressionLevel > Z_BEST_COMPRESSION) - compressionLevel = Z_BEST_COMPRESSION; +static void FIO_rust_gzip_readConsume(void* opaque, size_t consumed) +{ + AIO_ReadPool_consumeBytes((ReadPoolCtx_t*)opaque, consumed); +} - strm.zalloc = Z_NULL; - strm.zfree = Z_NULL; - strm.opaque = Z_NULL; +static void FIO_rust_gzip_writeAcquire(void* opaque, void** job, + unsigned char** buffer, size_t* bufferSize) +{ + IOJob_t* const writeJob = AIO_WritePool_acquireJob((WritePoolCtx_t*)opaque); + *job = writeJob; + *buffer = (unsigned char*)writeJob->buffer; + *bufferSize = writeJob->bufferSize; +} - { int const ret = deflateInit2(&strm, compressionLevel, Z_DEFLATED, +static void FIO_rust_gzip_writeEnqueue(void* opaque, void** job, + size_t usedBufferSize, + unsigned char** buffer, size_t* bufferSize) +{ + IOJob_t* writeJob = (IOJob_t*)*job; + writeJob->usedBufferSize = usedBufferSize; + AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob); + *job = writeJob; + *buffer = (unsigned char*)writeJob->buffer; + *bufferSize = writeJob->bufferSize; + (void)opaque; +} + +static void FIO_rust_gzip_writeRelease(void* opaque, void* job) +{ + AIO_WritePool_releaseIoJob((IOJob_t*)job); + (void)opaque; +} + +static void FIO_rust_gzip_sparseWriteEnd(void* opaque) +{ + AIO_WritePool_sparseWriteEnd((WritePoolCtx_t*)opaque); +} + +static int FIO_rust_gzip_zlibInit(void* opaque, int compressionLevel) +{ + z_stream* const strm = (z_stream*)opaque; + strm->zalloc = Z_NULL; + strm->zfree = Z_NULL; + strm->opaque = Z_NULL; + return deflateInit2(strm, compressionLevel, Z_DEFLATED, 15 /* maxWindowLogSize */ + 16 /* gzip only */, 8, Z_DEFAULT_STRATEGY); /* see https://www.zlib.net/manual.html */ - if (ret != Z_OK) { - EXM_THROW(71, "zstd: %s: deflateInit2 error %d \n", srcFileName, ret); - } } +} - writeJob = AIO_WritePool_acquireJob(ress->writeCtx); - strm.next_in = 0; - strm.avail_in = 0; - strm.next_out = (Bytef*)writeJob->buffer; - strm.avail_out = (uInt)writeJob->bufferSize; +static int FIO_rust_gzip_zlibDeflate(void* opaque, const unsigned char* input, + size_t inputSize, unsigned char* output, + size_t outputSize, int flush, + size_t* consumed, size_t* produced) +{ + z_stream* const strm = (z_stream*)opaque; + size_t const availBefore = inputSize; + size_t const outputBefore = outputSize; + int ret; - while (1) { - int ret; - if (strm.avail_in == 0) { - AIO_ReadPool_fillBuffer(ress->readCtx, ZSTD_CStreamInSize()); - if (ress->readCtx->srcBufferLoaded == 0) break; - inFileSize += ress->readCtx->srcBufferLoaded; - strm.next_in = (z_const unsigned char*)ress->readCtx->srcBuffer; - strm.avail_in = (uInt)ress->readCtx->srcBufferLoaded; - } + assert(inputSize <= UINT_MAX); + assert(outputSize <= UINT_MAX); + strm->next_in = (z_const unsigned char*)input; + strm->avail_in = (uInt)inputSize; + strm->next_out = output; + strm->avail_out = (uInt)outputSize; + ret = deflate(strm, flush); + *consumed = availBefore - strm->avail_in; + *produced = outputBefore - strm->avail_out; + return ret; +} - { - size_t const availBefore = strm.avail_in; - ret = deflate(&strm, Z_NO_FLUSH); - AIO_ReadPool_consumeBytes(ress->readCtx, availBefore - strm.avail_in); - } +static int FIO_rust_gzip_zlibEnd(void* opaque) +{ + return deflateEnd((z_stream*)opaque); +} - if (ret != Z_OK) - EXM_THROW(72, "zstd: %s: deflate error %d \n", srcFileName, ret); - { size_t const cSize = writeJob->bufferSize - strm.avail_out; - if (cSize) { - writeJob->usedBufferSize = cSize; - AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob); - outFileSize += cSize; - strm.next_out = (Bytef*)writeJob->buffer; - strm.avail_out = (uInt)writeJob->bufferSize; - } } - if (srcFileSize == UTIL_FILESIZE_UNKNOWN) { - DISPLAYUPDATE_PROGRESS( - "\rRead : %u MB ==> %.2f%% ", - (unsigned)(inFileSize>>20), - (double)outFileSize/(double)inFileSize*100) - } else { - DISPLAYUPDATE_PROGRESS( - "\rRead : %u / %u MB ==> %.2f%% ", - (unsigned)(inFileSize>>20), (unsigned)(srcFileSize>>20), - (double)outFileSize/(double)inFileSize*100); - } } - - while (1) { - int const ret = deflate(&strm, Z_FINISH); - { size_t const cSize = writeJob->bufferSize - strm.avail_out; - if (cSize) { - writeJob->usedBufferSize = cSize; - AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob); - outFileSize += cSize; - strm.next_out = (Bytef*)writeJob->buffer; - strm.avail_out = (uInt)writeJob->bufferSize; - } } - if (ret == Z_STREAM_END) break; - if (ret != Z_BUF_ERROR) - EXM_THROW(77, "zstd: %s: deflate error %d \n", srcFileName, ret); +static void FIO_rust_gzip_progress(void* opaque, U64 srcFileSize, + U64 inFileSize, U64 outFileSize) +{ + (void)opaque; + if (srcFileSize == UTIL_FILESIZE_UNKNOWN) { + DISPLAYUPDATE_PROGRESS( + "\rRead : %u MB ==> %.2f%% ", + (unsigned)(inFileSize>>20), + (double)outFileSize/(double)inFileSize*100) + } else { + DISPLAYUPDATE_PROGRESS( + "\rRead : %u / %u MB ==> %.2f%% ", + (unsigned)(inFileSize>>20), (unsigned)(srcFileSize>>20), + (double)outFileSize/(double)inFileSize*100); } - - { int const ret = deflateEnd(&strm); - if (ret != Z_OK) { - EXM_THROW(79, "zstd: %s: deflateEnd error %d \n", srcFileName, ret); - } } - *readsize = inFileSize; - AIO_WritePool_releaseIoJob(writeJob); - AIO_WritePool_sparseWriteEnd(ress->writeCtx); - return outFileSize; } #endif @@ -1691,8 +1752,47 @@ FIO_rust_compressGzipCallback(void* ress, const char* srcFileName, U64 srcFileSize, int compressionLevel, U64* readsize) { - return FIO_compressGzFrame((const cRess_t*)ress, srcFileName, srcFileSize, - compressionLevel, readsize); + const cRess_t* const ressPtr = (const cRess_t*)ress; + FIO_rust_gzip_compress_projection_t projection; + z_stream strm; + U64 compressedSize = 0; + int zlibResult = Z_OK; + int status; + + memset(&projection, 0, sizeof(projection)); + projection.readOpaque = (void*)ressPtr->readCtx; + projection.writeOpaque = (void*)ressPtr->writeCtx; + projection.zlibOpaque = &strm; + projection.readBufferSize = ZSTD_CStreamInSize(); + projection.readFill = FIO_rust_gzip_readFill; + projection.readConsume = FIO_rust_gzip_readConsume; + projection.writeAcquire = FIO_rust_gzip_writeAcquire; + projection.writeEnqueue = FIO_rust_gzip_writeEnqueue; + projection.writeRelease = FIO_rust_gzip_writeRelease; + projection.sparseWriteEnd = FIO_rust_gzip_sparseWriteEnd; + projection.zlibInit = FIO_rust_gzip_zlibInit; + projection.zlibDeflate = FIO_rust_gzip_zlibDeflate; + projection.zlibEnd = FIO_rust_gzip_zlibEnd; + projection.progress = FIO_rust_gzip_progress; + + status = FIO_rust_compressGzipFrame( + &projection, srcFileName, srcFileSize, compressionLevel, + readsize, &compressedSize, &zlibResult); + switch (status) { + case FIO_RUST_GZIP_OK: + return compressedSize; + case FIO_RUST_GZIP_INIT_ERROR: + EXM_THROW(71, "zstd: %s: deflateInit2 error %d \n", srcFileName, zlibResult); + case FIO_RUST_GZIP_DEFLATE_ERROR: + EXM_THROW(72, "zstd: %s: deflate error %d \n", srcFileName, zlibResult); + case FIO_RUST_GZIP_FINISH_ERROR: + EXM_THROW(77, "zstd: %s: deflate error %d \n", srcFileName, zlibResult); + case FIO_RUST_GZIP_END_ERROR: + EXM_THROW(79, "zstd: %s: deflateEnd error %d \n", srcFileName, zlibResult); + default: + assert(status == FIO_RUST_GZIP_INVALID_PROJECTION); + EXM_THROW(72, "zstd: %s: deflate error %d \n", srcFileName, zlibResult); + } } #endif diff --git a/rust/src/fileio_asyncio.rs b/rust/src/fileio_asyncio.rs index 965914f74..7c0bcb208 100644 --- a/rust/src/fileio_asyncio.rs +++ b/rust/src/fileio_asyncio.rs @@ -121,6 +121,65 @@ pub struct FIO_rust_compress_callbacks_t { pub display_status: Option, } +pub const FIO_RUST_GZIP_OK: c_int = 0; +pub const FIO_RUST_GZIP_INIT_ERROR: c_int = 1; +pub const FIO_RUST_GZIP_DEFLATE_ERROR: c_int = 2; +pub const FIO_RUST_GZIP_FINISH_ERROR: c_int = 3; +pub const FIO_RUST_GZIP_END_ERROR: c_int = 4; +pub const FIO_RUST_GZIP_INVALID_PROJECTION: c_int = 5; + +const FIO_RUST_GZIP_Z_OK: c_int = 0; +const FIO_RUST_GZIP_Z_STREAM_END: c_int = 1; +const FIO_RUST_GZIP_Z_BUF_ERROR: c_int = -5; +const FIO_RUST_GZIP_Z_NO_FLUSH: c_int = 0; +const FIO_RUST_GZIP_Z_FINISH: c_int = 4; +const FIO_RUST_GZIP_BEST_COMPRESSION: c_int = 9; + +pub type FIO_rust_gzip_read_fill_fn = + unsafe extern "C" fn(*mut c_void, usize, *mut *const u8, *mut usize); +pub type FIO_rust_gzip_read_consume_fn = unsafe extern "C" fn(*mut c_void, usize); +pub type FIO_rust_gzip_write_acquire_fn = + unsafe extern "C" fn(*mut c_void, *mut *mut c_void, *mut *mut u8, *mut usize); +pub type FIO_rust_gzip_write_enqueue_fn = + unsafe extern "C" fn(*mut c_void, *mut *mut c_void, usize, *mut *mut u8, *mut usize); +pub type FIO_rust_gzip_write_release_fn = unsafe extern "C" fn(*mut c_void, *mut c_void); +pub type FIO_rust_gzip_sparse_write_end_fn = unsafe extern "C" fn(*mut c_void); +pub type FIO_rust_gzip_zlib_init_fn = unsafe extern "C" fn(*mut c_void, c_int) -> c_int; +pub type FIO_rust_gzip_zlib_deflate_fn = unsafe extern "C" fn( + *mut c_void, + *const u8, + usize, + *mut u8, + usize, + c_int, + *mut usize, + *mut usize, +) -> c_int; +pub type FIO_rust_gzip_zlib_end_fn = unsafe extern "C" fn(*mut c_void) -> c_int; +pub type FIO_rust_gzip_progress_fn = unsafe extern "C" fn(*mut c_void, u64, u64, u64); + +/// Rust owns the gzip read/deflate/flush loop. C supplies opaque pool and +/// zlib operations so neither the private `cRess_t` nor `z_stream` layout +/// crosses this ABI. +#[repr(C)] +pub struct FIO_rust_gzip_compress_projection_t { + pub read_opaque: *mut c_void, + pub write_opaque: *mut c_void, + pub zlib_opaque: *mut c_void, + pub progress_opaque: *mut c_void, + pub read_buffer_size: usize, + pub read_fill: Option, + pub read_consume: Option, + pub write_acquire: Option, + pub write_enqueue: Option, + pub write_release: Option, + pub sparse_write_end: Option, + pub zlib_init: Option, + pub zlib_deflate: Option, + pub zlib_end: Option, + pub progress: Option, +} + /// C's `FIO_prefs_t` from `programs/fileio_types.h`. /// /// `fileio_prefs.rs` contains the same C layout for the preferences API. It @@ -1664,6 +1723,215 @@ pub unsafe extern "C" fn FIO_rust_compressFilenameInternal( FIO_RUST_COMPRESS_OK } +/// Compresses one gzip member through the C-owned zlib and asynchronous +/// resource callbacks. The loop mirrors the original `FIO_compressGzFrame` +/// sequencing: input is counted when a read-pool buffer is loaded, consumed +/// bytes are returned immediately after each zlib call, non-empty output is +/// flushed before progress is reported, and the final empty job is released +/// only after `deflateEnd()` succeeds. +#[no_mangle] +pub unsafe extern "C" fn FIO_rust_compressGzipFrame( + projection: *const FIO_rust_gzip_compress_projection_t, + src_file_name: *const c_char, + src_file_size: u64, + compression_level: c_int, + read_size: *mut u64, + compressed_size: *mut u64, + zlib_result: *mut c_int, +) -> c_int { + assert!(!projection.is_null()); + assert!(!src_file_name.is_null()); + assert!(!read_size.is_null()); + assert!(!compressed_size.is_null()); + assert!(!zlib_result.is_null()); + + unsafe { + *read_size = 0; + *compressed_size = 0; + *zlib_result = FIO_RUST_GZIP_Z_OK; + } + + let projection = unsafe { &*projection }; + let Some(read_fill) = projection.read_fill else { + return FIO_RUST_GZIP_INVALID_PROJECTION; + }; + let Some(read_consume) = projection.read_consume else { + return FIO_RUST_GZIP_INVALID_PROJECTION; + }; + let Some(write_acquire) = projection.write_acquire else { + return FIO_RUST_GZIP_INVALID_PROJECTION; + }; + let Some(write_enqueue) = projection.write_enqueue else { + return FIO_RUST_GZIP_INVALID_PROJECTION; + }; + let Some(write_release) = projection.write_release else { + return FIO_RUST_GZIP_INVALID_PROJECTION; + }; + let Some(sparse_write_end) = projection.sparse_write_end else { + return FIO_RUST_GZIP_INVALID_PROJECTION; + }; + let Some(zlib_init) = projection.zlib_init else { + return FIO_RUST_GZIP_INVALID_PROJECTION; + }; + let Some(zlib_deflate) = projection.zlib_deflate else { + return FIO_RUST_GZIP_INVALID_PROJECTION; + }; + let Some(zlib_end) = projection.zlib_end else { + return FIO_RUST_GZIP_INVALID_PROJECTION; + }; + let Some(progress) = projection.progress else { + return FIO_RUST_GZIP_INVALID_PROJECTION; + }; + + let compression_level = compression_level.min(FIO_RUST_GZIP_BEST_COMPRESSION); + let init_result = unsafe { zlib_init(projection.zlib_opaque, compression_level) }; + if init_result != FIO_RUST_GZIP_Z_OK { + unsafe { *zlib_result = init_result }; + return FIO_RUST_GZIP_INIT_ERROR; + } + + let mut job = ptr::null_mut::(); + let mut output = ptr::null_mut::(); + let mut output_size = 0_usize; + unsafe { + write_acquire( + projection.write_opaque, + &mut job, + &mut output, + &mut output_size, + ); + } + assert!(!job.is_null()); + assert!(!output.is_null() || output_size == 0); + + let mut input = ptr::null(); + let mut input_size = 0_usize; + let mut in_file_size = 0_u64; + let mut out_file_size = 0_u64; + + loop { + if input_size == 0 { + let mut loaded = 0_usize; + unsafe { + read_fill( + projection.read_opaque, + projection.read_buffer_size, + &mut input, + &mut loaded, + ); + } + if loaded == 0 { + break; + } + assert!(!input.is_null()); + input_size = loaded; + in_file_size = in_file_size.wrapping_add(loaded as u64); + } + + let mut consumed = 0_usize; + let mut produced = 0_usize; + let result = unsafe { + zlib_deflate( + projection.zlib_opaque, + input, + input_size, + output, + output_size, + FIO_RUST_GZIP_Z_NO_FLUSH, + &mut consumed, + &mut produced, + ) + }; + assert!(consumed <= input_size); + assert!(produced <= output_size); + unsafe { read_consume(projection.read_opaque, consumed) }; + if consumed != 0 { + input = unsafe { input.add(consumed) }; + } + input_size -= consumed; + + if result != FIO_RUST_GZIP_Z_OK { + unsafe { *zlib_result = result }; + return FIO_RUST_GZIP_DEFLATE_ERROR; + } + + if produced != 0 { + unsafe { + write_enqueue( + projection.write_opaque, + &mut job, + produced, + &mut output, + &mut output_size, + ); + } + out_file_size = out_file_size.wrapping_add(produced as u64); + } + unsafe { + progress( + projection.progress_opaque, + src_file_size, + in_file_size, + out_file_size, + ) + }; + } + + loop { + let mut consumed = 0_usize; + let mut produced = 0_usize; + let result = unsafe { + zlib_deflate( + projection.zlib_opaque, + ptr::null(), + 0, + output, + output_size, + FIO_RUST_GZIP_Z_FINISH, + &mut consumed, + &mut produced, + ) + }; + assert_eq!(consumed, 0); + assert!(produced <= output_size); + + if produced != 0 { + unsafe { + write_enqueue( + projection.write_opaque, + &mut job, + produced, + &mut output, + &mut output_size, + ); + } + out_file_size = out_file_size.wrapping_add(produced as u64); + } + + if result == FIO_RUST_GZIP_Z_STREAM_END { + break; + } + if result != FIO_RUST_GZIP_Z_BUF_ERROR { + unsafe { *zlib_result = result }; + return FIO_RUST_GZIP_FINISH_ERROR; + } + } + + let end_result = unsafe { zlib_end(projection.zlib_opaque) }; + if end_result != FIO_RUST_GZIP_Z_OK { + unsafe { *zlib_result = end_result }; + return FIO_RUST_GZIP_END_ERROR; + } + + unsafe { + *read_size = in_file_size; + *compressed_size = out_file_size; + write_release(projection.write_opaque, job); + sparse_write_end(projection.write_opaque); + } + FIO_RUST_GZIP_OK +} + /// Decompresses exactly one zstd frame using the existing asynchronous pools. /// /// C retains the frame dispatcher and all user-facing diagnostics. This ABI @@ -2386,6 +2654,279 @@ mod tests { assert_eq!(context.totalBytesOutput, 47); } + struct GzipProjectionState { + input: [u8; 5], + input_pos: usize, + output_buffer: [u8; 4], + read_fill_calls: usize, + consumed: Vec, + enqueue_sizes: Vec, + output_chunks: Vec>, + acquire_calls: usize, + release_calls: usize, + sparse_end_calls: usize, + init_levels: Vec, + deflate_flushes: Vec, + finish_calls: usize, + end_calls: usize, + progress: Vec<(u64, u64, u64)>, + init_result: c_int, + } + + impl Default for GzipProjectionState { + fn default() -> Self { + Self { + input: *b"abcde", + input_pos: 0, + output_buffer: [0; 4], + read_fill_calls: 0, + consumed: Vec::new(), + enqueue_sizes: Vec::new(), + output_chunks: Vec::new(), + acquire_calls: 0, + release_calls: 0, + sparse_end_calls: 0, + init_levels: Vec::new(), + deflate_flushes: Vec::new(), + finish_calls: 0, + end_calls: 0, + progress: Vec::new(), + init_result: FIO_RUST_GZIP_Z_OK, + } + } + } + + unsafe extern "C" fn gzip_test_read_fill( + opaque: *mut c_void, + requested: usize, + buffer: *mut *const u8, + loaded: *mut usize, + ) { + let state = unsafe { &mut *opaque.cast::() }; + state.read_fill_calls += 1; + let available = state.input.len() - state.input_pos; + let amount = available.min(requested); + unsafe { + *buffer = state.input.as_ptr().add(state.input_pos); + *loaded = amount; + } + } + + unsafe extern "C" fn gzip_test_read_consume(opaque: *mut c_void, amount: usize) { + let state = unsafe { &mut *opaque.cast::() }; + assert!(amount <= state.input.len() - state.input_pos); + state.input_pos += amount; + state.consumed.push(amount); + } + + unsafe extern "C" fn gzip_test_write_acquire( + opaque: *mut c_void, + job: *mut *mut c_void, + buffer: *mut *mut u8, + buffer_size: *mut usize, + ) { + let state = unsafe { &mut *opaque.cast::() }; + state.acquire_calls += 1; + unsafe { + *job = opaque; + *buffer = state.output_buffer.as_mut_ptr(); + *buffer_size = state.output_buffer.len(); + } + } + + unsafe extern "C" fn gzip_test_write_enqueue( + opaque: *mut c_void, + job: *mut *mut c_void, + used: usize, + buffer: *mut *mut u8, + buffer_size: *mut usize, + ) { + let state = unsafe { &mut *opaque.cast::() }; + assert_eq!(unsafe { *job }, opaque); + assert!(used <= state.output_buffer.len()); + state.enqueue_sizes.push(used); + state + .output_chunks + .push(state.output_buffer[..used].to_vec()); + unsafe { + *buffer = state.output_buffer.as_mut_ptr(); + *buffer_size = state.output_buffer.len(); + } + } + + unsafe extern "C" fn gzip_test_write_release(opaque: *mut c_void, job: *mut c_void) { + let state = unsafe { &mut *opaque.cast::() }; + assert_eq!(job, opaque); + state.release_calls += 1; + } + + unsafe extern "C" fn gzip_test_sparse_end(opaque: *mut c_void) { + let state = unsafe { &mut *opaque.cast::() }; + state.sparse_end_calls += 1; + } + + unsafe extern "C" fn gzip_test_zlib_init(opaque: *mut c_void, level: c_int) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + state.init_levels.push(level); + state.init_result + } + + unsafe extern "C" fn gzip_test_zlib_deflate( + opaque: *mut c_void, + _input: *const u8, + input_size: usize, + output: *mut u8, + output_size: usize, + flush: c_int, + consumed: *mut usize, + produced: *mut usize, + ) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + state.deflate_flushes.push(flush); + unsafe { + *consumed = 0; + *produced = 0; + } + if flush == FIO_RUST_GZIP_Z_NO_FLUSH { + assert!(input_size != 0); + assert!(output_size != 0); + unsafe { + *consumed = input_size.min(2); + *produced = 1; + *output = b'x'; + } + return FIO_RUST_GZIP_Z_OK; + } + + state.finish_calls += 1; + if state.finish_calls == 1 { + return FIO_RUST_GZIP_Z_BUF_ERROR; + } + assert!(output_size >= 2); + unsafe { + *produced = 2; + *output.add(0) = b'y'; + *output.add(1) = b'z'; + } + FIO_RUST_GZIP_Z_STREAM_END + } + + unsafe extern "C" fn gzip_test_zlib_end(opaque: *mut c_void) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + state.end_calls += 1; + FIO_RUST_GZIP_Z_OK + } + + unsafe extern "C" fn gzip_test_progress( + opaque: *mut c_void, + _src_file_size: u64, + in_file_size: u64, + out_file_size: u64, + ) { + let state = unsafe { &mut *opaque.cast::() }; + state + .progress + .push((_src_file_size, in_file_size, out_file_size)); + } + + fn gzip_test_projection( + state: &mut GzipProjectionState, + ) -> FIO_rust_gzip_compress_projection_t { + let opaque = (state as *mut GzipProjectionState).cast::(); + FIO_rust_gzip_compress_projection_t { + read_opaque: opaque, + write_opaque: opaque, + zlib_opaque: opaque, + progress_opaque: opaque, + read_buffer_size: 3, + read_fill: Some(gzip_test_read_fill), + read_consume: Some(gzip_test_read_consume), + write_acquire: Some(gzip_test_write_acquire), + write_enqueue: Some(gzip_test_write_enqueue), + write_release: Some(gzip_test_write_release), + sparse_write_end: Some(gzip_test_sparse_end), + zlib_init: Some(gzip_test_zlib_init), + zlib_deflate: Some(gzip_test_zlib_deflate), + zlib_end: Some(gzip_test_zlib_end), + progress: Some(gzip_test_progress), + } + } + + #[test] + fn gzip_projection_preserves_pool_order_and_accounting() { + let mut state = GzipProjectionState::default(); + let projection = gzip_test_projection(&mut state); + let mut read_size = 0; + let mut compressed_size = 0; + let mut zlib_result = 0; + + assert_eq!( + unsafe { + FIO_rust_compressGzipFrame( + &projection, + c"gzip-test".as_ptr(), + 5, + 99, + &mut read_size, + &mut compressed_size, + &mut zlib_result, + ) + }, + FIO_RUST_GZIP_OK + ); + assert_eq!(read_size, 5); + assert_eq!(compressed_size, 5); + assert_eq!(zlib_result, FIO_RUST_GZIP_Z_OK); + assert_eq!(state.init_levels, vec![FIO_RUST_GZIP_BEST_COMPRESSION]); + assert_eq!(state.consumed, vec![2, 1, 2]); + assert_eq!(state.enqueue_sizes, vec![1, 1, 1, 2]); + assert_eq!( + state.output_chunks, + vec![vec![b'x'], vec![b'x'], vec![b'x'], vec![b'y', b'z']] + ); + assert_eq!(state.deflate_flushes, vec![0, 0, 0, 4, 4]); + assert_eq!(state.finish_calls, 2); + assert_eq!(state.acquire_calls, 1); + assert_eq!(state.release_calls, 1); + assert_eq!(state.sparse_end_calls, 1); + assert_eq!(state.end_calls, 1); + assert_eq!(state.progress, vec![(5, 3, 1), (5, 3, 2), (5, 5, 3)]); + } + + #[test] + fn gzip_projection_propagates_init_error_before_pool_acquisition() { + let mut state = GzipProjectionState { + init_result: -7, + ..GzipProjectionState::default() + }; + let projection = gzip_test_projection(&mut state); + let mut read_size = 41; + let mut compressed_size = 43; + let mut zlib_result = 0; + + assert_eq!( + unsafe { + FIO_rust_compressGzipFrame( + &projection, + c"gzip-test".as_ptr(), + 5, + 1, + &mut read_size, + &mut compressed_size, + &mut zlib_result, + ) + }, + FIO_RUST_GZIP_INIT_ERROR + ); + assert_eq!(zlib_result, -7); + assert_eq!(read_size, 0); + assert_eq!(compressed_size, 0); + assert_eq!(state.init_levels, vec![1]); + assert_eq!(state.acquire_calls, 0); + assert_eq!(state.release_calls, 0); + assert_eq!(state.sparse_end_calls, 0); + } + #[test] fn classifies_all_cli_decompression_headers() { let is_mock_zstd = |buffer: &[u8]| buffer.starts_with(&[0x28, 0xB5, 0x2F, 0xFD]);