From 0b18e73d4b166c7232c0962cdc437e5d5fd7a337 Mon Sep 17 00:00:00 2001 From: ddidderr Date: Tue, 21 Jul 2026 22:06:59 +0200 Subject: [PATCH] feat(mt): move compression progress publication to Rust Move the intermediate worker progress state update out of the C callback and into a Rust-owned publication leaf. The projection keeps the private job descriptor and pthread objects in C while preserving the lock, counter update, diagnostic, signal, and unlock order, including size_t wrapping behavior. Test Plan: - git diff --cached --check - capped cargo check --tests - capped cargo clippy --tests -- -A clippy::manual-bits -D warnings - capped make -j1 - capped make -j1 -C tests test - standalone capped cargo test was attempted but its hybrid link lacks the pre-existing ZSTD_rust_dctx_trace_view, ZSTD_rust_dctx_view, and ZSTD_rust_block_context_init symbols --- lib/compress/zstdmt_compress.c | 67 +++++++++++++++-- rust/src/zstdmt_compress.rs | 130 +++++++++++++++++++++++++++++++++ 2 files changed, 190 insertions(+), 7 deletions(-) diff --git a/lib/compress/zstdmt_compress.c b/lib/compress/zstdmt_compress.c index b8be2e421..6532f9e98 100644 --- a/lib/compress/zstdmt_compress.c +++ b/lib/compress/zstdmt_compress.c @@ -269,6 +269,37 @@ typedef size_t (*ZSTDMT_compressionJobSetParameterFn)(void* opaque, int value); typedef void (*ZSTDMT_compressionJobApplySequencesFn)( void* opaque, void* cctx, void* sequences, size_t nbSequences); +typedef struct { + void* callbackContext; + size_t* cSize; + size_t* consumed; + ZSTDMT_compressionJobVoidFn lock; + ZSTDMT_compressionJobVoidFn signal; + ZSTDMT_compressionJobVoidFn unlock; + ZSTDMT_chunkProgressFn debug; +} ZSTDMT_RustCompressionJobProgressProjection; +typedef char ZSTDMT_compression_job_progress_projection_layout[ + (offsetof(ZSTDMT_RustCompressionJobProgressProjection, callbackContext) + == 0 + && offsetof(ZSTDMT_RustCompressionJobProgressProjection, cSize) + == sizeof(void*) + && offsetof(ZSTDMT_RustCompressionJobProgressProjection, consumed) + == sizeof(void*) + sizeof(size_t) + && offsetof(ZSTDMT_RustCompressionJobProgressProjection, lock) + == sizeof(void*) + 2 * sizeof(size_t) + && offsetof(ZSTDMT_RustCompressionJobProgressProjection, signal) + == sizeof(void*) + 2 * sizeof(size_t) + sizeof(void*) + && offsetof(ZSTDMT_RustCompressionJobProgressProjection, unlock) + == sizeof(void*) + 2 * sizeof(size_t) + 2 * sizeof(void*) + && offsetof(ZSTDMT_RustCompressionJobProgressProjection, debug) + == sizeof(void*) + 2 * sizeof(size_t) + 3 * sizeof(void*) + && sizeof(ZSTDMT_RustCompressionJobProgressProjection) + == sizeof(void*) + 2 * sizeof(size_t) + 4 * sizeof(void*)) + ? 1 : -1]; + +void ZSTDMT_rust_compressionJobProgress( + void* opaque, size_t cSize, size_t consumed); + typedef struct { void* callbackContext; size_t srcSize; @@ -1925,18 +1956,32 @@ typedef struct { unsigned frameChecksumNeeded; /* used only by mtctx */ } ZSTDMT_jobDescription; -static void ZSTDMT_compressionJobProgress(void* opaque, size_t cSize, size_t consumed) +static void ZSTDMT_compressionJobProgressLock(void* opaque) { ZSTDMT_jobDescription* const job = (ZSTDMT_jobDescription*)opaque; ZSTD_PTHREAD_MUTEX_LOCK(&job->job_mutex); - job->cSize += cSize; - job->consumed = consumed; - DEBUGLOG(5, "ZSTDMT_compressionJob: compress new block : cSize==%u bytes (total: %u)", - (U32)cSize, (U32)job->cSize); +} + +static void ZSTDMT_compressionJobProgressSignal(void* opaque) +{ + ZSTDMT_jobDescription* const job = (ZSTDMT_jobDescription*)opaque; ZSTD_pthread_cond_signal(&job->job_cond); /* warns some more data is ready to be flushed */ +} + +static void ZSTDMT_compressionJobProgressUnlock(void* opaque) +{ + ZSTDMT_jobDescription* const job = (ZSTDMT_jobDescription*)opaque; ZSTD_pthread_mutex_unlock(&job->job_mutex); } +static void ZSTDMT_compressionJobProgressDebug( + void* opaque, size_t cSize, size_t totalCSize) +{ + (void)opaque; + DEBUGLOG(5, "ZSTDMT_compressionJob: compress new block : cSize==%u bytes (total: %u)", + (U32)cSize, (U32)totalCSize); +} + typedef struct { ZSTDMT_jobDescription* job; ZSTD_CCtx_params jobParams; @@ -2147,8 +2192,9 @@ static ZSTDMT_chunkProcessResult ZSTDMT_compressionJobCompress( void* opaque, unsigned lastJob) { ZSTDMT_compressionJobState* const state = - (ZSTDMT_compressionJobState*)opaque; + (ZSTDMT_compressionJobState*)opaque; ZSTDMT_jobDescription* const job = state->job; + ZSTDMT_RustCompressionJobProgressProjection progress; size_t const chunkSize = 4*ZSTD_BLOCKSIZE_MAX; assert(lastJob == job->lastJob); @@ -2158,10 +2204,17 @@ static ZSTDMT_chunkProcessResult ZSTDMT_compressionJobCompress( assert(job->cSize == 0); assert(chunkSize > 0); assert((chunkSize & (chunkSize - 1)) == 0); /* chunkSize must be power of 2 for mask==(chunkSize-1) to work */ + progress.callbackContext = job; + progress.cSize = &job->cSize; + progress.consumed = &job->consumed; + progress.lock = ZSTDMT_compressionJobProgressLock; + progress.signal = ZSTDMT_compressionJobProgressSignal; + progress.unlock = ZSTDMT_compressionJobProgressUnlock; + progress.debug = ZSTDMT_compressionJobProgressDebug; return ZSTDMT_rust_compressJobChunks( state->cctx, job->src.start, job->src.size, state->dstBuff.start, state->dstBuff.capacity, chunkSize, lastJob, - job, ZSTDMT_compressionJobProgress); + &progress, ZSTDMT_rust_compressionJobProgress); } static void ZSTDMT_compressionJobTrace(void* opaque) diff --git a/rust/src/zstdmt_compress.rs b/rust/src/zstdmt_compress.rs index a81454e2c..eb40ef915 100644 --- a/rust/src/zstdmt_compress.rs +++ b/rust/src/zstdmt_compress.rs @@ -534,6 +534,77 @@ pub type ZSTDMT_compressionJobSetParameterFn = unsafe extern "C" fn(*mut c_void, pub type ZSTDMT_compressionJobApplySequencesFn = unsafe extern "C" fn(*mut c_void, *mut c_void, *mut c_void, usize); +/// Live fields and synchronization callbacks for one worker's intermediate +/// compression progress. The private job descriptor and pthread objects +/// remain in C; Rust owns the publication order and wrapping counter update. +#[repr(C)] +#[derive(Clone, Copy)] +pub struct ZSTDMT_compressionJobProgressProjection { + pub callbackContext: *mut c_void, + pub cSize: *mut usize, + pub consumed: *mut usize, + pub lock: Option, + pub signal: Option, + pub unlock: Option, + pub debug: Option, +} + +const _: () = { + assert!( + offset_of!(ZSTDMT_compressionJobProgressProjection, callbackContext) == 0 + ); + assert!( + offset_of!(ZSTDMT_compressionJobProgressProjection, cSize) == size_of::<*mut c_void>() + ); + assert!(offset_of!(ZSTDMT_compressionJobProgressProjection, consumed) + == size_of::<*mut c_void>() + size_of::()); + assert!(offset_of!(ZSTDMT_compressionJobProgressProjection, lock) + == size_of::<*mut c_void>() + 2 * size_of::()); + assert!(offset_of!(ZSTDMT_compressionJobProgressProjection, signal) + == size_of::<*mut c_void>() + 2 * size_of::() + size_of::()); + assert!(offset_of!(ZSTDMT_compressionJobProgressProjection, unlock) + == size_of::<*mut c_void>() + 2 * size_of::() + 2 * size_of::()); + assert!(offset_of!(ZSTDMT_compressionJobProgressProjection, debug) + == size_of::<*mut c_void>() + 2 * size_of::() + 3 * size_of::()); + assert!(size_of::() + == size_of::<*mut c_void>() + 2 * size_of::() + 4 * size_of::()); +}; + +/// Publish one intermediate worker block while preserving the C callback's +/// lock, update, diagnostic, signal, and unlock order. +#[no_mangle] +pub unsafe extern "C" fn ZSTDMT_rust_compressionJobProgress( + opaque: *mut c_void, + c_size: usize, + consumed: usize, +) { + let Some(state) = (unsafe { + opaque + .cast::() + .as_ref() + }) else { + return; + }; + let (Some(lock), Some(signal), Some(unlock)) = (state.lock, state.signal, state.unlock) + else { + return; + }; + if state.cSize.is_null() || state.consumed.is_null() { + return; + } + + unsafe { + lock(state.callbackContext); + *state.cSize = (*state.cSize).wrapping_add(c_size); + *state.consumed = consumed; + if let Some(debug) = state.debug { + debug(state.callbackContext, c_size, *state.cSize); + } + signal(state.callbackContext); + unlock(state.callbackContext); + } +} + #[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] struct ZSTDMT_compressionJobSequenceProjection { cctx: *mut c_void, @@ -6207,6 +6278,12 @@ mod tests { calls: Vec<(usize, usize)>, } + #[derive(Default)] + struct MockJobProgress { + events: Vec<&'static str>, + debug_calls: Vec<(usize, usize)>, + } + #[derive(Default)] struct MockCompressionJob { events: Vec<&'static str>, @@ -8349,6 +8426,28 @@ mod tests { progress.calls.push((c_size, consumed)); } + unsafe extern "C" fn mock_job_progress_lock(context: *mut c_void) { + unsafe { (*context.cast::()).events.push("lock") }; + } + + unsafe extern "C" fn mock_job_progress_signal(context: *mut c_void) { + unsafe { (*context.cast::()).events.push("signal") }; + } + + unsafe extern "C" fn mock_job_progress_unlock(context: *mut c_void) { + unsafe { (*context.cast::()).events.push("unlock") }; + } + + unsafe extern "C" fn mock_job_progress_debug( + context: *mut c_void, + c_size: usize, + total_c_size: usize, + ) { + let progress = unsafe { &mut *context.cast::() }; + progress.events.push("debug"); + progress.debug_calls.push((c_size, total_c_size)); + } + fn run_mock_chunk_loop( compressor: &mut MockChunkCompressor, progress: &mut MockChunkProgress, @@ -8375,6 +8474,37 @@ mod tests { } } + #[test] + fn compression_job_progress_preserves_publication_order_and_wraps_c_size() { + let mut c_size = usize::MAX - 1; + let mut consumed = 0; + let mut progress = MockJobProgress::default(); + let projection = ZSTDMT_compressionJobProgressProjection { + callbackContext: (&mut progress as *mut MockJobProgress).cast(), + cSize: &mut c_size, + consumed: &mut consumed, + lock: Some(mock_job_progress_lock), + signal: Some(mock_job_progress_signal), + unlock: Some(mock_job_progress_unlock), + debug: Some(mock_job_progress_debug), + }; + + unsafe { + ZSTDMT_rust_compressionJobProgress( + (&projection as *const ZSTDMT_compressionJobProgressProjection) + .cast_mut() + .cast(), + 3, + 17, + ) + }; + + assert_eq!(c_size, 1); + assert_eq!(consumed, 17); + assert_eq!(progress.events, ["lock", "debug", "signal", "unlock"]); + assert_eq!(progress.debug_calls, [(3, 1)]); + } + fn run_output_publication( output: &mut [u8], output_size: usize,