From 2442fddc388f72e2de96194570aa1f30354daace Mon Sep 17 00:00:00 2001 From: ddidderr Date: Mon, 20 Jul 2026 08:41:31 +0200 Subject: [PATCH] refactor(fileio): move separate-output routing policy to Rust Move the separate-destination mode decision into the Rust file-iteration boundary. The projection now carries the mirror-mode bit and one callback for each destination policy, so Rust selects the callback once before preserving the existing non-short-circuiting iteration and aggregate error behavior. Keep path construction, mirror-directory diagnostics, resource ownership, and the per-file compression leaf in the opaque C callbacks. Add C/Rust layout assertions and focused tests proving mirror/flat routing and ordered error aggregation. Test Plan: - `ulimit -v 41943040; cargo +nightly fmt --manifest-path rust/Cargo.toml --all -- --check` - `ulimit -v 41943040; CARGO_BUILD_JOBS=1 cargo clippy --manifest-path rust/Cargo.toml --all-targets -- -D warnings` - `ulimit -v 41943040; CARGO_BUILD_JOBS=1 cargo test --manifest-path rust/Cargo.toml --all-targets` (793 tests) - `ulimit -v 41943040; CARGO_BUILD_JOBS=1 cargo clippy --manifest-path rust/cli/Cargo.toml --all-targets -- -D warnings` - `ulimit -v 41943040; CARGO_BUILD_JOBS=1 cargo test --manifest-path rust/cli/Cargo.toml --all-targets` (184 tests) - `ulimit -v 41943040; make -j1` - `ulimit -v 41943040; make -j1 -C tests test` --- programs/fileio.c | 68 ++++++++----- rust/src/fileio_asyncio.rs | 193 ++++++++++++++++++++++++++++++++++--- 2 files changed, 221 insertions(+), 40 deletions(-) diff --git a/programs/fileio.c b/programs/fileio.c index de3e5d24c..c5ee9e5a7 100644 --- a/programs/fileio.c +++ b/programs/fileio.c @@ -1166,11 +1166,15 @@ int FIO_rust_compressMultipleFilenames( typedef int (*FIO_rust_compress_multiple_separate_file_fn)( void* opaque, const char* srcFileName); +/* Rust selects the destination-mode callback; both callbacks retain their + * private resources, path construction, diagnostics, and compression leaf. */ typedef struct { void* fCtx; const char** inFileNamesTable; void* opaque; - FIO_rust_compress_multiple_separate_file_fn compressFile; + int mirrorOutput; + FIO_rust_compress_multiple_separate_file_fn compressMirroredFile; + FIO_rust_compress_multiple_separate_file_fn compressFlatFile; } 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]; @@ -1180,12 +1184,22 @@ typedef char FIO_rust_compress_multiple_separate_input_names_offset[ 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) +typedef char FIO_rust_compress_multiple_separate_mode_offset[ + (offsetof(FIO_rust_compress_multiple_separate_projection_t, mirrorOutput) == 3 * sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_compress_multiple_separate_mirrored_callback_offset[ + (offsetof(FIO_rust_compress_multiple_separate_projection_t, compressMirroredFile) + == 4 * sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_compress_multiple_separate_flat_callback_offset[ + (offsetof(FIO_rust_compress_multiple_separate_projection_t, compressFlatFile) + == 5 * sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_compress_multiple_separate_int_size[ + (sizeof(int) <= sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_compress_multiple_separate_callback_size[ + (sizeof(FIO_rust_compress_multiple_separate_file_fn) == 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]; + == 6 * sizeof(void*)) ? 1 : -1]; int FIO_rust_compressMultipleSeparateFilenames( const FIO_rust_compress_multiple_separate_projection_t* projection); @@ -3517,30 +3531,34 @@ typedef struct { int compressionLevel; } FIO_rust_compress_multiple_separate_context_t; -static int FIO_rust_compressMultipleSeparateFileCallback(void* opaque, - const char* srcFileName) +static int FIO_rust_compressMultipleSeparateMirroredFileCallback(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); + char* const validMirroredDirName = UTIL_createMirroredDestDirName( + srcFileName, context->outMirroredRootDirName); + if (validMirroredDirName) { + const char* const dstFileName = FIO_determineCompressedName( + srcFileName, validMirroredDirName, context->suffix); + free(validMirroredDirName); + return FIO_compressFilename_srcFile( + context->fCtx, context->prefs, *context->ress, + dstFileName, srcFileName, context->compressionLevel); } + DISPLAYLEVEL(2, "zstd: --output-dir-mirror cannot compress '%s' into '%s' \n", + srcFileName, context->outMirroredRootDirName); + return 1; +} + +static int FIO_rust_compressMultipleSeparateFlatFileCallback(void* opaque, + const char* srcFileName) +{ + FIO_rust_compress_multiple_separate_context_t* const context = + (FIO_rust_compress_multiple_separate_context_t*)opaque; + const char* const dstFileName = FIO_determineCompressedName( + srcFileName, context->outDirName, context->suffix); return FIO_compressFilename_srcFile( context->fCtx, context->prefs, *context->ress, dstFileName, srcFileName, context->compressionLevel); @@ -3598,7 +3616,9 @@ int FIO_compressMultipleFilenames(FIO_ctx_t* const fCtx, }; FIO_rust_compress_multiple_separate_projection_t projection = { fCtx, inFileNamesTable, &callbackContext, - FIO_rust_compressMultipleSeparateFileCallback + outMirroredRootDirName != NULL, + FIO_rust_compressMultipleSeparateMirroredFileCallback, + FIO_rust_compressMultipleSeparateFlatFileCallback }; if (outMirroredRootDirName) UTIL_mirrorSourceFilesDirectories(inFileNamesTable, (unsigned)fCtx->nbFilesTotal, outMirroredRootDirName); diff --git a/rust/src/fileio_asyncio.rs b/rust/src/fileio_asyncio.rs index 09daf0e29..0ff6d9d9f 100644 --- a/rust/src/fileio_asyncio.rs +++ b/rust/src/fileio_asyncio.rs @@ -535,16 +535,18 @@ 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. +/// Rust owns separate-destination mode selection, file iteration, and +/// aggregate error handling. C retains destination-name construction, private +/// preferences/resources, diagnostics, and compression dispatch in the two +/// opaque per-source callbacks. #[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, + pub mirror_output: c_int, + pub compress_mirrored_file: Option, + pub compress_flat_file: Option, } const _: () = { assert!(std::mem::offset_of!(FIO_rust_compress_multiple_separate_projection_t, f_ctx) == 0); @@ -561,12 +563,25 @@ const _: () = { assert!( std::mem::offset_of!( FIO_rust_compress_multiple_separate_projection_t, - compress_file + mirror_output ) == 3 * size_of::() ); + assert!( + std::mem::offset_of!( + FIO_rust_compress_multiple_separate_projection_t, + compress_mirrored_file + ) == 4 * size_of::() + ); + assert!( + std::mem::offset_of!( + FIO_rust_compress_multiple_separate_projection_t, + compress_flat_file + ) == 5 * size_of::() + ); + assert!(size_of::() <= size_of::()); assert!(size_of::() == size_of::()); assert!( - size_of::() == 4 * size_of::() + size_of::() == 6 * size_of::() ); }; @@ -3114,10 +3129,10 @@ 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. +/// Select the destination-mode callback once, then iterate files that each +/// receive a separate destination. C 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, @@ -3126,9 +3141,15 @@ pub unsafe extern "C" fn FIO_rust_compressMultipleSeparateFilenames( 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 compress_file = if projection.mirror_output != 0 { + projection + .compress_mirrored_file + .expect("mirrored separate-destination compression callback is required") + } else { + projection + .compress_flat_file + .expect("flat separate-destination compression callback is required") + }; let f_ctx = projection.f_ctx.cast::(); let mut error = 0; @@ -6330,6 +6351,43 @@ mod tests { status } + #[derive(Default)] + struct SeparateCompressionRoutingState { + routes: Vec<&'static str>, + sources: Vec<*const c_char>, + statuses: Vec, + } + + fn record_separate_compression_route( + opaque: *mut c_void, + source: *const c_char, + route: &'static str, + ) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + let status = state + .statuses + .get(state.sources.len()) + .copied() + .unwrap_or(0); + state.routes.push(route); + state.sources.push(source); + status + } + + unsafe extern "C" fn record_mirrored_compression_file( + opaque: *mut c_void, + source: *const c_char, + ) -> c_int { + record_separate_compression_route(opaque, source, "mirror") + } + + unsafe extern "C" fn record_flat_compression_file( + opaque: *mut c_void, + source: *const c_char, + ) -> c_int { + record_separate_compression_route(opaque, source, "flat") + } + #[derive(Default)] struct MultipleDecompressionState { sources: Vec<*const c_char>, @@ -6367,6 +6425,105 @@ mod tests { status } + #[test] + fn compression_multiple_separate_destinations_selects_mirror_callback_only() { + let source_names = [c"one".as_ptr(), c"two".as_ptr()]; + let mut context = FIO_rust_compression_context_t { + nbFilesTotal: 2, + hasStdinInput: 0, + hasStdoutOutput: 0, + currFileIdx: 0, + nbFilesProcessed: 0, + totalBytesInput: 0, + totalBytesOutput: 0, + }; + let mut state = SeparateCompressionRoutingState::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 SeparateCompressionRoutingState).cast(), + mirror_output: 1, + compress_mirrored_file: Some(record_mirrored_compression_file), + compress_flat_file: Some(record_flat_compression_file), + }; + + assert_eq!( + unsafe { FIO_rust_compressMultipleSeparateFilenames(&projection) }, + 0 + ); + assert_eq!(state.routes, vec!["mirror", "mirror"]); + assert_eq!(state.sources, source_names); + assert_eq!(context.currFileIdx, 2); + assert_eq!(context.nbFilesProcessed, 2); + } + + #[test] + fn compression_multiple_separate_destinations_selects_flat_callback_only() { + let source_names = [c"one".as_ptr(), c"two".as_ptr()]; + let mut context = FIO_rust_compression_context_t { + nbFilesTotal: 2, + hasStdinInput: 0, + hasStdoutOutput: 0, + currFileIdx: 0, + nbFilesProcessed: 0, + totalBytesInput: 0, + totalBytesOutput: 0, + }; + let mut state = SeparateCompressionRoutingState::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 SeparateCompressionRoutingState).cast(), + mirror_output: 0, + compress_mirrored_file: Some(record_mirrored_compression_file), + compress_flat_file: Some(record_flat_compression_file), + }; + + assert_eq!( + unsafe { FIO_rust_compressMultipleSeparateFilenames(&projection) }, + 0 + ); + assert_eq!(state.routes, vec!["flat", "flat"]); + assert_eq!(state.sources, source_names); + assert_eq!(context.currFileIdx, 2); + assert_eq!(context.nbFilesProcessed, 2); + } + + #[test] + fn compression_multiple_separate_destinations_accumulates_selected_callback_errors_in_order() { + 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 = SeparateCompressionRoutingState { + statuses: vec![1, 4, 2], + ..SeparateCompressionRoutingState::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 SeparateCompressionRoutingState).cast(), + mirror_output: 1, + compress_mirrored_file: Some(record_mirrored_compression_file), + compress_flat_file: Some(record_flat_compression_file), + }; + + assert_eq!( + unsafe { FIO_rust_compressMultipleSeparateFilenames(&projection) }, + 1 | 4 | 2 + ); + assert_eq!(state.routes, vec!["mirror", "mirror", "mirror"]); + assert_eq!(state.sources, source_names); + assert_eq!(context.currFileIdx, 3); + assert_eq!(context.nbFilesProcessed, 4); + } + #[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()]; @@ -6387,7 +6544,9 @@ mod tests { 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), + mirror_output: 0, + compress_mirrored_file: Some(record_multiple_separate_compression_file), + compress_flat_file: Some(record_multiple_separate_compression_file), }; assert_eq!( @@ -6419,7 +6578,9 @@ mod tests { 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), + mirror_output: 0, + compress_mirrored_file: Some(record_multiple_separate_compression_file), + compress_flat_file: Some(record_multiple_separate_compression_file), }; assert_eq!(