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
This commit is contained in:
2026-07-19 10:52:31 +02:00
parent bd7fbc81fc
commit fbc80de7dc
2 changed files with 223 additions and 21 deletions
+150
View File
@@ -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::<FIO_rust_compress_multiple_projection_t>() == 5 * size_of::<usize>());
};
/// 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<FIO_rust_compress_multiple_separate_file_fn>,
}
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::<usize>()
);
assert!(
std::mem::offset_of!(FIO_rust_compress_multiple_separate_projection_t, opaque)
== 2 * size_of::<usize>()
);
assert!(
std::mem::offset_of!(
FIO_rust_compress_multiple_separate_projection_t,
compress_file
) == 3 * size_of::<usize>()
);
assert!(size_of::<FIO_rust_compress_multiple_separate_file_fn>() == size_of::<usize>());
assert!(
size_of::<FIO_rust_compress_multiple_separate_projection_t>() == 4 * size_of::<usize>()
);
};
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::<FIO_rust_compression_context_t>();
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::<MultipleCompressionState>() };
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()];