feat(cli): move separate decompression scheduling into Rust
Move the separate-destination multi-file decompression loop into the Rust projection so Rust owns file iteration, progress counters, and aggregate error handling. Keep destination-name construction, mirror setup, source opening, format dispatch, diagnostics, and source removal in the C callback boundary. Test Plan: - ulimit -v 41943040; CARGO_BUILD_JOBS=1 cargo test --manifest-path rust/Cargo.toml decompression_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:
+72
-20
@@ -1208,6 +1208,32 @@ typedef char FIO_rust_decompress_multiple_projection_size[
|
|||||||
int FIO_rust_decompressMultipleFilenames(
|
int FIO_rust_decompressMultipleFilenames(
|
||||||
const FIO_rust_decompress_multiple_projection_t* projection);
|
const FIO_rust_decompress_multiple_projection_t* projection);
|
||||||
|
|
||||||
|
typedef int (*FIO_rust_decompress_multiple_separate_file_fn)(
|
||||||
|
void* opaque, const char* srcFileName);
|
||||||
|
typedef struct {
|
||||||
|
void* fCtx;
|
||||||
|
const char** srcNamesTable;
|
||||||
|
void* opaque;
|
||||||
|
FIO_rust_decompress_multiple_separate_file_fn decompressFile;
|
||||||
|
} FIO_rust_decompress_multiple_separate_projection_t;
|
||||||
|
typedef char FIO_rust_decompress_multiple_separate_fctx_offset[
|
||||||
|
(offsetof(FIO_rust_decompress_multiple_separate_projection_t, fCtx) == 0) ? 1 : -1];
|
||||||
|
typedef char FIO_rust_decompress_multiple_separate_input_names_offset[
|
||||||
|
(offsetof(FIO_rust_decompress_multiple_separate_projection_t, srcNamesTable)
|
||||||
|
== sizeof(void*)) ? 1 : -1];
|
||||||
|
typedef char FIO_rust_decompress_multiple_separate_opaque_offset[
|
||||||
|
(offsetof(FIO_rust_decompress_multiple_separate_projection_t, opaque)
|
||||||
|
== 2 * sizeof(void*)) ? 1 : -1];
|
||||||
|
typedef char FIO_rust_decompress_multiple_separate_callback_offset[
|
||||||
|
(offsetof(FIO_rust_decompress_multiple_separate_projection_t, decompressFile)
|
||||||
|
== 3 * sizeof(void*)) ? 1 : -1];
|
||||||
|
typedef char FIO_rust_decompress_multiple_separate_projection_size[
|
||||||
|
(sizeof(FIO_rust_decompress_multiple_separate_projection_t)
|
||||||
|
== 3 * sizeof(void*)
|
||||||
|
+ sizeof(FIO_rust_decompress_multiple_separate_file_fn)) ? 1 : -1];
|
||||||
|
int FIO_rust_decompressMultipleSeparateFilenames(
|
||||||
|
const FIO_rust_decompress_multiple_separate_projection_t* projection);
|
||||||
|
|
||||||
enum {
|
enum {
|
||||||
FIO_RUST_ZSTD_OK = 0,
|
FIO_RUST_ZSTD_OK = 0,
|
||||||
FIO_RUST_ZSTD_COMPRESS_ERROR = 1,
|
FIO_RUST_ZSTD_COMPRESS_ERROR = 1,
|
||||||
@@ -3884,6 +3910,43 @@ static int FIO_rust_decompressMultipleFileCallback(void* opaque,
|
|||||||
context->fCtx, context->prefs, *context->ress, outFileName, srcFileName);
|
context->fCtx, context->prefs, *context->ress, outFileName, srcFileName);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
typedef struct {
|
||||||
|
FIO_ctx_t* fCtx;
|
||||||
|
FIO_prefs_t* prefs;
|
||||||
|
dRess_t* ress;
|
||||||
|
const char* outMirroredRootDirName;
|
||||||
|
const char* outDirName;
|
||||||
|
} FIO_rust_decompress_multiple_separate_context_t;
|
||||||
|
|
||||||
|
static const char* FIO_determineDstName(const char* srcFileName, const char* outDirName);
|
||||||
|
|
||||||
|
/* Keep destination-name construction, diagnostics, and decompression in C. */
|
||||||
|
static int FIO_rust_decompressMultipleSeparateFileCallback(void* opaque,
|
||||||
|
const char* srcFileName)
|
||||||
|
{
|
||||||
|
FIO_rust_decompress_multiple_separate_context_t* const context =
|
||||||
|
(FIO_rust_decompress_multiple_separate_context_t*)opaque;
|
||||||
|
const char* dstFileName = NULL;
|
||||||
|
|
||||||
|
if (context->outMirroredRootDirName) {
|
||||||
|
char* validMirroredDirName = UTIL_createMirroredDestDirName(
|
||||||
|
srcFileName, context->outMirroredRootDirName);
|
||||||
|
if (validMirroredDirName) {
|
||||||
|
dstFileName = FIO_determineDstName(srcFileName, validMirroredDirName);
|
||||||
|
free(validMirroredDirName);
|
||||||
|
} else {
|
||||||
|
DISPLAYLEVEL(2, "zstd: --output-dir-mirror cannot decompress '%s' into '%s'\n",
|
||||||
|
srcFileName, context->outMirroredRootDirName);
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
dstFileName = FIO_determineDstName(srcFileName, context->outDirName);
|
||||||
|
}
|
||||||
|
|
||||||
|
if (dstFileName == NULL) return 1;
|
||||||
|
return FIO_decompressSrcFile(
|
||||||
|
context->fCtx, context->prefs, *context->ress, dstFileName, srcFileName);
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
int FIO_decompressFilename(FIO_ctx_t* const fCtx, FIO_prefs_t* const prefs,
|
int FIO_decompressFilename(FIO_ctx_t* const fCtx, FIO_prefs_t* const prefs,
|
||||||
@@ -3953,7 +4016,6 @@ FIO_decompressMultipleFilenames(FIO_ctx_t* const fCtx,
|
|||||||
const char* outDirName, const char* outFileName,
|
const char* outDirName, const char* outFileName,
|
||||||
const char* dictFileName)
|
const char* dictFileName)
|
||||||
{
|
{
|
||||||
int status;
|
|
||||||
int error = 0;
|
int error = 0;
|
||||||
dRess_t ress = FIO_createDResources(prefs, dictFileName);
|
dRess_t ress = FIO_createDResources(prefs, dictFileName);
|
||||||
|
|
||||||
@@ -3981,28 +4043,18 @@ FIO_decompressMultipleFilenames(FIO_ctx_t* const fCtx,
|
|||||||
EXM_THROW(72, "Write error : %s : cannot properly close output file",
|
EXM_THROW(72, "Write error : %s : cannot properly close output file",
|
||||||
strerror(errno));
|
strerror(errno));
|
||||||
} else {
|
} else {
|
||||||
|
FIO_rust_decompress_multiple_separate_context_t callbackContext = {
|
||||||
|
fCtx, prefs, &ress, outMirroredRootDirName, outDirName
|
||||||
|
};
|
||||||
|
FIO_rust_decompress_multiple_separate_projection_t projection = {
|
||||||
|
fCtx, srcNamesTable, &callbackContext,
|
||||||
|
FIO_rust_decompressMultipleSeparateFileCallback
|
||||||
|
};
|
||||||
if (outMirroredRootDirName)
|
if (outMirroredRootDirName)
|
||||||
UTIL_mirrorSourceFilesDirectories(srcNamesTable, (unsigned)fCtx->nbFilesTotal, outMirroredRootDirName);
|
UTIL_mirrorSourceFilesDirectories(srcNamesTable, (unsigned)fCtx->nbFilesTotal, outMirroredRootDirName);
|
||||||
|
|
||||||
for (; fCtx->currFileIdx < fCtx->nbFilesTotal; fCtx->currFileIdx++) { /* create dstFileName */
|
error = FIO_rust_decompressMultipleSeparateFilenames(&projection);
|
||||||
const char* const srcFileName = srcNamesTable[fCtx->currFileIdx];
|
|
||||||
const char* dstFileName = NULL;
|
|
||||||
if (outMirroredRootDirName) {
|
|
||||||
char* validMirroredDirName = UTIL_createMirroredDestDirName(srcFileName, outMirroredRootDirName);
|
|
||||||
if (validMirroredDirName) {
|
|
||||||
dstFileName = FIO_determineDstName(srcFileName, validMirroredDirName);
|
|
||||||
free(validMirroredDirName);
|
|
||||||
} else {
|
|
||||||
DISPLAYLEVEL(2, "zstd: --output-dir-mirror cannot decompress '%s' into '%s'\n", srcFileName, outMirroredRootDirName);
|
|
||||||
}
|
|
||||||
} else {
|
|
||||||
dstFileName = FIO_determineDstName(srcFileName, outDirName);
|
|
||||||
}
|
|
||||||
if (dstFileName == NULL) { error=1; continue; }
|
|
||||||
status = FIO_decompressSrcFile(fCtx, prefs, ress, dstFileName, srcFileName);
|
|
||||||
if (!status) fCtx->nbFilesProcessed++;
|
|
||||||
error |= status;
|
|
||||||
}
|
|
||||||
if (outDirName)
|
if (outDirName)
|
||||||
FIO_checkFilenameCollisions(srcNamesTable , (unsigned)fCtx->nbFilesTotal);
|
FIO_checkFilenameCollisions(srcNamesTable , (unsigned)fCtx->nbFilesTotal);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -334,6 +334,8 @@ const _: () = {
|
|||||||
|
|
||||||
pub type FIO_rust_decompress_multiple_file_fn =
|
pub type FIO_rust_decompress_multiple_file_fn =
|
||||||
unsafe extern "C" fn(*mut c_void, *const c_char, *const c_char) -> c_int;
|
unsafe extern "C" fn(*mut c_void, *const c_char, *const c_char) -> c_int;
|
||||||
|
pub type FIO_rust_decompress_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
|
/// Rust owns only the shared-destination file iteration and aggregate error
|
||||||
/// handling. C retains source opening, format dispatch, diagnostics, and
|
/// handling. C retains source opening, format dispatch, diagnostics, and
|
||||||
@@ -368,6 +370,40 @@ const _: () = {
|
|||||||
assert!(size_of::<FIO_rust_decompress_multiple_projection_t>() == 5 * size_of::<usize>());
|
assert!(size_of::<FIO_rust_decompress_multiple_projection_t>() == 5 * size_of::<usize>());
|
||||||
};
|
};
|
||||||
|
|
||||||
|
/// Rust owns only the separate-destination file iteration and aggregate error
|
||||||
|
/// handling. C retains destination-name construction, source opening, format
|
||||||
|
/// dispatch, diagnostics, and source removal through the per-source callback.
|
||||||
|
#[repr(C)]
|
||||||
|
pub struct FIO_rust_decompress_multiple_separate_projection_t {
|
||||||
|
pub f_ctx: *mut c_void,
|
||||||
|
pub src_names_table: *const *const c_char,
|
||||||
|
pub opaque: *mut c_void,
|
||||||
|
pub decompress_file: Option<FIO_rust_decompress_multiple_separate_file_fn>,
|
||||||
|
}
|
||||||
|
const _: () = {
|
||||||
|
assert!(std::mem::offset_of!(FIO_rust_decompress_multiple_separate_projection_t, f_ctx) == 0);
|
||||||
|
assert!(
|
||||||
|
std::mem::offset_of!(
|
||||||
|
FIO_rust_decompress_multiple_separate_projection_t,
|
||||||
|
src_names_table
|
||||||
|
) == size_of::<usize>()
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
std::mem::offset_of!(FIO_rust_decompress_multiple_separate_projection_t, opaque)
|
||||||
|
== 2 * size_of::<usize>()
|
||||||
|
);
|
||||||
|
assert!(
|
||||||
|
std::mem::offset_of!(
|
||||||
|
FIO_rust_decompress_multiple_separate_projection_t,
|
||||||
|
decompress_file
|
||||||
|
) == 3 * size_of::<usize>()
|
||||||
|
);
|
||||||
|
assert!(size_of::<FIO_rust_decompress_multiple_separate_file_fn>() == size_of::<usize>());
|
||||||
|
assert!(
|
||||||
|
size_of::<FIO_rust_decompress_multiple_separate_projection_t>() == 4 * size_of::<usize>()
|
||||||
|
);
|
||||||
|
};
|
||||||
|
|
||||||
pub const FIO_RUST_ZSTD_OK: c_int = 0;
|
pub const FIO_RUST_ZSTD_OK: c_int = 0;
|
||||||
pub const FIO_RUST_ZSTD_COMPRESS_ERROR: c_int = 1;
|
pub const FIO_RUST_ZSTD_COMPRESS_ERROR: c_int = 1;
|
||||||
pub const FIO_RUST_ZSTD_INCOMPLETE_INPUT: c_int = 2;
|
pub const FIO_RUST_ZSTD_INCOMPLETE_INPUT: c_int = 2;
|
||||||
@@ -2344,6 +2380,41 @@ pub unsafe extern "C" fn FIO_rust_decompressMultipleFilenames(
|
|||||||
error
|
error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Iterate files that each receive a separate destination. The C callback
|
||||||
|
/// retains destination-name construction, source opening, format dispatch,
|
||||||
|
/// diagnostics, and `--rm`; Rust preserves the non-short-circuiting loop and
|
||||||
|
/// counter updates.
|
||||||
|
#[no_mangle]
|
||||||
|
pub unsafe extern "C" fn FIO_rust_decompressMultipleSeparateFilenames(
|
||||||
|
projection: *const FIO_rust_decompress_multiple_separate_projection_t,
|
||||||
|
) -> c_int {
|
||||||
|
assert!(!projection.is_null());
|
||||||
|
let projection = unsafe { &*projection };
|
||||||
|
assert!(!projection.f_ctx.is_null());
|
||||||
|
assert!(!projection.src_names_table.is_null());
|
||||||
|
let decompress_file = projection
|
||||||
|
.decompress_file
|
||||||
|
.expect("separate-destination decompression 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.src_names_table.add(file_index) };
|
||||||
|
let status = unsafe { decompress_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]
|
#[inline]
|
||||||
fn zstd_adapt_level(projection: &FIO_rust_zstd_adapt_projection_t, slower: bool) -> c_int {
|
fn zstd_adapt_level(projection: &FIO_rust_zstd_adapt_projection_t, slower: bool) -> c_int {
|
||||||
if slower {
|
if slower {
|
||||||
@@ -4668,6 +4739,20 @@ mod tests {
|
|||||||
status
|
status
|
||||||
}
|
}
|
||||||
|
|
||||||
|
unsafe extern "C" fn record_multiple_separate_decompression_file(
|
||||||
|
opaque: *mut c_void,
|
||||||
|
source: *const c_char,
|
||||||
|
) -> c_int {
|
||||||
|
let state = unsafe { &mut *opaque.cast::<MultipleDecompressionState>() };
|
||||||
|
let status = state
|
||||||
|
.statuses
|
||||||
|
.get(state.sources.len())
|
||||||
|
.copied()
|
||||||
|
.unwrap_or(0);
|
||||||
|
state.sources.push(source);
|
||||||
|
status
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn compression_multiple_separate_destinations_preserves_order_and_counters() {
|
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 source_names = [c"zero".as_ptr(), c"one".as_ptr(), c"two".as_ptr()];
|
||||||
@@ -4872,6 +4957,70 @@ mod tests {
|
|||||||
assert_eq!(context.nbFilesProcessed, 4);
|
assert_eq!(context.nbFilesProcessed, 4);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn decompression_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 = MultipleDecompressionState {
|
||||||
|
statuses: vec![0, 0],
|
||||||
|
..MultipleDecompressionState::default()
|
||||||
|
};
|
||||||
|
let projection = FIO_rust_decompress_multiple_separate_projection_t {
|
||||||
|
f_ctx: (&mut context as *mut FIO_rust_compression_context_t).cast(),
|
||||||
|
src_names_table: source_names.as_ptr(),
|
||||||
|
opaque: (&mut state as *mut MultipleDecompressionState).cast(),
|
||||||
|
decompress_file: Some(record_multiple_separate_decompression_file),
|
||||||
|
};
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
unsafe { FIO_rust_decompressMultipleSeparateFilenames(&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 decompression_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 = MultipleDecompressionState {
|
||||||
|
statuses: vec![1, 4, 2],
|
||||||
|
..MultipleDecompressionState::default()
|
||||||
|
};
|
||||||
|
let projection = FIO_rust_decompress_multiple_separate_projection_t {
|
||||||
|
f_ctx: (&mut context as *mut FIO_rust_compression_context_t).cast(),
|
||||||
|
src_names_table: source_names.as_ptr(),
|
||||||
|
opaque: (&mut state as *mut MultipleDecompressionState).cast(),
|
||||||
|
decompress_file: Some(record_multiple_separate_decompression_file),
|
||||||
|
};
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
unsafe { FIO_rust_decompressMultipleSeparateFilenames(&projection) },
|
||||||
|
1 | 4 | 2
|
||||||
|
);
|
||||||
|
assert_eq!(state.sources, source_names);
|
||||||
|
assert_eq!(context.currFileIdx, 3);
|
||||||
|
assert_eq!(context.nbFilesProcessed, 4);
|
||||||
|
}
|
||||||
|
|
||||||
struct ZstdProjectionState {
|
struct ZstdProjectionState {
|
||||||
input: [u8; 5],
|
input: [u8; 5],
|
||||||
input_pos: usize,
|
input_pos: usize,
|
||||||
|
|||||||
Reference in New Issue
Block a user