Files
zstd-rs/programs/fileio.c
T
ddidderr 05561c7494 feat(cli): move gzip codec loop into Rust
Move the gzip member read, zlib deflate, finish, output-enqueue, progress,
and accounting loop into Rust. Keep zlib, asynchronous pool, sparse-output,
and diagnostic ownership in C through an explicit projection so optional gzip
builds preserve the existing CLI behavior and error mapping. Add focused tests
for callback ordering, partial input consumption, finish retries, accounting,
and initialization failures.

Test Plan:
- cargo test --manifest-path rust/Cargo.toml --all-targets -- --test-threads=1
- cargo test --manifest-path rust/cli/Cargo.toml --all-targets -- --test-threads=1
- make -B -C programs -j2 zstd
- make -B -C tests -j2 test-cli-tests test-zstd
- FUZZERTEST=-T5s make -B -C tests -j2 test-fuzzer
- cargo clippy --manifest-path rust/Cargo.toml --tests -- -D warnings
2026-07-18 23:07:51 +02:00

3459 lines
135 KiB
C

/*
* Copyright (c) Meta Platforms, Inc. and affiliates.
* All rights reserved.
*
* This source code is licensed under both the BSD-style license (found in the
* LICENSE file in the root directory of this source tree) and the GPLv2 (found
* in the COPYING file in the root directory of this source tree).
* You may select, at your option, one of the above-listed licenses.
*/
/* *************************************
* Compiler Options
***************************************/
#ifdef _MSC_VER /* Visual */
# pragma warning(disable : 4127) /* disable: C4127: conditional expression is constant */
# pragma warning(disable : 4204) /* non-constant aggregate initializer */
#endif
#if defined(__MINGW32__) && !defined(_POSIX_SOURCE)
# define _POSIX_SOURCE 1 /* disable %llu warnings with MinGW on Windows */
#endif
/*-*************************************
* Includes
***************************************/
#include "platform.h" /* Large Files support, SET_BINARY_MODE */
#include "util.h" /* UTIL_getFileSize, UTIL_isRegularFile, UTIL_isSameFile */
#include <stdio.h> /* fprintf, open, fdopen, fread, _fileno, stdin, stdout */
#include <stdlib.h> /* malloc, free */
#include <string.h> /* strcmp, strlen */
#include <stddef.h> /* offsetof */
#include <time.h> /* clock_t, to measure process time */
#include <fcntl.h> /* O_WRONLY */
#include <assert.h>
#include <errno.h> /* errno */
#include <limits.h> /* INT_MAX */
#include <signal.h>
#include "timefn.h" /* UTIL_getTime, UTIL_clockSpanMicro */
#if defined (_MSC_VER)
# include <sys/stat.h>
# include <io.h>
#endif
#include "fileio.h"
#include "fileio_asyncio.h"
#include "fileio_common.h"
FIO_display_prefs_t g_display_prefs = {2, FIO_ps_auto};
UTIL_time_t g_displayClock = UTIL_TIME_INITIALIZER;
#define ZSTD_STATIC_LINKING_ONLY /* ZSTD_magicNumber, ZSTD_frameHeaderSize_max */
#include "../lib/zstd.h"
#include "../lib/zstd_errors.h" /* ZSTD_error_frameParameter_windowTooLarge */
#if defined(ZSTD_GZCOMPRESS) || defined(ZSTD_GZDECOMPRESS)
# include <zlib.h>
# if !defined(z_const)
# define z_const
# endif
#endif
#if defined(ZSTD_LZMACOMPRESS) || defined(ZSTD_LZMADECOMPRESS)
# include <lzma.h>
#endif
#define LZ4_MAGICNUMBER 0x184D2204
#if defined(ZSTD_LZ4COMPRESS) || defined(ZSTD_LZ4DECOMPRESS)
# define LZ4F_ENABLE_OBSOLETE_ENUMS
# include <lz4frame.h>
# include <lz4.h>
#endif
char const* FIO_zlibVersion(void)
{
#if defined(ZSTD_GZCOMPRESS) || defined(ZSTD_GZDECOMPRESS)
return zlibVersion();
#else
return "Unsupported";
#endif
}
char const* FIO_lz4Version(void)
{
#if defined(ZSTD_LZ4COMPRESS) || defined(ZSTD_LZ4DECOMPRESS)
/* LZ4_versionString() added in v1.7.3 */
# if LZ4_VERSION_NUMBER >= 10703
return LZ4_versionString();
# else
# define ZSTD_LZ4_VERSION LZ4_VERSION_MAJOR.LZ4_VERSION_MINOR.LZ4_VERSION_RELEASE
# define ZSTD_LZ4_VERSION_STRING ZSTD_EXPAND_AND_QUOTE(ZSTD_LZ4_VERSION)
return ZSTD_LZ4_VERSION_STRING;
# endif
#else
return "Unsupported";
#endif
}
char const* FIO_lzmaVersion(void)
{
#if defined(ZSTD_LZMACOMPRESS) || defined(ZSTD_LZMADECOMPRESS)
return lzma_version_string();
#else
return "Unsupported";
#endif
}
/*-*************************************
* Constants
***************************************/
#define ADAPT_WINDOWLOG_DEFAULT 23 /* 8 MB */
#define DICTSIZE_MAX (32 MB) /* protection against large input (attack scenario) */
#define FNSPACE 30
/* Default file permissions 0666 (modulated by umask) */
/* Temporary restricted file permissions are used when we're going to
* chmod/chown at the end of the operation. */
#if !defined(_WIN32)
/* These macros aren't defined on windows. */
#define DEFAULT_FILE_PERMISSIONS (S_IRUSR|S_IWUSR|S_IRGRP|S_IWGRP|S_IROTH|S_IWOTH)
#define TEMPORARY_FILE_PERMISSIONS (S_IRUSR|S_IWUSR)
#else
#define DEFAULT_FILE_PERMISSIONS (0666)
#define TEMPORARY_FILE_PERMISSIONS (0600)
#endif
/*-************************************
* Signal (Ctrl-C trapping)
**************************************/
static const char* g_artefact = NULL;
static void INThandler(int sig)
{
assert(sig==SIGINT); (void)sig;
#if !defined(_MSC_VER)
signal(sig, SIG_IGN); /* this invocation generates a buggy warning in Visual Studio */
#endif
if (g_artefact) {
assert(UTIL_isRegularFile(g_artefact));
remove(g_artefact);
}
DISPLAY("\n");
exit(2);
}
static void addHandler(char const* dstFileName)
{
if (UTIL_isRegularFile(dstFileName)) {
g_artefact = dstFileName;
signal(SIGINT, INThandler);
} else {
g_artefact = NULL;
}
}
/* Idempotent */
static void clearHandler(void)
{
if (g_artefact) signal(SIGINT, SIG_DFL);
g_artefact = NULL;
}
/*-*********************************************************
* Termination signal trapping (Print debug stack trace)
***********************************************************/
#if defined(__has_feature) && !defined(BACKTRACE_ENABLE) /* Clang compiler */
# if (__has_feature(address_sanitizer))
# define BACKTRACE_ENABLE 0
# endif /* __has_feature(address_sanitizer) */
#elif defined(__SANITIZE_ADDRESS__) && !defined(BACKTRACE_ENABLE) /* GCC compiler */
# define BACKTRACE_ENABLE 0
#endif
#if !defined(BACKTRACE_ENABLE)
/* automatic detector : backtrace enabled by default on linux+glibc and osx */
# if (defined(__linux__) && (defined(__GLIBC__) && !defined(__UCLIBC__))) \
|| (defined(__APPLE__) && defined(__MACH__))
# define BACKTRACE_ENABLE 1
# else
# define BACKTRACE_ENABLE 0
# endif
#endif
/* note : after this point, BACKTRACE_ENABLE is necessarily defined */
#if BACKTRACE_ENABLE
#include <execinfo.h> /* backtrace, backtrace_symbols */
#define MAX_STACK_FRAMES 50
static void ABRThandler(int sig) {
const char* name;
void* addrlist[MAX_STACK_FRAMES];
char** symbollist;
int addrlen, i;
switch (sig) {
case SIGABRT: name = "SIGABRT"; break;
case SIGFPE: name = "SIGFPE"; break;
case SIGILL: name = "SIGILL"; break;
case SIGINT: name = "SIGINT"; break;
case SIGSEGV: name = "SIGSEGV"; break;
default: name = "UNKNOWN";
}
DISPLAY("Caught %s signal, printing stack:\n", name);
/* Retrieve current stack addresses. */
addrlen = backtrace(addrlist, MAX_STACK_FRAMES);
if (addrlen == 0) {
DISPLAY("\n");
return;
}
/* Create readable strings to each frame. */
symbollist = backtrace_symbols(addrlist, addrlen);
/* Print the stack trace, excluding calls handling the signal. */
for (i = ZSTD_START_SYMBOLLIST_FRAME; i < addrlen; i++) {
DISPLAY("%s\n", symbollist[i]);
}
free(symbollist);
/* Reset and raise the signal so default handler runs. */
signal(sig, SIG_DFL);
raise(sig);
}
#endif
void FIO_addAbortHandler(void)
{
#if BACKTRACE_ENABLE
signal(SIGABRT, ABRThandler);
signal(SIGFPE, ABRThandler);
signal(SIGILL, ABRThandler);
signal(SIGSEGV, ABRThandler);
signal(SIGBUS, ABRThandler);
#endif
}
/*-*************************************
* Parameters: FIO_ctx_t
***************************************/
/* typedef'd to FIO_ctx_t within fileio.h */
struct FIO_ctx_s {
/* file i/o info */
int nbFilesTotal;
int hasStdinInput;
int hasStdoutOutput;
/* file i/o state */
int currFileIdx;
int nbFilesProcessed;
size_t totalBytesInput;
size_t totalBytesOutput;
};
/* Keep the Rust preference/context layouts in lock-step with these C
* definitions. The Rust side uses the same #[repr(C)] field order. */
typedef char FIO_rust_prefs_compression_type_offset[
(offsetof(FIO_prefs_t, compressionType) == 0) ? 1 : -1];
typedef char FIO_rust_prefs_stream_size_offset[
(offsetof(FIO_prefs_t, streamSrcSize) == 16 * sizeof(int)) ? 1 : -1];
typedef char FIO_rust_prefs_target_size_offset[
(offsetof(FIO_prefs_t, targetCBlockSize)
== 16 * sizeof(int) + sizeof(size_t)) ? 1 : -1];
typedef char FIO_rust_prefs_src_hint_offset[
(offsetof(FIO_prefs_t, srcSizeHint)
== 16 * sizeof(int) + 2 * sizeof(size_t)) ? 1 : -1];
typedef char FIO_rust_prefs_mmap_dict_offset[
(offsetof(FIO_prefs_t, mmapDict)
== 16 * sizeof(int) + 2 * sizeof(size_t) + 13 * sizeof(int)) ? 1 : -1];
typedef char FIO_rust_prefs_size[
(sizeof(FIO_prefs_t)
== 16 * sizeof(int) + 2 * sizeof(size_t) + 14 * sizeof(int)) ? 1 : -1];
typedef char FIO_rust_ctx_total_input_offset[
(offsetof(FIO_ctx_t, totalBytesInput)
== ((5 * sizeof(int) + sizeof(size_t) - 1) / sizeof(size_t)) * sizeof(size_t)) ? 1 : -1];
typedef char FIO_rust_ctx_total_output_offset[
(offsetof(FIO_ctx_t, totalBytesOutput)
== offsetof(FIO_ctx_t, totalBytesInput) + sizeof(size_t)) ? 1 : -1];
typedef char FIO_rust_ctx_size[
(sizeof(FIO_ctx_t)
== offsetof(FIO_ctx_t, totalBytesOutput) + sizeof(size_t)) ? 1 : -1];
typedef char FIO_rust_compression_params_strategy_offset[
(offsetof(ZSTD_compressionParameters, strategy)
== 6 * sizeof(unsigned)) ? 1 : -1];
typedef char FIO_rust_compression_params_size[
(sizeof(ZSTD_compressionParameters)
== 7 * sizeof(unsigned)) ? 1 : -1];
int FIO_rust_shouldDisplayFileSummary(const FIO_ctx_t* fCtx);
int FIO_rust_shouldDisplayMultipleFileSummary(const FIO_ctx_t* fCtx);
enum {
FIO_RUST_MULTI_FILES_ACTION_PROCEED = 0,
FIO_RUST_MULTI_FILES_ACTION_FATAL_STDOUT_REMOVE = 1,
FIO_RUST_MULTI_FILES_ACTION_FATAL_TEST_REMOVE = 2,
FIO_RUST_MULTI_FILES_ACTION_DISABLE_REMOVE = 3,
FIO_RUST_MULTI_FILES_ACTION_QUIET_ABORT = 4,
FIO_RUST_MULTI_FILES_ACTION_CONFIRM = 5
};
int FIO_rust_multiFilesConcatAction(int nbFilesTotal,
int hasStdoutOutput,
int testMode,
int hasOutputFile,
int removeSrcFile,
int overwrite,
int displayLevel,
int displayLevelCutoff);
static int FIO_shouldDisplayFileSummary(FIO_ctx_t const* fCtx)
{
return FIO_rust_shouldDisplayFileSummary(fCtx);
}
static int FIO_shouldDisplayMultipleFileSummary(FIO_ctx_t const* fCtx)
{
return FIO_rust_shouldDisplayMultipleFileSummary(fCtx);
}
/*-*************************************
* Parameters: Initialization
***************************************/
#define FIO_OVERLAP_LOG_NOTSET 9999
#define FIO_LDM_PARAM_NOTSET 9999
/* These policy and filesystem leaves are implemented by the Rust CLI archive.
* The declarations in fileio.h remain the C ABI shims while file-I/O
* orchestration and diagnostics stay here. */
int FIO_rust_checkFilenameCollisions(const char** filenameTable, unsigned nbFiles);
unsigned FIO_rust_highbit64(unsigned long long v);
unsigned long long FIO_rust_getLargestFileSize(const char** inFileNames, unsigned nbFiles);
int FIO_rust_adjustParamsForPatchFromMode(FIO_prefs_t* prefs,
ZSTD_compressionParameters* comprParams,
unsigned long long dictSize,
unsigned long long maxSrcFileSize,
ZSTD_compressionParameters cParams,
unsigned* fileWindowLog,
int* autoLdm,
int* optimalParser);
void FIO_rust_setInBuffer(ZSTD_inBuffer* output, const void* buf, size_t s, size_t pos);
void FIO_rust_setOutBuffer(ZSTD_outBuffer* output, void* buf, size_t s, size_t pos);
const char* FIO_rust_determineCompressedName(const char* srcFileName, const char* outDirName, const char* suffix);
const char* FIO_rust_determineDstName(const char* srcFileName, const char* outDirName,
const char* const* suffixList, const char* suffixListStr);
int FIO_rust_adjustMemLimitForPatchFromMode(FIO_prefs_t* prefs,
unsigned long long dictSize,
unsigned long long maxSrcFileSize);
int FIO_rust_getDictFileStat(const char* fileName, stat_t* statBuf);
enum {
FIO_RUST_OPEN_SRC_SUCCESS = 0,
FIO_RUST_OPEN_SRC_STAT_FAILED = 1,
FIO_RUST_OPEN_SRC_NON_REGULAR = 2,
FIO_RUST_OPEN_SRC_FOPEN_FAILED = 3,
};
int FIO_rust_openSrcFile(int allowBlockDevices,
const char* srcFileName,
stat_t* statbuf,
FILE** outFile);
enum {
FIO_RUST_OPEN_DST_SUCCESS = 0,
FIO_RUST_OPEN_DST_SETBUF_FAILED = 1,
FIO_RUST_OPEN_DST_TEST_MODE = 2,
FIO_RUST_OPEN_DST_STDOUT = 3,
FIO_RUST_OPEN_DST_SAME_FILE = 4,
FIO_RUST_OPEN_DST_EXISTING = 5,
FIO_RUST_OPEN_DST_NULL_DEVICE_REGULAR = 6,
FIO_RUST_OPEN_DST_OPEN_FAILED = 7,
};
int FIO_rust_openDstFile(int testMode,
int allowExisting,
const char* srcFileName,
const char* dstFileName,
int mode,
int* isDstRegFile,
FILE** outFile);
int FIO_rust_setDictBufferMalloc(const char* fileName,
unsigned long long expectedFileSize,
size_t maxSize,
void** buffer,
size_t* loadedSize);
#if (PLATFORM_POSIX_VERSION > 0) || defined(_MSC_VER) || defined(_WIN32)
enum {
FIO_DICT_MMAP_SUCCESS = 0,
FIO_DICT_MMAP_OPEN_FAILED = 1,
FIO_DICT_MMAP_TOO_LARGE = 2,
FIO_DICT_MMAP_FAILED = 3,
FIO_DICT_MMAP_VIEW_FAILED = 4,
};
int FIO_rust_setDictBufferMMap(const char* fileName,
unsigned long long expectedFileSize,
size_t maxSize,
void** buffer,
size_t* mappedSize,
void** dictHandle);
#endif
void FIO_rust_freeDict(int dictBufferType,
void** buffer,
size_t* bufferSize,
void** dictHandle);
int FIO_rust_removeFile(const char* path);
int FIO_rust_passThrough(ReadPoolCtx_t* readCtx, WritePoolCtx_t* writeCtx);
enum {
FIO_RUST_ZSTD_FRAME_OK = 0,
FIO_RUST_ZSTD_FRAME_DECODING_ERROR = 1,
FIO_RUST_ZSTD_FRAME_PREMATURE_END = 2,
};
typedef void (*FIO_rust_frame_progress_fn)(void* opaque,
const char* srcFileName,
U64 decodedSize);
int FIO_rust_decompressZstdFrames(void* fCtx,
void* dctx,
ReadPoolCtx_t* readCtx,
WritePoolCtx_t* writeCtx,
const char* srcFileName,
U64 alreadyDecoded,
U64* decodedSize,
size_t* zstdError,
FIO_rust_frame_progress_fn progress);
enum {
FIO_RUST_DECOMPRESS_OK = 0,
FIO_RUST_DECOMPRESS_PASS_THROUGH = 1,
FIO_RUST_DECOMPRESS_EMPTY_INPUT = 2,
FIO_RUST_DECOMPRESS_SHORT_INPUT = 3,
FIO_RUST_DECOMPRESS_GZIP_UNSUPPORTED = 4,
FIO_RUST_DECOMPRESS_LZMA_UNSUPPORTED = 5,
FIO_RUST_DECOMPRESS_LZ4_UNSUPPORTED = 6,
FIO_RUST_DECOMPRESS_FRAME_ERROR = 7,
FIO_RUST_DECOMPRESS_UNSUPPORTED_FORMAT = 8,
FIO_RUST_DECOMPRESS_PASS_THROUGH_ERROR = 9,
FIO_RUST_DECOMPRESS_ZSTD_UNSUPPORTED = 10
};
typedef int (*FIO_rust_decompress_frame_fn)(void* opaque,
const char* srcFileName,
U64 alreadyDecoded,
U64* frameSize,
size_t* errorCode,
int mode);
typedef int (*FIO_rust_pass_through_fn)(void* opaque);
typedef struct {
void* opaque;
FIO_rust_decompress_frame_fn decode_zstd;
FIO_rust_decompress_frame_fn decode_gzip;
FIO_rust_decompress_frame_fn decode_lzma;
FIO_rust_decompress_frame_fn decode_lz4;
FIO_rust_pass_through_fn pass_through;
} FIO_rust_decompress_callbacks_t;
int FIO_rust_decompressFrames(ReadPoolCtx_t* readCtx,
const char* srcFileName,
int passThrough,
U64* decodedSize,
const FIO_rust_decompress_callbacks_t* callbacks);
void FIO_rust_displayCompressionParameters(const FIO_prefs_t* prefs);
#ifdef ZSTD_LZ4COMPRESS
int FIO_rust_LZ4_GetBlockSize_FromBlockId(int id);
#endif
/*-*************************************
* Parameters: Display Options
***************************************/
/*-*************************************
* Parameters: Setters
***************************************/
/* FIO_prefs_t functions */
/*-*************************************
* Functions
***************************************/
/** FIO_removeFile() :
* @result : Unlink `fileName`, even if it's read-only */
static int FIO_removeFile(const char* path)
{
enum {
FIO_REMOVE_SUCCESS = 0,
FIO_REMOVE_STAT_FAILED = 1,
FIO_REMOVE_NON_REGULAR = 2,
};
int const status = FIO_rust_removeFile(path);
if (status == FIO_REMOVE_STAT_FAILED) {
DISPLAYLEVEL(2, "zstd: Failed to stat %s while trying to remove it\n", path);
return 0;
}
if (status == FIO_REMOVE_NON_REGULAR) {
DISPLAYLEVEL(2, "zstd: Refusing to remove non-regular file %s\n", path);
return 0;
}
return status == FIO_REMOVE_SUCCESS ? 0 : 1;
}
/** FIO_openSrcFile() :
* condition : `srcFileName` must be non-NULL. `prefs` may be NULL.
* @result : FILE* to `srcFileName`, or NULL if it fails */
static FILE* FIO_openSrcFile(const FIO_prefs_t* const prefs, const char* srcFileName, stat_t* statbuf)
{
int allowBlockDevices = prefs != NULL ? prefs->allowBlockDevices : 0;
assert(srcFileName != NULL);
assert(statbuf != NULL);
if (!strcmp (srcFileName, stdinmark)) {
DISPLAYLEVEL(4,"Using stdin for input \n");
SET_BINARY_MODE(stdin);
return stdin;
}
{ FILE* f = NULL;
int const status = FIO_rust_openSrcFile(allowBlockDevices,
srcFileName,
statbuf,
&f);
if (status == FIO_RUST_OPEN_SRC_STAT_FAILED) {
DISPLAYLEVEL(1, "zstd: can't stat %s : %s -- ignored \n",
srcFileName, strerror(errno));
return NULL;
}
if (status == FIO_RUST_OPEN_SRC_NON_REGULAR) {
DISPLAYLEVEL(1, "zstd: %s is not a regular file -- ignored \n",
srcFileName);
return NULL;
}
if (status == FIO_RUST_OPEN_SRC_FOPEN_FAILED) {
DISPLAYLEVEL(1, "zstd: %s: %s \n", srcFileName, strerror(errno));
return NULL;
}
assert(status == FIO_RUST_OPEN_SRC_SUCCESS);
return f;
}
}
/** FIO_openDstFile() :
* condition : `dstFileName` must be non-NULL.
* @result : FILE* to `dstFileName`, or NULL if it fails */
static FILE*
FIO_openDstFile(FIO_ctx_t* fCtx, FIO_prefs_t* const prefs,
const char* srcFileName, const char* dstFileName,
const int mode)
{
int isDstRegFile = 0;
int allowExisting = 0;
int sparseAdjusted = 0;
int status;
FILE* f = NULL;
status = FIO_rust_openDstFile(prefs->testMode, allowExisting,
srcFileName, dstFileName, mode,
&isDstRegFile, &f);
if (status == FIO_RUST_OPEN_DST_TEST_MODE)
return NULL; /* do not open file in test mode */
assert(dstFileName != NULL);
for (;;) {
if (status == FIO_RUST_OPEN_DST_STDOUT) {
DISPLAYLEVEL(4,"Using stdout for output \n");
SET_BINARY_MODE(stdout);
if (prefs->sparseFileSupport == 1) {
prefs->sparseFileSupport = 0;
DISPLAYLEVEL(4, "Sparse File Support is automatically disabled on stdout ; try --sparse \n");
}
return stdout;
}
if (status == FIO_RUST_OPEN_DST_SAME_FILE) {
DISPLAYLEVEL(1, "zstd: Refusing to open an output file which will overwrite the input file \n");
return NULL;
}
/* Keep the C-owned preference mutation and its diagnostics here. Do
* this only for the first classification: after an existing file is
* removed, the second Rust call must not reclassify it as a missing
* destination and emit a different sparse-mode message. */
if (!sparseAdjusted) {
if (prefs->sparseFileSupport == 1) {
prefs->sparseFileSupport = ZSTD_SPARSE_DEFAULT;
if (!isDstRegFile) {
prefs->sparseFileSupport = 0;
DISPLAYLEVEL(4, "Sparse File Support is disabled when output is not a file \n");
}
}
sparseAdjusted = 1;
}
if (status == FIO_RUST_OPEN_DST_NULL_DEVICE_REGULAR) {
EXM_THROW(40, "%s is unexpectedly categorized as a regular file",
dstFileName);
}
if (status == FIO_RUST_OPEN_DST_EXISTING) {
if (!prefs->overwrite) {
if (g_display_prefs.displayLevel <= 1) {
/* No interaction possible */
DISPLAYLEVEL(1, "zstd: %s already exists; not overwritten \n",
dstFileName);
return NULL;
}
DISPLAY("zstd: %s already exists; ", dstFileName);
if (UTIL_requireUserConfirmation("overwrite (y/n) ? ", "Not overwritten \n", "yY", fCtx->hasStdinInput))
return NULL;
}
/* Keep the existing C wrapper so its stat/non-regular diagnostics
* remain unchanged. The next Rust call opens with O_TRUNC. */
FIO_removeFile(dstFileName);
allowExisting = 1;
f = NULL;
status = FIO_rust_openDstFile(prefs->testMode, allowExisting,
srcFileName, dstFileName, mode,
&isDstRegFile, &f);
continue;
}
if (status == FIO_RUST_OPEN_DST_OPEN_FAILED) {
DISPLAYLEVEL(1, "zstd: %s: %s\n", dstFileName, strerror(errno));
return NULL;
}
if (status == FIO_RUST_OPEN_DST_SETBUF_FAILED)
DISPLAYLEVEL(2, "Warning: setvbuf failed for %s\n", dstFileName);
assert(status == FIO_RUST_OPEN_DST_SUCCESS
|| status == FIO_RUST_OPEN_DST_SETBUF_FAILED);
return f;
}
}
/* FIO_getDictFileStat() :
*/
static void FIO_getDictFileStat(const char* fileName, stat_t* dictFileStat) {
enum {
FIO_DICT_STAT_SUCCESS = 0,
FIO_DICT_STAT_FAILED = 1,
FIO_DICT_STAT_NON_REGULAR = 2,
};
int status;
assert(dictFileStat != NULL);
if (fileName == NULL) return;
status = FIO_rust_getDictFileStat(fileName, dictFileStat);
if (status == FIO_DICT_STAT_FAILED) {
EXM_THROW(31, "Stat failed on dictionary file %s: %s", fileName, strerror(errno));
}
if (status == FIO_DICT_STAT_NON_REGULAR) {
EXM_THROW(32, "Dictionary %s must be a regular file.", fileName);
}
assert(status == FIO_DICT_STAT_SUCCESS);
}
/* FIO_setDictBufferMalloc() :
* allocates a buffer, pointed by `dict->dictBuffer`,
* loads `filename` content into it, up to DICTSIZE_MAX bytes.
* @return : loaded size
* if fileName==NULL, returns 0 and a NULL pointer
*/
static size_t FIO_setDictBufferMalloc(FIO_Dict_t* dict, const char* fileName, FIO_prefs_t* const prefs, stat_t* dictFileStat)
{
U64 fileSize;
size_t loadedSize = 0;
size_t const dictSizeMax = prefs->patchFromMode ? prefs->memLimit : DICTSIZE_MAX;
void** bufferPtr = &dict->dictBuffer;
enum {
FIO_DICT_LOAD_SUCCESS = 0,
FIO_DICT_LOAD_OPEN_FAILED = 1,
FIO_DICT_LOAD_TOO_LARGE = 2,
FIO_DICT_LOAD_ALLOCATION_FAILED = 3,
FIO_DICT_LOAD_READ_FAILED = 4,
};
int status;
assert(bufferPtr != NULL);
assert(dictFileStat != NULL);
*bufferPtr = NULL;
if (fileName == NULL) return 0;
DISPLAYLEVEL(4,"Loading %s as dictionary \n", fileName);
fileSize = UTIL_getFileSizeStat(dictFileStat);
if (fileSize > dictSizeMax) {
EXM_THROW(34, "Dictionary file %s is too large (> %u bytes)",
fileName, (unsigned)dictSizeMax); /* avoid extreme cases */
}
status = FIO_rust_setDictBufferMalloc(fileName,
(unsigned long long)fileSize,
dictSizeMax,
bufferPtr,
&loadedSize);
if (status == FIO_DICT_LOAD_SUCCESS) return loadedSize;
if (status == FIO_DICT_LOAD_OPEN_FAILED) {
EXM_THROW(33, "Couldn't open dictionary %s: %s", fileName, strerror(errno));
}
if (status == FIO_DICT_LOAD_TOO_LARGE) {
EXM_THROW(34, "Dictionary file %s is too large (> %u bytes)",
fileName, (unsigned)dictSizeMax); /* avoid extreme cases */
}
if (status == FIO_DICT_LOAD_ALLOCATION_FAILED) {
EXM_THROW(34, "%s", strerror(errno));
}
if (status == FIO_DICT_LOAD_READ_FAILED) {
EXM_THROW(35, "Error reading dictionary file %s : %s",
fileName, strerror(errno));
}
assert(0); /* unexpected Rust status */
return 0;
}
#if (PLATFORM_POSIX_VERSION > 0)
static size_t FIO_setDictBufferMMap(FIO_Dict_t* dict, const char* fileName, FIO_prefs_t* const prefs, stat_t* dictFileStat)
{
U64 fileSize;
size_t const dictSizeMax = prefs->patchFromMode ? prefs->memLimit : DICTSIZE_MAX;
void** bufferPtr = &dict->dictBuffer;
int status;
assert(bufferPtr != NULL);
assert(dictFileStat != NULL);
*bufferPtr = NULL;
if (fileName == NULL) return 0;
DISPLAYLEVEL(4,"Loading %s as dictionary \n", fileName);
fileSize = UTIL_getFileSizeStat(dictFileStat);
if (fileSize > dictSizeMax) {
EXM_THROW(34, "Dictionary file %s is too large (> %u bytes)",
fileName, (unsigned)dictSizeMax); /* avoid extreme cases */
}
status = FIO_rust_setDictBufferMMap(fileName,
(unsigned long long)fileSize,
dictSizeMax,
bufferPtr,
&dict->dictBufferSize,
NULL);
if (status == FIO_DICT_MMAP_SUCCESS) return dict->dictBufferSize;
if (status == FIO_DICT_MMAP_OPEN_FAILED) {
EXM_THROW(33, "Couldn't open dictionary %s: %s", fileName, strerror(errno));
}
if (status == FIO_DICT_MMAP_TOO_LARGE) {
EXM_THROW(34, "Dictionary file %s is too large (> %u bytes)",
fileName, (unsigned)dictSizeMax); /* avoid extreme cases */
}
if (status == FIO_DICT_MMAP_FAILED) {
EXM_THROW(34, "%s", strerror(errno));
}
assert(0); /* unexpected Rust status */
return 0;
}
#elif defined(_MSC_VER) || defined(_WIN32)
static size_t FIO_setDictBufferMMap(FIO_Dict_t* dict, const char* fileName, FIO_prefs_t* const prefs, stat_t* dictFileStat)
{
U64 fileSize;
size_t const dictSizeMax = prefs->patchFromMode ? prefs->memLimit : DICTSIZE_MAX;
void** bufferPtr = &dict->dictBuffer;
int status;
assert(bufferPtr != NULL);
assert(dictFileStat != NULL);
*bufferPtr = NULL;
if (fileName == NULL) return 0;
DISPLAYLEVEL(4,"Loading %s as dictionary \n", fileName);
fileSize = UTIL_getFileSizeStat(dictFileStat);
if (fileSize > dictSizeMax) {
EXM_THROW(34, "Dictionary file %s is too large (> %u bytes)",
fileName, (unsigned)dictSizeMax); /* avoid extreme cases */
}
status = FIO_rust_setDictBufferMMap(fileName,
(unsigned long long)fileSize,
dictSizeMax,
bufferPtr,
&dict->dictBufferSize,
(void**)&dict->dictHandle);
if (status == FIO_DICT_MMAP_SUCCESS) return dict->dictBufferSize;
if (status == FIO_DICT_MMAP_OPEN_FAILED) {
EXM_THROW(33, "Couldn't open dictionary %s: %s", fileName, strerror(errno));
}
if (status == FIO_DICT_MMAP_TOO_LARGE) {
EXM_THROW(34, "Dictionary file %s is too large (> %u bytes)",
fileName, (unsigned)dictSizeMax); /* avoid extreme cases */
}
if (status == FIO_DICT_MMAP_FAILED) {
EXM_THROW(35, "Couldn't map dictionary %s: %s", fileName, strerror(errno));
}
if (status == FIO_DICT_MMAP_VIEW_FAILED) {
EXM_THROW(36, "%s", strerror(errno));
}
assert(0); /* unexpected Rust status */
return 0;
}
#else
static size_t FIO_setDictBufferMMap(FIO_Dict_t* dict, const char* fileName, FIO_prefs_t* const prefs, stat_t* dictFileStat)
{
return FIO_setDictBufferMalloc(dict, fileName, prefs, dictFileStat);
}
#endif
static void FIO_freeDict(FIO_Dict_t* dict) {
if (dict->dictBufferType != FIO_mallocDict
&& dict->dictBufferType != FIO_mmapDict) {
assert(0); /* Should not reach this case */
return;
}
FIO_rust_freeDict((int)dict->dictBufferType,
&dict->dictBuffer,
&dict->dictBufferSize,
#if defined(_MSC_VER) || defined(_WIN32)
(void**)&dict->dictHandle
#else
NULL
#endif
);
}
static void FIO_initDict(FIO_Dict_t* dict, const char* fileName, FIO_prefs_t* const prefs, stat_t* dictFileStat, FIO_dictBufferType_t dictBufferType) {
#if defined(_MSC_VER) || defined(_WIN32)
dict->dictHandle = NULL;
#endif
dict->dictBufferType = dictBufferType;
if (dict->dictBufferType == FIO_mallocDict) {
dict->dictBufferSize = FIO_setDictBufferMalloc(dict, fileName, prefs, dictFileStat);
} else if (dict->dictBufferType == FIO_mmapDict) {
dict->dictBufferSize = FIO_setDictBufferMMap(dict, fileName, prefs, dictFileStat);
} else {
assert(0); /* Should not reach this case */
}
}
/* FIO_checkFilenameCollisions() :
* Checks for and warns if there are any files that would have the same output path
*/
int FIO_checkFilenameCollisions(const char** filenameTable, unsigned nbFiles) {
return FIO_rust_checkFilenameCollisions(filenameTable, nbFiles);
}
/* FIO_highbit64() :
* gives position of highest bit.
* note : only works for v > 0 !
*/
static unsigned FIO_highbit64(unsigned long long v)
{
assert(v != 0);
return FIO_rust_highbit64(v);
}
static void FIO_adjustMemLimitForPatchFromMode(FIO_prefs_t* const prefs,
unsigned long long const dictSize,
unsigned long long const maxSrcFileSize)
{
enum {
FIO_PATCH_MEM_LIMIT_SUCCESS = 0,
FIO_PATCH_MEM_LIMIT_UNKNOWN_SIZE = 1,
FIO_PATCH_MEM_LIMIT_TOO_LARGE = 2
};
unsigned const maxWindowSize = (1U << ZSTD_WINDOWLOG_MAX);
int const status = FIO_rust_adjustMemLimitForPatchFromMode(prefs, dictSize, maxSrcFileSize);
if (status == FIO_PATCH_MEM_LIMIT_UNKNOWN_SIZE)
EXM_THROW(42, "Using --patch-from with stdin requires --stream-size");
if (status == FIO_PATCH_MEM_LIMIT_TOO_LARGE)
EXM_THROW(42, "Can't handle files larger than %u GB\n", maxWindowSize/(1 GB));
assert(status == FIO_PATCH_MEM_LIMIT_SUCCESS);
}
/* FIO_multiFilesConcatWarning() :
* This function handles logic when processing multiple files with -o or -c, displaying the appropriate warnings/prompts.
* Returns 1 if the console should abort, 0 if console should proceed.
*
* If output is stdout or test mode is active, check that `--rm` disabled.
*
* If there is just 1 file to process, zstd will proceed as usual.
* If each file get processed into its own separate destination file, proceed as usual.
*
* When multiple files are processed into a single output,
* display a warning message, then disable --rm if it's set.
*
* If -f is specified or if output is stdout, just proceed.
* If output is set with -o, prompt for confirmation.
*/
static int FIO_multiFilesConcatWarning(const FIO_ctx_t* fCtx, FIO_prefs_t* prefs, const char* outFileName, int displayLevelCutoff)
{
int action = FIO_rust_multiFilesConcatAction(
fCtx->nbFilesTotal, fCtx->hasStdoutOutput, prefs->testMode,
outFileName != NULL, prefs->removeSrcFile, prefs->overwrite,
g_display_prefs.displayLevel, displayLevelCutoff);
if (action == FIO_RUST_MULTI_FILES_ACTION_FATAL_STDOUT_REMOVE)
/* this should not happen ; hard fail, to protect user's data
* note: this should rather be an assert(), but we want to be certain that user's data will not be wiped out in case it nonetheless happen */
EXM_THROW(43, "It's not allowed to remove input files when processed output is piped to stdout. "
"This scenario is not supposed to be possible. "
"This is a programming error. File an issue for it to be fixed.");
if (action == FIO_RUST_MULTI_FILES_ACTION_FATAL_TEST_REMOVE)
/* this should not happen ; hard fail, to protect user's data
* note: this should rather be an assert(), but we want to be certain that user's data will not be wiped out in case it nonetheless happen */
EXM_THROW(43, "Test mode shall not remove input files! "
"This scenario is not supposed to be possible. "
"This is a programming error. File an issue for it to be fixed.");
if (fCtx->nbFilesTotal == 1) return 0;
assert(fCtx->nbFilesTotal > 1);
if (!outFileName || prefs->testMode) return 0;
if (fCtx->hasStdoutOutput) {
DISPLAYLEVEL(2, "zstd: WARNING: all input files will be processed and concatenated into stdout. \n");
} else {
DISPLAYLEVEL(2, "zstd: WARNING: all input files will be processed and concatenated into a single output file: %s \n", outFileName);
}
DISPLAYLEVEL(2, "The concatenated output CANNOT regenerate original file names nor directory structure. \n")
/* multi-input into single output : --rm is not allowed */
if (action == FIO_RUST_MULTI_FILES_ACTION_DISABLE_REMOVE) {
DISPLAYLEVEL(2, "Since it's a destructive operation, input files will not be removed. \n");
prefs->removeSrcFile = 0;
action = FIO_rust_multiFilesConcatAction(
fCtx->nbFilesTotal, fCtx->hasStdoutOutput, prefs->testMode,
1, 0, prefs->overwrite, g_display_prefs.displayLevel,
displayLevelCutoff);
}
if (action == FIO_RUST_MULTI_FILES_ACTION_PROCEED) return 0;
/* multiple files concatenated into single destination file using -o without -f */
if (action == FIO_RUST_MULTI_FILES_ACTION_QUIET_ABORT) {
/* quiet mode => no prompt => fail automatically */
DISPLAYLEVEL(1, "Concatenating multiple processed inputs into a single output loses file metadata. \n");
DISPLAYLEVEL(1, "Aborting. \n");
return 1;
}
/* normal mode => prompt */
assert(action == FIO_RUST_MULTI_FILES_ACTION_CONFIRM);
return UTIL_requireUserConfirmation("Proceed? (y/n): ", "Aborting...", "yY", fCtx->hasStdinInput);
}
static ZSTD_inBuffer setInBuffer(const void* buf, size_t s, size_t pos)
{
ZSTD_inBuffer i;
FIO_rust_setInBuffer(&i, buf, s, pos);
return i;
}
static ZSTD_outBuffer setOutBuffer(void* buf, size_t s, size_t pos)
{
ZSTD_outBuffer o;
FIO_rust_setOutBuffer(&o, buf, s, pos);
return o;
}
#ifndef ZSTD_NOCOMPRESS
/* **********************************************************************
* Compression
************************************************************************/
typedef struct {
FIO_Dict_t dict;
const char* dictFileName;
stat_t dictFileStat;
ZSTD_CStream* cctx;
WritePoolCtx_t *writeCtx;
ReadPoolCtx_t *readCtx;
} cRess_t;
enum {
FIO_RUST_COMPRESS_OK = 0,
FIO_RUST_COMPRESS_GZIP_UNSUPPORTED = 1,
FIO_RUST_COMPRESS_LZMA_UNSUPPORTED = 2,
FIO_RUST_COMPRESS_LZ4_UNSUPPORTED = 3,
FIO_RUST_COMPRESS_ZSTD_UNSUPPORTED = 4
};
typedef unsigned long long (*FIO_rust_compress_zstd_fn)(
void* fCtx, void* prefs, void* ress,
const char* srcFileName, U64 srcFileSize,
int compressionLevel, U64* readsize);
typedef unsigned long long (*FIO_rust_compress_gzip_fn)(
void* ress, const char* srcFileName, U64 srcFileSize,
int compressionLevel, U64* readsize);
typedef unsigned long long (*FIO_rust_compress_lzma_fn)(
void* ress, const char* srcFileName, U64 srcFileSize,
int compressionLevel, U64* readsize, int plainLzma);
typedef unsigned long long (*FIO_rust_compress_lz4_fn)(
void* ress, const char* srcFileName, U64 srcFileSize,
int compressionLevel, int checksumFlag, U64* readsize);
typedef void (*FIO_rust_compress_input_display_fn)(
void* opaque, const char* srcFileName, U64 fileSize);
typedef void (*FIO_rust_compress_status_display_fn)(
void* opaque, void* fCtx, const char* dstFileName,
const char* srcFileName, U64 readsize, U64 compressedfilesize);
typedef struct {
void* opaque;
FIO_rust_compress_zstd_fn compress_zstd;
FIO_rust_compress_gzip_fn compress_gzip;
FIO_rust_compress_lzma_fn compress_lzma;
FIO_rust_compress_lz4_fn compress_lz4;
FIO_rust_compress_input_display_fn display_input;
FIO_rust_compress_status_display_fn display_status;
} FIO_rust_compress_callbacks_t;
enum {
FIO_RUST_GZIP_OK = 0,
FIO_RUST_GZIP_INIT_ERROR = 1,
FIO_RUST_GZIP_DEFLATE_ERROR = 2,
FIO_RUST_GZIP_FINISH_ERROR = 3,
FIO_RUST_GZIP_END_ERROR = 4,
FIO_RUST_GZIP_INVALID_PROJECTION = 5
};
typedef void (*FIO_rust_gzip_read_fill_fn)(
void* opaque, size_t requested, const unsigned char** buffer, size_t* loaded);
typedef void (*FIO_rust_gzip_read_consume_fn)(void* opaque, size_t consumed);
typedef void (*FIO_rust_gzip_write_acquire_fn)(
void* opaque, void** job, unsigned char** buffer, size_t* bufferSize);
typedef void (*FIO_rust_gzip_write_enqueue_fn)(
void* opaque, void** job, size_t usedBufferSize,
unsigned char** buffer, size_t* bufferSize);
typedef void (*FIO_rust_gzip_write_release_fn)(void* opaque, void* job);
typedef void (*FIO_rust_gzip_sparse_write_end_fn)(void* opaque);
typedef int (*FIO_rust_gzip_zlib_init_fn)(void* opaque, int compressionLevel);
typedef int (*FIO_rust_gzip_zlib_deflate_fn)(
void* opaque, const unsigned char* input, size_t inputSize,
unsigned char* output, size_t outputSize, int flush,
size_t* consumed, size_t* produced);
typedef int (*FIO_rust_gzip_zlib_end_fn)(void* opaque);
typedef void (*FIO_rust_gzip_progress_fn)(
void* opaque, U64 srcFileSize, U64 inFileSize, U64 outFileSize);
typedef struct {
void* readOpaque;
void* writeOpaque;
void* zlibOpaque;
void* progressOpaque;
size_t readBufferSize;
FIO_rust_gzip_read_fill_fn readFill;
FIO_rust_gzip_read_consume_fn readConsume;
FIO_rust_gzip_write_acquire_fn writeAcquire;
FIO_rust_gzip_write_enqueue_fn writeEnqueue;
FIO_rust_gzip_write_release_fn writeRelease;
FIO_rust_gzip_sparse_write_end_fn sparseWriteEnd;
FIO_rust_gzip_zlib_init_fn zlibInit;
FIO_rust_gzip_zlib_deflate_fn zlibDeflate;
FIO_rust_gzip_zlib_end_fn zlibEnd;
FIO_rust_gzip_progress_fn progress;
} FIO_rust_gzip_compress_projection_t;
int FIO_rust_compressGzipFrame(
const FIO_rust_gzip_compress_projection_t* projection,
const char* srcFileName, U64 srcFileSize, int compressionLevel,
U64* readsize, U64* compressedSize, int* zlibResult);
int FIO_rust_compressFilenameInternal(
void* fCtx, FIO_prefs_t* prefs, void* ress,
const char* dstFileName, const char* srcFileName,
U64 fileSize, int compressionLevel,
const FIO_rust_compress_callbacks_t* callbacks);
static void FIO_adjustParamsForPatchFromMode(FIO_prefs_t* const prefs,
ZSTD_compressionParameters* comprParams,
unsigned long long const dictSize,
unsigned long long const maxSrcFileSize,
int cLevel)
{
enum {
FIO_PATCH_MEM_LIMIT_SUCCESS = 0,
FIO_PATCH_MEM_LIMIT_UNKNOWN_SIZE = 1,
FIO_PATCH_MEM_LIMIT_TOO_LARGE = 2
};
ZSTD_compressionParameters const cParams = ZSTD_getCParams(cLevel, (size_t)maxSrcFileSize, (size_t)dictSize);
unsigned fileWindowLog;
int autoLdm;
int optimalParser;
int const status = FIO_rust_adjustParamsForPatchFromMode(
prefs, comprParams, dictSize, maxSrcFileSize, cParams,
&fileWindowLog, &autoLdm, &optimalParser);
if (status == FIO_PATCH_MEM_LIMIT_UNKNOWN_SIZE)
EXM_THROW(42, "Using --patch-from with stdin requires --stream-size");
if (status == FIO_PATCH_MEM_LIMIT_TOO_LARGE) {
unsigned const maxWindowSize = (1U << ZSTD_WINDOWLOG_MAX);
EXM_THROW(42, "Can't handle files larger than %u GB\n", maxWindowSize/(1 GB));
}
assert(status == FIO_PATCH_MEM_LIMIT_SUCCESS);
if (fileWindowLog > ZSTD_WINDOWLOG_MAX)
DISPLAYLEVEL(1, "Max window log exceeded by file (compression ratio will suffer)\n");
if (autoLdm)
DISPLAYLEVEL(2, "long mode automatically triggered\n");
if (optimalParser) {
DISPLAYLEVEL(4, "[Optimal parser notes] Consider the following to improve patch size at the cost of speed:\n");
DISPLAYLEVEL(4, "- Set a larger targetLength (e.g. --zstd=targetLength=4096)\n");
DISPLAYLEVEL(4, "- Set a larger chainLog (e.g. --zstd=chainLog=%u)\n", ZSTD_CHAINLOG_MAX);
DISPLAYLEVEL(4, "- Set a larger LDM hashLog (e.g. --zstd=ldmHashLog=%u)\n", ZSTD_LDM_HASHLOG_MAX);
DISPLAYLEVEL(4, "- Set a smaller LDM rateLog (e.g. --zstd=ldmHashRateLog=%u)\n", ZSTD_LDM_HASHRATELOG_MIN);
DISPLAYLEVEL(4, "Also consider playing around with searchLog and hashLog\n");
}
}
static cRess_t FIO_createCResources(FIO_prefs_t* const prefs,
const char* dictFileName, unsigned long long const maxSrcFileSize,
int cLevel, ZSTD_compressionParameters comprParams) {
int useMMap = prefs->mmapDict == ZSTD_ps_enable;
int forceNoUseMMap = prefs->mmapDict == ZSTD_ps_disable;
FIO_dictBufferType_t dictBufferType;
cRess_t ress;
memset(&ress, 0, sizeof(ress));
DISPLAYLEVEL(6, "FIO_createCResources \n");
ress.cctx = ZSTD_createCCtx();
if (ress.cctx == NULL)
EXM_THROW(30, "allocation error (%s): can't create ZSTD_CCtx",
strerror(errno));
FIO_getDictFileStat(dictFileName, &ress.dictFileStat);
/* need to update memLimit before calling createDictBuffer
* because of memLimit check inside it */
if (prefs->patchFromMode) {
U64 const dictSize = UTIL_getFileSizeStat(&ress.dictFileStat);
unsigned long long const ssSize = (unsigned long long)prefs->streamSrcSize;
useMMap |= dictSize > prefs->memLimit;
FIO_adjustParamsForPatchFromMode(prefs, &comprParams, dictSize, ssSize > 0 ? ssSize : maxSrcFileSize, cLevel);
}
dictBufferType = (useMMap && !forceNoUseMMap) ? FIO_mmapDict : FIO_mallocDict;
FIO_initDict(&ress.dict, dictFileName, prefs, &ress.dictFileStat, dictBufferType); /* works with dictFileName==NULL */
ress.writeCtx = AIO_WritePool_create(prefs, ZSTD_CStreamOutSize());
ress.readCtx = AIO_ReadPool_create(prefs, ZSTD_CStreamInSize());
/* Advanced parameters, including dictionary */
if (dictFileName && (ress.dict.dictBuffer==NULL))
EXM_THROW(32, "allocation error : can't create dictBuffer");
ress.dictFileName = dictFileName;
if (prefs->adaptiveMode && !prefs->ldmFlag && !comprParams.windowLog)
comprParams.windowLog = ADAPT_WINDOWLOG_DEFAULT;
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_contentSizeFlag, prefs->contentSize) ); /* always enable content size when available (note: supposed to be default) */
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_dictIDFlag, prefs->dictIDFlag) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_checksumFlag, prefs->checksumFlag) );
/* compression level */
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_compressionLevel, cLevel) );
/* max compressed block size */
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_targetCBlockSize, (int)prefs->targetCBlockSize) );
/* source size hint */
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_srcSizeHint, (int)prefs->srcSizeHint) );
/* long distance matching */
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_enableLongDistanceMatching, prefs->ldmFlag) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_ldmHashLog, prefs->ldmHashLog) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_ldmMinMatch, prefs->ldmMinMatch) );
if (prefs->ldmBucketSizeLog != FIO_LDM_PARAM_NOTSET) {
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_ldmBucketSizeLog, prefs->ldmBucketSizeLog) );
}
if (prefs->ldmHashRateLog != FIO_LDM_PARAM_NOTSET) {
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_ldmHashRateLog, prefs->ldmHashRateLog) );
}
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_useRowMatchFinder, prefs->useRowMatchFinder));
/* compression parameters */
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_windowLog, (int)comprParams.windowLog) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_chainLog, (int)comprParams.chainLog) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_hashLog, (int)comprParams.hashLog) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_searchLog, (int)comprParams.searchLog) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_minMatch, (int)comprParams.minMatch) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_targetLength, (int)comprParams.targetLength) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_strategy, (int)comprParams.strategy) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_literalCompressionMode, (int)prefs->literalCompressionMode) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_enableDedicatedDictSearch, 1) );
/* multi-threading */
#ifdef ZSTD_MULTITHREAD
DISPLAYLEVEL(5,"set nb workers = %u \n", prefs->nbWorkers);
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_nbWorkers, prefs->nbWorkers) );
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_jobSize, prefs->blockSize) );
if (prefs->overlapLog != FIO_OVERLAP_LOG_NOTSET) {
DISPLAYLEVEL(3,"set overlapLog = %u \n", prefs->overlapLog);
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_overlapLog, prefs->overlapLog) );
}
CHECK( ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_rsyncable, prefs->rsyncable) );
#endif
/* dictionary */
if (prefs->patchFromMode) {
CHECK( ZSTD_CCtx_refPrefix(ress.cctx, ress.dict.dictBuffer, ress.dict.dictBufferSize) );
} else {
CHECK( ZSTD_CCtx_loadDictionary_byReference(ress.cctx, ress.dict.dictBuffer, ress.dict.dictBufferSize) );
}
return ress;
}
static void FIO_freeCResources(cRess_t* const ress)
{
FIO_freeDict(&(ress->dict));
AIO_WritePool_free(ress->writeCtx);
AIO_ReadPool_free(ress->readCtx);
ZSTD_freeCStream(ress->cctx); /* never fails */
}
#ifdef ZSTD_GZCOMPRESS
static void FIO_rust_gzip_readFill(void* opaque, size_t requested,
const unsigned char** buffer, size_t* loaded)
{
ReadPoolCtx_t* const readCtx = (ReadPoolCtx_t*)opaque;
AIO_ReadPool_fillBuffer(readCtx, requested);
*buffer = readCtx->srcBuffer;
*loaded = readCtx->srcBufferLoaded;
}
static void FIO_rust_gzip_readConsume(void* opaque, size_t consumed)
{
AIO_ReadPool_consumeBytes((ReadPoolCtx_t*)opaque, consumed);
}
static void FIO_rust_gzip_writeAcquire(void* opaque, void** job,
unsigned char** buffer, size_t* bufferSize)
{
IOJob_t* const writeJob = AIO_WritePool_acquireJob((WritePoolCtx_t*)opaque);
*job = writeJob;
*buffer = (unsigned char*)writeJob->buffer;
*bufferSize = writeJob->bufferSize;
}
static void FIO_rust_gzip_writeEnqueue(void* opaque, void** job,
size_t usedBufferSize,
unsigned char** buffer, size_t* bufferSize)
{
IOJob_t* writeJob = (IOJob_t*)*job;
writeJob->usedBufferSize = usedBufferSize;
AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob);
*job = writeJob;
*buffer = (unsigned char*)writeJob->buffer;
*bufferSize = writeJob->bufferSize;
(void)opaque;
}
static void FIO_rust_gzip_writeRelease(void* opaque, void* job)
{
AIO_WritePool_releaseIoJob((IOJob_t*)job);
(void)opaque;
}
static void FIO_rust_gzip_sparseWriteEnd(void* opaque)
{
AIO_WritePool_sparseWriteEnd((WritePoolCtx_t*)opaque);
}
static int FIO_rust_gzip_zlibInit(void* opaque, int compressionLevel)
{
z_stream* const strm = (z_stream*)opaque;
strm->zalloc = Z_NULL;
strm->zfree = Z_NULL;
strm->opaque = Z_NULL;
return deflateInit2(strm, compressionLevel, Z_DEFLATED,
15 /* maxWindowLogSize */ + 16 /* gzip only */,
8, Z_DEFAULT_STRATEGY); /* see https://www.zlib.net/manual.html */
}
static int FIO_rust_gzip_zlibDeflate(void* opaque, const unsigned char* input,
size_t inputSize, unsigned char* output,
size_t outputSize, int flush,
size_t* consumed, size_t* produced)
{
z_stream* const strm = (z_stream*)opaque;
size_t const availBefore = inputSize;
size_t const outputBefore = outputSize;
int ret;
assert(inputSize <= UINT_MAX);
assert(outputSize <= UINT_MAX);
strm->next_in = (z_const unsigned char*)input;
strm->avail_in = (uInt)inputSize;
strm->next_out = output;
strm->avail_out = (uInt)outputSize;
ret = deflate(strm, flush);
*consumed = availBefore - strm->avail_in;
*produced = outputBefore - strm->avail_out;
return ret;
}
static int FIO_rust_gzip_zlibEnd(void* opaque)
{
return deflateEnd((z_stream*)opaque);
}
static void FIO_rust_gzip_progress(void* opaque, U64 srcFileSize,
U64 inFileSize, U64 outFileSize)
{
(void)opaque;
if (srcFileSize == UTIL_FILESIZE_UNKNOWN) {
DISPLAYUPDATE_PROGRESS(
"\rRead : %u MB ==> %.2f%% ",
(unsigned)(inFileSize>>20),
(double)outFileSize/(double)inFileSize*100)
} else {
DISPLAYUPDATE_PROGRESS(
"\rRead : %u / %u MB ==> %.2f%% ",
(unsigned)(inFileSize>>20), (unsigned)(srcFileSize>>20),
(double)outFileSize/(double)inFileSize*100);
}
}
#endif
#ifdef ZSTD_LZMACOMPRESS
static unsigned long long
FIO_compressLzmaFrame(cRess_t* ress,
const char* srcFileName, U64 const srcFileSize,
int compressionLevel, U64* readsize, int plain_lzma)
{
unsigned long long inFileSize = 0, outFileSize = 0;
lzma_stream strm = LZMA_STREAM_INIT;
lzma_action action = LZMA_RUN;
lzma_ret ret;
IOJob_t *writeJob = NULL;
if (compressionLevel < 0) compressionLevel = 0;
if (compressionLevel > 9) compressionLevel = 9;
if (plain_lzma) {
lzma_options_lzma opt_lzma;
if (lzma_lzma_preset(&opt_lzma, compressionLevel))
EXM_THROW(81, "zstd: %s: lzma_lzma_preset error", srcFileName);
ret = lzma_alone_encoder(&strm, &opt_lzma); /* LZMA */
if (ret != LZMA_OK)
EXM_THROW(82, "zstd: %s: lzma_alone_encoder error %d", srcFileName, ret);
} else {
ret = lzma_easy_encoder(&strm, compressionLevel, LZMA_CHECK_CRC64); /* XZ */
if (ret != LZMA_OK)
EXM_THROW(83, "zstd: %s: lzma_easy_encoder error %d", srcFileName, ret);
}
writeJob =AIO_WritePool_acquireJob(ress->writeCtx);
strm.next_out = (BYTE*)writeJob->buffer;
strm.avail_out = writeJob->bufferSize;
strm.next_in = 0;
strm.avail_in = 0;
while (1) {
if (strm.avail_in == 0) {
size_t const inSize = AIO_ReadPool_fillBuffer(ress->readCtx, ZSTD_CStreamInSize());
if (ress->readCtx->srcBufferLoaded == 0) action = LZMA_FINISH;
inFileSize += inSize;
strm.next_in = (BYTE const*)ress->readCtx->srcBuffer;
strm.avail_in = ress->readCtx->srcBufferLoaded;
}
{
size_t const availBefore = strm.avail_in;
ret = lzma_code(&strm, action);
AIO_ReadPool_consumeBytes(ress->readCtx, availBefore - strm.avail_in);
}
if (ret != LZMA_OK && ret != LZMA_STREAM_END)
EXM_THROW(84, "zstd: %s: lzma_code encoding error %d", srcFileName, ret);
{ size_t const compBytes = writeJob->bufferSize - strm.avail_out;
if (compBytes) {
writeJob->usedBufferSize = compBytes;
AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob);
outFileSize += compBytes;
strm.next_out = (BYTE*)writeJob->buffer;
strm.avail_out = writeJob->bufferSize;
} }
if (srcFileSize == UTIL_FILESIZE_UNKNOWN)
DISPLAYUPDATE_PROGRESS("\rRead : %u MB ==> %.2f%%",
(unsigned)(inFileSize>>20),
(double)outFileSize/(double)inFileSize*100)
else
DISPLAYUPDATE_PROGRESS("\rRead : %u / %u MB ==> %.2f%%",
(unsigned)(inFileSize>>20), (unsigned)(srcFileSize>>20),
(double)outFileSize/(double)inFileSize*100);
if (ret == LZMA_STREAM_END) break;
}
lzma_end(&strm);
*readsize = inFileSize;
AIO_WritePool_releaseIoJob(writeJob);
AIO_WritePool_sparseWriteEnd(ress->writeCtx);
return outFileSize;
}
#endif
#ifdef ZSTD_LZ4COMPRESS
#if LZ4_VERSION_NUMBER <= 10600
#define LZ4F_blockLinked blockLinked
#define LZ4F_max64KB max64KB
#endif
static int FIO_LZ4_GetBlockSize_FromBlockId (int id)
{
return FIO_rust_LZ4_GetBlockSize_FromBlockId(id);
}
static unsigned long long
FIO_compressLz4Frame(cRess_t* ress,
const char* srcFileName, U64 const srcFileSize,
int compressionLevel, int checksumFlag,
U64* readsize)
{
const size_t blockSize = FIO_LZ4_GetBlockSize_FromBlockId(LZ4F_max64KB);
unsigned long long inFileSize = 0, outFileSize = 0;
LZ4F_preferences_t prefs;
LZ4F_compressionContext_t ctx;
IOJob_t* writeJob = AIO_WritePool_acquireJob(ress->writeCtx);
LZ4F_errorCode_t const errorCode = LZ4F_createCompressionContext(&ctx, LZ4F_VERSION);
if (LZ4F_isError(errorCode))
EXM_THROW(31, "zstd: failed to create lz4 compression context");
memset(&prefs, 0, sizeof(prefs));
assert(blockSize <= ress->readCtx->base.jobBufferSize);
/* autoflush off to mitigate a bug in lz4<=1.9.3 for compression level 12 */
prefs.autoFlush = 0;
prefs.compressionLevel = compressionLevel;
prefs.frameInfo.blockMode = LZ4F_blockLinked;
prefs.frameInfo.blockSizeID = LZ4F_max64KB;
prefs.frameInfo.contentChecksumFlag = (contentChecksum_t)checksumFlag;
#if LZ4_VERSION_NUMBER >= 10600
prefs.frameInfo.contentSize = (srcFileSize==UTIL_FILESIZE_UNKNOWN) ? 0 : srcFileSize;
#endif
assert(LZ4F_compressBound(blockSize, &prefs) <= writeJob->bufferSize);
{
size_t headerSize = LZ4F_compressBegin(ctx, writeJob->buffer, writeJob->bufferSize, &prefs);
if (LZ4F_isError(headerSize))
EXM_THROW(33, "File header generation failed : %s",
LZ4F_getErrorName(headerSize));
writeJob->usedBufferSize = headerSize;
AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob);
outFileSize += headerSize;
/* Read first block */
inFileSize += AIO_ReadPool_fillBuffer(ress->readCtx, blockSize);
/* Main Loop */
while (ress->readCtx->srcBufferLoaded) {
size_t inSize = MIN(blockSize, ress->readCtx->srcBufferLoaded);
size_t const outSize = LZ4F_compressUpdate(ctx, writeJob->buffer, writeJob->bufferSize,
ress->readCtx->srcBuffer, inSize, NULL);
if (LZ4F_isError(outSize))
EXM_THROW(35, "zstd: %s: lz4 compression failed : %s",
srcFileName, LZ4F_getErrorName(outSize));
outFileSize += outSize;
if (srcFileSize == UTIL_FILESIZE_UNKNOWN) {
DISPLAYUPDATE_PROGRESS("\rRead : %u MB ==> %.2f%%",
(unsigned)(inFileSize>>20),
(double)outFileSize/(double)inFileSize*100)
} else {
DISPLAYUPDATE_PROGRESS("\rRead : %u / %u MB ==> %.2f%%",
(unsigned)(inFileSize>>20), (unsigned)(srcFileSize>>20),
(double)outFileSize/(double)inFileSize*100);
}
/* Write Block */
writeJob->usedBufferSize = outSize;
AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob);
/* Read next block */
AIO_ReadPool_consumeBytes(ress->readCtx, inSize);
inFileSize += AIO_ReadPool_fillBuffer(ress->readCtx, blockSize);
}
/* End of Stream mark */
headerSize = LZ4F_compressEnd(ctx, writeJob->buffer, writeJob->bufferSize, NULL);
if (LZ4F_isError(headerSize))
EXM_THROW(38, "zstd: %s: lz4 end of file generation failed : %s",
srcFileName, LZ4F_getErrorName(headerSize));
writeJob->usedBufferSize = headerSize;
AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob);
outFileSize += headerSize;
}
*readsize = inFileSize;
LZ4F_freeCompressionContext(ctx);
AIO_WritePool_releaseIoJob(writeJob);
AIO_WritePool_sparseWriteEnd(ress->writeCtx);
return outFileSize;
}
#endif
static unsigned long long
FIO_compressZstdFrame(FIO_ctx_t* const fCtx,
FIO_prefs_t* const prefs,
const cRess_t* ressPtr,
const char* srcFileName, U64 fileSize,
int compressionLevel, U64* readsize)
{
cRess_t const ress = *ressPtr;
IOJob_t* writeJob = AIO_WritePool_acquireJob(ressPtr->writeCtx);
U64 compressedfilesize = 0;
ZSTD_EndDirective directive = ZSTD_e_continue;
U64 pledgedSrcSize = ZSTD_CONTENTSIZE_UNKNOWN;
/* stats */
ZSTD_frameProgression previous_zfp_update = { 0, 0, 0, 0, 0, 0 };
ZSTD_frameProgression previous_zfp_correction = { 0, 0, 0, 0, 0, 0 };
typedef enum { noChange, slower, faster } speedChange_e;
speedChange_e speedChange = noChange;
unsigned flushWaiting = 0;
unsigned inputPresented = 0;
unsigned inputBlocked = 0;
unsigned lastJobID = 0;
UTIL_time_t lastAdaptTime = UTIL_getTime();
U64 const adaptEveryMicro = REFRESH_RATE;
UTIL_HumanReadableSize_t const file_hrs = UTIL_makeHumanReadableSize(fileSize);
DISPLAYLEVEL(6, "compression using zstd format \n");
/* init */
if (fileSize != UTIL_FILESIZE_UNKNOWN) {
pledgedSrcSize = fileSize;
CHECK(ZSTD_CCtx_setPledgedSrcSize(ress.cctx, fileSize));
} else if (prefs->streamSrcSize > 0) {
/* unknown source size; use the declared stream size */
pledgedSrcSize = prefs->streamSrcSize;
CHECK( ZSTD_CCtx_setPledgedSrcSize(ress.cctx, prefs->streamSrcSize) );
}
{ int windowLog;
UTIL_HumanReadableSize_t windowSize;
CHECK(ZSTD_CCtx_getParameter(ress.cctx, ZSTD_c_windowLog, &windowLog));
if (windowLog == 0) {
if (prefs->ldmFlag) {
/* If long mode is set without a window size libzstd will set this size internally */
windowLog = ZSTD_WINDOWLOG_LIMIT_DEFAULT;
} else {
const ZSTD_compressionParameters cParams = ZSTD_getCParams(compressionLevel, fileSize, 0);
windowLog = (int)cParams.windowLog;
}
}
windowSize = UTIL_makeHumanReadableSize(MAX(1ULL, MIN(1ULL << windowLog, pledgedSrcSize)));
DISPLAYLEVEL(4, "Decompression will require %.*f%s of memory\n", windowSize.precision, windowSize.value, windowSize.suffix);
}
/* Main compression loop */
do {
size_t stillToFlush;
/* Fill input Buffer */
size_t const inSize = AIO_ReadPool_fillBuffer(ress.readCtx, ZSTD_CStreamInSize());
ZSTD_inBuffer inBuff = setInBuffer( ress.readCtx->srcBuffer, ress.readCtx->srcBufferLoaded, 0 );
DISPLAYLEVEL(6, "fread %u bytes from source \n", (unsigned)inSize);
*readsize += inSize;
if ((ress.readCtx->srcBufferLoaded == 0) || (*readsize == fileSize))
directive = ZSTD_e_end;
stillToFlush = 1;
while ((inBuff.pos != inBuff.size) /* input buffer must be entirely ingested */
|| (directive == ZSTD_e_end && stillToFlush != 0) ) {
size_t const oldIPos = inBuff.pos;
ZSTD_outBuffer outBuff = setOutBuffer( writeJob->buffer, writeJob->bufferSize, 0 );
size_t const toFlushNow = ZSTD_toFlushNow(ress.cctx);
CHECK_V(stillToFlush, ZSTD_compressStream2(ress.cctx, &outBuff, &inBuff, directive));
AIO_ReadPool_consumeBytes(ress.readCtx, inBuff.pos - oldIPos);
/* count stats */
inputPresented++;
if (oldIPos == inBuff.pos) inputBlocked++; /* input buffer is full and can't take any more : input speed is faster than consumption rate */
if (!toFlushNow) flushWaiting = 1;
/* Write compressed stream */
DISPLAYLEVEL(6, "ZSTD_compress_generic(end:%u) => input pos(%u)<=(%u)size ; output generated %u bytes \n",
(unsigned)directive, (unsigned)inBuff.pos, (unsigned)inBuff.size, (unsigned)outBuff.pos);
if (outBuff.pos) {
writeJob->usedBufferSize = outBuff.pos;
AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob);
compressedfilesize += outBuff.pos;
}
/* adaptive mode : statistics measurement and speed correction */
if (prefs->adaptiveMode && UTIL_clockSpanMicro(lastAdaptTime) > adaptEveryMicro) {
ZSTD_frameProgression const zfp = ZSTD_getFrameProgression(ress.cctx);
lastAdaptTime = UTIL_getTime();
/* check output speed */
if (zfp.currentJobID > 1) { /* only possible if nbWorkers >= 1 */
unsigned long long newlyProduced = zfp.produced - previous_zfp_update.produced;
unsigned long long newlyFlushed = zfp.flushed - previous_zfp_update.flushed;
assert(zfp.produced >= previous_zfp_update.produced);
assert(prefs->nbWorkers >= 1);
/* test if compression is blocked
* either because output is slow and all buffers are full
* or because input is slow and no job can start while waiting for at least one buffer to be filled.
* note : exclude starting part, since currentJobID > 1 */
if ( (zfp.consumed == previous_zfp_update.consumed) /* no data compressed : no data available, or no more buffer to compress to, OR compression is really slow (compression of a single block is slower than update rate)*/
&& (zfp.nbActiveWorkers == 0) /* confirmed : no compression ongoing */
) {
DISPLAYLEVEL(6, "all buffers full : compression stopped => slow down \n")
speedChange = slower;
}
previous_zfp_update = zfp;
if ( (newlyProduced > (newlyFlushed * 9 / 8)) /* compression produces more data than output can flush (though production can be spiky, due to work unit : (N==4)*block sizes) */
&& (flushWaiting == 0) /* flush speed was never slowed by lack of production, so it's operating at max capacity */
) {
DISPLAYLEVEL(6, "compression faster than flush (%llu > %llu), and flushed was never slowed down by lack of production => slow down \n", newlyProduced, newlyFlushed);
speedChange = slower;
}
flushWaiting = 0;
}
/* course correct only if there is at least one new job completed */
if (zfp.currentJobID > lastJobID) {
DISPLAYLEVEL(6, "compression level adaptation check \n")
/* check input speed */
if (zfp.currentJobID > (unsigned)(prefs->nbWorkers+1)) { /* warm up period, to fill all workers */
if (inputBlocked <= 0) {
DISPLAYLEVEL(6, "input is never blocked => input is slower than ingestion \n");
speedChange = slower;
} else if (speedChange == noChange) {
unsigned long long newlyIngested = zfp.ingested - previous_zfp_correction.ingested;
unsigned long long newlyConsumed = zfp.consumed - previous_zfp_correction.consumed;
unsigned long long newlyProduced = zfp.produced - previous_zfp_correction.produced;
unsigned long long newlyFlushed = zfp.flushed - previous_zfp_correction.flushed;
previous_zfp_correction = zfp;
assert(inputPresented > 0);
DISPLAYLEVEL(6, "input blocked %u/%u(%.2f) - ingested:%u vs %u:consumed - flushed:%u vs %u:produced \n",
inputBlocked, inputPresented, (double)inputBlocked/inputPresented*100,
(unsigned)newlyIngested, (unsigned)newlyConsumed,
(unsigned)newlyFlushed, (unsigned)newlyProduced);
if ( (inputBlocked > inputPresented / 8) /* input is waiting often, because input buffers is full : compression or output too slow */
&& (newlyFlushed * 33 / 32 > newlyProduced) /* flush everything that is produced */
&& (newlyIngested * 33 / 32 > newlyConsumed) /* input speed as fast or faster than compression speed */
) {
DISPLAYLEVEL(6, "recommend faster as in(%llu) >= (%llu)comp(%llu) <= out(%llu) \n",
newlyIngested, newlyConsumed, newlyProduced, newlyFlushed);
speedChange = faster;
}
}
inputBlocked = 0;
inputPresented = 0;
}
if (speedChange == slower) {
DISPLAYLEVEL(6, "slower speed , higher compression \n")
compressionLevel ++;
if (compressionLevel > ZSTD_maxCLevel()) compressionLevel = ZSTD_maxCLevel();
if (compressionLevel > prefs->maxAdaptLevel) compressionLevel = prefs->maxAdaptLevel;
compressionLevel += (compressionLevel == 0); /* skip 0 */
ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_compressionLevel, compressionLevel);
}
if (speedChange == faster) {
DISPLAYLEVEL(6, "faster speed , lighter compression \n")
compressionLevel --;
if (compressionLevel < prefs->minAdaptLevel) compressionLevel = prefs->minAdaptLevel;
compressionLevel -= (compressionLevel == 0); /* skip 0 */
ZSTD_CCtx_setParameter(ress.cctx, ZSTD_c_compressionLevel, compressionLevel);
}
speedChange = noChange;
lastJobID = zfp.currentJobID;
} /* if (zfp.currentJobID > lastJobID) */
} /* if (prefs->adaptiveMode && UTIL_clockSpanMicro(lastAdaptTime) > adaptEveryMicro) */
/* display notification */
if (SHOULD_DISPLAY_PROGRESS() && READY_FOR_UPDATE()) {
ZSTD_frameProgression const zfp = ZSTD_getFrameProgression(ress.cctx);
double const cShare = (double)zfp.produced / (double)(zfp.consumed + !zfp.consumed/*avoid div0*/) * 100;
UTIL_HumanReadableSize_t const buffered_hrs = UTIL_makeHumanReadableSize(zfp.ingested - zfp.consumed);
UTIL_HumanReadableSize_t const consumed_hrs = UTIL_makeHumanReadableSize(zfp.consumed);
UTIL_HumanReadableSize_t const produced_hrs = UTIL_makeHumanReadableSize(zfp.produced);
DELAY_NEXT_UPDATE();
/* display progress notifications */
DISPLAY_PROGRESS("\r%79s\r", ""); /* Clear out the current displayed line */
if (g_display_prefs.displayLevel >= 3) {
/* Verbose progress update */
DISPLAY_PROGRESS(
"(L%i) Buffered:%5.*f%s - Consumed:%5.*f%s - Compressed:%5.*f%s => %.2f%% ",
compressionLevel,
buffered_hrs.precision, buffered_hrs.value, buffered_hrs.suffix,
consumed_hrs.precision, consumed_hrs.value, consumed_hrs.suffix,
produced_hrs.precision, produced_hrs.value, produced_hrs.suffix,
cShare );
} else {
/* Require level 2 or forcibly displayed progress counter for summarized updates */
if (fCtx->nbFilesTotal > 1) {
size_t srcFileNameSize = strlen(srcFileName);
/* Ensure that the string we print is roughly the same size each time */
if (srcFileNameSize > 18) {
const char* truncatedSrcFileName = srcFileName + srcFileNameSize - 15;
DISPLAY_PROGRESS("Compress: %u/%u files. Current: ...%s ",
fCtx->currFileIdx+1, fCtx->nbFilesTotal, truncatedSrcFileName);
} else {
DISPLAY_PROGRESS("Compress: %u/%u files. Current: %*s ",
fCtx->currFileIdx+1, fCtx->nbFilesTotal, (int)(18-srcFileNameSize), srcFileName);
}
}
DISPLAY_PROGRESS("Read:%6.*f%4s ", consumed_hrs.precision, consumed_hrs.value, consumed_hrs.suffix);
if (fileSize != UTIL_FILESIZE_UNKNOWN)
DISPLAY_PROGRESS("/%6.*f%4s", file_hrs.precision, file_hrs.value, file_hrs.suffix);
DISPLAY_PROGRESS(" ==> %2.f%%", cShare);
}
} /* if (SHOULD_DISPLAY_PROGRESS() && READY_FOR_UPDATE()) */
} /* while ((inBuff.pos != inBuff.size) */
} while (directive != ZSTD_e_end);
if (fileSize != UTIL_FILESIZE_UNKNOWN && *readsize != fileSize) {
EXM_THROW(27, "Read error : Incomplete read : %llu / %llu B",
(unsigned long long)*readsize, (unsigned long long)fileSize);
}
AIO_WritePool_releaseIoJob(writeJob);
AIO_WritePool_sparseWriteEnd(ressPtr->writeCtx);
return compressedfilesize;
}
static unsigned long long
FIO_rust_compressZstdCallback(void* fCtx, void* prefs, void* ress,
const char* srcFileName, U64 srcFileSize,
int compressionLevel, U64* readsize)
{
return FIO_compressZstdFrame((FIO_ctx_t*)fCtx, (FIO_prefs_t*)prefs,
(const cRess_t*)ress, srcFileName,
srcFileSize, compressionLevel, readsize);
}
#ifdef ZSTD_GZCOMPRESS
static unsigned long long
FIO_rust_compressGzipCallback(void* ress, const char* srcFileName,
U64 srcFileSize, int compressionLevel,
U64* readsize)
{
const cRess_t* const ressPtr = (const cRess_t*)ress;
FIO_rust_gzip_compress_projection_t projection;
z_stream strm;
U64 compressedSize = 0;
int zlibResult = Z_OK;
int status;
memset(&projection, 0, sizeof(projection));
projection.readOpaque = (void*)ressPtr->readCtx;
projection.writeOpaque = (void*)ressPtr->writeCtx;
projection.zlibOpaque = &strm;
projection.readBufferSize = ZSTD_CStreamInSize();
projection.readFill = FIO_rust_gzip_readFill;
projection.readConsume = FIO_rust_gzip_readConsume;
projection.writeAcquire = FIO_rust_gzip_writeAcquire;
projection.writeEnqueue = FIO_rust_gzip_writeEnqueue;
projection.writeRelease = FIO_rust_gzip_writeRelease;
projection.sparseWriteEnd = FIO_rust_gzip_sparseWriteEnd;
projection.zlibInit = FIO_rust_gzip_zlibInit;
projection.zlibDeflate = FIO_rust_gzip_zlibDeflate;
projection.zlibEnd = FIO_rust_gzip_zlibEnd;
projection.progress = FIO_rust_gzip_progress;
status = FIO_rust_compressGzipFrame(
&projection, srcFileName, srcFileSize, compressionLevel,
readsize, &compressedSize, &zlibResult);
switch (status) {
case FIO_RUST_GZIP_OK:
return compressedSize;
case FIO_RUST_GZIP_INIT_ERROR:
EXM_THROW(71, "zstd: %s: deflateInit2 error %d \n", srcFileName, zlibResult);
case FIO_RUST_GZIP_DEFLATE_ERROR:
EXM_THROW(72, "zstd: %s: deflate error %d \n", srcFileName, zlibResult);
case FIO_RUST_GZIP_FINISH_ERROR:
EXM_THROW(77, "zstd: %s: deflate error %d \n", srcFileName, zlibResult);
case FIO_RUST_GZIP_END_ERROR:
EXM_THROW(79, "zstd: %s: deflateEnd error %d \n", srcFileName, zlibResult);
default:
assert(status == FIO_RUST_GZIP_INVALID_PROJECTION);
EXM_THROW(72, "zstd: %s: deflate error %d \n", srcFileName, zlibResult);
}
}
#endif
#ifdef ZSTD_LZMACOMPRESS
static unsigned long long
FIO_rust_compressLzmaCallback(void* ress, const char* srcFileName,
U64 srcFileSize, int compressionLevel,
U64* readsize, int plainLzma)
{
return FIO_compressLzmaFrame((cRess_t*)ress, srcFileName, srcFileSize,
compressionLevel, readsize, plainLzma);
}
#endif
#ifdef ZSTD_LZ4COMPRESS
static unsigned long long
FIO_rust_compressLz4Callback(void* ress, const char* srcFileName,
U64 srcFileSize, int compressionLevel,
int checksumFlag, U64* readsize)
{
return FIO_compressLz4Frame((cRess_t*)ress, srcFileName, srcFileSize,
compressionLevel, checksumFlag, readsize);
}
#endif
typedef struct {
UTIL_time_t timeStart;
clock_t cpuStart;
} FIO_compressDisplayContext_t;
static void
FIO_rust_compressInputDisplayCallback(void* opaque, const char* srcFileName,
U64 fileSize)
{
(void)opaque;
DISPLAYLEVEL(5, "%s: %llu bytes \n", srcFileName,
(unsigned long long)fileSize);
}
static void
FIO_rust_compressStatusDisplayCallback(void* opaque, void* fCtxOpaque,
const char* dstFileName,
const char* srcFileName,
U64 readsize, U64 compressedfilesize)
{
FIO_compressDisplayContext_t const* const displayCtx =
(const FIO_compressDisplayContext_t*)opaque;
FIO_ctx_t* const fCtx = (FIO_ctx_t*)fCtxOpaque;
DISPLAY_PROGRESS("\r%79s\r", "");
if (FIO_shouldDisplayFileSummary(fCtx)) {
UTIL_HumanReadableSize_t hr_isize =
UTIL_makeHumanReadableSize((U64)readsize);
UTIL_HumanReadableSize_t hr_osize =
UTIL_makeHumanReadableSize((U64)compressedfilesize);
if (readsize == 0) {
DISPLAY_SUMMARY("%-20s : (%6.*f%s => %6.*f%s, %s) \n",
srcFileName,
hr_isize.precision, hr_isize.value, hr_isize.suffix,
hr_osize.precision, hr_osize.value, hr_osize.suffix,
dstFileName);
} else {
DISPLAY_SUMMARY("%-20s :%6.2f%% (%6.*f%s => %6.*f%s, %s) \n",
srcFileName,
(double)compressedfilesize / (double)readsize * 100,
hr_isize.precision, hr_isize.value, hr_isize.suffix,
hr_osize.precision, hr_osize.value, hr_osize.suffix,
dstFileName);
}
}
/* Keep the original elapsed-time and CPU-load diagnostics in C so their
* formatting and clock source remain identical to the old loop. */
{ clock_t const cpuEnd = clock();
double const cpuLoad_s = (double)(cpuEnd - displayCtx->cpuStart) / CLOCKS_PER_SEC;
U64 const timeLength_ns = UTIL_clockSpanNano(displayCtx->timeStart);
double const timeLength_s = (double)timeLength_ns / 1000000000;
double const cpuLoad_pct = (cpuLoad_s / timeLength_s) * 100;
DISPLAYLEVEL(4, "%-20s : Completed in %.2f sec (cpu load : %.0f%%)\n",
srcFileName, timeLength_s, cpuLoad_pct);
}
}
/*! FIO_compressFilename_internal() :
* same as FIO_compressFilename_extRess(), with `ress.desFile` already opened.
* @return : 0 : compression completed correctly,
* 1 : missing or pb opening srcFileName
*/
static int
FIO_compressFilename_internal(FIO_ctx_t* const fCtx,
FIO_prefs_t* const prefs,
cRess_t ress,
const char* dstFileName, const char* srcFileName,
int compressionLevel)
{
FIO_compressDisplayContext_t displayCtx = {
UTIL_getTime(),
clock()
};
U64 const fileSize = UTIL_getFileSize(srcFileName);
FIO_rust_compress_callbacks_t callbacks;
int status;
memset(&callbacks, 0, sizeof(callbacks));
callbacks.opaque = &displayCtx;
callbacks.compress_zstd = FIO_rust_compressZstdCallback;
#ifdef ZSTD_GZCOMPRESS
callbacks.compress_gzip = FIO_rust_compressGzipCallback;
#endif
#ifdef ZSTD_LZMACOMPRESS
callbacks.compress_lzma = FIO_rust_compressLzmaCallback;
#endif
#ifdef ZSTD_LZ4COMPRESS
callbacks.compress_lz4 = FIO_rust_compressLz4Callback;
#endif
callbacks.display_input = FIO_rust_compressInputDisplayCallback;
callbacks.display_status = FIO_rust_compressStatusDisplayCallback;
status = FIO_rust_compressFilenameInternal(
fCtx, prefs, &ress, dstFileName, srcFileName, fileSize,
compressionLevel, &callbacks);
switch (status) {
case FIO_RUST_COMPRESS_OK:
return 0;
case FIO_RUST_COMPRESS_GZIP_UNSUPPORTED:
EXM_THROW(20, "zstd: %s: file cannot be compressed as gzip (zstd compiled without ZSTD_GZCOMPRESS) -- ignored \n",
srcFileName);
case FIO_RUST_COMPRESS_LZMA_UNSUPPORTED:
EXM_THROW(20, "zstd: %s: file cannot be compressed as xz/lzma (zstd compiled without ZSTD_LZMACOMPRESS) -- ignored \n",
srcFileName);
case FIO_RUST_COMPRESS_LZ4_UNSUPPORTED:
EXM_THROW(20, "zstd: %s: file cannot be compressed as lz4 (zstd compiled without ZSTD_LZ4COMPRESS) -- ignored \n",
srcFileName);
case FIO_RUST_COMPRESS_ZSTD_UNSUPPORTED:
default:
assert(status == FIO_RUST_COMPRESS_ZSTD_UNSUPPORTED);
EXM_THROW(20, "zstd: %s: file cannot be compressed as zstd -- ignored \n",
srcFileName);
}
}
/*! FIO_compressFilename_dstFile() :
* open dstFileName, or pass-through if ress.file != NULL,
* then start compression with FIO_compressFilename_internal().
* Manages source removal (--rm) and file permissions transfer.
* note : ress.srcFile must be != NULL,
* so reach this function through FIO_compressFilename_srcFile().
* @return : 0 : compression completed correctly,
* 1 : pb
*/
static int FIO_compressFilename_dstFile(FIO_ctx_t* const fCtx,
FIO_prefs_t* const prefs,
cRess_t ress,
const char* dstFileName,
const char* srcFileName,
const stat_t* srcFileStat,
int compressionLevel)
{
int closeDstFile = 0;
int result;
int transferStat = 0;
int dstFd = -1;
assert(AIO_ReadPool_getFile(ress.readCtx) != NULL);
if (AIO_WritePool_getFile(ress.writeCtx) == NULL) {
int dstFileInitialPermissions = DEFAULT_FILE_PERMISSIONS;
if ( strcmp (srcFileName, stdinmark)
&& strcmp (dstFileName, stdoutmark)
&& UTIL_isRegularFileStat(srcFileStat) ) {
transferStat = 1;
dstFileInitialPermissions = TEMPORARY_FILE_PERMISSIONS;
}
closeDstFile = 1;
DISPLAYLEVEL(6, "FIO_compressFilename_dstFile: opening dst: %s \n", dstFileName);
{ FILE *dstFile = FIO_openDstFile(fCtx, prefs, srcFileName, dstFileName, dstFileInitialPermissions);
if (dstFile==NULL) return 1; /* could not open dstFileName */
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_compressFilename_internal(fCtx, prefs, ress, dstFileName, srcFileName, compressionLevel);
if (closeDstFile) {
clearHandler();
if (transferStat) {
UTIL_setFDStat(dstFd, dstFileName, srcFileStat);
}
DISPLAYLEVEL(6, "FIO_compressFilename_dstFile: closing dst: %s \n", dstFileName);
if (AIO_WritePool_closeFile(ress.writeCtx)) { /* error closing file */
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 compression artefact; note don't do anything special if remove() fails */
}
}
return result;
}
/* List used to compare file extensions (used with --exclude-compressed flag)
* Different from the suffixList and should only apply to ZSTD compress operationResult
*/
static const char *compressedFileExtensions[] = {
ZSTD_EXTENSION,
TZSTD_EXTENSION,
GZ_EXTENSION,
TGZ_EXTENSION,
LZMA_EXTENSION,
XZ_EXTENSION,
TXZ_EXTENSION,
LZ4_EXTENSION,
TLZ4_EXTENSION,
".7z",
".aa3",
".aac",
".aar",
".ace",
".alac",
".ape",
".apk",
".apng",
".arc",
".archive",
".arj",
".ark",
".asf",
".avi",
".avif",
".ba",
".br",
".bz2",
".cab",
".cdx",
".chm",
".cr2",
".divx",
".dmg",
".dng",
".docm",
".docx",
".dotm",
".dotx",
".dsft",
".ear",
".eftx",
".emz",
".eot",
".epub",
".f4v",
".flac",
".flv",
".gho",
".gif",
".gifv",
".gnp",
".iso",
".jar",
".jpeg",
".jpg",
".jxl",
".lz",
".lzh",
".m4a",
".m4v",
".mkv",
".mov",
".mp2",
".mp3",
".mp4",
".mpa",
".mpc",
".mpe",
".mpeg",
".mpg",
".mpl",
".mpv",
".msi",
".odp",
".ods",
".odt",
".ogg",
".ogv",
".otp",
".ots",
".ott",
".pea",
".png",
".pptx",
".qt",
".rar",
".s7z",
".sfx",
".sit",
".sitx",
".sqx",
".svgz",
".swf",
".tbz2",
".tib",
".tlz",
".vob",
".war",
".webm",
".webp",
".wma",
".wmv",
".woff",
".woff2",
".wvl",
".xlsx",
".xpi",
".xps",
".zip",
".zipx",
".zoo",
".zpaq",
NULL
};
/*! FIO_compressFilename_srcFile() :
* @return : 0 : compression completed correctly,
* 1 : missing or pb opening srcFileName
*/
static int
FIO_compressFilename_srcFile(FIO_ctx_t* const fCtx,
FIO_prefs_t* const prefs,
cRess_t ress,
const char* dstFileName,
const char* srcFileName,
int compressionLevel)
{
int result;
FILE* srcFile;
stat_t srcFileStat;
U64 fileSize = UTIL_FILESIZE_UNKNOWN;
DISPLAYLEVEL(6, "FIO_compressFilename_srcFile: %s \n", srcFileName);
if (strcmp(srcFileName, stdinmark)) {
if (UTIL_stat(srcFileName, &srcFileStat)) {
/* failure to stat at all is handled during opening */
/* ensure src is not a directory */
if (UTIL_isDirectoryStat(&srcFileStat)) {
DISPLAYLEVEL(1, "zstd: %s is a directory -- ignored \n", srcFileName);
return 1;
}
/* ensure src is not the same as dict (if present) */
if (ress.dictFileName != NULL && UTIL_isSameFileStat(srcFileName, ress.dictFileName, &srcFileStat, &ress.dictFileStat)) {
DISPLAYLEVEL(1, "zstd: cannot use %s as an input file and dictionary \n", srcFileName);
return 1;
}
}
}
/* Check if "srcFile" is compressed. Only done if --exclude-compressed flag is used
* YES => ZSTD will skip compression of the file and will return 0.
* NO => ZSTD will resume with compress operation.
*/
if (prefs->excludeCompressedFiles == 1 && UTIL_isCompressedFile(srcFileName, compressedFileExtensions)) {
DISPLAYLEVEL(4, "File is already compressed : %s \n", srcFileName);
return 0;
}
srcFile = FIO_openSrcFile(prefs, srcFileName, &srcFileStat);
if (srcFile == NULL) return 1; /* srcFile could not be opened */
/* 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);
}
AIO_ReadPool_setFile(ress.readCtx, srcFile);
result = FIO_compressFilename_dstFile(
fCtx, prefs, ress,
dstFileName, srcFileName,
&srcFileStat, compressionLevel);
AIO_ReadPool_closeFile(ress.readCtx);
if ( prefs->removeSrcFile /* --rm */
&& result == 0 /* success */
&& strcmp(srcFileName, stdinmark) /* exception : don't erase 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))
EXM_THROW(1, "zstd: %s: %s", srcFileName, strerror(errno));
}
return result;
}
void FIO_displayCompressionParameters(const FIO_prefs_t* prefs)
{
assert(g_display_prefs.displayLevel >= 4);
FIO_rust_displayCompressionParameters(prefs);
}
int FIO_compressFilename(FIO_ctx_t* const fCtx, FIO_prefs_t* const prefs, const char* dstFileName,
const char* srcFileName, const char* dictFileName,
int compressionLevel, ZSTD_compressionParameters comprParams)
{
cRess_t ress = FIO_createCResources(prefs, dictFileName, UTIL_getFileSize(srcFileName), compressionLevel, comprParams);
int const result = FIO_compressFilename_srcFile(fCtx, prefs, ress, dstFileName, srcFileName, compressionLevel);
#define DISPLAY_LEVEL_DEFAULT 2
FIO_freeCResources(&ress);
return result;
}
/* FIO_determineCompressedName() :
* create a destination filename for compressed srcFileName.
* @return a pointer to it.
* This function never returns an error (it may abort() in case of pb)
*/
static const char*
FIO_determineCompressedName(const char* srcFileName, const char* outDirName, const char* suffix)
{
return FIO_rust_determineCompressedName(srcFileName, outDirName, suffix);
}
static unsigned long long FIO_getLargestFileSize(const char** inFileNames, unsigned nbFiles)
{
return FIO_rust_getLargestFileSize(inFileNames, nbFiles);
}
/* FIO_compressMultipleFilenames() :
* compress nbFiles files
* into either one destination (outFileName),
* or into one file each (outFileName == NULL, but suffix != NULL),
* or into a destination folder (specified with -O)
*/
int FIO_compressMultipleFilenames(FIO_ctx_t* const fCtx,
FIO_prefs_t* const prefs,
const char** inFileNamesTable,
const char* outMirroredRootDirName,
const char* outDirName,
const char* outFileName, const char* suffix,
const char* dictFileName, int compressionLevel,
ZSTD_compressionParameters comprParams)
{
int status;
int error = 0;
cRess_t ress = FIO_createCResources(prefs, dictFileName,
FIO_getLargestFileSize(inFileNamesTable, (unsigned)fCtx->nbFilesTotal),
compressionLevel, comprParams);
/* init */
assert(outFileName != NULL || suffix != NULL);
if (outFileName != NULL) { /* output into a single destination (stdout typically) */
FILE *dstFile;
if (FIO_multiFilesConcatWarning(fCtx, prefs, outFileName, 1 /* displayLevelCutoff */)) {
FIO_freeCResources(&ress);
return 1;
}
dstFile = FIO_openDstFile(fCtx, prefs, NULL, outFileName, DEFAULT_FILE_PERMISSIONS);
if (dstFile == NULL) { /* could not open outFileName */
error = 1;
} else {
AIO_WritePool_setFile(ress.writeCtx, dstFile);
for (; fCtx->currFileIdx < fCtx->nbFilesTotal; ++fCtx->currFileIdx) {
status = FIO_compressFilename_srcFile(fCtx, prefs, ress, outFileName, inFileNamesTable[fCtx->currFileIdx], compressionLevel);
if (!status) fCtx->nbFilesProcessed++;
error |= status;
}
if (AIO_WritePool_closeFile(ress.writeCtx))
EXM_THROW(29, "Write error (%s) : cannot properly close %s",
strerror(errno), outFileName);
}
} else {
if (outMirroredRootDirName)
UTIL_mirrorSourceFilesDirectories(inFileNamesTable, (unsigned)fCtx->nbFilesTotal, outMirroredRootDirName);
for (; fCtx->currFileIdx < fCtx->nbFilesTotal; ++fCtx->currFileIdx) {
const char* const srcFileName = inFileNamesTable[fCtx->currFileIdx];
const char* dstFileName = NULL;
if (outMirroredRootDirName) {
char* validMirroredDirName = UTIL_createMirroredDestDirName(srcFileName, outMirroredRootDirName);
if (validMirroredDirName) {
dstFileName = FIO_determineCompressedName(srcFileName, validMirroredDirName, suffix);
free(validMirroredDirName);
} else {
DISPLAYLEVEL(2, "zstd: --output-dir-mirror cannot compress '%s' into '%s' \n", srcFileName, outMirroredRootDirName);
error=1;
continue;
}
} else {
dstFileName = FIO_determineCompressedName(srcFileName, outDirName, suffix); /* cannot fail */
}
status = FIO_compressFilename_srcFile(fCtx, prefs, ress, dstFileName, srcFileName, compressionLevel);
if (!status) fCtx->nbFilesProcessed++;
error |= status;
}
if (outDirName)
FIO_checkFilenameCollisions(inFileNamesTable , (unsigned)fCtx->nbFilesTotal);
}
if (FIO_shouldDisplayMultipleFileSummary(fCtx)) {
UTIL_HumanReadableSize_t hr_isize = UTIL_makeHumanReadableSize((U64) fCtx->totalBytesInput);
UTIL_HumanReadableSize_t hr_osize = UTIL_makeHumanReadableSize((U64) fCtx->totalBytesOutput);
DISPLAY_PROGRESS("\r%79s\r", "");
if (fCtx->totalBytesInput == 0) {
DISPLAY_SUMMARY("%3d files compressed : (%6.*f%4s => %6.*f%4s)\n",
fCtx->nbFilesProcessed,
hr_isize.precision, hr_isize.value, hr_isize.suffix,
hr_osize.precision, hr_osize.value, hr_osize.suffix);
} else {
DISPLAY_SUMMARY("%3d files compressed : %.2f%% (%6.*f%4s => %6.*f%4s)\n",
fCtx->nbFilesProcessed,
(double)fCtx->totalBytesOutput/((double)fCtx->totalBytesInput)*100,
hr_isize.precision, hr_isize.value, hr_isize.suffix,
hr_osize.precision, hr_osize.value, hr_osize.suffix);
}
}
FIO_freeCResources(&ress);
return error;
}
#endif /* #ifndef ZSTD_NOCOMPRESS */
#ifndef ZSTD_NODECOMPRESS
/* **************************************************************************
* Decompression
***************************************************************************/
typedef struct {
FIO_Dict_t dict;
ZSTD_DStream* dctx;
WritePoolCtx_t *writeCtx;
ReadPoolCtx_t *readCtx;
} dRess_t;
static dRess_t FIO_createDResources(FIO_prefs_t* const prefs, const char* dictFileName)
{
int useMMap = prefs->mmapDict == ZSTD_ps_enable;
int forceNoUseMMap = prefs->mmapDict == ZSTD_ps_disable;
stat_t statbuf;
dRess_t ress;
memset(&statbuf, 0, sizeof(statbuf));
memset(&ress, 0, sizeof(ress));
FIO_getDictFileStat(dictFileName, &statbuf);
if (prefs->patchFromMode){
U64 const dictSize = UTIL_getFileSizeStat(&statbuf);
useMMap |= dictSize > prefs->memLimit;
FIO_adjustMemLimitForPatchFromMode(prefs, dictSize, 0 /* just use the dict size */);
}
/* Allocation */
ress.dctx = ZSTD_createDStream();
if (ress.dctx==NULL)
EXM_THROW(60, "Error: %s : can't create ZSTD_DStream", strerror(errno));
CHECK( ZSTD_DCtx_setMaxWindowSize(ress.dctx, prefs->memLimit) );
CHECK( ZSTD_DCtx_setParameter(ress.dctx, ZSTD_d_forceIgnoreChecksum, !prefs->checksumFlag));
/* dictionary */
{
FIO_dictBufferType_t dictBufferType = (useMMap && !forceNoUseMMap) ? FIO_mmapDict : FIO_mallocDict;
FIO_initDict(&ress.dict, dictFileName, prefs, &statbuf, dictBufferType);
CHECK(ZSTD_DCtx_reset(ress.dctx, ZSTD_reset_session_only) );
if (prefs->patchFromMode){
CHECK(ZSTD_DCtx_refPrefix(ress.dctx, ress.dict.dictBuffer, ress.dict.dictBufferSize));
} else {
CHECK(ZSTD_DCtx_loadDictionary_byReference(ress.dctx, ress.dict.dictBuffer, ress.dict.dictBufferSize));
}
}
ress.writeCtx = AIO_WritePool_create(prefs, ZSTD_DStreamOutSize());
ress.readCtx = AIO_ReadPool_create(prefs, ZSTD_DStreamInSize());
return ress;
}
static void FIO_freeDResources(dRess_t ress)
{
FIO_freeDict(&(ress.dict));
CHECK( ZSTD_freeDStream(ress.dctx) );
AIO_WritePool_free(ress.writeCtx);
AIO_ReadPool_free(ress.readCtx);
}
/* FIO_passThrough() : just copy input into output, for compatibility with gzip -df mode
* @return : 0 (no error) */
static int FIO_passThrough(dRess_t *ress)
{
return FIO_rust_passThrough(ress->readCtx, ress->writeCtx);
}
/* FIO_zstdErrorHelp() :
* detailed error message when requested window size is too large */
static void
FIO_zstdErrorHelp(const FIO_prefs_t* const prefs,
const dRess_t* ress,
size_t err,
const char* srcFileName)
{
ZSTD_FrameHeader header;
/* Help message only for one specific error */
if (ZSTD_getErrorCode(err) != ZSTD_error_frameParameter_windowTooLarge)
return;
/* Try to decode the frame header */
err = ZSTD_getFrameHeader(&header, ress->readCtx->srcBuffer, ress->readCtx->srcBufferLoaded);
if (err == 0) {
unsigned long long const windowSize = header.windowSize;
unsigned const windowLog = FIO_highbit64(windowSize) + ((windowSize & (windowSize - 1)) != 0);
assert(prefs->memLimit > 0);
DISPLAYLEVEL(1, "%s : Window size larger than maximum : %llu > %u \n",
srcFileName, windowSize, prefs->memLimit);
if (windowLog <= ZSTD_WINDOWLOG_MAX) {
unsigned const windowMB = (unsigned)((windowSize >> 20) + ((windowSize & ((1 MB) - 1)) != 0));
assert(windowSize < (U64)(1ULL << 52)); /* ensure now overflow for windowMB */
DISPLAYLEVEL(1, "%s : Use --long=%u or --memory=%uMB \n",
srcFileName, windowLog, windowMB);
return;
} }
DISPLAYLEVEL(1, "%s : Window log larger than ZSTD_WINDOWLOG_MAX=%u; not supported \n",
srcFileName, ZSTD_WINDOWLOG_MAX);
}
/** FIO_decompressFrame() :
* @return : size of decoded zstd frame, or an error code
*/
#define FIO_ERROR_FRAME_DECODING ((unsigned long long)(-2))
static void
FIO_decompressZstdFrameProgress(void* const opaque,
const char* const srcFileName,
U64 const decodedSize)
{
FIO_ctx_t* const fCtx = (FIO_ctx_t*)opaque;
const char* srcFName20 = srcFileName;
size_t const srcFileLength = strlen(srcFileName);
UTIL_HumanReadableSize_t const hrs = UTIL_makeHumanReadableSize(decodedSize);
/* display last 20 characters only when not --verbose */
if ((srcFileLength>20) && (g_display_prefs.displayLevel<3))
srcFName20 += srcFileLength-20;
if (fCtx->nbFilesTotal > 1) {
DISPLAYUPDATE_PROGRESS(
"\rDecompress: %2u/%2u files. Current: %s : %.*f%s... ",
fCtx->currFileIdx+1, fCtx->nbFilesTotal, srcFName20,
hrs.precision, hrs.value, hrs.suffix);
} else {
DISPLAYUPDATE_PROGRESS("\r%-20.20s : %.*f%s... ",
srcFName20, hrs.precision, hrs.value, hrs.suffix);
}
}
static unsigned long long
FIO_decompressZstdFrames(FIO_ctx_t* const fCtx, dRess_t* ress,
const FIO_prefs_t* const prefs,
const char* srcFileName,
U64 alreadyDecoded) /* for multi-frames streams */
{
U64 decodedSize = 0;
size_t zstdError = 0;
int const status = FIO_rust_decompressZstdFrames(
fCtx, ress->dctx, ress->readCtx, ress->writeCtx, srcFileName,
alreadyDecoded, &decodedSize, &zstdError,
FIO_decompressZstdFrameProgress);
if (status == FIO_RUST_ZSTD_FRAME_OK)
return decodedSize;
if (status == FIO_RUST_ZSTD_FRAME_DECODING_ERROR) {
DISPLAYLEVEL(1, "%s : Decoding error (36) : %s \n",
srcFileName, ZSTD_getErrorName(zstdError));
FIO_zstdErrorHelp(prefs, ress, zstdError, srcFileName);
return FIO_ERROR_FRAME_DECODING;
}
if (status == FIO_RUST_ZSTD_FRAME_PREMATURE_END) {
DISPLAYLEVEL(1, "%s : Read error (39) : premature end \n",
srcFileName);
return FIO_ERROR_FRAME_DECODING;
}
assert(0);
return FIO_ERROR_FRAME_DECODING;
}
#ifdef ZSTD_GZDECOMPRESS
static unsigned long long
FIO_decompressGzFrame(dRess_t* ress, const char* srcFileName)
{
unsigned long long outFileSize = 0;
z_stream strm;
int flush = Z_NO_FLUSH;
int decodingError = 0;
IOJob_t *writeJob = NULL;
strm.zalloc = Z_NULL;
strm.zfree = Z_NULL;
strm.opaque = Z_NULL;
strm.next_in = 0;
strm.avail_in = 0;
/* see https://www.zlib.net/manual.html */
if (inflateInit2(&strm, 15 /* maxWindowLogSize */ + 16 /* gzip only */) != Z_OK)
return FIO_ERROR_FRAME_DECODING;
writeJob = AIO_WritePool_acquireJob(ress->writeCtx);
strm.next_out = (Bytef*)writeJob->buffer;
strm.avail_out = (uInt)writeJob->bufferSize;
strm.avail_in = (uInt)ress->readCtx->srcBufferLoaded;
strm.next_in = (z_const unsigned char*)ress->readCtx->srcBuffer;
for ( ; ; ) {
int ret;
if (strm.avail_in == 0) {
AIO_ReadPool_consumeAndRefill(ress->readCtx);
if (ress->readCtx->srcBufferLoaded == 0) flush = Z_FINISH;
strm.next_in = (z_const unsigned char*)ress->readCtx->srcBuffer;
strm.avail_in = (uInt)ress->readCtx->srcBufferLoaded;
}
ret = inflate(&strm, flush);
if (ret == Z_BUF_ERROR) {
DISPLAYLEVEL(1, "zstd: %s: premature gz end \n", srcFileName);
decodingError = 1; break;
}
if (ret != Z_OK && ret != Z_STREAM_END) {
DISPLAYLEVEL(1, "zstd: %s: inflate error %d \n", srcFileName, ret);
decodingError = 1; break;
}
{ size_t const decompBytes = writeJob->bufferSize - strm.avail_out;
if (decompBytes) {
writeJob->usedBufferSize = decompBytes;
AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob);
outFileSize += decompBytes;
strm.next_out = (Bytef*)writeJob->buffer;
strm.avail_out = (uInt)writeJob->bufferSize;
}
}
if (ret == Z_STREAM_END) break;
}
AIO_ReadPool_consumeBytes(ress->readCtx, ress->readCtx->srcBufferLoaded - strm.avail_in);
if ( (inflateEnd(&strm) != Z_OK) /* release resources ; error detected */
&& (decodingError==0) ) {
DISPLAYLEVEL(1, "zstd: %s: inflateEnd error \n", srcFileName);
decodingError = 1;
}
AIO_WritePool_releaseIoJob(writeJob);
AIO_WritePool_sparseWriteEnd(ress->writeCtx);
return decodingError ? FIO_ERROR_FRAME_DECODING : outFileSize;
}
#endif
#ifdef ZSTD_LZMADECOMPRESS
static unsigned long long
FIO_decompressLzmaFrame(dRess_t* ress,
const char* srcFileName, int plain_lzma)
{
unsigned long long outFileSize = 0;
lzma_stream strm = LZMA_STREAM_INIT;
lzma_action action = LZMA_RUN;
lzma_ret initRet;
int decodingError = 0;
IOJob_t *writeJob = NULL;
strm.next_in = 0;
strm.avail_in = 0;
if (plain_lzma) {
initRet = lzma_alone_decoder(&strm, UINT64_MAX); /* LZMA */
} else {
initRet = lzma_stream_decoder(&strm, UINT64_MAX, 0); /* XZ */
}
if (initRet != LZMA_OK) {
DISPLAYLEVEL(1, "zstd: %s: %s error %d \n",
plain_lzma ? "lzma_alone_decoder" : "lzma_stream_decoder",
srcFileName, initRet);
return FIO_ERROR_FRAME_DECODING;
}
writeJob = AIO_WritePool_acquireJob(ress->writeCtx);
strm.next_out = (BYTE*)writeJob->buffer;
strm.avail_out = writeJob->bufferSize;
strm.next_in = (BYTE const*)ress->readCtx->srcBuffer;
strm.avail_in = ress->readCtx->srcBufferLoaded;
for ( ; ; ) {
lzma_ret ret;
if (strm.avail_in == 0) {
AIO_ReadPool_consumeAndRefill(ress->readCtx);
if (ress->readCtx->srcBufferLoaded == 0) action = LZMA_FINISH;
strm.next_in = (BYTE const*)ress->readCtx->srcBuffer;
strm.avail_in = ress->readCtx->srcBufferLoaded;
}
ret = lzma_code(&strm, action);
if (ret == LZMA_BUF_ERROR) {
DISPLAYLEVEL(1, "zstd: %s: premature lzma end \n", srcFileName);
decodingError = 1; break;
}
if (ret != LZMA_OK && ret != LZMA_STREAM_END) {
DISPLAYLEVEL(1, "zstd: %s: lzma_code decoding error %d \n",
srcFileName, ret);
decodingError = 1; break;
}
{ size_t const decompBytes = writeJob->bufferSize - strm.avail_out;
if (decompBytes) {
writeJob->usedBufferSize = decompBytes;
AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob);
outFileSize += decompBytes;
strm.next_out = (BYTE*)writeJob->buffer;
strm.avail_out = writeJob->bufferSize;
} }
if (ret == LZMA_STREAM_END) break;
}
AIO_ReadPool_consumeBytes(ress->readCtx, ress->readCtx->srcBufferLoaded - strm.avail_in);
lzma_end(&strm);
AIO_WritePool_releaseIoJob(writeJob);
AIO_WritePool_sparseWriteEnd(ress->writeCtx);
return decodingError ? FIO_ERROR_FRAME_DECODING : outFileSize;
}
#endif
#ifdef ZSTD_LZ4DECOMPRESS
static unsigned long long
FIO_decompressLz4Frame(dRess_t* ress, const char* srcFileName)
{
unsigned long long filesize = 0;
LZ4F_errorCode_t nextToLoad = 4;
LZ4F_decompressionContext_t dCtx;
LZ4F_errorCode_t const errorCode = LZ4F_createDecompressionContext(&dCtx, LZ4F_VERSION);
int decodingError = 0;
IOJob_t *writeJob = NULL;
if (LZ4F_isError(errorCode)) {
DISPLAYLEVEL(1, "zstd: failed to create lz4 decompression context \n");
return FIO_ERROR_FRAME_DECODING;
}
writeJob = AIO_WritePool_acquireJob(ress->writeCtx);
/* Main Loop */
for (;nextToLoad;) {
size_t pos = 0;
size_t decodedBytes = writeJob->bufferSize;
int fullBufferDecoded = 0;
/* Read input */
AIO_ReadPool_fillBuffer(ress->readCtx, nextToLoad);
if(!ress->readCtx->srcBufferLoaded) break; /* reached end of file */
while ((pos < ress->readCtx->srcBufferLoaded) || fullBufferDecoded) { /* still to read, or still to flush */
/* Decode Input (at least partially) */
size_t remaining = ress->readCtx->srcBufferLoaded - pos;
decodedBytes = writeJob->bufferSize;
nextToLoad = LZ4F_decompress(dCtx, writeJob->buffer, &decodedBytes, (char*)(ress->readCtx->srcBuffer)+pos,
&remaining, NULL);
if (LZ4F_isError(nextToLoad)) {
DISPLAYLEVEL(1, "zstd: %s: lz4 decompression error : %s \n",
srcFileName, LZ4F_getErrorName(nextToLoad));
decodingError = 1; nextToLoad = 0; break;
}
pos += remaining;
assert(pos <= ress->readCtx->srcBufferLoaded);
fullBufferDecoded = decodedBytes == writeJob->bufferSize;
/* Write Block */
if (decodedBytes) {
UTIL_HumanReadableSize_t hrs;
writeJob->usedBufferSize = decodedBytes;
AIO_WritePool_enqueueAndReacquireWriteJob(&writeJob);
filesize += decodedBytes;
hrs = UTIL_makeHumanReadableSize(filesize);
DISPLAYUPDATE_PROGRESS("\rDecompressed : %.*f%s ", hrs.precision, hrs.value, hrs.suffix);
}
if (!nextToLoad) break;
}
AIO_ReadPool_consumeBytes(ress->readCtx, pos);
}
if (nextToLoad!=0) {
DISPLAYLEVEL(1, "zstd: %s: unfinished lz4 stream \n", srcFileName);
decodingError=1;
}
LZ4F_freeDecompressionContext(dCtx);
AIO_WritePool_releaseIoJob(writeJob);
AIO_WritePool_sparseWriteEnd(ress->writeCtx);
return decodingError ? FIO_ERROR_FRAME_DECODING : filesize;
}
#endif
/* Rust owns the mixed-format loop, while these adapters keep each codec and
* the private dRess_t layout in this C translation unit. */
typedef struct {
FIO_ctx_t* fCtx;
dRess_t* ress;
const FIO_prefs_t* prefs;
} FIO_rust_decompression_projection_t;
static int FIO_rust_decompressZstdFrameCallback(void* opaque,
const char* srcFileName,
U64 alreadyDecoded,
U64* frameSize,
size_t* errorCode,
int mode)
{
FIO_rust_decompression_projection_t* const projection =
(FIO_rust_decompression_projection_t*)opaque;
(void)mode;
*errorCode = 0;
*frameSize = FIO_decompressZstdFrames(projection->fCtx,
projection->ress,
projection->prefs,
srcFileName,
alreadyDecoded);
return *frameSize == FIO_ERROR_FRAME_DECODING;
}
#ifdef ZSTD_GZDECOMPRESS
static int FIO_rust_decompressGzFrameCallback(void* opaque,
const char* srcFileName,
U64 alreadyDecoded,
U64* frameSize,
size_t* errorCode,
int mode)
{
FIO_rust_decompression_projection_t* const projection =
(FIO_rust_decompression_projection_t*)opaque;
(void)alreadyDecoded;
(void)mode;
*errorCode = 0;
*frameSize = FIO_decompressGzFrame(projection->ress, srcFileName);
return *frameSize == FIO_ERROR_FRAME_DECODING;
}
#endif
#ifdef ZSTD_LZMADECOMPRESS
static int FIO_rust_decompressLzmaFrameCallback(void* opaque,
const char* srcFileName,
U64 alreadyDecoded,
U64* frameSize,
size_t* errorCode,
int mode)
{
FIO_rust_decompression_projection_t* const projection =
(FIO_rust_decompression_projection_t*)opaque;
(void)alreadyDecoded;
*errorCode = 0;
*frameSize = FIO_decompressLzmaFrame(projection->ress, srcFileName, mode);
return *frameSize == FIO_ERROR_FRAME_DECODING;
}
#endif
#ifdef ZSTD_LZ4DECOMPRESS
static int FIO_rust_decompressLz4FrameCallback(void* opaque,
const char* srcFileName,
U64 alreadyDecoded,
U64* frameSize,
size_t* errorCode,
int mode)
{
FIO_rust_decompression_projection_t* const projection =
(FIO_rust_decompression_projection_t*)opaque;
(void)alreadyDecoded;
(void)mode;
*errorCode = 0;
*frameSize = FIO_decompressLz4Frame(projection->ress, srcFileName);
return *frameSize == FIO_ERROR_FRAME_DECODING;
}
#endif
static int FIO_rust_decompressPassThroughCallback(void* opaque)
{
FIO_rust_decompression_projection_t* const projection =
(FIO_rust_decompression_projection_t*)opaque;
return FIO_passThrough(projection->ress);
}
/** FIO_decompressFrames() :
* Find and decode frames inside srcFile
* srcFile presumed opened and valid
* @return : 0 : OK
* 1 : error
*/
static int FIO_decompressFrames(FIO_ctx_t* const fCtx,
dRess_t ress, const FIO_prefs_t* const prefs,
const char* dstFileName, const char* srcFileName)
{
U64 filesize = 0;
int passThrough = prefs->passThrough;
FIO_rust_decompression_projection_t projection;
FIO_rust_decompress_callbacks_t callbacks;
int status;
if (passThrough == -1) {
/* If pass-through mode is not explicitly enabled or disabled,
* default to the legacy behavior of enabling it if we are writing
* to stdout with the overwrite flag enabled.
*/
passThrough = prefs->overwrite && !strcmp(dstFileName, stdoutmark);
}
assert(passThrough == 0 || passThrough == 1);
projection.fCtx = fCtx;
projection.ress = &ress;
projection.prefs = prefs;
memset(&callbacks, 0, sizeof(callbacks));
callbacks.opaque = &projection;
callbacks.decode_zstd = FIO_rust_decompressZstdFrameCallback;
#ifdef ZSTD_GZDECOMPRESS
callbacks.decode_gzip = FIO_rust_decompressGzFrameCallback;
#endif
#ifdef ZSTD_LZMADECOMPRESS
callbacks.decode_lzma = FIO_rust_decompressLzmaFrameCallback;
#endif
#ifdef ZSTD_LZ4DECOMPRESS
callbacks.decode_lz4 = FIO_rust_decompressLz4FrameCallback;
#endif
callbacks.pass_through = FIO_rust_decompressPassThroughCallback;
status = FIO_rust_decompressFrames(ress.readCtx, srcFileName, passThrough,
&filesize, &callbacks);
switch (status) {
case FIO_RUST_DECOMPRESS_OK:
break;
case FIO_RUST_DECOMPRESS_PASS_THROUGH:
return 0;
case FIO_RUST_DECOMPRESS_EMPTY_INPUT:
DISPLAYLEVEL(1, "zstd: %s: unexpected end of file \n", srcFileName);
return 1;
case FIO_RUST_DECOMPRESS_SHORT_INPUT:
DISPLAYLEVEL(1, "zstd: %s: unknown header \n", srcFileName);
return 1;
case FIO_RUST_DECOMPRESS_GZIP_UNSUPPORTED:
DISPLAYLEVEL(1, "zstd: %s: gzip file cannot be uncompressed (zstd compiled without HAVE_ZLIB) -- ignored \n", srcFileName);
return 1;
case FIO_RUST_DECOMPRESS_LZMA_UNSUPPORTED:
DISPLAYLEVEL(1, "zstd: %s: xz/lzma file cannot be uncompressed (zstd compiled without HAVE_LZMA) -- ignored \n", srcFileName);
return 1;
case FIO_RUST_DECOMPRESS_LZ4_UNSUPPORTED:
DISPLAYLEVEL(1, "zstd: %s: lz4 file cannot be uncompressed (zstd compiled without HAVE_LZ4) -- ignored \n", srcFileName);
return 1;
case FIO_RUST_DECOMPRESS_UNSUPPORTED_FORMAT:
DISPLAYLEVEL(1, "zstd: %s: unsupported format \n", srcFileName);
return 1;
case FIO_RUST_DECOMPRESS_FRAME_ERROR:
case FIO_RUST_DECOMPRESS_PASS_THROUGH_ERROR:
case FIO_RUST_DECOMPRESS_ZSTD_UNSUPPORTED:
return 1;
default:
assert(0);
return 1;
}
/* Final Status */
fCtx->totalBytesOutput += (size_t)filesize;
DISPLAY_PROGRESS("\r%79s\r", "");
if (FIO_shouldDisplayFileSummary(fCtx))
DISPLAY_SUMMARY("%-20s: %llu bytes \n", srcFileName,
(unsigned long long)filesize);
return 0;
}
/** 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;
stat_t srcFileStat;
int result;
U64 fileSize = UTIL_FILESIZE_UNKNOWN;
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;
/* 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);
}
AIO_ReadPool_setFile(ress.readCtx, srcFile);
result = FIO_decompressDstFile(fCtx, prefs, ress, dstFileName, srcFileName, &srcFileStat);
AIO_ReadPool_setFile(ress.readCtx, NULL);
/* Close file */
if (fclose(srcFile)) {
DISPLAYLEVEL(1, "zstd: %s: %s \n", srcFileName, strerror(errno)); /* error should not happen */
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;
}
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);
FIO_freeDResources(ress);
return decodingError;
}
static const char *suffixList[] = {
ZSTD_EXTENSION,
TZSTD_EXTENSION,
#ifndef ZSTD_NODECOMPRESS
ZSTD_ALT_EXTENSION,
#endif
#ifdef ZSTD_GZDECOMPRESS
GZ_EXTENSION,
TGZ_EXTENSION,
#endif
#ifdef ZSTD_LZMADECOMPRESS
LZMA_EXTENSION,
XZ_EXTENSION,
TXZ_EXTENSION,
#endif
#ifdef ZSTD_LZ4DECOMPRESS
LZ4_EXTENSION,
TLZ4_EXTENSION,
#endif
NULL
};
static const char *suffixListStr =
ZSTD_EXTENSION "/" TZSTD_EXTENSION
#ifdef ZSTD_GZDECOMPRESS
"/" GZ_EXTENSION "/" TGZ_EXTENSION
#endif
#ifdef ZSTD_LZMADECOMPRESS
"/" LZMA_EXTENSION "/" XZ_EXTENSION "/" TXZ_EXTENSION
#endif
#ifdef ZSTD_LZ4DECOMPRESS
"/" LZ4_EXTENSION "/" TLZ4_EXTENSION
#endif
;
/* FIO_determineDstName() :
* create a destination filename from a srcFileName.
* @return a pointer to it.
* @return == NULL if there is an error */
static const char*
FIO_determineDstName(const char* srcFileName, const char* outDirName)
{
return FIO_rust_determineDstName(srcFileName, outDirName, suffixList, suffixListStr);
}
int
FIO_decompressMultipleFilenames(FIO_ctx_t* const fCtx,
FIO_prefs_t* const prefs,
const char** srcNamesTable,
const char* outMirroredRootDirName,
const char* outDirName, const char* outFileName,
const char* dictFileName)
{
int status;
int error = 0;
dRess_t ress = FIO_createDResources(prefs, dictFileName);
if (outFileName) {
if (FIO_multiFilesConcatWarning(fCtx, prefs, outFileName, 1 /* displayLevelCutoff */)) {
FIO_freeDResources(ress);
return 1;
}
if (!prefs->testMode) {
FILE* dstFile = FIO_openDstFile(fCtx, prefs, NULL, outFileName, DEFAULT_FILE_PERMISSIONS);
if (dstFile == 0) EXM_THROW(19, "cannot open %s", outFileName);
AIO_WritePool_setFile(ress.writeCtx, dstFile);
}
for (; fCtx->currFileIdx < fCtx->nbFilesTotal; fCtx->currFileIdx++) {
status = FIO_decompressSrcFile(fCtx, prefs, ress, outFileName, srcNamesTable[fCtx->currFileIdx]);
if (!status) fCtx->nbFilesProcessed++;
error |= status;
}
if ((!prefs->testMode) && (AIO_WritePool_closeFile(ress.writeCtx)))
EXM_THROW(72, "Write error : %s : cannot properly close output file",
strerror(errno));
} else {
if (outMirroredRootDirName)
UTIL_mirrorSourceFilesDirectories(srcNamesTable, (unsigned)fCtx->nbFilesTotal, outMirroredRootDirName);
for (; fCtx->currFileIdx < fCtx->nbFilesTotal; fCtx->currFileIdx++) { /* create dstFileName */
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)
FIO_checkFilenameCollisions(srcNamesTable , (unsigned)fCtx->nbFilesTotal);
}
if (FIO_shouldDisplayMultipleFileSummary(fCtx)) {
DISPLAY_PROGRESS("\r%79s\r", "");
DISPLAY_SUMMARY("%d files decompressed : %6llu bytes total \n",
fCtx->nbFilesProcessed, (unsigned long long)fCtx->totalBytesOutput);
}
FIO_freeDResources(ress);
return error;
}
/* **************************************************************************
* .zst file info (--list command)
***************************************************************************/
typedef struct {
U64 decompressedSize;
U64 compressedSize;
U64 windowSize;
int numActualFrames;
int numSkippableFrames;
int decompUnavailable;
int usesCheck;
BYTE checksum[4];
U32 nbFiles;
unsigned dictID;
} fileInfo_t;
typedef char FIO_rust_file_info_compressed_size_offset[
(offsetof(fileInfo_t, compressedSize) == sizeof(U64)) ? 1 : -1];
typedef char FIO_rust_file_info_window_size_offset[
(offsetof(fileInfo_t, windowSize) == 2 * sizeof(U64)) ? 1 : -1];
typedef char FIO_rust_file_info_frame_counts_offset[
(offsetof(fileInfo_t, numActualFrames) == 3 * sizeof(U64)) ? 1 : -1];
typedef char FIO_rust_file_info_checksum_offset[
(offsetof(fileInfo_t, checksum) == 3 * sizeof(U64) + 4 * sizeof(int)) ? 1 : -1];
typedef char FIO_rust_file_info_size[
(sizeof(fileInfo_t) == (sizeof(void*) == 8
? 7 * sizeof(U64)
: 13 * sizeof(U32))) ? 1 : -1];
void FIO_rust_addFInfo(fileInfo_t* output, fileInfo_t const* fi1, fileInfo_t const* fi2);
static fileInfo_t FIO_addFInfo(fileInfo_t fi1, fileInfo_t fi2)
{
fileInfo_t total;
FIO_rust_addFInfo(&total, &fi1, &fi2);
return total;
}
typedef enum {
info_success=0,
info_frame_error=1,
info_not_zstd=2,
info_file_error=3,
info_truncated_input=4
} InfoError;
enum {
FIO_RUST_ANALYZE_DIAG_SEEKED_PAST_FILE = 0,
FIO_RUST_ANALYZE_DIAG_INCOMPLETE_FRAME = 1,
FIO_RUST_ANALYZE_DIAG_RAN_OUT_OF_FRAMES = 2,
FIO_RUST_ANALYZE_DIAG_DECODE_FRAME_HEADER = 3,
FIO_RUST_ANALYZE_DIAG_FRAME_HEADER_SIZE = 4,
FIO_RUST_ANALYZE_DIAG_MOVE_TO_FRAME_HEADER_END = 5,
FIO_RUST_ANALYZE_DIAG_BLOCK_HEADER = 6,
FIO_RUST_ANALYZE_DIAG_UNSUPPORTED_BLOCK_TYPE = 7,
FIO_RUST_ANALYZE_DIAG_SKIP_BLOCK = 8,
FIO_RUST_ANALYZE_DIAG_CHECKSUM = 9,
FIO_RUST_ANALYZE_DIAG_SKIP_FRAME = 10,
FIO_RUST_ANALYZE_DIAG_MIXED_DICTIONARY_IDS = 11
};
int FIO_rust_analyzeFrames(fileInfo_t* info, FILE* srcFile);
void FIO_rust_analyzeFrames_display(int diagnostic,
unsigned long long filePosition,
unsigned long long fileSize);
#define ERROR_IF(c,n,...) { \
if (c) { \
DISPLAYLEVEL(1, __VA_ARGS__); \
DISPLAYLEVEL(1, " \n"); \
return n; \
} \
}
void
FIO_rust_analyzeFrames_display(int diagnostic,
unsigned long long filePosition,
unsigned long long fileSize)
{
switch (diagnostic) {
case FIO_RUST_ANALYZE_DIAG_SEEKED_PAST_FILE:
DISPLAYLEVEL(1,
"Error: seeked to position %llu, which is beyond file size of %llu\n",
filePosition, fileSize);
break;
case FIO_RUST_ANALYZE_DIAG_INCOMPLETE_FRAME:
DISPLAYLEVEL(1, "Error: reached end of file with incomplete frame");
break;
case FIO_RUST_ANALYZE_DIAG_RAN_OUT_OF_FRAMES:
DISPLAYLEVEL(1, "Error: did not reach end of file but ran out of frames");
break;
case FIO_RUST_ANALYZE_DIAG_DECODE_FRAME_HEADER:
DISPLAYLEVEL(1, "Error: could not decode frame header");
break;
case FIO_RUST_ANALYZE_DIAG_FRAME_HEADER_SIZE:
DISPLAYLEVEL(1, "Error: could not determine frame header size");
break;
case FIO_RUST_ANALYZE_DIAG_MOVE_TO_FRAME_HEADER_END:
DISPLAYLEVEL(1, "Error: could not move to end of frame header");
break;
case FIO_RUST_ANALYZE_DIAG_BLOCK_HEADER:
DISPLAYLEVEL(1, "Error while reading block header");
break;
case FIO_RUST_ANALYZE_DIAG_UNSUPPORTED_BLOCK_TYPE:
DISPLAYLEVEL(1, "Error: unsupported block type");
break;
case FIO_RUST_ANALYZE_DIAG_SKIP_BLOCK:
DISPLAYLEVEL(1, "Error: could not skip to end of block");
break;
case FIO_RUST_ANALYZE_DIAG_CHECKSUM:
DISPLAYLEVEL(1, "Error: could not read checksum");
break;
case FIO_RUST_ANALYZE_DIAG_SKIP_FRAME:
DISPLAYLEVEL(1, "Error: could not find end of skippable frame");
break;
case FIO_RUST_ANALYZE_DIAG_MIXED_DICTIONARY_IDS:
DISPLAY("WARNING: File contains multiple frames with different dictionary IDs. Showing dictID 0 instead");
return;
default:
assert(0);
return;
}
DISPLAYLEVEL(1, " \n");
}
static InfoError
FIO_analyzeFrames(fileInfo_t* info, FILE* const srcFile)
{
return (InfoError)FIO_rust_analyzeFrames(info, srcFile);
}
static InfoError
getFileInfo_fileConfirmed(fileInfo_t* info, const char* inFileName)
{
InfoError status;
stat_t srcFileStat;
FILE* const srcFile = FIO_openSrcFile(NULL, inFileName, &srcFileStat);
ERROR_IF(srcFile == NULL, info_file_error, "Error: could not open source file %s", inFileName);
info->compressedSize = UTIL_getFileSizeStat(&srcFileStat);
status = FIO_analyzeFrames(info, srcFile);
fclose(srcFile);
info->nbFiles = 1;
return status;
}
/** getFileInfo() :
* Reads information from file, stores in *info
* @return : InfoError status
*/
static InfoError
getFileInfo(fileInfo_t* info, const char* srcFileName)
{
ERROR_IF(!UTIL_isRegularFile(srcFileName),
info_file_error, "Error : %s is not a file", srcFileName);
return getFileInfo_fileConfirmed(info, srcFileName);
}
static void
displayInfo(const char* inFileName, const fileInfo_t* info, int displayLevel)
{
UTIL_HumanReadableSize_t const window_hrs = UTIL_makeHumanReadableSize(info->windowSize);
UTIL_HumanReadableSize_t const compressed_hrs = UTIL_makeHumanReadableSize(info->compressedSize);
UTIL_HumanReadableSize_t const decompressed_hrs = UTIL_makeHumanReadableSize(info->decompressedSize);
double const ratio = (info->compressedSize == 0) ? 0 : ((double)info->decompressedSize)/(double)info->compressedSize;
const char* const checkString = (info->usesCheck ? "XXH64" : "None");
if (displayLevel <= 2) {
if (!info->decompUnavailable) {
DISPLAYOUT("%6d %5d %6.*f%4s %8.*f%4s %5.3f %5s %s\n",
info->numSkippableFrames + info->numActualFrames,
info->numSkippableFrames,
compressed_hrs.precision, compressed_hrs.value, compressed_hrs.suffix,
decompressed_hrs.precision, decompressed_hrs.value, decompressed_hrs.suffix,
ratio, checkString, inFileName);
} else {
DISPLAYOUT("%6d %5d %6.*f%4s %5s %s\n",
info->numSkippableFrames + info->numActualFrames,
info->numSkippableFrames,
compressed_hrs.precision, compressed_hrs.value, compressed_hrs.suffix,
checkString, inFileName);
}
} else {
DISPLAYOUT("%s \n", inFileName);
DISPLAYOUT("# Zstandard Frames: %d\n", info->numActualFrames);
if (info->numSkippableFrames)
DISPLAYOUT("# Skippable Frames: %d\n", info->numSkippableFrames);
DISPLAYOUT("DictID: %u\n", info->dictID);
DISPLAYOUT("Window Size: %.*f%s (%llu B)\n",
window_hrs.precision, window_hrs.value, window_hrs.suffix,
(unsigned long long)info->windowSize);
DISPLAYOUT("Compressed Size: %.*f%s (%llu B)\n",
compressed_hrs.precision, compressed_hrs.value, compressed_hrs.suffix,
(unsigned long long)info->compressedSize);
if (!info->decompUnavailable) {
DISPLAYOUT("Decompressed Size: %.*f%s (%llu B)\n",
decompressed_hrs.precision, decompressed_hrs.value, decompressed_hrs.suffix,
(unsigned long long)info->decompressedSize);
DISPLAYOUT("Ratio: %.4f\n", ratio);
}
if (info->usesCheck && info->numActualFrames == 1) {
DISPLAYOUT("Check: %s %02x%02x%02x%02x\n", checkString,
info->checksum[3], info->checksum[2],
info->checksum[1], info->checksum[0]
);
} else {
DISPLAYOUT("Check: %s\n", checkString);
}
DISPLAYOUT("\n");
}
}
static int
FIO_listFile(fileInfo_t* total, const char* inFileName, int displayLevel)
{
fileInfo_t info;
memset(&info, 0, sizeof(info));
{ InfoError const error = getFileInfo(&info, inFileName);
switch (error) {
case info_frame_error:
/* display error, but provide output */
DISPLAYLEVEL(1, "Error while parsing \"%s\" \n", inFileName);
break;
case info_not_zstd:
DISPLAYOUT("File \"%s\" not compressed by zstd \n", inFileName);
if (displayLevel > 2) DISPLAYOUT("\n");
return 1;
case info_file_error:
/* error occurred while opening the file */
if (displayLevel > 2) DISPLAYOUT("\n");
return 1;
case info_truncated_input:
DISPLAYOUT("File \"%s\" is truncated \n", inFileName);
if (displayLevel > 2) DISPLAYOUT("\n");
return 1;
case info_success:
default:
break;
}
displayInfo(inFileName, &info, displayLevel);
*total = FIO_addFInfo(*total, info);
assert(error == info_success || error == info_frame_error);
return (int)error;
}
}
int FIO_listMultipleFiles(unsigned numFiles, const char** filenameTable, int displayLevel)
{
/* ensure no specified input is stdin (needs fseek() capability) */
{ unsigned u;
for (u=0; u<numFiles;u++) {
ERROR_IF(!strcmp (filenameTable[u], stdinmark),
1, "zstd: --list does not support reading from standard input");
} }
if (numFiles == 0) {
if (!UTIL_isConsole(stdin)) {
DISPLAYLEVEL(1, "zstd: --list does not support reading from standard input \n");
}
DISPLAYLEVEL(1, "No files given \n");
return 1;
}
if (displayLevel <= 2) {
DISPLAYOUT("Frames Skips Compressed Uncompressed Ratio Check Filename\n");
}
{ int error = 0;
fileInfo_t total;
memset(&total, 0, sizeof(total));
total.usesCheck = 1;
/* --list each file, and check for any error */
{ unsigned u;
for (u=0; u<numFiles;u++) {
error |= FIO_listFile(&total, filenameTable[u], displayLevel);
} }
if (numFiles > 1 && displayLevel <= 2) { /* display total */
UTIL_HumanReadableSize_t const compressed_hrs = UTIL_makeHumanReadableSize(total.compressedSize);
UTIL_HumanReadableSize_t const decompressed_hrs = UTIL_makeHumanReadableSize(total.decompressedSize);
double const ratio = (total.compressedSize == 0) ? 0 : ((double)total.decompressedSize)/(double)total.compressedSize;
const char* const checkString = (total.usesCheck ? "XXH64" : "");
DISPLAYOUT("----------------------------------------------------------------- \n");
if (total.decompUnavailable) {
DISPLAYOUT("%6d %5d %6.*f%4s %5s %u files\n",
total.numSkippableFrames + total.numActualFrames,
total.numSkippableFrames,
compressed_hrs.precision, compressed_hrs.value, compressed_hrs.suffix,
checkString, (unsigned)total.nbFiles);
} else {
DISPLAYOUT("%6d %5d %6.*f%4s %8.*f%4s %5.3f %5s %u files\n",
total.numSkippableFrames + total.numActualFrames,
total.numSkippableFrames,
compressed_hrs.precision, compressed_hrs.value, compressed_hrs.suffix,
decompressed_hrs.precision, decompressed_hrs.value, decompressed_hrs.suffix,
ratio, checkString, (unsigned)total.nbFiles);
} }
return error;
}
}
#endif /* #ifndef ZSTD_NODECOMPRESS */