feat(cli): move compression format dispatch into Rust
`FIO_compressFilename_internal()` still selected the zstd, gzip, xz/lzma, and lz4 paths in a C-owned format switch and updated aggregate byte accounting around those callbacks. That left the mixed-format CLI dispatch loop outside the Rust file-I/O orchestration even though the pools and decompression path were already projected there. Add a callback table for the codec leaves and move format selection, optional codec status handling, post-success accounting, and progress-hook ordering to Rust. The existing C codec implementations, diagnostics, elapsed-time formatting, and resource ownership remain unchanged; unsupported optional formats return explicit statuses so the C wrapper preserves its original EXM_THROW messages. Test Plan: - `cargo test --manifest-path rust/Cargo.toml --all-targets -- --test-threads=1` -- 473 passed. - Native `test-cli-tests` -- all 41 passed. - Native `test-zstd` -- passed, including async zstd/gzip/xz/lzma/lz4 round trips. - `make -B -C programs -j2 zstd` and `make -B -C tests/fuzz -j2 all` -- passed.
This commit is contained in:
@@ -75,6 +75,52 @@ pub struct FIO_rust_decompress_callbacks_t {
|
||||
pub pass_through: Option<FIO_rust_pass_through_fn>,
|
||||
}
|
||||
|
||||
pub const FIO_RUST_COMPRESS_OK: c_int = 0;
|
||||
pub const FIO_RUST_COMPRESS_GZIP_UNSUPPORTED: c_int = 1;
|
||||
pub const FIO_RUST_COMPRESS_LZMA_UNSUPPORTED: c_int = 2;
|
||||
pub const FIO_RUST_COMPRESS_LZ4_UNSUPPORTED: c_int = 3;
|
||||
pub const FIO_RUST_COMPRESS_ZSTD_UNSUPPORTED: c_int = 4;
|
||||
|
||||
const FIO_ZSTD_COMPRESSION: c_int = 0;
|
||||
const FIO_GZIP_COMPRESSION: c_int = 1;
|
||||
const FIO_XZ_COMPRESSION: c_int = 2;
|
||||
const FIO_LZMA_COMPRESSION: c_int = 3;
|
||||
const FIO_LZ4_COMPRESSION: c_int = 4;
|
||||
|
||||
pub type FIO_rust_compress_zstd_fn = unsafe extern "C" fn(
|
||||
*mut c_void,
|
||||
*mut c_void,
|
||||
*mut c_void,
|
||||
*const c_char,
|
||||
u64,
|
||||
c_int,
|
||||
*mut u64,
|
||||
) -> u64;
|
||||
pub type FIO_rust_compress_gzip_fn =
|
||||
unsafe extern "C" fn(*mut c_void, *const c_char, u64, c_int, *mut u64) -> u64;
|
||||
pub type FIO_rust_compress_lzma_fn =
|
||||
unsafe extern "C" fn(*mut c_void, *const c_char, u64, c_int, *mut u64, c_int) -> u64;
|
||||
pub type FIO_rust_compress_lz4_fn =
|
||||
unsafe extern "C" fn(*mut c_void, *const c_char, u64, c_int, c_int, *mut u64) -> u64;
|
||||
pub type FIO_rust_compress_input_display_fn = unsafe extern "C" fn(*mut c_void, *const c_char, u64);
|
||||
pub type FIO_rust_compress_status_display_fn =
|
||||
unsafe extern "C" fn(*mut c_void, *mut c_void, *const c_char, *const c_char, u64, u64);
|
||||
|
||||
/// C-owned compression codecs and CLI display hooks used by the Rust file
|
||||
/// compression selector. The codec callbacks deliberately receive opaque
|
||||
/// resource pointers: Rust owns the selection/accounting loop but never
|
||||
/// assumes the private `cRess_t` layout.
|
||||
#[repr(C)]
|
||||
pub struct FIO_rust_compress_callbacks_t {
|
||||
pub opaque: *mut c_void,
|
||||
pub compress_zstd: Option<FIO_rust_compress_zstd_fn>,
|
||||
pub compress_gzip: Option<FIO_rust_compress_gzip_fn>,
|
||||
pub compress_lzma: Option<FIO_rust_compress_lzma_fn>,
|
||||
pub compress_lz4: Option<FIO_rust_compress_lz4_fn>,
|
||||
pub display_input: Option<FIO_rust_compress_input_display_fn>,
|
||||
pub display_status: Option<FIO_rust_compress_status_display_fn>,
|
||||
}
|
||||
|
||||
/// C's `FIO_prefs_t` from `programs/fileio_types.h`.
|
||||
///
|
||||
/// `fileio_prefs.rs` contains the same C layout for the preferences API. It
|
||||
@@ -145,6 +191,20 @@ pub struct IOJob_t {
|
||||
pub offset: u64,
|
||||
}
|
||||
|
||||
/// The compression selector updates only these public counters in C's
|
||||
/// private `FIO_ctx_s`. The field order is kept explicit so no C-owned
|
||||
/// context implementation details cross the callback boundary.
|
||||
#[repr(C)]
|
||||
struct FIO_rust_compression_context_t {
|
||||
nbFilesTotal: c_int,
|
||||
hasStdinInput: c_int,
|
||||
hasStdoutOutput: c_int,
|
||||
currFileIdx: c_int,
|
||||
nbFilesProcessed: c_int,
|
||||
totalBytesInput: usize,
|
||||
totalBytesOutput: usize,
|
||||
}
|
||||
|
||||
type PoolFunction = unsafe extern "C" fn(*mut c_void);
|
||||
|
||||
#[cfg(not(test))]
|
||||
@@ -1437,6 +1497,173 @@ pub unsafe extern "C" fn FIO_rust_passThrough(
|
||||
0
|
||||
}
|
||||
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
enum CompressionFormat {
|
||||
Zstd,
|
||||
Gzip,
|
||||
Lzma { plain_lzma: c_int },
|
||||
Lz4 { checksum: c_int },
|
||||
}
|
||||
|
||||
fn select_compression_format(compression_type: c_int, checksum: c_int) -> CompressionFormat {
|
||||
match compression_type {
|
||||
FIO_GZIP_COMPRESSION => CompressionFormat::Gzip,
|
||||
FIO_XZ_COMPRESSION => CompressionFormat::Lzma { plain_lzma: 0 },
|
||||
FIO_LZMA_COMPRESSION => CompressionFormat::Lzma { plain_lzma: 1 },
|
||||
FIO_LZ4_COMPRESSION => CompressionFormat::Lz4 { checksum },
|
||||
FIO_ZSTD_COMPRESSION => CompressionFormat::Zstd,
|
||||
_ => CompressionFormat::Zstd,
|
||||
}
|
||||
}
|
||||
|
||||
unsafe fn run_compression_callback(
|
||||
f_ctx: *mut c_void,
|
||||
prefs: *mut FIO_prefs_t,
|
||||
ress: *mut c_void,
|
||||
src_file_name: *const c_char,
|
||||
file_size: u64,
|
||||
compression_level: c_int,
|
||||
callbacks: &FIO_rust_compress_callbacks_t,
|
||||
) -> Result<(u64, u64), c_int> {
|
||||
let format = select_compression_format(unsafe { (*prefs).compressionType }, unsafe {
|
||||
(*prefs).checksumFlag
|
||||
});
|
||||
let mut read_size = 0_u64;
|
||||
|
||||
let compressed_size = match format {
|
||||
CompressionFormat::Zstd => {
|
||||
let Some(callback) = callbacks.compress_zstd else {
|
||||
return Err(FIO_RUST_COMPRESS_ZSTD_UNSUPPORTED);
|
||||
};
|
||||
unsafe {
|
||||
callback(
|
||||
f_ctx,
|
||||
prefs.cast(),
|
||||
ress,
|
||||
src_file_name,
|
||||
file_size,
|
||||
compression_level,
|
||||
&mut read_size,
|
||||
)
|
||||
}
|
||||
}
|
||||
CompressionFormat::Gzip => {
|
||||
let Some(callback) = callbacks.compress_gzip else {
|
||||
return Err(FIO_RUST_COMPRESS_GZIP_UNSUPPORTED);
|
||||
};
|
||||
unsafe {
|
||||
callback(
|
||||
ress,
|
||||
src_file_name,
|
||||
file_size,
|
||||
compression_level,
|
||||
&mut read_size,
|
||||
)
|
||||
}
|
||||
}
|
||||
CompressionFormat::Lzma { plain_lzma } => {
|
||||
let Some(callback) = callbacks.compress_lzma else {
|
||||
return Err(FIO_RUST_COMPRESS_LZMA_UNSUPPORTED);
|
||||
};
|
||||
unsafe {
|
||||
callback(
|
||||
ress,
|
||||
src_file_name,
|
||||
file_size,
|
||||
compression_level,
|
||||
&mut read_size,
|
||||
plain_lzma,
|
||||
)
|
||||
}
|
||||
}
|
||||
CompressionFormat::Lz4 { checksum } => {
|
||||
let Some(callback) = callbacks.compress_lz4 else {
|
||||
return Err(FIO_RUST_COMPRESS_LZ4_UNSUPPORTED);
|
||||
};
|
||||
unsafe {
|
||||
callback(
|
||||
ress,
|
||||
src_file_name,
|
||||
file_size,
|
||||
compression_level,
|
||||
checksum,
|
||||
&mut read_size,
|
||||
)
|
||||
}
|
||||
}
|
||||
};
|
||||
|
||||
Ok((read_size, compressed_size))
|
||||
}
|
||||
|
||||
/// Runs the compression-side format selection and final file accounting.
|
||||
///
|
||||
/// The codec implementations remain in C and are selected through the
|
||||
/// explicit callback table above. Rust updates the C-visible aggregate byte
|
||||
/// counters only after a codec returns, then asks C to perform the existing
|
||||
/// progress, summary, and elapsed-time display. Missing optional codecs are
|
||||
/// returned as status values so the C wrapper can preserve its exact
|
||||
/// `EXM_THROW()` diagnostics.
|
||||
#[no_mangle]
|
||||
pub unsafe extern "C" fn FIO_rust_compressFilenameInternal(
|
||||
f_ctx: *mut c_void,
|
||||
prefs: *mut FIO_prefs_t,
|
||||
ress: *mut c_void,
|
||||
dst_file_name: *const c_char,
|
||||
src_file_name: *const c_char,
|
||||
file_size: u64,
|
||||
compression_level: c_int,
|
||||
callbacks: *const FIO_rust_compress_callbacks_t,
|
||||
) -> c_int {
|
||||
assert!(!f_ctx.is_null());
|
||||
assert!(!prefs.is_null());
|
||||
assert!(!ress.is_null());
|
||||
assert!(!dst_file_name.is_null());
|
||||
assert!(!src_file_name.is_null());
|
||||
assert!(!callbacks.is_null());
|
||||
|
||||
let callbacks = unsafe { &*callbacks };
|
||||
if let Some(display) = callbacks.display_input {
|
||||
unsafe { display(callbacks.opaque, src_file_name, file_size) };
|
||||
}
|
||||
|
||||
let (read_size, compressed_size) = match unsafe {
|
||||
run_compression_callback(
|
||||
f_ctx,
|
||||
prefs,
|
||||
ress,
|
||||
src_file_name,
|
||||
file_size,
|
||||
compression_level,
|
||||
callbacks,
|
||||
)
|
||||
} {
|
||||
Ok(sizes) => sizes,
|
||||
Err(status) => return status,
|
||||
};
|
||||
|
||||
let context = unsafe { &mut *f_ctx.cast::<FIO_rust_compression_context_t>() };
|
||||
context.totalBytesInput = context.totalBytesInput.wrapping_add(read_size as usize);
|
||||
context.totalBytesOutput = context
|
||||
.totalBytesOutput
|
||||
.wrapping_add(compressed_size as usize);
|
||||
|
||||
if let Some(display) = callbacks.display_status {
|
||||
unsafe {
|
||||
display(
|
||||
callbacks.opaque,
|
||||
f_ctx,
|
||||
dst_file_name,
|
||||
src_file_name,
|
||||
read_size,
|
||||
compressed_size,
|
||||
)
|
||||
};
|
||||
}
|
||||
|
||||
FIO_RUST_COMPRESS_OK
|
||||
}
|
||||
|
||||
/// Decompresses exactly one zstd frame using the existing asynchronous pools.
|
||||
///
|
||||
/// C retains the frame dispatcher and all user-facing diagnostics. This ABI
|
||||
@@ -1896,6 +2123,269 @@ mod tests {
|
||||
prefs
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct CompressionCallbackState {
|
||||
zstd_calls: usize,
|
||||
gzip_calls: usize,
|
||||
lzma_calls: usize,
|
||||
lz4_calls: usize,
|
||||
lzma_plain: c_int,
|
||||
lz4_checksum: c_int,
|
||||
}
|
||||
|
||||
unsafe extern "C" fn record_compress_zstd(
|
||||
_f_ctx: *mut c_void,
|
||||
_prefs: *mut c_void,
|
||||
ress: *mut c_void,
|
||||
_src_file_name: *const c_char,
|
||||
_file_size: u64,
|
||||
_compression_level: c_int,
|
||||
read_size: *mut u64,
|
||||
) -> u64 {
|
||||
let state = unsafe { &mut *ress.cast::<CompressionCallbackState>() };
|
||||
state.zstd_calls += 1;
|
||||
unsafe { *read_size = 11 };
|
||||
7
|
||||
}
|
||||
|
||||
unsafe extern "C" fn record_compress_gzip(
|
||||
ress: *mut c_void,
|
||||
_src_file_name: *const c_char,
|
||||
_file_size: u64,
|
||||
_compression_level: c_int,
|
||||
read_size: *mut u64,
|
||||
) -> u64 {
|
||||
let state = unsafe { &mut *ress.cast::<CompressionCallbackState>() };
|
||||
state.gzip_calls += 1;
|
||||
unsafe { *read_size = 12 };
|
||||
8
|
||||
}
|
||||
|
||||
unsafe extern "C" fn record_compress_lzma(
|
||||
ress: *mut c_void,
|
||||
_src_file_name: *const c_char,
|
||||
_file_size: u64,
|
||||
_compression_level: c_int,
|
||||
read_size: *mut u64,
|
||||
plain_lzma: c_int,
|
||||
) -> u64 {
|
||||
let state = unsafe { &mut *ress.cast::<CompressionCallbackState>() };
|
||||
state.lzma_calls += 1;
|
||||
state.lzma_plain = plain_lzma;
|
||||
unsafe { *read_size = 13 };
|
||||
9
|
||||
}
|
||||
|
||||
unsafe extern "C" fn record_compress_lz4(
|
||||
ress: *mut c_void,
|
||||
_src_file_name: *const c_char,
|
||||
_file_size: u64,
|
||||
_compression_level: c_int,
|
||||
checksum: c_int,
|
||||
read_size: *mut u64,
|
||||
) -> u64 {
|
||||
let state = unsafe { &mut *ress.cast::<CompressionCallbackState>() };
|
||||
state.lz4_calls += 1;
|
||||
state.lz4_checksum = checksum;
|
||||
unsafe { *read_size = 14 };
|
||||
10
|
||||
}
|
||||
|
||||
fn compression_test_callbacks() -> FIO_rust_compress_callbacks_t {
|
||||
FIO_rust_compress_callbacks_t {
|
||||
opaque: ptr::null_mut(),
|
||||
compress_zstd: Some(record_compress_zstd),
|
||||
compress_gzip: Some(record_compress_gzip),
|
||||
compress_lzma: Some(record_compress_lzma),
|
||||
compress_lz4: Some(record_compress_lz4),
|
||||
display_input: None,
|
||||
display_status: None,
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct CompressionDisplayState {
|
||||
input_sizes: Vec<u64>,
|
||||
status_sizes: Vec<(u64, u64)>,
|
||||
}
|
||||
|
||||
unsafe extern "C" fn record_compress_input(
|
||||
opaque: *mut c_void,
|
||||
_src_file_name: *const c_char,
|
||||
file_size: u64,
|
||||
) {
|
||||
let state = unsafe { &mut *opaque.cast::<CompressionDisplayState>() };
|
||||
state.input_sizes.push(file_size);
|
||||
}
|
||||
|
||||
unsafe extern "C" fn record_compress_status(
|
||||
opaque: *mut c_void,
|
||||
_f_ctx: *mut c_void,
|
||||
_dst_file_name: *const c_char,
|
||||
_src_file_name: *const c_char,
|
||||
read_size: u64,
|
||||
compressed_size: u64,
|
||||
) {
|
||||
let state = unsafe { &mut *opaque.cast::<CompressionDisplayState>() };
|
||||
state.status_sizes.push((read_size, compressed_size));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compression_selection_dispatches_each_format_and_forwards_options() {
|
||||
let source = c"source";
|
||||
let destination = c"destination";
|
||||
let cases = [
|
||||
(99, 1, 0_usize, 0, 11_u64, 7_u64),
|
||||
(FIO_ZSTD_COMPRESSION, 1, 0, 0, 11, 7),
|
||||
(FIO_GZIP_COMPRESSION, 1, 1, 0, 12, 8),
|
||||
(FIO_XZ_COMPRESSION, 1, 2, 0, 13, 9),
|
||||
(FIO_LZMA_COMPRESSION, 1, 2, 1, 13, 9),
|
||||
(FIO_LZ4_COMPRESSION, 2, 3, 0, 14, 10),
|
||||
];
|
||||
|
||||
for (
|
||||
compression_type,
|
||||
checksum,
|
||||
expected_codec,
|
||||
expected_plain,
|
||||
expected_read,
|
||||
expected_output,
|
||||
) in cases
|
||||
{
|
||||
let mut prefs: FIO_prefs_t = unsafe { std::mem::zeroed() };
|
||||
prefs.compressionType = compression_type;
|
||||
prefs.checksumFlag = checksum;
|
||||
let mut context = FIO_rust_compression_context_t {
|
||||
nbFilesTotal: 1,
|
||||
hasStdinInput: 0,
|
||||
hasStdoutOutput: 0,
|
||||
currFileIdx: 0,
|
||||
nbFilesProcessed: 0,
|
||||
totalBytesInput: 100,
|
||||
totalBytesOutput: 200,
|
||||
};
|
||||
let callbacks = compression_test_callbacks();
|
||||
let mut state = CompressionCallbackState::default();
|
||||
|
||||
assert_eq!(
|
||||
unsafe {
|
||||
FIO_rust_compressFilenameInternal(
|
||||
(&mut context as *mut FIO_rust_compression_context_t).cast(),
|
||||
&mut prefs,
|
||||
(&mut state as *mut CompressionCallbackState).cast(),
|
||||
destination.as_ptr(),
|
||||
source.as_ptr(),
|
||||
123,
|
||||
5,
|
||||
&callbacks,
|
||||
)
|
||||
},
|
||||
FIO_RUST_COMPRESS_OK
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
[
|
||||
state.zstd_calls,
|
||||
state.gzip_calls,
|
||||
state.lzma_calls,
|
||||
state.lz4_calls,
|
||||
][expected_codec],
|
||||
1
|
||||
);
|
||||
assert_eq!(
|
||||
state.zstd_calls + state.gzip_calls + state.lzma_calls + state.lz4_calls,
|
||||
1
|
||||
);
|
||||
assert_eq!(state.lzma_plain, expected_plain);
|
||||
assert_eq!(
|
||||
state.lz4_checksum,
|
||||
if expected_codec == 3 { checksum } else { 0 }
|
||||
);
|
||||
assert_eq!(context.totalBytesInput, 100 + expected_read as usize);
|
||||
assert_eq!(context.totalBytesOutput, 200 + expected_output as usize);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compression_accounting_calls_display_hooks_after_codec_success() {
|
||||
let source = c"source";
|
||||
let destination = c"destination";
|
||||
let mut prefs: FIO_prefs_t = unsafe { std::mem::zeroed() };
|
||||
prefs.compressionType = FIO_ZSTD_COMPRESSION;
|
||||
let mut context = FIO_rust_compression_context_t {
|
||||
nbFilesTotal: 1,
|
||||
hasStdinInput: 0,
|
||||
hasStdoutOutput: 0,
|
||||
currFileIdx: 0,
|
||||
nbFilesProcessed: 0,
|
||||
totalBytesInput: 0,
|
||||
totalBytesOutput: 0,
|
||||
};
|
||||
let mut codec_state = CompressionCallbackState::default();
|
||||
let mut display_state = CompressionDisplayState::default();
|
||||
let mut callbacks = compression_test_callbacks();
|
||||
callbacks.opaque = (&mut display_state as *mut CompressionDisplayState).cast();
|
||||
callbacks.display_input = Some(record_compress_input);
|
||||
callbacks.display_status = Some(record_compress_status);
|
||||
|
||||
assert_eq!(
|
||||
unsafe {
|
||||
FIO_rust_compressFilenameInternal(
|
||||
(&mut context as *mut FIO_rust_compression_context_t).cast(),
|
||||
&mut prefs,
|
||||
(&mut codec_state as *mut CompressionCallbackState).cast(),
|
||||
destination.as_ptr(),
|
||||
source.as_ptr(),
|
||||
123,
|
||||
5,
|
||||
&callbacks,
|
||||
)
|
||||
},
|
||||
FIO_RUST_COMPRESS_OK
|
||||
);
|
||||
assert_eq!(display_state.input_sizes, vec![123]);
|
||||
assert_eq!(display_state.status_sizes, vec![(11, 7)]);
|
||||
assert_eq!(context.totalBytesInput, 11);
|
||||
assert_eq!(context.totalBytesOutput, 7);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn compression_missing_optional_codec_leaves_accounting_unchanged() {
|
||||
let source = c"source";
|
||||
let destination = c"destination";
|
||||
let mut prefs: FIO_prefs_t = unsafe { std::mem::zeroed() };
|
||||
prefs.compressionType = FIO_GZIP_COMPRESSION;
|
||||
let mut context = FIO_rust_compression_context_t {
|
||||
nbFilesTotal: 1,
|
||||
hasStdinInput: 0,
|
||||
hasStdoutOutput: 0,
|
||||
currFileIdx: 0,
|
||||
nbFilesProcessed: 0,
|
||||
totalBytesInput: 31,
|
||||
totalBytesOutput: 47,
|
||||
};
|
||||
let mut callbacks = compression_test_callbacks();
|
||||
callbacks.compress_gzip = None;
|
||||
|
||||
assert_eq!(
|
||||
unsafe {
|
||||
FIO_rust_compressFilenameInternal(
|
||||
(&mut context as *mut FIO_rust_compression_context_t).cast(),
|
||||
&mut prefs,
|
||||
ptr::dangling_mut::<c_void>(),
|
||||
destination.as_ptr(),
|
||||
source.as_ptr(),
|
||||
123,
|
||||
5,
|
||||
&callbacks,
|
||||
)
|
||||
},
|
||||
FIO_RUST_COMPRESS_GZIP_UNSUPPORTED
|
||||
);
|
||||
assert_eq!(context.totalBytesInput, 31);
|
||||
assert_eq!(context.totalBytesOutput, 47);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn classifies_all_cli_decompression_headers() {
|
||||
let is_mock_zstd = |buffer: &[u8]| buffer.starts_with(&[0x28, 0xB5, 0x2F, 0xFD]);
|
||||
|
||||
Reference in New Issue
Block a user