feat(cli): move zstd compression loop into Rust
The CLI's zstd-format compression path previously kept the asynchronous read, compressStream2, output-job, flush, accounting, adaptive-policy, and progress loop together in fileio.c. That made the high-level stream orchestration another large C-owned surface even though the Rust fileio module already owned the neighboring format loops. Move the format-independent zstd stream loop and read/output accounting into FIO_rust_compressZstdFrame. The C adapter now projects read-pool, write-pool, codec, and policy operations through narrow callbacks. C retains the private ZSTD_CCtx interaction, adaptive-level policy, memory diagnostics, progress formatting, and CLI error mapping, so no private C layout crosses into Rust. Focused seam tests cover input/output ordering, final flush behavior, codec error propagation, and incomplete known-size input handling. Test Plan: - `cargo test --manifest-path rust/Cargo.toml --all-targets -- --test-threads=1` -- 513 passed under the 40 GiB virtual-memory cap. - `cargo test --manifest-path rust/cli/Cargo.toml --all-targets -- --test-threads=1` -- 169 passed under the cap. - `cargo clippy` lib/benches/tests for `rust` and `rust/cli`, with `-D warnings`, and nightly formatting -- passed. - `make -B -C lib -j1 lib` and `make -B -C programs -j1 zstd` -- passed. - Native CLI, full zstd, fuzzer, zstream, and decode-corpus targets -- passed serially under the cap. GPG signing was attempted but unavailable because no pinentry process was available; this repository's preceding commits are unsigned, so this commit uses the explicit unsigned fallback.
This commit is contained in:
+375
-221
@@ -1125,6 +1125,55 @@ typedef struct {
|
||||
FIO_rust_compress_status_display_fn display_status;
|
||||
} FIO_rust_compress_callbacks_t;
|
||||
|
||||
enum {
|
||||
FIO_RUST_ZSTD_OK = 0,
|
||||
FIO_RUST_ZSTD_COMPRESS_ERROR = 1,
|
||||
FIO_RUST_ZSTD_INCOMPLETE_INPUT = 2,
|
||||
FIO_RUST_ZSTD_INVALID_PROJECTION = 3
|
||||
};
|
||||
|
||||
typedef size_t (*FIO_rust_zstd_read_fill_fn)(
|
||||
void* opaque, size_t requested,
|
||||
const unsigned char** buffer, size_t* loaded);
|
||||
typedef void (*FIO_rust_zstd_read_consume_fn)(void* opaque, size_t consumed);
|
||||
typedef void (*FIO_rust_zstd_write_acquire_fn)(
|
||||
void* opaque, void** job, unsigned char** buffer, size_t* bufferSize);
|
||||
typedef void (*FIO_rust_zstd_write_enqueue_fn)(
|
||||
void* opaque, void** job, size_t usedBufferSize,
|
||||
unsigned char** buffer, size_t* bufferSize);
|
||||
typedef void (*FIO_rust_zstd_write_release_fn)(void* opaque, void* job);
|
||||
typedef void (*FIO_rust_zstd_sparse_write_end_fn)(void* opaque);
|
||||
typedef int (*FIO_rust_zstd_compress_stream_fn)(
|
||||
void* opaque, const char* srcFileName, int directive,
|
||||
const unsigned char* input, size_t inputSize, size_t inputPos,
|
||||
unsigned char* output, size_t outputSize,
|
||||
size_t* inputPosAfter, size_t* outputProduced,
|
||||
size_t* toFlushNow, size_t* zstdResult);
|
||||
typedef void (*FIO_rust_zstd_iteration_fn)(
|
||||
void* opaque, const char* srcFileName, int* compressionLevel,
|
||||
size_t oldInputPos, size_t newInputPos, size_t toFlushNow);
|
||||
|
||||
typedef struct {
|
||||
void* readOpaque;
|
||||
void* writeOpaque;
|
||||
void* codecOpaque;
|
||||
void* policyOpaque;
|
||||
size_t readBufferSize;
|
||||
FIO_rust_zstd_read_fill_fn readFill;
|
||||
FIO_rust_zstd_read_consume_fn readConsume;
|
||||
FIO_rust_zstd_write_acquire_fn writeAcquire;
|
||||
FIO_rust_zstd_write_enqueue_fn writeEnqueue;
|
||||
FIO_rust_zstd_write_release_fn writeRelease;
|
||||
FIO_rust_zstd_sparse_write_end_fn sparseWriteEnd;
|
||||
FIO_rust_zstd_compress_stream_fn compressStream;
|
||||
FIO_rust_zstd_iteration_fn iteration;
|
||||
} FIO_rust_zstd_compress_projection_t;
|
||||
|
||||
int FIO_rust_compressZstdFrame(
|
||||
const FIO_rust_zstd_compress_projection_t* projection,
|
||||
const char* srcFileName, U64 srcFileSize, int compressionLevel,
|
||||
U64* readsize, U64* compressedSize, size_t* zstdResult);
|
||||
|
||||
enum {
|
||||
FIO_RUST_GZIP_OK = 0,
|
||||
FIO_RUST_GZIP_INIT_ERROR = 1,
|
||||
@@ -1913,241 +1962,266 @@ FIO_compressLz4Frame(cRess_t* ress,
|
||||
}
|
||||
#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)
|
||||
typedef enum {
|
||||
FIO_rust_zstd_noChange,
|
||||
FIO_rust_zstd_slower,
|
||||
FIO_rust_zstd_faster
|
||||
} FIO_rust_zstd_speed_change_e;
|
||||
|
||||
typedef struct {
|
||||
FIO_ctx_t* fCtx;
|
||||
FIO_prefs_t* prefs;
|
||||
ZSTD_CCtx* cctx;
|
||||
ZSTD_frameProgression previousZfpUpdate;
|
||||
ZSTD_frameProgression previousZfpCorrection;
|
||||
FIO_rust_zstd_speed_change_e speedChange;
|
||||
unsigned flushWaiting;
|
||||
unsigned inputPresented;
|
||||
unsigned inputBlocked;
|
||||
unsigned lastJobID;
|
||||
UTIL_time_t lastAdaptTime;
|
||||
U64 srcFileSize;
|
||||
UTIL_HumanReadableSize_t fileHrs;
|
||||
} FIO_rust_zstd_projection_context_t;
|
||||
|
||||
static size_t FIO_rust_zstd_readFill(void* opaque, size_t requested,
|
||||
const unsigned char** buffer, size_t* loaded)
|
||||
{
|
||||
cRess_t const ress = *ressPtr;
|
||||
IOJob_t* writeJob = AIO_WritePool_acquireJob(ressPtr->writeCtx);
|
||||
ReadPoolCtx_t* const readCtx = (ReadPoolCtx_t*)opaque;
|
||||
size_t const added = AIO_ReadPool_fillBuffer(readCtx, requested);
|
||||
*buffer = readCtx->srcBuffer;
|
||||
*loaded = readCtx->srcBufferLoaded;
|
||||
DISPLAYLEVEL(6, "fread %u bytes from source \n", (unsigned)added);
|
||||
return added;
|
||||
}
|
||||
|
||||
U64 compressedfilesize = 0;
|
||||
ZSTD_EndDirective directive = ZSTD_e_continue;
|
||||
U64 pledgedSrcSize = ZSTD_CONTENTSIZE_UNKNOWN;
|
||||
static void FIO_rust_zstd_readConsume(void* opaque, size_t consumed)
|
||||
{
|
||||
AIO_ReadPool_consumeBytes((ReadPoolCtx_t*)opaque, consumed);
|
||||
}
|
||||
|
||||
/* 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;
|
||||
static void FIO_rust_zstd_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;
|
||||
}
|
||||
|
||||
UTIL_HumanReadableSize_t const file_hrs = UTIL_makeHumanReadableSize(fileSize);
|
||||
static void FIO_rust_zstd_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;
|
||||
}
|
||||
|
||||
DISPLAYLEVEL(6, "compression using zstd format \n");
|
||||
static void FIO_rust_zstd_writeRelease(void* opaque, void* job)
|
||||
{
|
||||
AIO_WritePool_releaseIoJob((IOJob_t*)job);
|
||||
(void)opaque;
|
||||
}
|
||||
|
||||
/* 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) );
|
||||
static void FIO_rust_zstd_sparseWriteEnd(void* opaque)
|
||||
{
|
||||
AIO_WritePool_sparseWriteEnd((WritePoolCtx_t*)opaque);
|
||||
}
|
||||
|
||||
static int FIO_rust_zstd_compressStream(
|
||||
void* opaque, const char* srcFileName, int directive,
|
||||
const unsigned char* input, size_t inputSize, size_t inputPos,
|
||||
unsigned char* output, size_t outputSize,
|
||||
size_t* inputPosAfter, size_t* outputProduced,
|
||||
size_t* toFlushNow, size_t* zstdResult)
|
||||
{
|
||||
FIO_rust_zstd_projection_context_t* const context =
|
||||
(FIO_rust_zstd_projection_context_t*)opaque;
|
||||
ZSTD_inBuffer inBuff = setInBuffer(input, inputSize, inputPos);
|
||||
ZSTD_outBuffer outBuff = setOutBuffer(output, outputSize, 0);
|
||||
size_t const toFlush = ZSTD_toFlushNow(context->cctx);
|
||||
size_t const result = ZSTD_compressStream2(
|
||||
context->cctx, &outBuff, &inBuff, (ZSTD_EndDirective)directive);
|
||||
|
||||
*inputPosAfter = inBuff.pos;
|
||||
*outputProduced = outBuff.pos;
|
||||
*toFlushNow = toFlush;
|
||||
*zstdResult = result;
|
||||
if (!ZSTD_isError(result)) {
|
||||
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);
|
||||
}
|
||||
(void)srcFileName;
|
||||
return ZSTD_isError(result);
|
||||
}
|
||||
|
||||
{ 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;
|
||||
static void FIO_rust_zstd_iteration(void* opaque, const char* srcFileName,
|
||||
int* compressionLevel,
|
||||
size_t oldInputPos, size_t newInputPos,
|
||||
size_t toFlushNow)
|
||||
{
|
||||
FIO_rust_zstd_projection_context_t* const context =
|
||||
(FIO_rust_zstd_projection_context_t*)opaque;
|
||||
FIO_prefs_t* const prefs = context->prefs;
|
||||
FIO_ctx_t* const fCtx = context->fCtx;
|
||||
|
||||
context->inputPresented++;
|
||||
if (oldInputPos == newInputPos) context->inputBlocked++;
|
||||
if (!toFlushNow) context->flushWaiting = 1;
|
||||
|
||||
/* Adaptive mode remains in C so its policy and private preference/context
|
||||
* access stay identical while Rust owns the surrounding stream loop. */
|
||||
if (prefs->adaptiveMode &&
|
||||
UTIL_clockSpanMicro(context->lastAdaptTime) > REFRESH_RATE) {
|
||||
ZSTD_frameProgression const zfp = ZSTD_getFrameProgression(context->cctx);
|
||||
|
||||
context->lastAdaptTime = UTIL_getTime();
|
||||
|
||||
/* check output speed */
|
||||
if (zfp.currentJobID > 1) { /* only possible if nbWorkers >= 1 */
|
||||
unsigned long long const newlyProduced =
|
||||
zfp.produced - context->previousZfpUpdate.produced;
|
||||
unsigned long long const newlyFlushed =
|
||||
zfp.flushed - context->previousZfpUpdate.flushed;
|
||||
assert(zfp.produced >= context->previousZfpUpdate.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 == context->previousZfpUpdate.consumed)
|
||||
&& (zfp.nbActiveWorkers == 0)) {
|
||||
DISPLAYLEVEL(6, "all buffers full : compression stopped => slow down \n")
|
||||
context->speedChange = FIO_rust_zstd_slower;
|
||||
}
|
||||
|
||||
context->previousZfpUpdate = zfp;
|
||||
|
||||
if ((newlyProduced > (newlyFlushed * 9 / 8))
|
||||
&& (context->flushWaiting == 0)) {
|
||||
DISPLAYLEVEL(6, "compression faster than flush (%llu > %llu), and flushed was never slowed down by lack of production => slow down \n",
|
||||
newlyProduced, newlyFlushed);
|
||||
context->speedChange = FIO_rust_zstd_slower;
|
||||
}
|
||||
context->flushWaiting = 0;
|
||||
}
|
||||
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;
|
||||
/* course correct only if there is at least one new job completed */
|
||||
if (zfp.currentJobID > context->lastJobID) {
|
||||
DISPLAYLEVEL(6, "compression level adaptation check \n")
|
||||
|
||||
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;
|
||||
/* check input speed */
|
||||
if (zfp.currentJobID > (unsigned)(prefs->nbWorkers+1)) {
|
||||
if (context->inputBlocked <= 0) {
|
||||
DISPLAYLEVEL(6, "input is never blocked => input is slower than ingestion \n");
|
||||
context->speedChange = FIO_rust_zstd_slower;
|
||||
} else if (context->speedChange == FIO_rust_zstd_noChange) {
|
||||
unsigned long long const newlyIngested =
|
||||
zfp.ingested - context->previousZfpCorrection.ingested;
|
||||
unsigned long long const newlyConsumed =
|
||||
zfp.consumed - context->previousZfpCorrection.consumed;
|
||||
unsigned long long const newlyProduced =
|
||||
zfp.produced - context->previousZfpCorrection.produced;
|
||||
unsigned long long const newlyFlushed =
|
||||
zfp.flushed - context->previousZfpCorrection.flushed;
|
||||
context->previousZfpCorrection = zfp;
|
||||
assert(context->inputPresented > 0);
|
||||
DISPLAYLEVEL(6, "input blocked %u/%u(%.2f) - ingested:%u vs %u:consumed - flushed:%u vs %u:produced \n",
|
||||
context->inputBlocked, context->inputPresented,
|
||||
(double)context->inputBlocked/context->inputPresented*100,
|
||||
(unsigned)newlyIngested, (unsigned)newlyConsumed,
|
||||
(unsigned)newlyFlushed, (unsigned)newlyProduced);
|
||||
if ((context->inputBlocked > context->inputPresented / 8)
|
||||
&& (newlyFlushed * 33 / 32 > newlyProduced)
|
||||
&& (newlyIngested * 33 / 32 > newlyConsumed)) {
|
||||
DISPLAYLEVEL(6, "recommend faster as in(%llu) >= (%llu)comp(%llu) <= out(%llu) \n",
|
||||
newlyIngested, newlyConsumed,
|
||||
newlyProduced, newlyFlushed);
|
||||
context->speedChange = FIO_rust_zstd_faster;
|
||||
}
|
||||
}
|
||||
context->inputBlocked = 0;
|
||||
context->inputPresented = 0;
|
||||
}
|
||||
|
||||
/* 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);
|
||||
if (context->speedChange == FIO_rust_zstd_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);
|
||||
ZSTD_CCtx_setParameter(context->cctx,
|
||||
ZSTD_c_compressionLevel,
|
||||
*compressionLevel);
|
||||
}
|
||||
if (context->speedChange == FIO_rust_zstd_faster) {
|
||||
DISPLAYLEVEL(6, "faster speed , lighter compression \n")
|
||||
(*compressionLevel)--;
|
||||
if (*compressionLevel < prefs->minAdaptLevel)
|
||||
*compressionLevel = prefs->minAdaptLevel;
|
||||
*compressionLevel -= (*compressionLevel == 0);
|
||||
ZSTD_CCtx_setParameter(context->cctx,
|
||||
ZSTD_c_compressionLevel,
|
||||
*compressionLevel);
|
||||
}
|
||||
context->speedChange = FIO_rust_zstd_noChange;
|
||||
context->lastJobID = zfp.currentJobID;
|
||||
}
|
||||
}
|
||||
|
||||
AIO_WritePool_releaseIoJob(writeJob);
|
||||
AIO_WritePool_sparseWriteEnd(ressPtr->writeCtx);
|
||||
/* Keep progress formatting and the frame-progression query in C. */
|
||||
if (SHOULD_DISPLAY_PROGRESS() && READY_FOR_UPDATE()) {
|
||||
ZSTD_frameProgression const zfp = ZSTD_getFrameProgression(context->cctx);
|
||||
double const cShare = (double)zfp.produced /
|
||||
(double)(zfp.consumed + !zfp.consumed) * 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);
|
||||
|
||||
return compressedfilesize;
|
||||
DELAY_NEXT_UPDATE();
|
||||
DISPLAY_PROGRESS("\r%79s\r", "");
|
||||
if (g_display_prefs.displayLevel >= 3) {
|
||||
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 {
|
||||
if (fCtx->nbFilesTotal > 1) {
|
||||
size_t const srcFileNameSize = strlen(srcFileName);
|
||||
if (srcFileNameSize > 18) {
|
||||
const char* const 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 (context->srcFileSize != UTIL_FILESIZE_UNKNOWN)
|
||||
DISPLAY_PROGRESS("/%6.*f%4s", context->fileHrs.precision,
|
||||
context->fileHrs.value, context->fileHrs.suffix);
|
||||
DISPLAY_PROGRESS(" ==> %2.f%%", cShare);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
static unsigned long long
|
||||
@@ -2155,9 +2229,89 @@ 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);
|
||||
FIO_ctx_t* const fCtxPtr = (FIO_ctx_t*)fCtx;
|
||||
FIO_prefs_t* const prefsPtr = (FIO_prefs_t*)prefs;
|
||||
cRess_t const* const ressPtr = (const cRess_t*)ress;
|
||||
FIO_rust_zstd_projection_context_t context;
|
||||
FIO_rust_zstd_compress_projection_t projection;
|
||||
U64 compressedSize = 0;
|
||||
U64 pledgedSrcSize = ZSTD_CONTENTSIZE_UNKNOWN;
|
||||
size_t zstdResult = 0;
|
||||
int status;
|
||||
|
||||
memset(&context, 0, sizeof(context));
|
||||
context.fCtx = fCtxPtr;
|
||||
context.prefs = prefsPtr;
|
||||
context.cctx = ressPtr->cctx;
|
||||
context.lastAdaptTime = UTIL_getTime();
|
||||
context.srcFileSize = srcFileSize;
|
||||
context.fileHrs = UTIL_makeHumanReadableSize(srcFileSize);
|
||||
|
||||
memset(&projection, 0, sizeof(projection));
|
||||
projection.readOpaque = (void*)ressPtr->readCtx;
|
||||
projection.writeOpaque = (void*)ressPtr->writeCtx;
|
||||
projection.codecOpaque = &context;
|
||||
projection.policyOpaque = &context;
|
||||
projection.readBufferSize = ZSTD_CStreamInSize();
|
||||
projection.readFill = FIO_rust_zstd_readFill;
|
||||
projection.readConsume = FIO_rust_zstd_readConsume;
|
||||
projection.writeAcquire = FIO_rust_zstd_writeAcquire;
|
||||
projection.writeEnqueue = FIO_rust_zstd_writeEnqueue;
|
||||
projection.writeRelease = FIO_rust_zstd_writeRelease;
|
||||
projection.sparseWriteEnd = FIO_rust_zstd_sparseWriteEnd;
|
||||
projection.compressStream = FIO_rust_zstd_compressStream;
|
||||
projection.iteration = FIO_rust_zstd_iteration;
|
||||
|
||||
DISPLAYLEVEL(6, "compression using zstd format \n");
|
||||
|
||||
/* Keep pledged-size and memory diagnostics in C while Rust owns the
|
||||
* surrounding asynchronous stream loop. */
|
||||
if (srcFileSize != UTIL_FILESIZE_UNKNOWN) {
|
||||
pledgedSrcSize = srcFileSize;
|
||||
CHECK(ZSTD_CCtx_setPledgedSrcSize(ressPtr->cctx, srcFileSize));
|
||||
} else if (prefsPtr->streamSrcSize > 0) {
|
||||
/* unknown source size; use the declared stream size */
|
||||
pledgedSrcSize = prefsPtr->streamSrcSize;
|
||||
CHECK(ZSTD_CCtx_setPledgedSrcSize(ressPtr->cctx, prefsPtr->streamSrcSize));
|
||||
}
|
||||
|
||||
{ int windowLog;
|
||||
UTIL_HumanReadableSize_t windowSize;
|
||||
CHECK(ZSTD_CCtx_getParameter(ressPtr->cctx, ZSTD_c_windowLog, &windowLog));
|
||||
if (windowLog == 0) {
|
||||
if (prefsPtr->ldmFlag) {
|
||||
/* If long mode is set without a window size libzstd will set this size internally */
|
||||
windowLog = ZSTD_WINDOWLOG_LIMIT_DEFAULT;
|
||||
} else {
|
||||
ZSTD_compressionParameters const cParams =
|
||||
ZSTD_getCParams(compressionLevel, srcFileSize, 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);
|
||||
}
|
||||
|
||||
status = FIO_rust_compressZstdFrame(
|
||||
&projection, srcFileName, srcFileSize, compressionLevel,
|
||||
readsize, &compressedSize, &zstdResult);
|
||||
switch (status) {
|
||||
case FIO_RUST_ZSTD_OK:
|
||||
return compressedSize;
|
||||
case FIO_RUST_ZSTD_COMPRESS_ERROR:
|
||||
DISPLAYLEVEL(5, "%s \n",
|
||||
"ZSTD_compressStream2(ress.cctx, &outBuff, &inBuff, directive)");
|
||||
EXM_THROW(11, "%s", ZSTD_getErrorName(zstdResult));
|
||||
case FIO_RUST_ZSTD_INCOMPLETE_INPUT:
|
||||
EXM_THROW(27, "Read error : Incomplete read : %llu / %llu B",
|
||||
(unsigned long long)*readsize,
|
||||
(unsigned long long)srcFileSize);
|
||||
default:
|
||||
assert(status == FIO_RUST_ZSTD_INVALID_PROJECTION);
|
||||
EXM_THROW(11, "zstd compression projection is invalid");
|
||||
}
|
||||
}
|
||||
|
||||
#ifdef ZSTD_GZCOMPRESS
|
||||
|
||||
Reference in New Issue
Block a user