refactor(cli): move adaptive feedback policy to Rust

The zstd file-I/O loop still kept adaptive compression's state machine in C: progression deltas, refresh gating, job-completion checks, input counters, speed decisions, and level updates were interleaved with private FIO and ZSTD state. That left orchestration policy behind the existing scalar Rust predicates.

Move the adaptive state and callback order into Rust. C now supplies scalar progression snapshots, the clock gate, exact diagnostics, and the ignored CCtx parameter setter through callbacks; private FIO_prefs_t, ZSTD_CCtx, ZSTD_frameProgression, clocks, and progress formatting remain C-owned. Preserve the previous-progression publication before backlog evaluation, input counter reset points, and serial/MT level-clamp behavior.

Test Plan:

- cargo +nightly fmt --manifest-path rust/Cargo.toml --all -- --check

- git diff --check

- ulimit -v 41943040; CARGO_BUILD_JOBS=1 make -j1
This commit is contained in:
2026-07-21 11:32:06 +02:00
parent 1a3971ff10
commit 1554c5aacb
2 changed files with 594 additions and 159 deletions
+166 -156
View File
@@ -1552,6 +1552,34 @@ void FIO_rust_zstd_compressStreamDisplay(
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 FIO_rust_zstd_adapt_projection_s
FIO_rust_zstd_adapt_projection_t;
typedef struct {
U64 ingested;
U64 consumed;
U64 produced;
U64 flushed;
unsigned currentJobID;
unsigned nbActiveWorkers;
} FIO_rust_zstd_progression_t;
typedef char FIO_rust_zstd_progression_layout[
(offsetof(FIO_rust_zstd_progression_t, ingested) == 0
&& offsetof(FIO_rust_zstd_progression_t, consumed) == sizeof(U64)
&& offsetof(FIO_rust_zstd_progression_t, produced) == 2 * sizeof(U64)
&& offsetof(FIO_rust_zstd_progression_t, flushed) == 3 * sizeof(U64)
&& offsetof(FIO_rust_zstd_progression_t, currentJobID) == 4 * sizeof(U64)
&& offsetof(FIO_rust_zstd_progression_t, nbActiveWorkers)
== 4 * sizeof(U64) + sizeof(unsigned)
&& sizeof(FIO_rust_zstd_progression_t)
== 4 * sizeof(U64) + 2 * sizeof(unsigned)) ? 1 : -1];
typedef int (*FIO_rust_zstd_adaptive_refresh_fn)(void* opaque);
typedef void (*FIO_rust_zstd_adaptive_progression_fn)(
void* opaque, FIO_rust_zstd_progression_t* progression);
typedef void (*FIO_rust_zstd_adaptive_set_parameter_fn)(
void* opaque, int compressionLevel);
typedef void (*FIO_rust_zstd_adaptive_diagnostic_fn)(
void* opaque, int diagnostic,
const FIO_rust_zstd_adapt_projection_t* projection);
typedef struct {
void* readOpaque;
@@ -1568,7 +1596,33 @@ typedef struct {
FIO_rust_zstd_compress_stream_fn compressStream;
FIO_rust_zstd_iteration_fn iteration;
FIO_rust_zstd_compress_display_fn compressStreamDisplay;
int adaptiveMode;
int nbWorkers;
int minAdaptLevel;
int maxAdaptLevel;
int maxCLevel;
FIO_rust_zstd_adaptive_refresh_fn adaptiveRefresh;
FIO_rust_zstd_adaptive_progression_fn adaptiveProgression;
FIO_rust_zstd_adaptive_set_parameter_fn adaptiveSetParameter;
FIO_rust_zstd_adaptive_diagnostic_fn adaptiveDiagnostic;
} FIO_rust_zstd_compress_projection_t;
typedef char FIO_rust_zstd_compress_projection_layout[
(offsetof(FIO_rust_zstd_compress_projection_t, adaptiveMode)
== 14 * sizeof(void*)
&& offsetof(FIO_rust_zstd_compress_projection_t, nbWorkers)
== 14 * sizeof(void*) + sizeof(int)
&& offsetof(FIO_rust_zstd_compress_projection_t, maxCLevel)
== 14 * sizeof(void*) + 4 * sizeof(int)
&& offsetof(FIO_rust_zstd_compress_projection_t, adaptiveRefresh)
== ((14 * sizeof(void*) + 5 * sizeof(int) + sizeof(void*) - 1)
/ sizeof(void*)) * sizeof(void*)
&& offsetof(FIO_rust_zstd_compress_projection_t, adaptiveDiagnostic)
== ((14 * sizeof(void*) + 5 * sizeof(int) + sizeof(void*) - 1)
/ sizeof(void*)) * sizeof(void*) + 3 * sizeof(void*)
&& sizeof(FIO_rust_zstd_compress_projection_t)
== ((14 * sizeof(void*) + 5 * sizeof(int) + sizeof(void*) - 1)
/ sizeof(void*)) * sizeof(void*) + 4 * sizeof(void*))
? 1 : -1];
int FIO_rust_compressZstdFrame(
const FIO_rust_zstd_compress_projection_t* projection,
@@ -2621,15 +2675,9 @@ FIO_compressLz4Frame(cRess_t* ress,
}
#endif
typedef enum {
FIO_rust_zstd_noChange,
FIO_rust_zstd_slower,
FIO_rust_zstd_faster
} FIO_rust_zstd_speed_change_e;
/* Rust owns only the scalar adaptive predicates. Keep the FIO context,
/* Rust owns scalar adaptive state and policy. Keep the FIO context,
* preference, and ZSTD progression layouts private to this translation unit. */
typedef struct {
struct FIO_rust_zstd_adapt_projection_s {
U64 consumed;
U64 previousConsumed;
unsigned nbActiveWorkers;
@@ -2644,7 +2692,7 @@ typedef struct {
int minAdaptLevel;
int maxAdaptLevel;
int maxCLevel;
} FIO_rust_zstd_adapt_projection_t;
};
typedef char FIO_rust_zstd_adapt_projection_layout[
(offsetof(FIO_rust_zstd_adapt_projection_t, consumed) == 0
&& offsetof(FIO_rust_zstd_adapt_projection_t, previousConsumed)
@@ -2678,28 +2726,19 @@ typedef char FIO_rust_zstd_adapt_projection_layout[
? 1 : -1];
enum {
FIO_RUST_ZSTD_ADAPT_OUTPUT_BLOCKED,
FIO_RUST_ZSTD_ADAPT_OUTPUT_BACKLOG,
FIO_RUST_ZSTD_ADAPT_INPUT_STARVATION,
FIO_RUST_ZSTD_ADAPT_BLOCKED_INPUT,
FIO_RUST_ZSTD_ADAPT_LEVEL_SLOWER,
FIO_RUST_ZSTD_ADAPT_LEVEL_FASTER
FIO_RUST_ZSTD_ADAPT_DIAG_OUTPUT_BLOCKED,
FIO_RUST_ZSTD_ADAPT_DIAG_OUTPUT_BACKLOG,
FIO_RUST_ZSTD_ADAPT_DIAG_CHECK,
FIO_RUST_ZSTD_ADAPT_DIAG_INPUT_STARVATION,
FIO_RUST_ZSTD_ADAPT_DIAG_INPUT_STATS,
FIO_RUST_ZSTD_ADAPT_DIAG_RECOMMEND_FASTER,
FIO_RUST_ZSTD_ADAPT_DIAG_SLOWER_LEVEL,
FIO_RUST_ZSTD_ADAPT_DIAG_FASTER_LEVEL
};
int FIO_rust_zstd_adapt(
int policy, const FIO_rust_zstd_adapt_projection_t* projection);
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;
@@ -2754,6 +2793,95 @@ static void FIO_rust_zstd_sparseWriteEnd(void* opaque)
AIO_WritePool_sparseWriteEnd((WritePoolCtx_t*)opaque);
}
static int FIO_rust_zstd_adaptiveRefresh(void* opaque)
{
FIO_rust_zstd_projection_context_t* const context =
(FIO_rust_zstd_projection_context_t*)opaque;
if (UTIL_clockSpanMicro(context->lastAdaptTime) > REFRESH_RATE) {
context->lastAdaptTime = UTIL_getTime();
return 1;
}
return 0;
}
static void FIO_rust_zstd_adaptiveProgression(
void* opaque, FIO_rust_zstd_progression_t* projection)
{
FIO_rust_zstd_projection_context_t* const context =
(FIO_rust_zstd_projection_context_t*)opaque;
ZSTD_frameProgression const zfp = ZSTD_getFrameProgression(context->cctx);
assert(projection != NULL);
projection->ingested = zfp.ingested;
projection->consumed = zfp.consumed;
projection->produced = zfp.produced;
projection->flushed = zfp.flushed;
projection->currentJobID = zfp.currentJobID;
projection->nbActiveWorkers = zfp.nbActiveWorkers;
}
static void FIO_rust_zstd_adaptiveSetParameter(void* opaque,
int compressionLevel)
{
FIO_rust_zstd_projection_context_t* const context =
(FIO_rust_zstd_projection_context_t*)opaque;
(void)ZSTD_CCtx_setParameter(context->cctx,
ZSTD_c_compressionLevel,
compressionLevel);
}
static void FIO_rust_zstd_adaptiveDiagnostic(
void* opaque, int diagnostic,
const FIO_rust_zstd_adapt_projection_t* projection)
{
(void)opaque;
switch (diagnostic) {
case FIO_RUST_ZSTD_ADAPT_DIAG_OUTPUT_BLOCKED:
DISPLAYLEVEL(6, "all buffers full : compression stopped => slow down \n")
break;
case FIO_RUST_ZSTD_ADAPT_DIAG_OUTPUT_BACKLOG:
assert(projection != NULL);
DISPLAYLEVEL(6, "compression faster than flush (%llu > %llu), and flushed was never slowed down by lack of production => slow down \n",
(unsigned long long)projection->newlyProduced,
(unsigned long long)projection->newlyFlushed);
break;
case FIO_RUST_ZSTD_ADAPT_DIAG_CHECK:
DISPLAYLEVEL(6, "compression level adaptation check \n")
break;
case FIO_RUST_ZSTD_ADAPT_DIAG_INPUT_STARVATION:
DISPLAYLEVEL(6, "input is never blocked => input is slower than ingestion \n");
break;
case FIO_RUST_ZSTD_ADAPT_DIAG_INPUT_STATS:
assert(projection != NULL);
assert(projection->inputPresented > 0);
DISPLAYLEVEL(6, "input blocked %u/%u(%.2f) - ingested:%u vs %u:consumed - flushed:%u vs %u:produced \n",
projection->inputBlocked, projection->inputPresented,
(double)projection->inputBlocked/
projection->inputPresented*100,
(unsigned)projection->newlyIngested,
(unsigned)projection->newlyConsumed,
(unsigned)projection->newlyFlushed,
(unsigned)projection->newlyProduced);
break;
case FIO_RUST_ZSTD_ADAPT_DIAG_RECOMMEND_FASTER:
assert(projection != NULL);
DISPLAYLEVEL(6, "recommend faster as in(%llu) >= (%llu)comp(%llu) <= out(%llu) \n",
(unsigned long long)projection->newlyIngested,
(unsigned long long)projection->newlyConsumed,
(unsigned long long)projection->newlyProduced,
(unsigned long long)projection->newlyFlushed);
break;
case FIO_RUST_ZSTD_ADAPT_DIAG_SLOWER_LEVEL:
DISPLAYLEVEL(6, "slower speed , higher compression \n")
break;
case FIO_RUST_ZSTD_ADAPT_DIAG_FASTER_LEVEL:
DISPLAYLEVEL(6, "faster speed , lighter compression \n")
break;
default:
assert(0);
break;
}
}
void FIO_rust_zstd_compressStreamDisplay(
int directive, size_t inputPos, size_t inputSize, size_t outputProduced)
{
@@ -2769,136 +2897,10 @@ static void FIO_rust_zstd_iteration(void* opaque, const char* srcFileName,
{
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;
FIO_rust_zstd_adapt_projection_t adaptProjection;
context->inputPresented++;
if (oldInputPos == newInputPos) context->inputBlocked++;
if (!toFlushNow) context->flushWaiting = 1;
/* C retains the clock gate, progression snapshots, diagnostics, and
* private context mutations; Rust supplies only scalar policy decisions. */
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 */
memset(&adaptProjection, 0, sizeof(adaptProjection));
adaptProjection.consumed = zfp.consumed;
adaptProjection.previousConsumed = context->previousZfpUpdate.consumed;
adaptProjection.nbActiveWorkers = zfp.nbActiveWorkers;
if (FIO_rust_zstd_adapt(
FIO_RUST_ZSTD_ADAPT_OUTPUT_BLOCKED, &adaptProjection)
== FIO_rust_zstd_slower) {
DISPLAYLEVEL(6, "all buffers full : compression stopped => slow down \n")
context->speedChange = FIO_rust_zstd_slower;
}
context->previousZfpUpdate = zfp;
adaptProjection.newlyProduced = newlyProduced;
adaptProjection.newlyFlushed = newlyFlushed;
adaptProjection.flushWaiting = context->flushWaiting;
if (FIO_rust_zstd_adapt(
FIO_RUST_ZSTD_ADAPT_OUTPUT_BACKLOG, &adaptProjection)
== FIO_rust_zstd_slower) {
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;
}
/* course correct only if there is at least one new job completed */
if (zfp.currentJobID > context->lastJobID) {
DISPLAYLEVEL(6, "compression level adaptation check \n")
/* check input speed */
if (zfp.currentJobID > (unsigned)(prefs->nbWorkers+1)) {
memset(&adaptProjection, 0, sizeof(adaptProjection));
adaptProjection.inputBlocked = context->inputBlocked;
if (FIO_rust_zstd_adapt(
FIO_RUST_ZSTD_ADAPT_INPUT_STARVATION, &adaptProjection)
== FIO_rust_zstd_slower) {
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);
adaptProjection.inputBlocked = context->inputBlocked;
adaptProjection.inputPresented = context->inputPresented;
adaptProjection.newlyIngested = newlyIngested;
adaptProjection.newlyConsumed = newlyConsumed;
adaptProjection.newlyProduced = newlyProduced;
adaptProjection.newlyFlushed = newlyFlushed;
if (FIO_rust_zstd_adapt(
FIO_RUST_ZSTD_ADAPT_BLOCKED_INPUT, &adaptProjection)
== FIO_rust_zstd_faster) {
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;
}
if (context->speedChange == FIO_rust_zstd_slower) {
DISPLAYLEVEL(6, "slower speed , higher compression \n")
memset(&adaptProjection, 0, sizeof(adaptProjection));
adaptProjection.compressionLevel = *compressionLevel;
adaptProjection.maxAdaptLevel = prefs->maxAdaptLevel;
adaptProjection.maxCLevel = ZSTD_maxCLevel();
*compressionLevel = FIO_rust_zstd_adapt(
FIO_RUST_ZSTD_ADAPT_LEVEL_SLOWER, &adaptProjection);
ZSTD_CCtx_setParameter(context->cctx,
ZSTD_c_compressionLevel,
*compressionLevel);
}
if (context->speedChange == FIO_rust_zstd_faster) {
DISPLAYLEVEL(6, "faster speed , lighter compression \n")
memset(&adaptProjection, 0, sizeof(adaptProjection));
adaptProjection.compressionLevel = *compressionLevel;
adaptProjection.minAdaptLevel = prefs->minAdaptLevel;
*compressionLevel = FIO_rust_zstd_adapt(
FIO_RUST_ZSTD_ADAPT_LEVEL_FASTER, &adaptProjection);
ZSTD_CCtx_setParameter(context->cctx,
ZSTD_c_compressionLevel,
*compressionLevel);
}
context->speedChange = FIO_rust_zstd_noChange;
context->lastJobID = zfp.currentJobID;
}
}
(void)oldInputPos;
(void)newInputPos;
(void)toFlushNow;
/* Keep progress formatting and the frame-progression query in C. */
if (SHOULD_DISPLAY_PROGRESS() && READY_FOR_UPDATE()) {
@@ -2964,7 +2966,6 @@ FIO_rust_compressZstdCallback(void* fCtx, void* prefs, void* ress,
memset(&context, 0, sizeof(context));
context.fCtx = fCtxPtr;
context.prefs = prefsPtr;
context.cctx = ressPtr->cctx;
context.lastAdaptTime = UTIL_getTime();
context.srcFileSize = srcFileSize;
@@ -2985,6 +2986,15 @@ FIO_rust_compressZstdCallback(void* fCtx, void* prefs, void* ress,
projection.compressStream = FIO_rust_zstd_compressStream;
projection.iteration = FIO_rust_zstd_iteration;
projection.compressStreamDisplay = FIO_rust_zstd_compressStreamDisplay;
projection.adaptiveMode = prefsPtr->adaptiveMode;
projection.nbWorkers = prefsPtr->nbWorkers;
projection.minAdaptLevel = prefsPtr->minAdaptLevel;
projection.maxAdaptLevel = prefsPtr->maxAdaptLevel;
projection.maxCLevel = ZSTD_maxCLevel();
projection.adaptiveRefresh = FIO_rust_zstd_adaptiveRefresh;
projection.adaptiveProgression = FIO_rust_zstd_adaptiveProgression;
projection.adaptiveSetParameter = FIO_rust_zstd_adaptiveSetParameter;
projection.adaptiveDiagnostic = FIO_rust_zstd_adaptiveDiagnostic;
DISPLAYLEVEL(6, "compression using zstd format \n");