From 45d54a71ad0b0a8f9dee0517b340600f627c9797 Mon Sep 17 00:00:00 2001 From: ddidderr Date: Sat, 18 Jul 2026 23:28:36 +0200 Subject: [PATCH] feat(mt): move compression job creation policy into Rust Move the multithreaded compression job-creation decision tree into a Rust projection while keeping C-owned job descriptors, input buffers, pools, synchronization, worker callbacks, and terminal empty-block serialization behind explicit callbacks. The Rust policy now preserves ring-table full checks, prepared-job retry state, prefix advancement, frame checksum rules, and the non-first empty terminal job path without crossing private C layouts. Test Plan: - cargo test --manifest-path rust/Cargo.toml --all-targets -- --test-threads=1 - cargo test --manifest-path rust/Cargo.toml --no-default-features --features compression,decompression,dict-builder,legacy-v01,legacy-v02,legacy-v03,legacy-v04,legacy-v05,legacy-v06,legacy-v07 --all-targets -- --test-threads=1 - cargo clippy --manifest-path rust/Cargo.toml --all-targets -- -D warnings - make -B -C lib -j2 lib - make -B -C tests -j2 test-zstd test-pool test-legacy test-invalidDictionaries --- lib/compress/zstdmt_compress.c | 208 +++++++++++++------ rust/src/zstdmt_compress.rs | 357 +++++++++++++++++++++++++++++++++ 2 files changed, 499 insertions(+), 66 deletions(-) diff --git a/lib/compress/zstdmt_compress.c b/lib/compress/zstdmt_compress.c index 90aa7aebe..c37a2b7de 100644 --- a/lib/compress/zstdmt_compress.c +++ b/lib/compress/zstdmt_compress.c @@ -202,6 +202,56 @@ ZSTDMT_RustEmptyBlockResult ZSTDMT_rust_writeLastEmptyBlock( const ZSTDMT_RustEmptyBlockJobProjection* projection, void* opaque, ZSTDMT_bufferGetFn getBuffer); +typedef struct { + unsigned doneJobID; + unsigned nextJobID; + unsigned jobIDMask; + unsigned jobReady; + const void* srcStart; + size_t srcSize; + size_t inBuffFilled; + const void* prefixStart; + size_t prefixSize; + size_t targetPrefixSize; + unsigned endFrame; + unsigned checksumFlag; +} ZSTDMT_RustCreateJobProjection; +typedef struct { + const void* srcStart; + size_t srcSize; + const void* prefixStart; + size_t prefixSize; + const void* nextPrefixStart; + size_t nextPrefixSize; + size_t roundBuffPosDelta; + unsigned jobNumber; + unsigned firstJob; + unsigned lastJob; + unsigned frameChecksumNeeded; + unsigned clearChecksumFlag; +} ZSTDMT_RustJobInitialization; +typedef struct { + size_t returnCode; + unsigned action; + unsigned jobID; + unsigned jobNumber; + unsigned nextJobID; + unsigned jobReady; +} ZSTDMT_RustCreateJobResult; +enum { + ZSTDMT_CREATE_JOB_TABLE_FULL = 0, + ZSTDMT_CREATE_JOB_POST = 1, + ZSTDMT_CREATE_JOB_EMPTY = 2 +}; +typedef void (*ZSTDMT_prepareJobFn)(void* opaque, unsigned jobID, + const ZSTDMT_RustJobInitialization* init); +typedef void (*ZSTDMT_writeEmptyJobFn)(void* opaque, unsigned jobID); +typedef int (*ZSTDMT_tryAddJobFn)(void* opaque, unsigned jobID); +ZSTDMT_RustCreateJobResult ZSTDMT_rust_createCompressionJob( + const ZSTDMT_RustCreateJobProjection* projection, + void* opaque, ZSTDMT_prepareJobFn prepareJob, + ZSTDMT_writeEmptyJobFn writeEmptyJob, ZSTDMT_tryAddJobFn tryAddJob); + typedef struct { rawSeq* seq; size_t pos; @@ -1477,82 +1527,108 @@ static void ZSTDMT_writeLastEmptyBlock(ZSTDMT_jobDescription* job) assert(!ZSTD_isError(job->cSize)); } +static void ZSTDMT_prepareCompressionJob( + void* opaque, unsigned jobID, + const ZSTDMT_RustJobInitialization* initialization) +{ + ZSTDMT_CCtx* const mtctx = (ZSTDMT_CCtx*)opaque; + ZSTDMT_jobDescription* const job = &mtctx->jobs[jobID]; + BYTE const* const src = (const BYTE*)initialization->srcStart; + + DEBUGLOG(5, "ZSTDMT_createCompressionJob: preparing job %u to compress %u bytes with %u preload ", + initialization->jobNumber, (U32)initialization->srcSize, + (U32)initialization->prefixSize); + assert(mtctx->inBuff.filled >= initialization->srcSize); + job->src = (Range){ src, initialization->srcSize }; + job->prefix = (Range){ + (const BYTE*)initialization->prefixStart, initialization->prefixSize + }; + job->consumed = 0; + job->cSize = 0; + job->params = mtctx->params; + job->cdict = initialization->firstJob ? mtctx->cdict : NULL; + job->fullFrameSize = mtctx->frameContentSize; + job->dstBuff = g_nullBuffer; + job->cctxPool = mtctx->cctxPool; + job->bufPool = mtctx->bufPool; + job->seqPool = mtctx->seqPool; + job->serial = &mtctx->serial; + job->jobID = initialization->jobNumber; + job->firstJob = initialization->firstJob; + job->lastJob = initialization->lastJob; + job->frameChecksumNeeded = initialization->frameChecksumNeeded; + job->dstFlushed = 0; + + /* Update the round buffer position and clear the input buffer to be reset. */ + mtctx->roundBuff.pos += initialization->roundBuffPosDelta; + mtctx->inBuff.buffer = g_nullBuffer; + mtctx->inBuff.filled = 0; + mtctx->inBuff.prefix = (Range){ + (const BYTE*)initialization->nextPrefixStart, + initialization->nextPrefixSize + }; + if (initialization->lastJob) { + mtctx->frameEnded = 1; + if (initialization->clearChecksumFlag) + mtctx->params.fParams.checksumFlag = 0; + } +} + +static void ZSTDMT_writeEmptyCompressionJob(void* opaque, unsigned jobID) +{ + ZSTDMT_CCtx* const mtctx = (ZSTDMT_CCtx*)opaque; + ZSTDMT_writeLastEmptyBlock(&mtctx->jobs[jobID]); +} + +static int ZSTDMT_tryAddCompressionJob(void* opaque, unsigned jobID) +{ + ZSTDMT_CCtx* const mtctx = (ZSTDMT_CCtx*)opaque; + return POOL_tryAdd(mtctx->factory, ZSTDMT_compressionJob, &mtctx->jobs[jobID]); +} + static size_t ZSTDMT_createCompressionJob(ZSTDMT_CCtx* mtctx, size_t srcSize, ZSTD_EndDirective endOp) { - unsigned const jobID = mtctx->nextJobID & mtctx->jobIDMask; - int const endFrame = (endOp == ZSTD_e_end); + ZSTDMT_RustCreateJobProjection const projection = { + mtctx->doneJobID, + mtctx->nextJobID, + mtctx->jobIDMask, + mtctx->jobReady, + mtctx->inBuff.buffer.start, + srcSize, + mtctx->inBuff.filled, + mtctx->inBuff.prefix.start, + mtctx->inBuff.prefix.size, + mtctx->targetPrefixSize, + (unsigned)(endOp == ZSTD_e_end), + (unsigned)mtctx->params.fParams.checksumFlag + }; + ZSTDMT_RustCreateJobResult const result = ZSTDMT_rust_createCompressionJob( + &projection, mtctx, ZSTDMT_prepareCompressionJob, + ZSTDMT_writeEmptyCompressionJob, ZSTDMT_tryAddCompressionJob); - if (mtctx->nextJobID > mtctx->doneJobID + mtctx->jobIDMask) { + if (result.action == ZSTDMT_CREATE_JOB_TABLE_FULL) { DEBUGLOG(5, "ZSTDMT_createCompressionJob: will not create new job : table is full"); assert((mtctx->nextJobID & mtctx->jobIDMask) == (mtctx->doneJobID & mtctx->jobIDMask)); - return 0; + return result.returnCode; } - if (!mtctx->jobReady) { - BYTE const* src = (BYTE const*)mtctx->inBuff.buffer.start; - DEBUGLOG(5, "ZSTDMT_createCompressionJob: preparing job %u to compress %u bytes with %u preload ", - mtctx->nextJobID, (U32)srcSize, (U32)mtctx->inBuff.prefix.size); - mtctx->jobs[jobID].src.start = src; - mtctx->jobs[jobID].src.size = srcSize; - assert(mtctx->inBuff.filled >= srcSize); - mtctx->jobs[jobID].prefix = mtctx->inBuff.prefix; - mtctx->jobs[jobID].consumed = 0; - mtctx->jobs[jobID].cSize = 0; - mtctx->jobs[jobID].params = mtctx->params; - mtctx->jobs[jobID].cdict = mtctx->nextJobID==0 ? mtctx->cdict : NULL; - mtctx->jobs[jobID].fullFrameSize = mtctx->frameContentSize; - mtctx->jobs[jobID].dstBuff = g_nullBuffer; - mtctx->jobs[jobID].cctxPool = mtctx->cctxPool; - mtctx->jobs[jobID].bufPool = mtctx->bufPool; - mtctx->jobs[jobID].seqPool = mtctx->seqPool; - mtctx->jobs[jobID].serial = &mtctx->serial; - mtctx->jobs[jobID].jobID = mtctx->nextJobID; - mtctx->jobs[jobID].firstJob = (mtctx->nextJobID==0); - mtctx->jobs[jobID].lastJob = endFrame; - mtctx->jobs[jobID].frameChecksumNeeded = mtctx->params.fParams.checksumFlag && endFrame && (mtctx->nextJobID>0); - mtctx->jobs[jobID].dstFlushed = 0; - - /* Update the round buffer pos and clear the input buffer to be reset */ - mtctx->roundBuff.pos += srcSize; - mtctx->inBuff.buffer = g_nullBuffer; - mtctx->inBuff.filled = 0; - /* Set the prefix for next job */ - if (!endFrame) { - size_t const newPrefixSize = MIN(srcSize, mtctx->targetPrefixSize); - mtctx->inBuff.prefix.start = src + srcSize - newPrefixSize; - mtctx->inBuff.prefix.size = newPrefixSize; - } else { /* endFrame==1 => no need for another input buffer */ - mtctx->inBuff.prefix = kNullRange; - mtctx->frameEnded = endFrame; - if (mtctx->nextJobID == 0) { - /* single job exception : checksum is already calculated directly within worker thread */ - mtctx->params.fParams.checksumFlag = 0; - } } - - if ( (srcSize == 0) - && (mtctx->nextJobID>0)/*single job must also write frame header*/ ) { - DEBUGLOG(5, "ZSTDMT_createCompressionJob: creating a last empty block to end frame"); - assert(endOp == ZSTD_e_end); /* only possible case : need to end the frame with an empty last block */ - ZSTDMT_writeLastEmptyBlock(mtctx->jobs + jobID); - mtctx->nextJobID++; - return 0; - } + mtctx->nextJobID = result.nextJobID; + mtctx->jobReady = result.jobReady; + if (result.action == ZSTDMT_CREATE_JOB_EMPTY) { + DEBUGLOG(5, "ZSTDMT_createCompressionJob: creating a last empty block to end frame"); + assert(endOp == ZSTD_e_end); /* only possible case : need to end the frame with an empty last block */ + return result.returnCode; } DEBUGLOG(5, "ZSTDMT_createCompressionJob: posting job %u : %u bytes (end:%u, jobNb == %u (mod:%u))", - mtctx->nextJobID, - (U32)mtctx->jobs[jobID].src.size, - mtctx->jobs[jobID].lastJob, - mtctx->nextJobID, - jobID); - if (POOL_tryAdd(mtctx->factory, ZSTDMT_compressionJob, &mtctx->jobs[jobID])) { - mtctx->nextJobID++; - mtctx->jobReady = 0; - } else { - DEBUGLOG(5, "ZSTDMT_createCompressionJob: no worker available for job %u", mtctx->nextJobID); - mtctx->jobReady = 1; - } - return 0; + result.jobNumber, + (U32)mtctx->jobs[result.jobID].src.size, + mtctx->jobs[result.jobID].lastJob, + result.jobNumber, + result.jobID); + if (result.jobReady) + DEBUGLOG(5, "ZSTDMT_createCompressionJob: no worker available for job %u", result.jobNumber); + return result.returnCode; } diff --git a/rust/src/zstdmt_compress.rs b/rust/src/zstdmt_compress.rs index 605b34da2..0ce5e9d60 100644 --- a/rust/src/zstdmt_compress.rs +++ b/rust/src/zstdmt_compress.rs @@ -240,6 +240,207 @@ pub unsafe extern "C" fn ZSTDMT_rust_writeLastEmptyBlock( unsafe { write_last_empty_block_with(projection, || get_buffer(opaque)) } } +const CREATE_JOB_TABLE_FULL: c_uint = 0; +const CREATE_JOB_POST: c_uint = 1; +const CREATE_JOB_EMPTY: c_uint = 2; + +/// Scalar MT state supplied by C for one compression-job creation attempt. +/// +/// The job descriptor, input-buffer ownership, and synchronization objects +/// remain private to C. Rust owns the ring-capacity and scheduling policy over +/// this scalar view. +#[repr(C)] +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub struct ZSTDMT_createJobProjection { + pub doneJobID: c_uint, + pub nextJobID: c_uint, + pub jobIDMask: c_uint, + pub jobReady: c_uint, + pub srcStart: *const c_void, + pub srcSize: usize, + pub inBuffFilled: usize, + pub prefixStart: *const c_void, + pub prefixSize: usize, + pub targetPrefixSize: usize, + pub endFrame: c_uint, + pub checksumFlag: c_uint, +} + +/// The scalar fields needed by C to initialize its private job descriptor and +/// the surrounding input state after Rust chooses to prepare a new job. +#[repr(C)] +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub struct ZSTDMT_jobInitialization { + pub srcStart: *const c_void, + pub srcSize: usize, + pub prefixStart: *const c_void, + pub prefixSize: usize, + pub nextPrefixStart: *const c_void, + pub nextPrefixSize: usize, + pub roundBuffPosDelta: usize, + pub jobNumber: c_uint, + pub firstJob: c_uint, + pub lastJob: c_uint, + pub frameChecksumNeeded: c_uint, + pub clearChecksumFlag: c_uint, +} + +#[repr(C)] +#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)] +pub struct ZSTDMT_createJobResult { + pub returnCode: usize, + pub action: c_uint, + pub jobID: c_uint, + pub jobNumber: c_uint, + pub nextJobID: c_uint, + pub jobReady: c_uint, +} + +pub type ZSTDMT_prepareJobFn = unsafe extern "C" fn( + opaque: *mut c_void, + job_id: c_uint, + initialization: *const ZSTDMT_jobInitialization, +); +pub type ZSTDMT_writeEmptyJobFn = unsafe extern "C" fn(opaque: *mut c_void, job_id: c_uint); +pub type ZSTDMT_tryAddJobFn = unsafe extern "C" fn(opaque: *mut c_void, job_id: c_uint) -> c_int; + +#[inline] +fn create_job_table_full(projection: ZSTDMT_createJobProjection) -> ZSTDMT_createJobResult { + ZSTDMT_createJobResult { + action: CREATE_JOB_TABLE_FULL, + jobID: projection.nextJobID & projection.jobIDMask, + jobNumber: projection.nextJobID, + nextJobID: projection.nextJobID, + jobReady: projection.jobReady, + ..ZSTDMT_createJobResult::default() + } +} + +#[inline] +fn create_job_initialization( + projection: ZSTDMT_createJobProjection, +) -> (c_uint, ZSTDMT_jobInitialization) { + let job_number = projection.nextJobID; + let first_job = (job_number == 0) as c_uint; + let last_job = (projection.endFrame != 0) as c_uint; + let frame_checksum_needed = + (projection.checksumFlag != 0 && last_job != 0 && job_number > 0) as c_uint; + + let (next_prefix_start, next_prefix_size) = if last_job != 0 { + (ptr::null(), 0) + } else { + let next_prefix_size = projection.srcSize.min(projection.targetPrefixSize); + let next_prefix_start = projection + .srcStart + .cast::() + .wrapping_add(projection.srcSize - next_prefix_size) + .cast(); + (next_prefix_start, next_prefix_size) + }; + + ( + projection.nextJobID & projection.jobIDMask, + ZSTDMT_jobInitialization { + srcStart: projection.srcStart, + srcSize: projection.srcSize, + prefixStart: projection.prefixStart, + prefixSize: projection.prefixSize, + nextPrefixStart: next_prefix_start, + nextPrefixSize: next_prefix_size, + roundBuffPosDelta: projection.srcSize, + jobNumber: job_number, + firstJob: first_job, + lastJob: last_job, + frameChecksumNeeded: frame_checksum_needed, + clearChecksumFlag: (last_job != 0 && first_job != 0) as c_uint, + }, + ) +} + +#[inline] +fn create_compression_job_with( + projection: ZSTDMT_createJobProjection, + mut prepare_job: P, + mut write_empty_job: E, + mut try_add_job: T, +) -> ZSTDMT_createJobResult +where + P: FnMut(c_uint, &ZSTDMT_jobInitialization), + E: FnMut(c_uint), + T: FnMut(c_uint) -> bool, +{ + /* Match the original unsigned C comparison, including wraparound. */ + if projection.nextJobID > projection.doneJobID.wrapping_add(projection.jobIDMask) { + return create_job_table_full(projection); + } + + let job_id = projection.nextJobID & projection.jobIDMask; + let mut result = ZSTDMT_createJobResult { + action: CREATE_JOB_POST, + jobID: job_id, + jobNumber: projection.nextJobID, + nextJobID: projection.nextJobID, + jobReady: projection.jobReady, + ..ZSTDMT_createJobResult::default() + }; + + if projection.jobReady == 0 { + debug_assert!(projection.inBuffFilled >= projection.srcSize); + let (_, initialization) = create_job_initialization(projection); + prepare_job(job_id, &initialization); + + /* A non-first empty job is represented by the terminal empty block, + * not by a worker submission. */ + if projection.srcSize == 0 && projection.nextJobID > 0 { + write_empty_job(job_id); + result.action = CREATE_JOB_EMPTY; + result.nextJobID = projection.nextJobID.wrapping_add(1); + result.jobReady = 0; + return result; + } + } + + if try_add_job(job_id) { + result.nextJobID = projection.nextJobID.wrapping_add(1); + result.jobReady = 0; + } else { + result.jobReady = 1; + } + result +} + +/// Apply MT scalar scheduling policy while C retains all private descriptor, +/// pool, mutex, condition-variable, and worker-callback operations. +#[cfg(not(test))] +#[no_mangle] +pub unsafe extern "C" fn ZSTDMT_rust_createCompressionJob( + projection: *const ZSTDMT_createJobProjection, + opaque: *mut c_void, + prepareJob: Option, + writeEmptyJob: Option, + tryAddJob: Option, +) -> ZSTDMT_createJobResult { + let Some(projection) = (unsafe { projection.as_ref() }).copied() else { + return ZSTDMT_createJobResult::default(); + }; + let (Some(prepare_job), Some(write_empty_job), Some(try_add_job)) = + (prepareJob, writeEmptyJob, tryAddJob) + else { + return ZSTDMT_createJobResult::default(); + }; + + create_compression_job_with( + projection, + |job_id, initialization| unsafe { + prepare_job(opaque, job_id, initialization); + }, + |job_id| unsafe { + write_empty_job(opaque, job_id); + }, + |job_id| unsafe { try_add_job(opaque, job_id) != 0 }, + ) +} + #[inline] fn invalid_flush_publication( output_pos: usize, @@ -2155,6 +2356,162 @@ mod tests { assert_eq!(output, original); } + #[test] + fn create_job_rejects_full_table_without_callbacks() { + let projection = ZSTDMT_createJobProjection { + doneJobID: 0, + nextJobID: 8, + jobIDMask: 7, + ..ZSTDMT_createJobProjection::default() + }; + let result = create_compression_job_with( + projection, + |_job_id, _initialization| panic!("table-full job must not be prepared"), + |_job_id| panic!("table-full job must not be written"), + |_job_id| panic!("table-full job must not be posted"), + ); + + assert_eq!(result.action, CREATE_JOB_TABLE_FULL); + assert_eq!(result.jobID, 0); + assert_eq!(result.jobNumber, 8); + assert_eq!(result.nextJobID, 8); + assert_eq!(result.jobReady, 0); + } + + #[test] + fn create_job_projects_ordinary_initialization_and_prefix_advance() { + let source = [0u8; 8]; + let prefix = [1u8; 3]; + let projection = ZSTDMT_createJobProjection { + doneJobID: 0, + nextJobID: 3, + jobIDMask: 7, + srcStart: source.as_ptr().cast(), + srcSize: source.len(), + inBuffFilled: source.len(), + prefixStart: prefix.as_ptr().cast(), + prefixSize: prefix.len(), + targetPrefixSize: 4, + checksumFlag: 1, + ..ZSTDMT_createJobProjection::default() + }; + let mut initialization = None; + let mut posted = None; + let result = create_compression_job_with( + projection, + |job_id, value| { + assert_eq!(job_id, 3); + initialization = Some(*value); + }, + |_job_id| panic!("ordinary job must not use the empty-block path"), + |job_id| { + posted = Some(job_id); + true + }, + ); + + let initialization = initialization.expect("job should be initialized"); + assert_eq!(initialization.srcStart, source.as_ptr().cast()); + assert_eq!(initialization.srcSize, source.len()); + assert_eq!(initialization.prefixStart, prefix.as_ptr().cast()); + assert_eq!(initialization.prefixSize, prefix.len()); + assert_eq!( + initialization.nextPrefixStart, + source.as_ptr().wrapping_add(4).cast() + ); + assert_eq!(initialization.nextPrefixSize, 4); + assert_eq!(initialization.roundBuffPosDelta, source.len()); + assert_eq!(initialization.jobNumber, 3); + assert_eq!(initialization.firstJob, 0); + assert_eq!(initialization.lastJob, 0); + assert_eq!(initialization.frameChecksumNeeded, 0); + assert_eq!(posted, Some(3)); + assert_eq!(result.action, CREATE_JOB_POST); + assert_eq!(result.jobID, 3); + assert_eq!(result.nextJobID, 4); + assert_eq!(result.jobReady, 0); + } + + #[test] + fn create_job_projects_terminal_empty_frame() { + let prefix = [2u8; 2]; + let projection = ZSTDMT_createJobProjection { + doneJobID: 0, + nextJobID: 2, + jobIDMask: 7, + inBuffFilled: 0, + prefixStart: prefix.as_ptr().cast(), + prefixSize: prefix.len(), + endFrame: 1, + checksumFlag: 1, + ..ZSTDMT_createJobProjection::default() + }; + let mut initialization = None; + let mut empty_job = None; + let result = create_compression_job_with( + projection, + |job_id, value| { + assert_eq!(job_id, 2); + initialization = Some(*value); + }, + |job_id| empty_job = Some(job_id), + |_job_id| panic!("terminal empty job must not be posted"), + ); + + let initialization = initialization.expect("empty job should be initialized"); + assert_eq!(initialization.srcSize, 0); + assert_eq!(initialization.prefixStart, prefix.as_ptr().cast()); + assert_eq!(initialization.prefixSize, prefix.len()); + assert!(initialization.nextPrefixStart.is_null()); + assert_eq!(initialization.nextPrefixSize, 0); + assert_eq!(initialization.roundBuffPosDelta, 0); + assert_eq!(initialization.jobNumber, 2); + assert_eq!(initialization.firstJob, 0); + assert_eq!(initialization.lastJob, 1); + assert_eq!(initialization.frameChecksumNeeded, 1); + assert_eq!(initialization.clearChecksumFlag, 0); + assert_eq!(empty_job, Some(2)); + assert_eq!(result.action, CREATE_JOB_EMPTY); + assert_eq!(result.nextJobID, 3); + assert_eq!(result.jobReady, 0); + } + + #[test] + fn create_job_preserves_pool_success_and_failure_state() { + let projection = ZSTDMT_createJobProjection { + doneJobID: 0, + nextJobID: 4, + jobIDMask: 7, + jobReady: 1, + ..ZSTDMT_createJobProjection::default() + }; + let success = create_compression_job_with( + projection, + |_job_id, _initialization| panic!("prepared job must not be reinitialized"), + |_job_id| panic!("prepared job must not use the empty-block path"), + |job_id| { + assert_eq!(job_id, 4); + true + }, + ); + assert_eq!(success.action, CREATE_JOB_POST); + assert_eq!(success.nextJobID, 5); + assert_eq!(success.jobReady, 0); + + let failure = create_compression_job_with( + projection, + |_job_id, _initialization| panic!("prepared job must not be reinitialized"), + |_job_id| panic!("prepared job must not use the empty-block path"), + |job_id| { + assert_eq!(job_id, 4); + false + }, + ); + assert_eq!(failure.action, CREATE_JOB_POST); + assert_eq!(failure.nextJobID, 4); + assert_eq!(failure.jobReady, 1); + } + #[test] fn flush_state_machine_preserves_offsets_and_completes_job() { let job = [0x60u8, 0x61, 0x62, 0x63, 0x64, 0x65];