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
This commit is contained in:
2026-07-21 22:06:59 +02:00
parent f89ae75898
commit 0b18e73d4b
2 changed files with 190 additions and 7 deletions
+60 -7
View File
@@ -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)
+130
View File
@@ -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<ZSTDMT_compressionJobVoidFn>,
pub signal: Option<ZSTDMT_compressionJobVoidFn>,
pub unlock: Option<ZSTDMT_compressionJobVoidFn>,
pub debug: Option<ZSTDMT_chunkProgressFn>,
}
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::<usize>());
assert!(offset_of!(ZSTDMT_compressionJobProgressProjection, lock)
== size_of::<*mut c_void>() + 2 * size_of::<usize>());
assert!(offset_of!(ZSTDMT_compressionJobProgressProjection, signal)
== size_of::<*mut c_void>() + 2 * size_of::<usize>() + size_of::<usize>());
assert!(offset_of!(ZSTDMT_compressionJobProgressProjection, unlock)
== size_of::<*mut c_void>() + 2 * size_of::<usize>() + 2 * size_of::<usize>());
assert!(offset_of!(ZSTDMT_compressionJobProgressProjection, debug)
== size_of::<*mut c_void>() + 2 * size_of::<usize>() + 3 * size_of::<usize>());
assert!(size_of::<ZSTDMT_compressionJobProgressProjection>()
== size_of::<*mut c_void>() + 2 * size_of::<usize>() + 4 * size_of::<usize>());
};
/// 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::<ZSTDMT_compressionJobProgressProjection>()
.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::<MockJobProgress>()).events.push("lock") };
}
unsafe extern "C" fn mock_job_progress_signal(context: *mut c_void) {
unsafe { (*context.cast::<MockJobProgress>()).events.push("signal") };
}
unsafe extern "C" fn mock_job_progress_unlock(context: *mut c_void) {
unsafe { (*context.cast::<MockJobProgress>()).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::<MockJobProgress>() };
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,