feat(cli): move gzip codec loop into Rust

Move the gzip member read, zlib deflate, finish, output-enqueue, progress,
and accounting loop into Rust. Keep zlib, asynchronous pool, sparse-output,
and diagnostic ownership in C through an explicit projection so optional gzip
builds preserve the existing CLI behavior and error mapping. Add focused tests
for callback ordering, partial input consumption, finish retries, accounting,
and initialization failures.

Test Plan:
- cargo test --manifest-path rust/Cargo.toml --all-targets -- --test-threads=1
- cargo test --manifest-path rust/cli/Cargo.toml --all-targets -- --test-threads=1
- make -B -C programs -j2 zstd
- make -B -C tests -j2 test-cli-tests test-zstd
- FUZZERTEST=-T5s make -B -C tests -j2 test-fuzzer
- cargo clippy --manifest-path rust/Cargo.toml --tests -- -D warnings
This commit is contained in:
2026-07-18 23:07:51 +02:00
parent 575360c6b1
commit 05561c7494
2 changed files with 722 additions and 81 deletions
+181 -81
View File
@@ -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