feat(mt): move compression-job finish policy into Rust
The MT worker previously reported errors through one C callback and left a combined C finish callback responsible for serial completion, private resource release, output-size publication, consumed-size publication, and signaling. That kept the branch and lifecycle policy on the C side of the existing Rust stage scheduler. Add a narrow finish projection whose C callbacks expose only those private leaves and the source-size scalar. Rust now publishes an error before serial completion on every failed resource or codec stage, releases the sequence and CCtx resources in the original order, publishes the final block size only on success, then publishes the consumed size and signals the job condition. The failed path therefore normalizes any codec-reported last-block size to zero without changing the C-owned mutex, descriptor, pool, or context layouts. Focused tests cover successful final-block publication, resource failure, codec failure, cleanup ordering, consumed-size publication, and error-path normalization. Test Plan: - `rustfmt --check --edition 2021 rust/src/zstdmt_compress.rs` -- passed - `cc -fsyntax-only -Werror=incompatible-pointer-types -Ilib -Ilib/common -Ilib/compress -Ilib/decompress -Ilib/dict -Ilib/legacy lib/compress/zstdmt_compress.c` -- passed - `git diff --check -- lib/compress/zstdmt_compress.c rust/src/zstdmt_compress.rs` and `git diff --cached --check` -- passed - Rust unit tests were added but not run because the request prohibited Cargo and heavy commands
This commit is contained in:
+103
-21
@@ -234,8 +234,39 @@ 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);
|
||||
typedef void (*ZSTDMT_compressionJobSizeFn)(void* opaque, size_t size);
|
||||
|
||||
typedef struct {
|
||||
void* callbackContext;
|
||||
size_t srcSize;
|
||||
ZSTDMT_compressionJobVoidFn ensureFinished;
|
||||
ZSTDMT_compressionJobSizeFn publishError;
|
||||
ZSTDMT_compressionJobSizeFn publishLastBlockSize;
|
||||
ZSTDMT_compressionJobVoidFn releaseSeq;
|
||||
ZSTDMT_compressionJobVoidFn releaseCCtx;
|
||||
ZSTDMT_compressionJobSizeFn publishConsumed;
|
||||
ZSTDMT_compressionJobVoidFn signal;
|
||||
} ZSTDMT_RustCompressionJobFinishProjection;
|
||||
typedef char ZSTDMT_compression_job_finish_projection_layout[
|
||||
(offsetof(ZSTDMT_RustCompressionJobFinishProjection, callbackContext) == 0
|
||||
&& offsetof(ZSTDMT_RustCompressionJobFinishProjection, srcSize)
|
||||
== sizeof(void*)
|
||||
&& offsetof(ZSTDMT_RustCompressionJobFinishProjection, ensureFinished)
|
||||
== 2 * sizeof(void*)
|
||||
&& offsetof(ZSTDMT_RustCompressionJobFinishProjection, publishError)
|
||||
== 3 * sizeof(void*)
|
||||
&& offsetof(ZSTDMT_RustCompressionJobFinishProjection, publishLastBlockSize)
|
||||
== 4 * sizeof(void*)
|
||||
&& offsetof(ZSTDMT_RustCompressionJobFinishProjection, releaseSeq)
|
||||
== 5 * sizeof(void*)
|
||||
&& offsetof(ZSTDMT_RustCompressionJobFinishProjection, releaseCCtx)
|
||||
== 6 * sizeof(void*)
|
||||
&& offsetof(ZSTDMT_RustCompressionJobFinishProjection, publishConsumed)
|
||||
== 7 * sizeof(void*)
|
||||
&& offsetof(ZSTDMT_RustCompressionJobFinishProjection, signal)
|
||||
== 8 * sizeof(void*)
|
||||
&& sizeof(ZSTDMT_RustCompressionJobFinishProjection) == 9 * sizeof(void*))
|
||||
? 1 : -1];
|
||||
|
||||
typedef struct {
|
||||
int checksumFlag;
|
||||
@@ -255,6 +286,7 @@ ZSTDMT_RustCompressionJobParameters ZSTDMT_rust_prepareCompressionJobParameters(
|
||||
|
||||
void ZSTDMT_rust_compressionJob(
|
||||
const ZSTDMT_RustCompressionJobProjection* projection,
|
||||
const ZSTDMT_RustCompressionJobFinishProjection* finishProjection,
|
||||
void* opaque,
|
||||
ZSTDMT_compressionJobStepFn acquireResources,
|
||||
ZSTDMT_compressionJobVoidFn prepareParameters,
|
||||
@@ -262,9 +294,7 @@ void ZSTDMT_rust_compressionJob(
|
||||
ZSTDMT_compressionJobStepFn beginJob,
|
||||
ZSTDMT_compressionJobVoidFn applySequences,
|
||||
ZSTDMT_compressionJobCompressFn compressJob,
|
||||
ZSTDMT_compressionJobVoidFn traceJob,
|
||||
ZSTDMT_compressionJobErrorFn setError,
|
||||
ZSTDMT_compressionJobFinishFn finishJob);
|
||||
ZSTDMT_compressionJobVoidFn traceJob);
|
||||
|
||||
ZSTDMT_chunkProcessResult ZSTDMT_rust_compressJobChunks(
|
||||
ZSTD_CCtx* cctx, const void* src, size_t srcSize,
|
||||
@@ -1764,7 +1794,19 @@ static void ZSTDMT_compressionJobTrace(void* opaque)
|
||||
ZSTD_CCtx_trace(state->cctx, 0);
|
||||
}
|
||||
|
||||
static void ZSTDMT_compressionJobSetError(void* opaque, size_t error)
|
||||
static void ZSTDMT_compressionJobEnsureFinished(void* opaque)
|
||||
{
|
||||
ZSTDMT_compressionJobState* const state =
|
||||
(ZSTDMT_compressionJobState*)opaque;
|
||||
ZSTDMT_jobDescription* const job = state->job;
|
||||
|
||||
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);
|
||||
}
|
||||
|
||||
static void ZSTDMT_compressionJobPublishError(void* opaque, size_t error)
|
||||
{
|
||||
ZSTDMT_compressionJobState* const state =
|
||||
(ZSTDMT_compressionJobState*)opaque;
|
||||
@@ -1775,24 +1817,54 @@ static void ZSTDMT_compressionJobSetError(void* opaque, size_t error)
|
||||
ZSTD_pthread_mutex_unlock(&job->job_mutex);
|
||||
}
|
||||
|
||||
static void ZSTDMT_compressionJobFinish(void* opaque, size_t lastCBlockSize)
|
||||
static void ZSTDMT_compressionJobPublishLastBlockSize(void* opaque, size_t lastCBlockSize)
|
||||
{
|
||||
ZSTDMT_compressionJobState* const state =
|
||||
(ZSTDMT_compressionJobState*)opaque;
|
||||
ZSTDMT_jobDescription* const job = state->job;
|
||||
|
||||
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, state->rawSeqStore);
|
||||
ZSTDMT_releaseCCtx(job->cctxPool, state->cctx);
|
||||
/* report */
|
||||
ZSTD_PTHREAD_MUTEX_LOCK(&job->job_mutex);
|
||||
if (ZSTD_isError(job->cSize)) assert(lastCBlockSize == 0);
|
||||
job->cSize += lastCBlockSize;
|
||||
job->consumed = job->src.size; /* when job->consumed == job->src.size , compression job is presumed completed */
|
||||
ZSTD_pthread_mutex_unlock(&job->job_mutex);
|
||||
}
|
||||
|
||||
static void ZSTDMT_compressionJobReleaseSeq(void* opaque)
|
||||
{
|
||||
ZSTDMT_compressionJobState* const state =
|
||||
(ZSTDMT_compressionJobState*)opaque;
|
||||
ZSTDMT_jobDescription* const job = state->job;
|
||||
|
||||
ZSTDMT_releaseSeq(job->seqPool, state->rawSeqStore);
|
||||
}
|
||||
|
||||
static void ZSTDMT_compressionJobReleaseCCtx(void* opaque)
|
||||
{
|
||||
ZSTDMT_compressionJobState* const state =
|
||||
(ZSTDMT_compressionJobState*)opaque;
|
||||
ZSTDMT_jobDescription* const job = state->job;
|
||||
|
||||
ZSTDMT_releaseCCtx(job->cctxPool, state->cctx);
|
||||
}
|
||||
|
||||
static void ZSTDMT_compressionJobPublishConsumed(void* opaque, size_t srcSize)
|
||||
{
|
||||
ZSTDMT_compressionJobState* const state =
|
||||
(ZSTDMT_compressionJobState*)opaque;
|
||||
ZSTDMT_jobDescription* const job = state->job;
|
||||
|
||||
assert(srcSize == job->src.size);
|
||||
ZSTD_PTHREAD_MUTEX_LOCK(&job->job_mutex);
|
||||
job->consumed = srcSize; /* when job->consumed == job->src.size , compression job is presumed completed */
|
||||
ZSTD_pthread_mutex_unlock(&job->job_mutex);
|
||||
}
|
||||
|
||||
static void ZSTDMT_compressionJobSignal(void* opaque)
|
||||
{
|
||||
ZSTDMT_compressionJobState* const state =
|
||||
(ZSTDMT_compressionJobState*)opaque;
|
||||
ZSTDMT_jobDescription* const job = state->job;
|
||||
|
||||
ZSTD_PTHREAD_MUTEX_LOCK(&job->job_mutex);
|
||||
ZSTD_pthread_cond_signal(&job->job_cond);
|
||||
ZSTD_pthread_mutex_unlock(&job->job_mutex);
|
||||
}
|
||||
@@ -1803,6 +1875,7 @@ static void ZSTDMT_compressionJob(void* jobDescription)
|
||||
{
|
||||
ZSTDMT_jobDescription* const job = (ZSTDMT_jobDescription*)jobDescription;
|
||||
ZSTDMT_compressionJobState state;
|
||||
ZSTDMT_RustCompressionJobFinishProjection finishProjection;
|
||||
ZSTDMT_RustCompressionJobProjection const projection = {
|
||||
job->firstJob,
|
||||
job->lastJob,
|
||||
@@ -1813,18 +1886,27 @@ static void ZSTDMT_compressionJob(void* jobDescription)
|
||||
state.job = job;
|
||||
state.jobParams = job->params; /* do not modify job->params ! copy it, modify the copy */
|
||||
state.dstBuff = job->dstBuff;
|
||||
finishProjection = (ZSTDMT_RustCompressionJobFinishProjection){
|
||||
&state,
|
||||
job->src.size,
|
||||
ZSTDMT_compressionJobEnsureFinished,
|
||||
ZSTDMT_compressionJobPublishError,
|
||||
ZSTDMT_compressionJobPublishLastBlockSize,
|
||||
ZSTDMT_compressionJobReleaseSeq,
|
||||
ZSTDMT_compressionJobReleaseCCtx,
|
||||
ZSTDMT_compressionJobPublishConsumed,
|
||||
ZSTDMT_compressionJobSignal
|
||||
};
|
||||
|
||||
ZSTDMT_rust_compressionJob(
|
||||
&projection, &state,
|
||||
&projection, &finishProjection, &state,
|
||||
ZSTDMT_compressionJobAcquireResources,
|
||||
ZSTDMT_compressionJobPrepareParameters,
|
||||
ZSTDMT_compressionJobGenerateSequences,
|
||||
ZSTDMT_compressionJobBegin,
|
||||
ZSTDMT_compressionJobApplySequences,
|
||||
ZSTDMT_compressionJobCompress,
|
||||
ZSTDMT_compressionJobTrace,
|
||||
ZSTDMT_compressionJobSetError,
|
||||
ZSTDMT_compressionJobFinish);
|
||||
ZSTDMT_compressionJobTrace);
|
||||
}
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user