From fbc80de7dc4a3764ef5c3fdcfb4e913fe03e1bbd Mon Sep 17 00:00:00 2001 From: ddidderr Date: Sun, 19 Jul 2026 10:52:31 +0200 Subject: [PATCH] feat(cli): move separate-output scheduling into Rust Move the separate-destination branch of FIO_compressMultipleFilenames into a Rust-owned scheduler. Preserve source ordering, currFileIdx advancement, successful-file counting, aggregate error handling, and post-loop collision checks while retaining destination construction and private file/resource work in C callbacks. Test Plan: - ulimit -v 41943040; CARGO_BUILD_JOBS=1 cargo test --manifest-path rust/Cargo.toml compression_multiple_separate -- --test-threads=1 - ulimit -v 41943040; make -B -C programs -j1 zstd - ulimit -v 41943040; make -B -C tests -j1 test-cli-tests --- programs/fileio.c | 94 +++++++++++++++++------ rust/src/fileio_asyncio.rs | 150 +++++++++++++++++++++++++++++++++++++ 2 files changed, 223 insertions(+), 21 deletions(-) diff --git a/programs/fileio.c b/programs/fileio.c index 046159273..497c7cf85 100644 --- a/programs/fileio.c +++ b/programs/fileio.c @@ -1154,6 +1154,31 @@ typedef char FIO_rust_compress_multiple_projection_size[ int FIO_rust_compressMultipleFilenames( const FIO_rust_compress_multiple_projection_t* projection); +typedef int (*FIO_rust_compress_multiple_separate_file_fn)( + void* opaque, const char* srcFileName); +typedef struct { + void* fCtx; + const char** inFileNamesTable; + void* opaque; + FIO_rust_compress_multiple_separate_file_fn compressFile; +} FIO_rust_compress_multiple_separate_projection_t; +typedef char FIO_rust_compress_multiple_separate_fctx_offset[ + (offsetof(FIO_rust_compress_multiple_separate_projection_t, fCtx) == 0) ? 1 : -1]; +typedef char FIO_rust_compress_multiple_separate_input_names_offset[ + (offsetof(FIO_rust_compress_multiple_separate_projection_t, inFileNamesTable) + == sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_compress_multiple_separate_opaque_offset[ + (offsetof(FIO_rust_compress_multiple_separate_projection_t, opaque) + == 2 * sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_compress_multiple_separate_callback_offset[ + (offsetof(FIO_rust_compress_multiple_separate_projection_t, compressFile) + == 3 * sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_compress_multiple_separate_projection_size[ + (sizeof(FIO_rust_compress_multiple_separate_projection_t) + == 3 * sizeof(void*) + sizeof(FIO_rust_compress_multiple_separate_file_fn)) ? 1 : -1]; +int FIO_rust_compressMultipleSeparateFilenames( + const FIO_rust_compress_multiple_separate_projection_t* projection); + enum { FIO_RUST_ZSTD_OK = 0, FIO_RUST_ZSTD_COMPRESS_ERROR = 1, @@ -2944,6 +2969,45 @@ static int FIO_rust_compressMultipleFileCallback(void* opaque, dstFileName, srcFileName, context->compressionLevel); } +typedef struct { + FIO_ctx_t* fCtx; + FIO_prefs_t* prefs; + cRess_t* ress; + const char* outMirroredRootDirName; + const char* outDirName; + const char* suffix; + int compressionLevel; +} FIO_rust_compress_multiple_separate_context_t; + +static int FIO_rust_compressMultipleSeparateFileCallback(void* opaque, + const char* srcFileName) +{ + FIO_rust_compress_multiple_separate_context_t* const context = + (FIO_rust_compress_multiple_separate_context_t*)opaque; + const char* dstFileName; + + if (context->outMirroredRootDirName) { + char* const validMirroredDirName = UTIL_createMirroredDestDirName( + srcFileName, context->outMirroredRootDirName); + if (validMirroredDirName) { + dstFileName = FIO_determineCompressedName( + srcFileName, validMirroredDirName, context->suffix); + free(validMirroredDirName); + } else { + DISPLAYLEVEL(2, "zstd: --output-dir-mirror cannot compress '%s' into '%s' \n", + srcFileName, context->outMirroredRootDirName); + return 1; + } + } else { + dstFileName = FIO_determineCompressedName( + srcFileName, context->outDirName, context->suffix); + } + + return FIO_compressFilename_srcFile( + context->fCtx, context->prefs, *context->ress, + dstFileName, srcFileName, context->compressionLevel); +} + /* FIO_compressMultipleFilenames() : * compress nbFiles files * into either one destination (outFileName), @@ -2959,7 +3023,6 @@ int FIO_compressMultipleFilenames(FIO_ctx_t* const fCtx, const char* dictFileName, int compressionLevel, ZSTD_compressionParameters comprParams) { - int status; int error = 0; cRess_t ress = FIO_createCResources(prefs, dictFileName, FIO_getLargestFileSize(inFileNamesTable, (unsigned)fCtx->nbFilesTotal), @@ -2991,29 +3054,18 @@ int FIO_compressMultipleFilenames(FIO_ctx_t* const fCtx, strerror(errno), outFileName); } } else { + FIO_rust_compress_multiple_separate_context_t callbackContext = { + fCtx, prefs, &ress, outMirroredRootDirName, + outDirName, suffix, compressionLevel + }; + FIO_rust_compress_multiple_separate_projection_t projection = { + fCtx, inFileNamesTable, &callbackContext, + FIO_rust_compressMultipleSeparateFileCallback + }; if (outMirroredRootDirName) UTIL_mirrorSourceFilesDirectories(inFileNamesTable, (unsigned)fCtx->nbFilesTotal, outMirroredRootDirName); - for (; fCtx->currFileIdx < fCtx->nbFilesTotal; ++fCtx->currFileIdx) { - const char* const srcFileName = inFileNamesTable[fCtx->currFileIdx]; - const char* dstFileName = NULL; - if (outMirroredRootDirName) { - char* validMirroredDirName = UTIL_createMirroredDestDirName(srcFileName, outMirroredRootDirName); - if (validMirroredDirName) { - dstFileName = FIO_determineCompressedName(srcFileName, validMirroredDirName, suffix); - free(validMirroredDirName); - } else { - DISPLAYLEVEL(2, "zstd: --output-dir-mirror cannot compress '%s' into '%s' \n", srcFileName, outMirroredRootDirName); - error=1; - continue; - } - } else { - dstFileName = FIO_determineCompressedName(srcFileName, outDirName, suffix); /* cannot fail */ - } - status = FIO_compressFilename_srcFile(fCtx, prefs, ress, dstFileName, srcFileName, compressionLevel); - if (!status) fCtx->nbFilesProcessed++; - error |= status; - } + error = FIO_rust_compressMultipleSeparateFilenames(&projection); if (outDirName) FIO_checkFilenameCollisions(inFileNamesTable , (unsigned)fCtx->nbFilesTotal); diff --git a/rust/src/fileio_asyncio.rs b/rust/src/fileio_asyncio.rs index 4e90bef23..f4071bfe2 100644 --- a/rust/src/fileio_asyncio.rs +++ b/rust/src/fileio_asyncio.rs @@ -261,6 +261,8 @@ pub struct FIO_rust_compress_callbacks_t { pub type FIO_rust_compress_multiple_file_fn = unsafe extern "C" fn(*mut c_void, *const c_char, *const c_char) -> c_int; +pub type FIO_rust_compress_multiple_separate_file_fn = + unsafe extern "C" fn(*mut c_void, *const c_char) -> c_int; /// Rust owns only the shared-destination file iteration and aggregate error /// handling. C retains the private preferences/resources and supplies one @@ -295,6 +297,41 @@ const _: () = { assert!(size_of::() == 5 * size_of::()); }; +/// Rust owns only the separate-destination file iteration and aggregate error +/// handling. C retains destination-name construction, private +/// preferences/resources, diagnostics, and compression dispatch behind the +/// per-source callback. +#[repr(C)] +pub struct FIO_rust_compress_multiple_separate_projection_t { + pub f_ctx: *mut c_void, + pub input_file_names: *const *const c_char, + pub opaque: *mut c_void, + pub compress_file: Option, +} +const _: () = { + assert!(std::mem::offset_of!(FIO_rust_compress_multiple_separate_projection_t, f_ctx) == 0); + assert!( + std::mem::offset_of!( + FIO_rust_compress_multiple_separate_projection_t, + input_file_names + ) == size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_compress_multiple_separate_projection_t, opaque) + == 2 * size_of::() + ); + assert!( + std::mem::offset_of!( + FIO_rust_compress_multiple_separate_projection_t, + compress_file + ) == 3 * size_of::() + ); + assert!(size_of::() == size_of::()); + assert!( + size_of::() == 4 * size_of::() + ); +}; + pub const FIO_RUST_ZSTD_OK: c_int = 0; pub const FIO_RUST_ZSTD_COMPRESS_ERROR: c_int = 1; pub const FIO_RUST_ZSTD_INCOMPLETE_INPUT: c_int = 2; @@ -2200,6 +2237,41 @@ pub unsafe extern "C" fn FIO_rust_compressMultipleFilenames( error } +/// Iterate files that each receive a separate destination. The C callback +/// retains destination-name construction, resource state, diagnostics, and +/// compression dispatch; Rust preserves the original non-short-circuiting +/// loop and its counter updates. +#[no_mangle] +pub unsafe extern "C" fn FIO_rust_compressMultipleSeparateFilenames( + projection: *const FIO_rust_compress_multiple_separate_projection_t, +) -> c_int { + assert!(!projection.is_null()); + let projection = unsafe { &*projection }; + assert!(!projection.f_ctx.is_null()); + assert!(!projection.input_file_names.is_null()); + let compress_file = projection + .compress_file + .expect("separate-destination compression callback is required"); + + let f_ctx = projection.f_ctx.cast::(); + let mut error = 0; + while unsafe { (*f_ctx).currFileIdx < (*f_ctx).nbFilesTotal } { + let file_index = unsafe { (*f_ctx).currFileIdx as usize }; + let src_file_name = unsafe { *projection.input_file_names.add(file_index) }; + let status = unsafe { compress_file(projection.opaque, src_file_name) }; + if status == 0 { + unsafe { + (*f_ctx).nbFilesProcessed = (*f_ctx).nbFilesProcessed.wrapping_add(1); + } + } + error |= status; + unsafe { + (*f_ctx).currFileIdx = (*f_ctx).currFileIdx.wrapping_add(1); + } + } + error +} + #[inline] fn zstd_adapt_level(projection: &FIO_rust_zstd_adapt_projection_t, slower: bool) -> c_int { if slower { @@ -4487,6 +4559,84 @@ mod tests { status } + unsafe extern "C" fn record_multiple_separate_compression_file( + opaque: *mut c_void, + source: *const c_char, + ) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + let status = state + .statuses + .get(state.sources.len()) + .copied() + .unwrap_or(0); + state.sources.push(source); + status + } + + #[test] + fn compression_multiple_separate_destinations_preserves_order_and_counters() { + let source_names = [c"zero".as_ptr(), c"one".as_ptr(), c"two".as_ptr()]; + let mut context = FIO_rust_compression_context_t { + nbFilesTotal: 3, + hasStdinInput: 0, + hasStdoutOutput: 0, + currFileIdx: 1, + nbFilesProcessed: 5, + totalBytesInput: 0, + totalBytesOutput: 0, + }; + let mut state = MultipleCompressionState { + statuses: vec![0, 0], + ..MultipleCompressionState::default() + }; + let projection = FIO_rust_compress_multiple_separate_projection_t { + f_ctx: (&mut context as *mut FIO_rust_compression_context_t).cast(), + input_file_names: source_names.as_ptr(), + opaque: (&mut state as *mut MultipleCompressionState).cast(), + compress_file: Some(record_multiple_separate_compression_file), + }; + + assert_eq!( + unsafe { FIO_rust_compressMultipleSeparateFilenames(&projection) }, + 0 + ); + assert_eq!(state.sources, vec![source_names[1], source_names[2]]); + assert_eq!(context.currFileIdx, 3); + assert_eq!(context.nbFilesProcessed, 7); + } + + #[test] + fn compression_multiple_separate_destinations_accumulates_errors_without_short_circuiting() { + let source_names = [c"one".as_ptr(), c"two".as_ptr(), c"three".as_ptr()]; + let mut context = FIO_rust_compression_context_t { + nbFilesTotal: 3, + hasStdinInput: 0, + hasStdoutOutput: 0, + currFileIdx: 0, + nbFilesProcessed: 4, + totalBytesInput: 0, + totalBytesOutput: 0, + }; + let mut state = MultipleCompressionState { + statuses: vec![1, 4, 2], + ..MultipleCompressionState::default() + }; + let projection = FIO_rust_compress_multiple_separate_projection_t { + f_ctx: (&mut context as *mut FIO_rust_compression_context_t).cast(), + input_file_names: source_names.as_ptr(), + opaque: (&mut state as *mut MultipleCompressionState).cast(), + compress_file: Some(record_multiple_separate_compression_file), + }; + + assert_eq!( + unsafe { FIO_rust_compressMultipleSeparateFilenames(&projection) }, + 1 | 4 | 2 + ); + assert_eq!(state.sources, source_names); + assert_eq!(context.currFileIdx, 3); + assert_eq!(context.nbFilesProcessed, 4); + } + #[test] fn compression_multiple_shared_destination_preserves_order_and_counters() { let source_names = [c"zero".as_ptr(), c"one".as_ptr(), c"two".as_ptr()];