From 1655302473a749db548798716bad3ca60fff9b38 Mon Sep 17 00:00:00 2001 From: ddidderr Date: Sun, 19 Jul 2026 09:14:50 +0200 Subject: [PATCH] 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. --- lib/compress/zstdmt_compress.c | 278 ++++++++++++++++++++++--------- rust/src/zstdmt_compress.rs | 295 +++++++++++++++++++++++++++++++++ 2 files changed, 493 insertions(+), 80 deletions(-) diff --git a/lib/compress/zstdmt_compress.c b/lib/compress/zstdmt_compress.c index 6529e2301..0875c7640 100644 --- a/lib/compress/zstdmt_compress.c +++ b/lib/compress/zstdmt_compress.c @@ -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 ===== */ diff --git a/rust/src/zstdmt_compress.rs b/rust/src/zstdmt_compress.rs index 56f3d090c..5065a9abd 100644 --- a/rust/src/zstdmt_compress.rs +++ b/rust/src/zstdmt_compress.rs @@ -56,6 +56,25 @@ pub struct ZSTDMT_chunkProcessResult { pub lastBlockSize: usize, } +/// Scalar job state used by the Rust compression-job scheduler. +/// +/// The job descriptor, pools, synchronization, and codec state remain +/// private to C. Rust uses only the frame-position flags to choose the +/// sequencing and non-first-job header stages. +#[repr(C)] +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub struct ZSTDMT_compressionJobProjection { + pub firstJob: c_uint, + pub lastJob: c_uint, +} + +pub type ZSTDMT_compressionJobStepFn = unsafe extern "C" fn(*mut c_void) -> usize; +pub type ZSTDMT_compressionJobVoidFn = unsafe extern "C" fn(*mut c_void); +pub type ZSTDMT_compressionJobCompressFn = + unsafe extern "C" fn(*mut c_void, c_uint) -> ZSTDMT_chunkProcessResult; +pub type ZSTDMT_compressionJobErrorFn = unsafe extern "C" fn(*mut c_void, usize); +pub type ZSTDMT_compressionJobFinishFn = unsafe extern "C" fn(*mut c_void, usize); + #[repr(C)] #[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] pub struct ZSTDMT_flushPublicationResult { @@ -560,6 +579,135 @@ pub unsafe extern "C" fn ZSTDMT_rust_createCompressionJob( ) } +/// Run the high-level worker-job sequence while C owns all codec and +/// synchronization operations behind callbacks. +/// +/// Resource acquisition and every codec-facing stage can fail with a zstd +/// error. Rust stops at the first such error, reports it before the common +/// C-owned cleanup callback, and passes the final block size only on the +/// successful compression path. +#[inline] +fn compression_job_with( + projection: ZSTDMT_compressionJobProjection, + mut acquire_resources: A, + mut prepare_parameters: P, + mut generate_sequences: S, + mut begin_job: B, + mut apply_sequences: Q, + mut write_frame_header: H, + mut compress_job: C, + mut trace_job: T, + mut set_error: E, + mut finish_job: F, +) where + A: FnMut() -> usize, + P: FnMut(), + S: FnMut(), + B: FnMut() -> usize, + Q: FnMut(), + H: FnMut() -> usize, + C: FnMut(c_uint) -> ZSTDMT_chunkProcessResult, + T: FnMut(), + E: FnMut(usize), + F: FnMut(usize), +{ + let mut error = acquire_resources(); + let mut last_block_size = 0; + + if !ERR_isError(error) { + prepare_parameters(); + generate_sequences(); + error = begin_job(); + } + + if !ERR_isError(error) { + apply_sequences(); + if projection.firstJob == 0 { + error = write_frame_header(); + } + } + + if !ERR_isError(error) { + let result = compress_job(projection.lastJob); + if ERR_isError(result.error) { + error = result.error; + } else { + last_block_size = result.lastBlockSize; + trace_job(); + } + } + + if ERR_isError(error) { + set_error(error); + last_block_size = 0; + } + finish_job(last_block_size); +} + +/// C ABI entry point for the worker-job orchestration. C supplies callbacks +/// that keep the private descriptor, pools, mutexes, and codec operations on +/// the C side of this narrow projection. +#[cfg(not(test))] +#[no_mangle] +pub unsafe extern "C" fn ZSTDMT_rust_compressionJob( + projection: *const ZSTDMT_compressionJobProjection, + opaque: *mut c_void, + acquireResources: Option, + prepareParameters: Option, + generateSequences: Option, + beginJob: Option, + applySequences: Option, + writeFrameHeader: Option, + compressJob: Option, + traceJob: Option, + setError: Option, + finishJob: Option, +) { + let Some(projection) = (unsafe { projection.as_ref() }).copied() else { + return; + }; + let ( + Some(acquire_resources), + Some(prepare_parameters), + Some(generate_sequences), + Some(begin_job), + Some(apply_sequences), + Some(write_frame_header), + Some(compress_job), + Some(trace_job), + Some(set_error), + Some(finish_job), + ) = ( + acquireResources, + prepareParameters, + generateSequences, + beginJob, + applySequences, + writeFrameHeader, + compressJob, + traceJob, + setError, + finishJob, + ) + else { + return; + }; + + compression_job_with( + projection, + || unsafe { acquire_resources(opaque) }, + || unsafe { prepare_parameters(opaque) }, + || unsafe { generate_sequences(opaque) }, + || unsafe { begin_job(opaque) }, + || unsafe { apply_sequences(opaque) }, + || unsafe { write_frame_header(opaque) }, + |last_job| unsafe { compress_job(opaque, last_job) }, + || unsafe { trace_job(opaque) }, + |error| unsafe { set_error(opaque, error) }, + |last_block_size| unsafe { finish_job(opaque, last_block_size) }, + ); +} + #[inline] fn invalid_flush_publication( output_pos: usize, @@ -2556,6 +2704,7 @@ pub unsafe extern "C" fn ZSTDMT_rust_cctx_pool_release(pool: *mut RustCCtxPool, mod tests { use super::*; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; + use std::{cell::RefCell, rc::Rc}; const DEFAULT_MEM: ZstdCustomMem = ZstdCustomMem { customAlloc: None, @@ -2580,6 +2729,152 @@ mod tests { calls: Vec<(usize, usize)>, } + #[derive(Default)] + struct MockCompressionJob { + events: Vec<&'static str>, + errors: Vec, + finished: Vec, + } + + fn record_compression_job_event(state: &Rc>, event: &'static str) { + state.borrow_mut().events.push(event); + } + + #[test] + fn compression_job_runs_first_job_stages_and_reports_final_block() { + let state = Rc::new(RefCell::new(MockCompressionJob::default())); + let acquire_state = Rc::clone(&state); + let prepare_state = Rc::clone(&state); + let sequence_state = Rc::clone(&state); + let begin_state = Rc::clone(&state); + let apply_state = Rc::clone(&state); + let compress_state = Rc::clone(&state); + let trace_state = Rc::clone(&state); + let finish_state = Rc::clone(&state); + + compression_job_with( + ZSTDMT_compressionJobProjection { + firstJob: 1, + lastJob: 1, + }, + move || { + record_compression_job_event(&acquire_state, "acquire"); + 0 + }, + move || record_compression_job_event(&prepare_state, "prepare"), + move || record_compression_job_event(&sequence_state, "sequences"), + move || { + record_compression_job_event(&begin_state, "begin"); + 0 + }, + move || record_compression_job_event(&apply_state, "apply"), + || panic!("first jobs do not rewrite a frame header"), + move |last_job| { + assert_eq!(last_job, 1); + record_compression_job_event(&compress_state, "compress"); + ZSTDMT_chunkProcessResult { + error: 0, + lastBlockSize: 7, + } + }, + move || record_compression_job_event(&trace_state, "trace"), + |_error| panic!("success must not report an error"), + move |last_block_size| { + record_compression_job_event(&finish_state, "finish"); + finish_state.borrow_mut().finished.push(last_block_size); + }, + ); + + let state = state.borrow(); + assert_eq!( + state.events, + vec![ + "acquire", + "prepare", + "sequences", + "begin", + "apply", + "compress", + "trace", + "finish" + ] + ); + assert!(state.errors.is_empty()); + assert_eq!(state.finished, vec![7]); + } + + #[test] + fn compression_job_stops_on_non_first_chunk_error_and_cleans_up() { + let state = Rc::new(RefCell::new(MockCompressionJob::default())); + let acquire_state = Rc::clone(&state); + let prepare_state = Rc::clone(&state); + let sequence_state = Rc::clone(&state); + let begin_state = Rc::clone(&state); + let apply_state = Rc::clone(&state); + let header_state = Rc::clone(&state); + let compress_state = Rc::clone(&state); + let error_state = Rc::clone(&state); + let finish_state = Rc::clone(&state); + let expected_error = ERROR(ZstdErrorCode::DstSizeTooSmall); + + compression_job_with( + ZSTDMT_compressionJobProjection { + firstJob: 0, + lastJob: 1, + }, + move || { + record_compression_job_event(&acquire_state, "acquire"); + 0 + }, + move || record_compression_job_event(&prepare_state, "prepare"), + move || record_compression_job_event(&sequence_state, "sequences"), + move || { + record_compression_job_event(&begin_state, "begin"); + 0 + }, + move || record_compression_job_event(&apply_state, "apply"), + move || { + record_compression_job_event(&header_state, "header"); + 0 + }, + move |last_job| { + assert_eq!(last_job, 1); + record_compression_job_event(&compress_state, "compress"); + ZSTDMT_chunkProcessResult { + error: expected_error, + lastBlockSize: 0, + } + }, + || panic!("a failed chunk must not be traced"), + move |error| { + record_compression_job_event(&error_state, "error"); + error_state.borrow_mut().errors.push(error); + }, + move |last_block_size| { + record_compression_job_event(&finish_state, "finish"); + finish_state.borrow_mut().finished.push(last_block_size); + }, + ); + + let state = state.borrow(); + assert_eq!( + state.events, + vec![ + "acquire", + "prepare", + "sequences", + "begin", + "apply", + "header", + "compress", + "error", + "finish" + ] + ); + assert_eq!(state.errors, vec![expected_error]); + assert_eq!(state.finished, vec![0]); + } + unsafe extern "C" fn mock_compress_continue( cctx: *mut c_void, _dst: *mut c_void,