feat(compress): move MT job orchestration into Rust

Move the high-level ZSTDMT compression-job stage sequence into Rust: resource
acquisition, per-job parameter preparation, serial sequence handling, context
initialization, external-sequence application, non-first frame-header repair,
chunk compression/error routing, tracing, and common finalization. Keep C-owned
job descriptors, pools, synchronization, codec contexts, serial state, and
cleanup behind callbacks so private worker state does not cross the boundary.

Test Plan:
- Rust library all-target tests: 517 passed, including MT job-order tests.
- Rust legacy feature matrix: 572 passed.
- Rust and CLI clippy, nightly fmt, native CLI tests (41), and library smoke.
- Native test-zstd, bounded fuzzer (319), zstream (152 + 297), and decode
  corpus (1,647) all passed, including multi-GiB and MT round trips.
- All heavy checks ran serially with CARGO_BUILD_JOBS=1 or make -j1 and
  ulimit -v 41943040 (40 GiB virtual memory).
- Commit is intentionally unsigned because configured GPG pinentry was
  unavailable and hung during the signing attempt.
This commit is contained in:
2026-07-19 09:14:50 +02:00
parent 47b0da7ad7
commit 1655302473
2 changed files with 493 additions and 80 deletions
+198 -80
View File
@@ -121,6 +121,11 @@ typedef struct {
size_t lastBlockSize;
} ZSTDMT_chunkProcessResult;
typedef struct {
unsigned firstJob;
unsigned lastJob;
} ZSTDMT_RustCompressionJobProjection;
typedef struct {
int status;
size_t toFlush;
@@ -129,6 +134,26 @@ typedef struct {
} ZSTDMT_flushPublicationResult;
typedef void (*ZSTDMT_chunkProgressFn)(void* opaque, size_t cSize, size_t consumed);
typedef size_t (*ZSTDMT_compressionJobStepFn)(void* opaque);
typedef void (*ZSTDMT_compressionJobVoidFn)(void* opaque);
typedef ZSTDMT_chunkProcessResult (*ZSTDMT_compressionJobCompressFn)(
void* opaque, unsigned lastJob);
typedef void (*ZSTDMT_compressionJobErrorFn)(void* opaque, size_t error);
typedef void (*ZSTDMT_compressionJobFinishFn)(void* opaque, size_t lastBlockSize);
void ZSTDMT_rust_compressionJob(
const ZSTDMT_RustCompressionJobProjection* projection,
void* opaque,
ZSTDMT_compressionJobStepFn acquireResources,
ZSTDMT_compressionJobVoidFn prepareParameters,
ZSTDMT_compressionJobVoidFn generateSequences,
ZSTDMT_compressionJobStepFn beginJob,
ZSTDMT_compressionJobVoidFn applySequences,
ZSTDMT_compressionJobStepFn writeFrameHeader,
ZSTDMT_compressionJobCompressFn compressJob,
ZSTDMT_compressionJobVoidFn traceJob,
ZSTDMT_compressionJobErrorFn setError,
ZSTDMT_compressionJobFinishFn finishJob);
ZSTDMT_chunkProcessResult ZSTDMT_rust_compressJobChunks(
ZSTD_CCtx* cctx, const void* src, size_t srcSize,
@@ -913,114 +938,177 @@ static void ZSTDMT_compressionJobProgress(void* opaque, size_t cSize, size_t con
ZSTD_pthread_mutex_unlock(&job->job_mutex);
}
#define JOB_ERROR(e) \
do { \
ZSTD_PTHREAD_MUTEX_LOCK(&job->job_mutex); \
job->cSize = e; \
ZSTD_pthread_mutex_unlock(&job->job_mutex); \
goto _endJob; \
} while (0)
typedef struct {
ZSTDMT_jobDescription* job;
ZSTD_CCtx_params jobParams;
ZSTD_CCtx* cctx;
RawSeqStore_t rawSeqStore;
Buffer dstBuff;
} ZSTDMT_compressionJobState;
/* ZSTDMT_compressionJob() is a POOL_function type */
static void ZSTDMT_compressionJob(void* jobDescription)
static size_t ZSTDMT_compressionJobAcquireResources(void* opaque)
{
ZSTDMT_jobDescription* const job = (ZSTDMT_jobDescription*)jobDescription;
ZSTD_CCtx_params jobParams = job->params; /* do not modify job->params ! copy it, modify the copy */
ZSTD_CCtx* const cctx = ZSTDMT_getCCtx(job->cctxPool);
RawSeqStore_t rawSeqStore = ZSTDMT_getSeq(job->seqPool);
Buffer dstBuff = job->dstBuff;
size_t lastCBlockSize = 0;
ZSTDMT_compressionJobState* const state =
(ZSTDMT_compressionJobState*)opaque;
ZSTDMT_jobDescription* const job = state->job;
state->cctx = ZSTDMT_getCCtx(job->cctxPool);
state->rawSeqStore = ZSTDMT_getSeq(job->seqPool);
state->dstBuff = job->dstBuff;
DEBUGLOG(5, "ZSTDMT_compressionJob: job %u", job->jobID);
/* resources */
if (cctx==NULL) JOB_ERROR(ERROR(memory_allocation));
if (dstBuff.start == NULL) { /* streaming job : doesn't provide a dstBuffer */
dstBuff = ZSTDMT_getBuffer(job->bufPool);
if (dstBuff.start==NULL) JOB_ERROR(ERROR(memory_allocation));
job->dstBuff = dstBuff; /* this value can be read in ZSTDMT_flush, when it copies the whole job */
if (state->cctx == NULL) return ERROR(memory_allocation);
if (state->dstBuff.start == NULL) {
state->dstBuff = ZSTDMT_getBuffer(job->bufPool);
if (state->dstBuff.start == NULL) return ERROR(memory_allocation);
job->dstBuff = state->dstBuff;
}
if (jobParams.ldmParams.enableLdm == ZSTD_ps_enable && rawSeqStore.seq == NULL)
JOB_ERROR(ERROR(memory_allocation));
if (state->jobParams.ldmParams.enableLdm == ZSTD_ps_enable &&
state->rawSeqStore.seq == NULL)
return ERROR(memory_allocation);
return 0;
}
static void ZSTDMT_compressionJobPrepareParameters(void* opaque)
{
ZSTDMT_compressionJobState* const state =
(ZSTDMT_compressionJobState*)opaque;
ZSTDMT_jobDescription* const job = state->job;
/* Don't compute the checksum for chunks, since we compute it externally,
* but write it in the header.
*/
if (job->jobID != 0) jobParams.fParams.checksumFlag = 0;
/* Don't run LDM for the chunks, since we handle it externally */
jobParams.ldmParams.enableLdm = ZSTD_ps_disable;
* but write it in the header. */
if (job->jobID != 0) state->jobParams.fParams.checksumFlag = 0;
/* Don't run LDM for the chunks, since we handle it externally. */
state->jobParams.ldmParams.enableLdm = ZSTD_ps_disable;
/* Correct nbWorkers to 0. */
jobParams.nbWorkers = 0;
state->jobParams.nbWorkers = 0;
}
static void ZSTDMT_compressionJobGenerateSequences(void* opaque)
{
ZSTDMT_compressionJobState* const state =
(ZSTDMT_compressionJobState*)opaque;
ZSTDMT_jobDescription* const job = state->job;
/* init */
/* Perform serial step as early as possible. */
ZSTDMT_serialState_genSequences(job->serial, &state->rawSeqStore,
job->src, job->jobID);
}
/* Perform serial step as early as possible */
ZSTDMT_serialState_genSequences(job->serial, &rawSeqStore, job->src, job->jobID);
static size_t ZSTDMT_compressionJobBegin(void* opaque)
{
ZSTDMT_compressionJobState* const state =
(ZSTDMT_compressionJobState*)opaque;
ZSTDMT_jobDescription* const job = state->job;
if (job->cdict) {
size_t const initError = ZSTD_compressBegin_advanced_internal(cctx, NULL, 0, ZSTD_dct_auto, ZSTD_dtlm_fast, job->cdict, &jobParams, job->fullFrameSize);
size_t const initError = ZSTD_compressBegin_advanced_internal(
state->cctx, NULL, 0, ZSTD_dct_auto, ZSTD_dtlm_fast, job->cdict,
&state->jobParams, job->fullFrameSize);
assert(job->firstJob); /* only allowed for first job */
if (ZSTD_isError(initError)) JOB_ERROR(initError);
} else {
U64 const pledgedSrcSize = job->firstJob ? job->fullFrameSize : job->src.size;
{ size_t const forceWindowError = ZSTD_CCtxParams_setParameter(&jobParams, ZSTD_c_forceMaxWindow, !job->firstJob);
if (ZSTD_isError(forceWindowError)) JOB_ERROR(forceWindowError);
}
return initError;
}
{ U64 const pledgedSrcSize = job->firstJob ? job->fullFrameSize : job->src.size;
size_t const forceWindowError = ZSTD_CCtxParams_setParameter(
&state->jobParams, ZSTD_c_forceMaxWindow, !job->firstJob);
if (ZSTD_isError(forceWindowError)) return forceWindowError;
if (!job->firstJob) {
size_t const err = ZSTD_CCtxParams_setParameter(&jobParams, ZSTD_c_deterministicRefPrefix, 0);
if (ZSTD_isError(err)) JOB_ERROR(err);
size_t const err = ZSTD_CCtxParams_setParameter(
&state->jobParams, ZSTD_c_deterministicRefPrefix, 0);
if (ZSTD_isError(err)) return err;
}
DEBUGLOG(6, "ZSTDMT_compressionJob: job %u: loading prefix of size %zu", job->jobID, job->prefix.size);
{ size_t const initError = ZSTD_compressBegin_advanced_internal(cctx,
job->prefix.start, job->prefix.size, ZSTD_dct_rawContent,
ZSTD_dtlm_fast,
NULL, /*cdict*/
&jobParams, pledgedSrcSize);
if (ZSTD_isError(initError)) JOB_ERROR(initError);
} }
/* External Sequences can only be applied after CCtx initialization */
ZSTDMT_serialState_applySequences(job->serial, cctx, &rawSeqStore);
if (!job->firstJob) { /* flush and overwrite frame header when it's not first job */
size_t const hSize = ZSTD_compressContinue_public(cctx, dstBuff.start, dstBuff.capacity, job->src.start, 0);
if (ZSTD_isError(hSize)) JOB_ERROR(hSize);
DEBUGLOG(5, "ZSTDMT_compressionJob: flush and overwrite %u bytes of frame header (not first job)", (U32)hSize);
ZSTD_invalidateRepCodes(cctx);
return ZSTD_compressBegin_advanced_internal(
state->cctx, job->prefix.start, job->prefix.size,
ZSTD_dct_rawContent, ZSTD_dtlm_fast, NULL, /*cdict*/
&state->jobParams, pledgedSrcSize);
}
}
static void ZSTDMT_compressionJobApplySequences(void* opaque)
{
ZSTDMT_compressionJobState* const state =
(ZSTDMT_compressionJobState*)opaque;
ZSTDMT_jobDescription* const job = state->job;
/* External Sequences can only be applied after CCtx initialization. */
ZSTDMT_serialState_applySequences(job->serial, state->cctx,
&state->rawSeqStore);
}
static size_t ZSTDMT_compressionJobWriteFrameHeader(void* opaque)
{
ZSTDMT_compressionJobState* const state =
(ZSTDMT_compressionJobState*)opaque;
ZSTDMT_jobDescription* const job = state->job;
size_t const hSize = ZSTD_compressContinue_public(
state->cctx, state->dstBuff.start, state->dstBuff.capacity,
job->src.start, 0);
if (ZSTD_isError(hSize)) return hSize;
DEBUGLOG(5, "ZSTDMT_compressionJob: flush and overwrite %u bytes of frame header (not first job)", (U32)hSize);
ZSTD_invalidateRepCodes(state->cctx);
return hSize;
}
static ZSTDMT_chunkProcessResult ZSTDMT_compressionJobCompress(
void* opaque, unsigned lastJob)
{
ZSTDMT_compressionJobState* const state =
(ZSTDMT_compressionJobState*)opaque;
ZSTDMT_jobDescription* const job = state->job;
size_t const chunkSize = 4*ZSTD_BLOCKSIZE_MAX;
assert(lastJob == job->lastJob);
if (sizeof(size_t) > sizeof(int)) assert(job->src.size < ((size_t)INT_MAX) * chunkSize); /* check overflow */
DEBUGLOG(5, "ZSTDMT_compressionJob: compress %u bytes in %zu blocks",
(U32)job->src.size, (job->src.size + (chunkSize-1)) / chunkSize);
assert(job->cSize == 0);
assert(chunkSize > 0);
assert((chunkSize & (chunkSize - 1)) == 0); /* chunkSize must be power of 2 for mask==(chunkSize-1) to work */
return ZSTDMT_rust_compressJobChunks(
state->cctx, job->src.start, job->src.size,
state->dstBuff.start, state->dstBuff.capacity, chunkSize, lastJob,
job, ZSTDMT_compressionJobProgress);
}
static void ZSTDMT_compressionJobTrace(void* opaque)
{
ZSTDMT_compressionJobState* const state =
(ZSTDMT_compressionJobState*)opaque;
ZSTDMT_jobDescription* const job = state->job;
/* compress the entire job by smaller chunks, for better granularity */
{ size_t const chunkSize = 4*ZSTD_BLOCKSIZE_MAX;
if (sizeof(size_t) > sizeof(int)) assert(job->src.size < ((size_t)INT_MAX) * chunkSize); /* check overflow */
DEBUGLOG(5, "ZSTDMT_compressionJob: compress %u bytes in %zu blocks",
(U32)job->src.size, (job->src.size + (chunkSize-1)) / chunkSize);
assert(job->cSize == 0);
assert(chunkSize > 0);
assert((chunkSize & (chunkSize - 1)) == 0); /* chunkSize must be power of 2 for mask==(chunkSize-1) to work */
{ ZSTDMT_chunkProcessResult const result = ZSTDMT_rust_compressJobChunks(
cctx, job->src.start, job->src.size,
dstBuff.start, dstBuff.capacity, chunkSize, job->lastJob,
job, ZSTDMT_compressionJobProgress);
if (ZSTD_isError(result.error)) JOB_ERROR(result.error);
lastCBlockSize = result.lastBlockSize;
}
}
if (!job->firstJob) {
/* Double check that we don't have an ext-dict, because then our
* repcode invalidation doesn't work.
*/
assert(!ZSTD_window_hasExtDict(cctx->blockState.matchState.window));
* repcode invalidation doesn't work. */
assert(!ZSTD_window_hasExtDict(state->cctx->blockState.matchState.window));
}
ZSTD_CCtx_trace(cctx, 0);
ZSTD_CCtx_trace(state->cctx, 0);
}
static void ZSTDMT_compressionJobSetError(void* opaque, size_t error)
{
ZSTDMT_compressionJobState* const state =
(ZSTDMT_compressionJobState*)opaque;
ZSTDMT_jobDescription* const job = state->job;
ZSTD_PTHREAD_MUTEX_LOCK(&job->job_mutex);
job->cSize = error;
ZSTD_pthread_mutex_unlock(&job->job_mutex);
}
static void ZSTDMT_compressionJobFinish(void* opaque, size_t lastCBlockSize)
{
ZSTDMT_compressionJobState* const state =
(ZSTDMT_compressionJobState*)opaque;
ZSTDMT_jobDescription* const job = state->job;
_endJob:
ZSTDMT_serialState_ensureFinished(job->serial, job->jobID, job->cSize);
if (job->prefix.size > 0)
DEBUGLOG(5, "Finished with prefix: %zx", (size_t)job->prefix.start);
DEBUGLOG(5, "Finished with source: %zx", (size_t)job->src.start);
/* release resources */
ZSTDMT_releaseSeq(job->seqPool, rawSeqStore);
ZSTDMT_releaseCCtx(job->cctxPool, cctx);
ZSTDMT_releaseSeq(job->seqPool, state->rawSeqStore);
ZSTDMT_releaseCCtx(job->cctxPool, state->cctx);
/* report */
ZSTD_PTHREAD_MUTEX_LOCK(&job->job_mutex);
if (ZSTD_isError(job->cSize)) assert(lastCBlockSize == 0);
@@ -1030,6 +1118,36 @@ _endJob:
ZSTD_pthread_mutex_unlock(&job->job_mutex);
}
/* ZSTDMT_compressionJob() is a POOL_function type. Rust owns the stage
* ordering; these callbacks retain the private job and codec operations. */
static void ZSTDMT_compressionJob(void* jobDescription)
{
ZSTDMT_jobDescription* const job = (ZSTDMT_jobDescription*)jobDescription;
ZSTDMT_compressionJobState state;
ZSTDMT_RustCompressionJobProjection const projection = {
job->firstJob,
job->lastJob
};
ZSTD_memset(&state, 0, sizeof(state));
state.job = job;
state.jobParams = job->params; /* do not modify job->params ! copy it, modify the copy */
state.dstBuff = job->dstBuff;
ZSTDMT_rust_compressionJob(
&projection, &state,
ZSTDMT_compressionJobAcquireResources,
ZSTDMT_compressionJobPrepareParameters,
ZSTDMT_compressionJobGenerateSequences,
ZSTDMT_compressionJobBegin,
ZSTDMT_compressionJobApplySequences,
ZSTDMT_compressionJobWriteFrameHeader,
ZSTDMT_compressionJobCompress,
ZSTDMT_compressionJobTrace,
ZSTDMT_compressionJobSetError,
ZSTDMT_compressionJobFinish);
}
/* ------------------------------------------ */
/* ===== Multi-threaded compression ===== */