diff --git a/programs/fileio.c b/programs/fileio.c index 8e8b21802..49a1ee3c7 100644 --- a/programs/fileio.c +++ b/programs/fileio.c @@ -1234,6 +1234,91 @@ typedef char FIO_rust_decompress_multiple_separate_projection_size[ int FIO_rust_decompressMultipleSeparateFilenames( const FIO_rust_decompress_multiple_separate_projection_t* projection); +enum { + FIO_RUST_DECOMPRESS_SRC_OPEN_OK = 0, +}; +typedef int (*FIO_rust_decompress_src_open_fn)(void* opaque, const char* srcFileName, + U64* fileSize, int* sourceIsRegular); +typedef void (*FIO_rust_decompress_src_async_fn)(void* opaque, int asyncMode); +typedef void (*FIO_rust_decompress_src_attach_fn)(void* opaque); +typedef int (*FIO_rust_decompress_src_close_fn)(void* opaque); +typedef int (*FIO_rust_decompress_dst_open_fn)(void* opaque, + const char* srcFileName, + const char* dstFileName, + int transferStat, int* dstFd); +typedef void (*FIO_rust_decompress_dst_attach_fn)(void* opaque); +typedef void (*FIO_rust_decompress_handler_fn)(void* opaque); +typedef int (*FIO_rust_decompress_file_fn)(void* opaque, + const char* dstFileName, + const char* srcFileName); +typedef void (*FIO_rust_decompress_dst_stat_fn)(void* opaque, int dstFd, + const char* dstFileName); +typedef int (*FIO_rust_decompress_dst_close_fn)(void* opaque); +typedef void (*FIO_rust_decompress_dst_remove_fn)(void* opaque, const char* dstFileName); +typedef int (*FIO_rust_decompress_src_remove_fn)(void* opaque, const char* srcFileName); + +/* Rust owns per-file source/destination ordering; the callback context keeps + * dRess_t, FILE/pool handles, diagnostics, and the format dispatcher private. */ +typedef struct { + void* opaque; + const char* dstFileName; + const char* srcFileName; + int destinationAlreadyOpen; + int testMode; + int sourceIsStdin; + int destinationIsStdout; + int removeSource; + FIO_rust_decompress_src_open_fn openSource; + FIO_rust_decompress_src_async_fn setAsync; + FIO_rust_decompress_src_attach_fn attachSource; + FIO_rust_decompress_src_attach_fn detachSource; + FIO_rust_decompress_src_close_fn closeSource; + FIO_rust_decompress_dst_open_fn openDestination; + FIO_rust_decompress_dst_attach_fn attachDestination; + FIO_rust_decompress_handler_fn addHandler; + FIO_rust_decompress_file_fn decompress; + FIO_rust_decompress_handler_fn clearHandler; + FIO_rust_decompress_dst_stat_fn setFDStat; + FIO_rust_decompress_dst_close_fn closeDestination; + FIO_rust_decompress_handler_fn utimeDestination; + FIO_rust_decompress_dst_remove_fn removeDestination; + FIO_rust_decompress_src_remove_fn removeSourceFile; +} FIO_rust_decompress_file_projection_t; +typedef char FIO_rust_decompress_file_opaque_offset[ + (offsetof(FIO_rust_decompress_file_projection_t, opaque) == 0) ? 1 : -1]; +typedef char FIO_rust_decompress_file_dst_name_offset[ + (offsetof(FIO_rust_decompress_file_projection_t, dstFileName) + == sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_decompress_file_src_name_offset[ + (offsetof(FIO_rust_decompress_file_projection_t, srcFileName) + == 2 * sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_decompress_file_destination_open_offset[ + (offsetof(FIO_rust_decompress_file_projection_t, destinationAlreadyOpen) + == 3 * sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_decompress_file_test_mode_offset[ + (offsetof(FIO_rust_decompress_file_projection_t, testMode) + == 3 * sizeof(void*) + sizeof(int)) ? 1 : -1]; +typedef char FIO_rust_decompress_file_stdin_offset[ + (offsetof(FIO_rust_decompress_file_projection_t, sourceIsStdin) + == 3 * sizeof(void*) + 2 * sizeof(int)) ? 1 : -1]; +typedef char FIO_rust_decompress_file_stdout_offset[ + (offsetof(FIO_rust_decompress_file_projection_t, destinationIsStdout) + == 3 * sizeof(void*) + 3 * sizeof(int)) ? 1 : -1]; +typedef char FIO_rust_decompress_file_remove_source_offset[ + (offsetof(FIO_rust_decompress_file_projection_t, removeSource) + == 3 * sizeof(void*) + 4 * sizeof(int)) ? 1 : -1]; +typedef char FIO_rust_decompress_file_callback_offset[ + (offsetof(FIO_rust_decompress_file_projection_t, openSource) + == ((3 * sizeof(void*) + 5 * sizeof(int) + sizeof(void*) - 1) + / sizeof(void*)) * sizeof(void*)) ? 1 : -1]; +typedef char FIO_rust_decompress_file_projection_size[ + (sizeof(FIO_rust_decompress_file_projection_t) + == ((3 * sizeof(void*) + 5 * sizeof(int) + sizeof(void*) - 1) + / sizeof(void*)) * sizeof(void*) + + 15 * sizeof(FIO_rust_decompress_src_open_fn)) ? 1 : -1]; +int FIO_rust_decompressFilename( + const FIO_rust_decompress_file_projection_t* projection); + enum { FIO_RUST_ZSTD_OK = 0, FIO_RUST_ZSTD_COMPRESS_ERROR = 1, @@ -3896,132 +3981,209 @@ static int FIO_decompressFrames(FIO_ctx_t* const fCtx, return FIO_rust_finishDecompressFrames(status, srcFileName, filesize, &callbacks); } -/** FIO_decompressDstFile() : - open `dstFileName`, or pass-through if writeCtx's file is already != 0, - then start decompression process (FIO_decompressFrames()). - @return : 0 : OK - 1 : operation aborted -*/ -static int FIO_decompressDstFile(FIO_ctx_t* const fCtx, - FIO_prefs_t* const prefs, - dRess_t ress, - const char* dstFileName, - const char* srcFileName, - const stat_t* srcFileStat) -{ - int result; - int releaseDstFile = 0; - int transferStat = 0; - int dstFd = 0; - - if ((AIO_WritePool_getFile(ress.writeCtx) == NULL) && (prefs->testMode == 0)) { - FILE *dstFile; - int dstFilePermissions = DEFAULT_FILE_PERMISSIONS; - if ( strcmp(srcFileName, stdinmark) /* special case : don't transfer permissions from stdin */ - && strcmp(dstFileName, stdoutmark) - && UTIL_isRegularFileStat(srcFileStat) ) { - transferStat = 1; - dstFilePermissions = TEMPORARY_FILE_PERMISSIONS; - } - - releaseDstFile = 1; - - dstFile = FIO_openDstFile(fCtx, prefs, srcFileName, dstFileName, dstFilePermissions); - if (dstFile==NULL) return 1; - dstFd = fileno(dstFile); - AIO_WritePool_setFile(ress.writeCtx, dstFile); - - /* Must only be added after FIO_openDstFile() succeeds. - * Otherwise we may delete the destination file if it already exists, - * and the user presses Ctrl-C when asked if they wish to overwrite. - */ - addHandler(dstFileName); - } - - result = FIO_decompressFrames(fCtx, ress, prefs, dstFileName, srcFileName); - - if (releaseDstFile) { - clearHandler(); - - if (transferStat) { - UTIL_setFDStat(dstFd, dstFileName, srcFileStat); - } - - if (AIO_WritePool_closeFile(ress.writeCtx)) { - DISPLAYLEVEL(1, "zstd: %s: %s \n", dstFileName, strerror(errno)); - result = 1; - } - - if (transferStat) { - UTIL_utime(dstFileName, srcFileStat); - } - - if ( (result != 0) /* operation failure */ - && strcmp(dstFileName, stdoutmark) /* special case : don't remove() stdout */ - ) { - FIO_removeFile(dstFileName); /* remove decompression artefact; note: don't do anything special if remove() fails */ - } - } - - return result; -} - - -/** FIO_decompressSrcFile() : - Open `srcFileName`, transfer control to decompressDstFile() - @return : 0 : OK - 1 : error -*/ -static int FIO_decompressSrcFile(FIO_ctx_t* const fCtx, FIO_prefs_t* const prefs, dRess_t ress, const char* dstFileName, const char* srcFileName) -{ - FILE* srcFile; +typedef struct { + FIO_ctx_t* fCtx; + FIO_prefs_t* prefs; + dRess_t* ress; + const char* dstFileName; + const char* srcFileName; stat_t srcFileStat; - int result; - U64 fileSize = UTIL_FILESIZE_UNKNOWN; + FILE* srcFile; + FILE* dstFile; +} FIO_rust_decompress_file_context_t; + +static int FIO_rust_decompressSourceOpen(void* opaque, const char* srcFileName, + U64* fileSize, int* sourceIsRegular) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; if (UTIL_isDirectory(srcFileName)) { DISPLAYLEVEL(1, "zstd: %s is a directory -- ignored \n", srcFileName); return 1; } - srcFile = FIO_openSrcFile(prefs, srcFileName, &srcFileStat); - if (srcFile==NULL) return 1; + context->srcFile = FIO_openSrcFile(context->prefs, srcFileName, &context->srcFileStat); + if (context->srcFile == NULL) + return 1; - /* Don't use AsyncIO for small files */ - if (strcmp(srcFileName, stdinmark)) /* Stdin doesn't have stats */ - fileSize = UTIL_getFileSizeStat(&srcFileStat); - if(fileSize != UTIL_FILESIZE_UNKNOWN && fileSize < ZSTD_BLOCKSIZE_MAX * 3) { - AIO_ReadPool_setAsync(ress.readCtx, 0); - AIO_WritePool_setAsync(ress.writeCtx, 0); - } else { - AIO_ReadPool_setAsync(ress.readCtx, 1); - AIO_WritePool_setAsync(ress.writeCtx, 1); - } + *fileSize = strcmp(srcFileName, stdinmark) + ? UTIL_getFileSizeStat(&context->srcFileStat) + : UTIL_FILESIZE_UNKNOWN; + *sourceIsRegular = strcmp(srcFileName, stdinmark) + ? UTIL_isRegularFileStat(&context->srcFileStat) + : 0; + return FIO_RUST_DECOMPRESS_SRC_OPEN_OK; +} - AIO_ReadPool_setFile(ress.readCtx, srcFile); +static void FIO_rust_decompressSourceSetAsync(void* opaque, int asyncMode) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + AIO_ReadPool_setAsync(context->ress->readCtx, asyncMode); + AIO_WritePool_setAsync(context->ress->writeCtx, asyncMode); +} - result = FIO_decompressDstFile(fCtx, prefs, ress, dstFileName, srcFileName, &srcFileStat); +static void FIO_rust_decompressSourceAttach(void* opaque) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + AIO_ReadPool_setFile(context->ress->readCtx, context->srcFile); +} - AIO_ReadPool_setFile(ress.readCtx, NULL); +static void FIO_rust_decompressSourceDetach(void* opaque) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + AIO_ReadPool_setFile(context->ress->readCtx, NULL); +} - /* Close file */ - if (fclose(srcFile)) { - DISPLAYLEVEL(1, "zstd: %s: %s \n", srcFileName, strerror(errno)); /* error should not happen */ +static int FIO_rust_decompressSourceClose(void* opaque) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + int const result = fclose(context->srcFile); + context->srcFile = NULL; + if (result) { + DISPLAYLEVEL(1, "zstd: %s: %s \n", context->srcFileName, strerror(errno)); return 1; } - if ( prefs->removeSrcFile /* --rm */ - && (result==0) /* decompression successful */ - && strcmp(srcFileName, stdinmark) ) /* not stdin */ { - /* We must clear the handler, since after this point calling it would - * delete both the source and destination files. - */ - clearHandler(); - if (FIO_removeFile(srcFileName)) { - /* failed to remove src file */ - DISPLAYLEVEL(1, "zstd: %s: %s \n", srcFileName, strerror(errno)); - return 1; - } } - return result; + return 0; +} + +static int FIO_rust_decompressDestinationOpen(void* opaque, + const char* srcFileName, + const char* dstFileName, + int transferStat, int* dstFd) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + int const permissions = transferStat + ? TEMPORARY_FILE_PERMISSIONS + : DEFAULT_FILE_PERMISSIONS; + + context->dstFile = FIO_openDstFile( + context->fCtx, context->prefs, srcFileName, dstFileName, permissions); + if (context->dstFile == NULL) + return 1; + *dstFd = fileno(context->dstFile); + return 0; +} + +static void FIO_rust_decompressDestinationAttach(void* opaque) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + AIO_WritePool_setFile(context->ress->writeCtx, context->dstFile); +} + +static void FIO_rust_decompressAddHandler(void* opaque) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + addHandler(context->dstFileName); +} + +static int FIO_rust_decompressFile(void* opaque, + const char* dstFileName, + const char* srcFileName) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + return FIO_decompressFrames( + context->fCtx, *context->ress, context->prefs, dstFileName, srcFileName); +} + +static void FIO_rust_decompressClearHandler(void* opaque) +{ + (void)opaque; + clearHandler(); +} + +static void FIO_rust_decompressSetFDStat(void* opaque, int dstFd, + const char* dstFileName) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + UTIL_setFDStat(dstFd, dstFileName, &context->srcFileStat); +} + +static int FIO_rust_decompressDestinationClose(void* opaque) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + int const result = AIO_WritePool_closeFile(context->ress->writeCtx); + context->dstFile = NULL; + if (result) { + DISPLAYLEVEL(1, "zstd: %s: %s \n", context->dstFileName, strerror(errno)); + return 1; + } + return 0; +} + +static void FIO_rust_decompressDestinationUtime(void* opaque) +{ + FIO_rust_decompress_file_context_t* const context = + (FIO_rust_decompress_file_context_t*)opaque; + UTIL_utime(context->dstFileName, &context->srcFileStat); +} + +static void FIO_rust_decompressDestinationRemove(void* opaque, const char* dstFileName) +{ + (void)opaque; + (void)FIO_removeFile(dstFileName); +} + +static int FIO_rust_decompressSourceRemove(void* opaque, const char* srcFileName) +{ + (void)opaque; + if (FIO_removeFile(srcFileName)) { + DISPLAYLEVEL(1, "zstd: %s: %s \n", srcFileName, strerror(errno)); + return 1; + } + return 0; +} + +static int FIO_decompressFile(FIO_ctx_t* const fCtx, FIO_prefs_t* const prefs, + dRess_t* const ress, const char* dstFileName, + const char* srcFileName) +{ + FIO_rust_decompress_file_context_t context; + FIO_rust_decompress_file_projection_t projection; + + memset(&context, 0, sizeof(context)); + context.fCtx = fCtx; + context.prefs = prefs; + context.ress = ress; + context.dstFileName = dstFileName; + context.srcFileName = srcFileName; + + memset(&projection, 0, sizeof(projection)); + projection.opaque = &context; + projection.dstFileName = dstFileName; + projection.srcFileName = srcFileName; + projection.destinationAlreadyOpen = AIO_WritePool_getFile(ress->writeCtx) != NULL; + projection.testMode = prefs->testMode; + projection.sourceIsStdin = !strcmp(srcFileName, stdinmark); + projection.destinationIsStdout = !strcmp(dstFileName, stdoutmark); + projection.removeSource = prefs->removeSrcFile; + projection.openSource = FIO_rust_decompressSourceOpen; + projection.setAsync = FIO_rust_decompressSourceSetAsync; + projection.attachSource = FIO_rust_decompressSourceAttach; + projection.detachSource = FIO_rust_decompressSourceDetach; + projection.closeSource = FIO_rust_decompressSourceClose; + projection.openDestination = FIO_rust_decompressDestinationOpen; + projection.attachDestination = FIO_rust_decompressDestinationAttach; + projection.addHandler = FIO_rust_decompressAddHandler; + projection.decompress = FIO_rust_decompressFile; + projection.clearHandler = FIO_rust_decompressClearHandler; + projection.setFDStat = FIO_rust_decompressSetFDStat; + projection.closeDestination = FIO_rust_decompressDestinationClose; + projection.utimeDestination = FIO_rust_decompressDestinationUtime; + projection.removeDestination = FIO_rust_decompressDestinationRemove; + projection.removeSourceFile = FIO_rust_decompressSourceRemove; + + return FIO_rust_decompressFilename(&projection); } typedef struct { @@ -4030,15 +4192,16 @@ typedef struct { dRess_t* ress; } FIO_rust_decompress_multiple_context_t; -/* Keep source opening, format dispatch, diagnostics, and --rm in C. */ +/* Keep private source/destination operations, format dispatch, and + * diagnostics in C; Rust owns per-file ordering and removal policy. */ static int FIO_rust_decompressMultipleFileCallback(void* opaque, const char* outFileName, const char* srcFileName) { FIO_rust_decompress_multiple_context_t* const context = (FIO_rust_decompress_multiple_context_t*)opaque; - return FIO_decompressSrcFile( - context->fCtx, context->prefs, *context->ress, outFileName, srcFileName); + return FIO_decompressFile( + context->fCtx, context->prefs, context->ress, outFileName, srcFileName); } typedef struct { @@ -4051,7 +4214,8 @@ typedef struct { static const char* FIO_determineDstName(const char* srcFileName, const char* outDirName); -/* Keep destination-name construction, diagnostics, and decompression in C. */ +/* Keep destination-name construction, diagnostics, and private operations in + * C; Rust owns per-file ordering and removal policy. */ static int FIO_rust_decompressMultipleSeparateFileCallback(void* opaque, const char* srcFileName) { @@ -4074,8 +4238,8 @@ static int FIO_rust_decompressMultipleSeparateFileCallback(void* opaque, } if (dstFileName == NULL) return 1; - return FIO_decompressSrcFile( - context->fCtx, context->prefs, *context->ress, dstFileName, srcFileName); + return FIO_decompressFile( + context->fCtx, context->prefs, context->ress, dstFileName, srcFileName); } @@ -4084,11 +4248,9 @@ int FIO_decompressFilename(FIO_ctx_t* const fCtx, FIO_prefs_t* const prefs, const char* dstFileName, const char* srcFileName, const char* dictFileName) { - dRess_t const ress = FIO_createDResources(prefs, dictFileName); - - int const decodingError = FIO_decompressSrcFile(fCtx, prefs, ress, dstFileName, srcFileName); - - + dRess_t ress = FIO_createDResources(prefs, dictFileName); + int const decodingError = FIO_decompressFile( + fCtx, prefs, &ress, dstFileName, srcFileName); FIO_freeDResources(ress); return decodingError; diff --git a/rust/src/fileio_asyncio.rs b/rust/src/fileio_asyncio.rs index fc4e0ee31..a5d45d34d 100644 --- a/rust/src/fileio_asyncio.rs +++ b/rust/src/fileio_asyncio.rs @@ -510,6 +510,159 @@ const _: () = { ); }; +pub const FIO_RUST_DECOMPRESS_SRC_OPEN_OK: c_int = 0; + +const FIO_RUST_DECOMPRESS_SRC_ASYNC_THRESHOLD: u64 = (1 << 17) * 3; +const FIO_RUST_DECOMPRESS_SRC_UNKNOWN_SIZE: u64 = u64::MAX; + +pub type FIO_rust_decompress_src_open_fn = + unsafe extern "C" fn(*mut c_void, *const c_char, *mut u64, *mut c_int) -> c_int; +pub type FIO_rust_decompress_src_async_fn = unsafe extern "C" fn(*mut c_void, c_int); +pub type FIO_rust_decompress_src_attach_fn = unsafe extern "C" fn(*mut c_void); +pub type FIO_rust_decompress_src_close_fn = unsafe extern "C" fn(*mut c_void) -> c_int; +pub type FIO_rust_decompress_dst_open_fn = + unsafe extern "C" fn(*mut c_void, *const c_char, *const c_char, c_int, *mut c_int) -> c_int; +pub type FIO_rust_decompress_dst_attach_fn = unsafe extern "C" fn(*mut c_void); +pub type FIO_rust_decompress_handler_fn = unsafe extern "C" fn(*mut c_void); +pub type FIO_rust_decompress_file_fn = + unsafe extern "C" fn(*mut c_void, *const c_char, *const c_char) -> c_int; +pub type FIO_rust_decompress_dst_stat_fn = unsafe extern "C" fn(*mut c_void, c_int, *const c_char); +pub type FIO_rust_decompress_dst_close_fn = unsafe extern "C" fn(*mut c_void) -> c_int; +pub type FIO_rust_decompress_dst_remove_fn = unsafe extern "C" fn(*mut c_void, *const c_char); +pub type FIO_rust_decompress_src_remove_fn = + unsafe extern "C" fn(*mut c_void, *const c_char) -> c_int; + +/// C-owned resources and private I/O handles are projected as callbacks while +/// Rust owns the per-file decompression ordering and result policy. +#[repr(C)] +pub struct FIO_rust_decompress_file_projection_t { + pub opaque: *mut c_void, + pub dst_file_name: *const c_char, + pub src_file_name: *const c_char, + pub destination_already_open: c_int, + pub test_mode: c_int, + pub source_is_stdin: c_int, + pub destination_is_stdout: c_int, + pub remove_source: c_int, + pub open_source: Option, + pub set_async: Option, + pub attach_source: Option, + pub detach_source: Option, + pub close_source: Option, + pub open_destination: Option, + pub attach_destination: Option, + pub add_handler: Option, + pub decompress: Option, + pub clear_handler: Option, + pub set_fd_stat: Option, + pub close_destination: Option, + pub utime_destination: Option, + pub remove_destination: Option, + pub remove_source_file: Option, +} + +const FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET: usize = + (3 * size_of::() + 5 * size_of::() + size_of::() - 1) + & !(size_of::() - 1); + +const _: () = { + assert!(std::mem::offset_of!(FIO_rust_decompress_file_projection_t, opaque) == 0); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, dst_file_name) + == size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, src_file_name) + == 2 * size_of::() + ); + assert!( + std::mem::offset_of!( + FIO_rust_decompress_file_projection_t, + destination_already_open + ) == 3 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, test_mode) + == 3 * size_of::() + size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, source_is_stdin) + == 3 * size_of::() + 2 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, destination_is_stdout) + == 3 * size_of::() + 3 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, remove_source) + == 3 * size_of::() + 4 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, open_source) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, set_async) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, attach_source) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 2 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, detach_source) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 3 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, close_source) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 4 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, open_destination) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 5 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, attach_destination) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 6 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, add_handler) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 7 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, decompress) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 8 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, clear_handler) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 9 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, set_fd_stat) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 10 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, close_destination) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 11 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, utime_destination) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 12 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, remove_destination) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 13 * size_of::() + ); + assert!( + std::mem::offset_of!(FIO_rust_decompress_file_projection_t, remove_source_file) + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 14 * size_of::() + ); + assert!( + size_of::() + == FIO_RUST_DECOMPRESS_FILE_CALLBACK_OFFSET + 15 * 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; @@ -2463,6 +2616,177 @@ pub unsafe extern "C" fn FIO_rust_compressFilenameSrcFile( result } +unsafe fn decompress_destination( + projection: &FIO_rust_decompress_file_projection_t, + source_is_regular: c_int, +) -> c_int { + let decompress = projection + .decompress + .expect("decompression callback is required"); + let release_destination = projection.destination_already_open == 0 && projection.test_mode == 0; + + if !release_destination { + return unsafe { + decompress( + projection.opaque, + projection.dst_file_name, + projection.src_file_name, + ) + }; + } + + let transfer_stat = projection.source_is_stdin == 0 + && projection.destination_is_stdout == 0 + && source_is_regular != 0; + let mut destination_fd = 0; + let open_destination = projection + .open_destination + .expect("destination open callback is required"); + let open_status = unsafe { + open_destination( + projection.opaque, + projection.src_file_name, + projection.dst_file_name, + c_int::from(transfer_stat), + &mut destination_fd, + ) + }; + if open_status != 0 { + return 1; + } + + let attach_destination = projection + .attach_destination + .expect("destination attachment callback is required"); + unsafe { attach_destination(projection.opaque) }; + + let add_handler = projection + .add_handler + .expect("destination handler callback is required"); + unsafe { add_handler(projection.opaque) }; + + let mut result = unsafe { + decompress( + projection.opaque, + projection.dst_file_name, + projection.src_file_name, + ) + }; + + let clear_handler = projection + .clear_handler + .expect("handler-clear callback is required"); + unsafe { clear_handler(projection.opaque) }; + + if transfer_stat { + let set_fd_stat = projection + .set_fd_stat + .expect("destination stat callback is required"); + unsafe { set_fd_stat(projection.opaque, destination_fd, projection.dst_file_name) }; + } + + let close_destination = projection + .close_destination + .expect("destination close callback is required"); + if unsafe { close_destination(projection.opaque) } != 0 { + result = 1; + } + + if transfer_stat { + let utime_destination = projection + .utime_destination + .expect("destination timestamp callback is required"); + unsafe { utime_destination(projection.opaque) }; + } + + if result != 0 && projection.destination_is_stdout == 0 { + let remove_destination = projection + .remove_destination + .expect("destination removal callback is required"); + unsafe { remove_destination(projection.opaque, projection.dst_file_name) }; + } + + result +} + +/// Runs the source and destination policy around C-owned decompression +/// resources and format callbacks. The callback order mirrors +/// `FIO_decompressSrcFile()` and `FIO_decompressDstFile()` exactly. +#[no_mangle] +pub unsafe extern "C" fn FIO_rust_decompressFilename( + projection: *const FIO_rust_decompress_file_projection_t, +) -> c_int { + assert!(!projection.is_null()); + let projection = unsafe { &*projection }; + assert!(!projection.opaque.is_null()); + assert!(!projection.dst_file_name.is_null()); + assert!(!projection.src_file_name.is_null()); + + let open_source = projection + .open_source + .expect("source open callback is required"); + let mut file_size = FIO_RUST_DECOMPRESS_SRC_UNKNOWN_SIZE; + let mut source_is_regular = 0; + let open_status = unsafe { + open_source( + projection.opaque, + projection.src_file_name, + &mut file_size, + &mut source_is_regular, + ) + }; + if open_status != FIO_RUST_DECOMPRESS_SRC_OPEN_OK { + return 1; + } + + let async_mode = if file_size != FIO_RUST_DECOMPRESS_SRC_UNKNOWN_SIZE + && file_size < FIO_RUST_DECOMPRESS_SRC_ASYNC_THRESHOLD + { + 0 + } else { + 1 + }; + let set_async = projection + .set_async + .expect("async-selection callback is required"); + unsafe { set_async(projection.opaque, async_mode) }; + + let attach_source = projection + .attach_source + .expect("source attachment callback is required"); + unsafe { attach_source(projection.opaque) }; + + let mut result = unsafe { decompress_destination(projection, source_is_regular) }; + + let detach_source = projection + .detach_source + .expect("source detachment callback is required"); + unsafe { detach_source(projection.opaque) }; + + let close_source = projection + .close_source + .expect("source close callback is required"); + if unsafe { close_source(projection.opaque) } != 0 { + return 1; + } + + if projection.remove_source != 0 && result == 0 && projection.source_is_stdin == 0 { + let clear_handler = projection + .clear_handler + .expect("handler-clear callback is required"); + unsafe { clear_handler(projection.opaque) }; + + let remove_source = projection + .remove_source_file + .expect("source removal callback is required"); + if unsafe { remove_source(projection.opaque, projection.src_file_name) } != 0 { + result = 1; + } + } + + result +} + /// Iterate the files that share one already-open destination. The C /// callback retains source validation, resource state, diagnostics, and /// compression dispatch; Rust preserves the original non-short-circuiting @@ -2540,9 +2864,9 @@ pub unsafe extern "C" fn FIO_rust_compressMultipleSeparateFilenames( error } -/// Iterate files that share one already-open decompression destination. The -/// C callback retains source opening, format dispatch, diagnostics, and -/// `--rm`; Rust preserves the non-short-circuiting loop and counter updates. +/// Iterate files that share one already-open decompression destination. The C +/// callback retains private source/destination operations, format dispatch, +/// and diagnostics; Rust owns per-file ordering and removal policy. #[no_mangle] pub unsafe extern "C" fn FIO_rust_decompressMultipleFilenames( projection: *const FIO_rust_decompress_multiple_projection_t, @@ -2577,9 +2901,8 @@ pub unsafe extern "C" fn FIO_rust_decompressMultipleFilenames( } /// 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. +/// retains destination-name construction, private I/O, format dispatch, and +/// diagnostics; Rust owns per-file ordering and removal policy. #[no_mangle] pub unsafe extern "C" fn FIO_rust_decompressMultipleSeparateFilenames( projection: *const FIO_rust_decompress_multiple_separate_projection_t, @@ -4689,6 +5012,279 @@ mod tests { ); } + const DECOMPRESS_POLICY_OPEN_SOURCE: u8 = 1; + const DECOMPRESS_POLICY_ASYNC: u8 = 2; + const DECOMPRESS_POLICY_ATTACH_SOURCE: u8 = 3; + const DECOMPRESS_POLICY_OPEN_DESTINATION: u8 = 4; + const DECOMPRESS_POLICY_ATTACH_DESTINATION: u8 = 5; + const DECOMPRESS_POLICY_ADD_HANDLER: u8 = 6; + const DECOMPRESS_POLICY_DECOMPRESS: u8 = 7; + const DECOMPRESS_POLICY_CLEAR_HANDLER: u8 = 8; + const DECOMPRESS_POLICY_SET_FD_STAT: u8 = 9; + const DECOMPRESS_POLICY_CLOSE_DESTINATION: u8 = 10; + const DECOMPRESS_POLICY_UTIME_DESTINATION: u8 = 11; + const DECOMPRESS_POLICY_REMOVE_DESTINATION: u8 = 12; + const DECOMPRESS_POLICY_DETACH_SOURCE: u8 = 13; + const DECOMPRESS_POLICY_CLOSE_SOURCE: u8 = 14; + const DECOMPRESS_POLICY_REMOVE_SOURCE: u8 = 15; + + #[derive(Default)] + struct DecompressPolicyState { + events: Vec, + async_modes: Vec, + transfers: Vec, + destination_fds: Vec, + source_file_size: u64, + source_is_regular: c_int, + source_open_status: c_int, + source_close_status: c_int, + destination_open_status: c_int, + destination_close_status: c_int, + decompress_status: c_int, + remove_source_status: c_int, + } + + unsafe extern "C" fn decompress_policy_open_source( + opaque: *mut c_void, + _src_file_name: *const c_char, + file_size: *mut u64, + source_is_regular: *mut c_int, + ) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_OPEN_SOURCE); + unsafe { + *file_size = state.source_file_size; + *source_is_regular = state.source_is_regular; + } + state.source_open_status + } + + unsafe extern "C" fn decompress_policy_async(opaque: *mut c_void, async_mode: c_int) { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_ASYNC); + state.async_modes.push(async_mode); + } + + unsafe extern "C" fn decompress_policy_attach_source(opaque: *mut c_void) { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_ATTACH_SOURCE); + } + + unsafe extern "C" fn decompress_policy_open_destination( + opaque: *mut c_void, + _src_file_name: *const c_char, + _dst_file_name: *const c_char, + transfer_stat: c_int, + dst_fd: *mut c_int, + ) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_OPEN_DESTINATION); + state.transfers.push(transfer_stat); + unsafe { *dst_fd = 37 }; + state.destination_open_status + } + + unsafe extern "C" fn decompress_policy_attach_destination(opaque: *mut c_void) { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_ATTACH_DESTINATION); + } + + unsafe extern "C" fn decompress_policy_add_handler(opaque: *mut c_void) { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_ADD_HANDLER); + } + + unsafe extern "C" fn decompress_policy_decompress( + opaque: *mut c_void, + _dst_file_name: *const c_char, + _src_file_name: *const c_char, + ) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_DECOMPRESS); + state.decompress_status + } + + unsafe extern "C" fn decompress_policy_clear_handler(opaque: *mut c_void) { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_CLEAR_HANDLER); + } + + unsafe extern "C" fn decompress_policy_set_fd_stat( + opaque: *mut c_void, + dst_fd: c_int, + _dst_file_name: *const c_char, + ) { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_SET_FD_STAT); + state.destination_fds.push(dst_fd); + } + + unsafe extern "C" fn decompress_policy_close_destination(opaque: *mut c_void) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_CLOSE_DESTINATION); + state.destination_close_status + } + + unsafe extern "C" fn decompress_policy_utime_destination(opaque: *mut c_void) { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_UTIME_DESTINATION); + } + + unsafe extern "C" fn decompress_policy_remove_destination( + opaque: *mut c_void, + _dst_file_name: *const c_char, + ) { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_REMOVE_DESTINATION); + } + + unsafe extern "C" fn decompress_policy_detach_source(opaque: *mut c_void) { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_DETACH_SOURCE); + } + + unsafe extern "C" fn decompress_policy_close_source(opaque: *mut c_void) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_CLOSE_SOURCE); + state.source_close_status + } + + unsafe extern "C" fn decompress_policy_remove_source( + opaque: *mut c_void, + _src_file_name: *const c_char, + ) -> c_int { + let state = unsafe { &mut *opaque.cast::() }; + state.events.push(DECOMPRESS_POLICY_REMOVE_SOURCE); + state.remove_source_status + } + + fn decompress_policy_projection( + state: &mut DecompressPolicyState, + destination_already_open: c_int, + test_mode: c_int, + source_is_stdin: c_int, + destination_is_stdout: c_int, + remove_source: c_int, + ) -> FIO_rust_decompress_file_projection_t { + FIO_rust_decompress_file_projection_t { + opaque: (state as *mut DecompressPolicyState).cast(), + dst_file_name: c"destination".as_ptr(), + src_file_name: c"source".as_ptr(), + destination_already_open, + test_mode, + source_is_stdin, + destination_is_stdout, + remove_source, + open_source: Some(decompress_policy_open_source), + set_async: Some(decompress_policy_async), + attach_source: Some(decompress_policy_attach_source), + detach_source: Some(decompress_policy_detach_source), + close_source: Some(decompress_policy_close_source), + open_destination: Some(decompress_policy_open_destination), + attach_destination: Some(decompress_policy_attach_destination), + add_handler: Some(decompress_policy_add_handler), + decompress: Some(decompress_policy_decompress), + clear_handler: Some(decompress_policy_clear_handler), + set_fd_stat: Some(decompress_policy_set_fd_stat), + close_destination: Some(decompress_policy_close_destination), + utime_destination: Some(decompress_policy_utime_destination), + remove_destination: Some(decompress_policy_remove_destination), + remove_source_file: Some(decompress_policy_remove_source), + } + } + + #[test] + fn decompress_file_policy_orders_cleanup_and_successful_source_removal() { + let mut state = DecompressPolicyState { + source_file_size: FIO_RUST_DECOMPRESS_SRC_ASYNC_THRESHOLD - 1, + source_is_regular: 1, + ..DecompressPolicyState::default() + }; + let projection = decompress_policy_projection(&mut state, 0, 0, 0, 0, 1); + + assert_eq!(unsafe { FIO_rust_decompressFilename(&projection) }, 0); + assert_eq!(state.async_modes, vec![0]); + assert_eq!(state.transfers, vec![1]); + assert_eq!(state.destination_fds, vec![37]); + assert_eq!( + state.events, + vec![ + DECOMPRESS_POLICY_OPEN_SOURCE, + DECOMPRESS_POLICY_ASYNC, + DECOMPRESS_POLICY_ATTACH_SOURCE, + DECOMPRESS_POLICY_OPEN_DESTINATION, + DECOMPRESS_POLICY_ATTACH_DESTINATION, + DECOMPRESS_POLICY_ADD_HANDLER, + DECOMPRESS_POLICY_DECOMPRESS, + DECOMPRESS_POLICY_CLEAR_HANDLER, + DECOMPRESS_POLICY_SET_FD_STAT, + DECOMPRESS_POLICY_CLOSE_DESTINATION, + DECOMPRESS_POLICY_UTIME_DESTINATION, + DECOMPRESS_POLICY_DETACH_SOURCE, + DECOMPRESS_POLICY_CLOSE_SOURCE, + DECOMPRESS_POLICY_CLEAR_HANDLER, + DECOMPRESS_POLICY_REMOVE_SOURCE, + ] + ); + } + + #[test] + fn decompress_file_policy_removes_failed_destination_and_skips_source_removal() { + let mut state = DecompressPolicyState { + source_file_size: FIO_RUST_DECOMPRESS_SRC_ASYNC_THRESHOLD, + decompress_status: 1, + destination_close_status: 1, + ..DecompressPolicyState::default() + }; + let projection = decompress_policy_projection(&mut state, 0, 0, 0, 0, 1); + + assert_eq!(unsafe { FIO_rust_decompressFilename(&projection) }, 1); + assert_eq!(state.async_modes, vec![1]); + assert_eq!(state.transfers, vec![0]); + assert!(!state.events.contains(&DECOMPRESS_POLICY_REMOVE_SOURCE)); + assert_eq!( + state.events, + vec![ + DECOMPRESS_POLICY_OPEN_SOURCE, + DECOMPRESS_POLICY_ASYNC, + DECOMPRESS_POLICY_ATTACH_SOURCE, + DECOMPRESS_POLICY_OPEN_DESTINATION, + DECOMPRESS_POLICY_ATTACH_DESTINATION, + DECOMPRESS_POLICY_ADD_HANDLER, + DECOMPRESS_POLICY_DECOMPRESS, + DECOMPRESS_POLICY_CLEAR_HANDLER, + DECOMPRESS_POLICY_CLOSE_DESTINATION, + DECOMPRESS_POLICY_REMOVE_DESTINATION, + DECOMPRESS_POLICY_DETACH_SOURCE, + DECOMPRESS_POLICY_CLOSE_SOURCE, + ] + ); + } + + #[test] + fn decompress_file_policy_does_not_remove_source_after_close_error() { + let mut state = DecompressPolicyState { + source_file_size: FIO_RUST_DECOMPRESS_SRC_UNKNOWN_SIZE, + source_close_status: 1, + ..DecompressPolicyState::default() + }; + let projection = decompress_policy_projection(&mut state, 1, 0, 1, 1, 1); + + assert_eq!(unsafe { FIO_rust_decompressFilename(&projection) }, 1); + assert_eq!(state.async_modes, vec![1]); + assert_eq!( + state.events, + vec![ + DECOMPRESS_POLICY_OPEN_SOURCE, + DECOMPRESS_POLICY_ASYNC, + DECOMPRESS_POLICY_ATTACH_SOURCE, + DECOMPRESS_POLICY_DECOMPRESS, + DECOMPRESS_POLICY_DETACH_SOURCE, + DECOMPRESS_POLICY_CLOSE_SOURCE, + ] + ); + } + fn run_zstd_adapt(policy: c_int, projection: &FIO_rust_zstd_adapt_projection_t) -> c_int { unsafe { FIO_rust_zstd_adapt(policy, projection) } }