feat(compress): move external producer success path into Rust

Move block-level external sequence producer invocation, result
post-processing, source-length validation, and transfer ordering into a
Rust-owned ABI leaf. Retain the private CCtx transfer callback and C-side
fallback/block-compressor selection, preserving producer-error fallback
eligibility and direct invalid-sequence/transfer failures.

Test Plan:
- cargo test --manifest-path rust/Cargo.toml --lib
- cargo clippy --manifest-path rust/Cargo.toml --all-targets -- -D warnings
- make -B -C programs -j1 zstd
- make -C tests -j1 test-zstream ZSTREAM_TESTTIME=-T1s
- focused external_sequence_producer unit tests
This commit is contained in:
2026-07-19 14:41:09 +02:00
parent 724c7c60fa
commit 9a68253710
2 changed files with 545 additions and 58 deletions
+460 -4
View File
@@ -37,8 +37,9 @@ use crate::zstd_compress_stats::{
ZSTD_rust_convertSequencesNoRepcodes, ZSTD_rust_copyBlockSequences,
ZSTD_rust_countSeqStoreLiteralsBytes, ZSTD_rust_countSeqStoreMatchBytes,
ZSTD_rust_deriveSeqStoreChunk, ZSTD_rust_determineBlockSize, ZSTD_rust_entropyCompressSeqStore,
ZSTD_rust_entropyCompressSeqStore_internal, ZSTD_rust_finalizeOffBase,
ZSTD_rust_get1BlockSummary, ZSTD_rust_isRLE, ZSTD_rust_maybeRLE, ZSTD_rust_resetSeqStore,
ZSTD_rust_entropyCompressSeqStore_internal, ZSTD_rust_fastSequenceLengthSum,
ZSTD_rust_finalizeOffBase, ZSTD_rust_get1BlockSummary, ZSTD_rust_isRLE, ZSTD_rust_maybeRLE,
ZSTD_rust_postProcessSequenceProducerResult, ZSTD_rust_resetSeqStore,
ZSTD_rust_seqStore_resolveOffCodes, ZSTD_rust_storeLastLiterals,
ZSTD_rust_transferSequencesNoDelim, ZSTD_rust_transferSequencesWBlockDelim,
ZSTD_rust_validateSeqStore, ZSTD_LLT_LITERAL_LENGTH, ZSTD_LLT_MATCH_LENGTH,
@@ -184,8 +185,8 @@ type BuildSeqStoreSelectFn = unsafe extern "C" fn(
/// Explicit projection for the sequence-store builder.
///
/// Rust owns the threshold/reset/repcode/literal-store orchestration. The
/// callbacks retain the private matchfinder, LDM, and external sequence
/// producer operations in C without passing `ZSTD_CCtx` across the ABI.
/// callbacks retain the private matchfinder, LDM, and external-producer
/// fallback operations in C without passing `ZSTD_CCtx` across the ABI.
#[repr(C)]
pub struct ZSTD_rust_buildSeqStoreState {
seq_store: *mut SeqStore_t,
@@ -227,6 +228,181 @@ const _: () = {
);
};
type ExternalSequenceProducerFn = unsafe extern "C" fn(
*mut c_void,
*mut ZSTD_Sequence,
usize,
*const c_void,
usize,
*const c_void,
usize,
c_int,
usize,
) -> usize;
type ExternalSequenceTransferFn = unsafe extern "C" fn(
*mut c_void,
*mut ZSTD_SequencePosition,
*const ZSTD_Sequence,
usize,
*const c_void,
usize,
) -> usize;
/// Projection for the successful block-level external sequence-producer path.
///
/// Rust owns producer invocation, result post-processing, length validation,
/// and transfer ordering. The transfer callback retains the private CCtx
/// match-state projection; C owns fallback selection after this leaf returns.
#[repr(C)]
pub struct ZSTD_rust_externalSequenceProducerState {
callback_context: *mut c_void,
producer_state: *mut c_void,
producer: Option<ExternalSequenceProducerFn>,
ext_seq_buf: *mut ZSTD_Sequence,
ext_seq_buf_capacity: *const usize,
src: *const c_void,
src_size: *const usize,
compression_level: *const c_int,
window_size: *const usize,
transfer: Option<ExternalSequenceTransferFn>,
external_seq_count: *mut usize,
seq_store_complete: *mut c_int,
allow_fallback: *mut c_int,
}
const _: () = {
assert!(offset_of!(ZSTD_rust_externalSequenceProducerState, callback_context) == 0);
assert!(
offset_of!(ZSTD_rust_externalSequenceProducerState, producer_state) == size_of::<usize>()
);
assert!(
offset_of!(ZSTD_rust_externalSequenceProducerState, producer) == 2 * size_of::<usize>()
);
assert!(
offset_of!(ZSTD_rust_externalSequenceProducerState, ext_seq_buf) == 3 * size_of::<usize>()
);
assert!(
offset_of!(
ZSTD_rust_externalSequenceProducerState,
ext_seq_buf_capacity
) == 4 * size_of::<usize>()
);
assert!(offset_of!(ZSTD_rust_externalSequenceProducerState, src) == 5 * size_of::<usize>());
assert!(
offset_of!(ZSTD_rust_externalSequenceProducerState, src_size) == 6 * size_of::<usize>()
);
assert!(
offset_of!(ZSTD_rust_externalSequenceProducerState, compression_level)
== 7 * size_of::<usize>()
);
assert!(
offset_of!(ZSTD_rust_externalSequenceProducerState, window_size) == size_of::<[usize; 8]>()
);
assert!(
offset_of!(ZSTD_rust_externalSequenceProducerState, transfer) == 9 * size_of::<usize>()
);
assert!(
offset_of!(ZSTD_rust_externalSequenceProducerState, external_seq_count)
== 10 * size_of::<usize>()
);
assert!(
offset_of!(ZSTD_rust_externalSequenceProducerState, seq_store_complete)
== 11 * size_of::<usize>()
);
assert!(
offset_of!(ZSTD_rust_externalSequenceProducerState, allow_fallback)
== 12 * size_of::<usize>()
);
assert!(size_of::<ZSTD_rust_externalSequenceProducerState>() == 13 * size_of::<usize>());
};
/// Try an external sequence producer; C retains only the fallback decision.
#[no_mangle]
pub unsafe extern "C" fn ZSTD_rust_tryExternalSequenceProducer(
state: *const ZSTD_rust_externalSequenceProducerState,
) -> usize {
if state.is_null() {
return ERROR(ZstdErrorCode::Generic);
}
let state = unsafe { &*state };
let Some(producer) = state.producer else {
return ERROR(ZstdErrorCode::Generic);
};
let Some(transfer) = state.transfer else {
return ERROR(ZstdErrorCode::Generic);
};
if state.ext_seq_buf.is_null()
|| state.ext_seq_buf_capacity.is_null()
|| state.src.is_null()
|| state.src_size.is_null()
|| state.compression_level.is_null()
|| state.window_size.is_null()
|| state.external_seq_count.is_null()
|| state.seq_store_complete.is_null()
|| state.allow_fallback.is_null()
{
return ERROR(ZstdErrorCode::Generic);
}
unsafe {
*state.seq_store_complete = 0;
*state.allow_fallback = 0;
}
let ext_seq_buf_capacity = unsafe { *state.ext_seq_buf_capacity };
let src_size = unsafe { *state.src_size };
let nb_external_seqs = unsafe {
producer(
state.producer_state,
state.ext_seq_buf,
ext_seq_buf_capacity,
state.src,
src_size,
ptr::null(),
0,
*state.compression_level,
*state.window_size,
)
};
unsafe { *state.external_seq_count = nb_external_seqs };
let nb_post_processed_seqs = unsafe {
ZSTD_rust_postProcessSequenceProducerResult(
state.ext_seq_buf,
nb_external_seqs,
ext_seq_buf_capacity,
src_size,
)
};
if ERR_isError(nb_post_processed_seqs) {
unsafe { *state.allow_fallback = 1 };
return nb_post_processed_seqs;
}
let seq_len_sum =
unsafe { ZSTD_rust_fastSequenceLengthSum(state.ext_seq_buf, nb_post_processed_seqs) };
if seq_len_sum > src_size {
return ERROR(ZstdErrorCode::ExternalSequencesInvalid);
}
let mut seq_pos = ZSTD_SequencePosition::default();
let transfer_result = unsafe {
transfer(
state.callback_context,
&mut seq_pos,
state.ext_seq_buf,
nb_post_processed_seqs,
state.src,
src_size,
)
};
if ERR_isError(transfer_result) {
return transfer_result;
}
unsafe { *state.seq_store_complete = 1 };
0
}
/// Explicit projection of the state used by `ZSTD_compress_frameChunk`.
///
/// The Rust side owns the per-frame block loop and its savings/dispatch
@@ -6038,6 +6214,286 @@ mod tests {
const ZSTD_BTOPT: c_int = 7;
const ZSTD_BTULTRA2: c_int = 9;
#[derive(Default)]
struct ExternalSequenceProducerProbe {
events: Vec<&'static str>,
sequences: Vec<ZSTD_Sequence>,
producer_result: usize,
producer_dict: *const c_void,
producer_dict_size: usize,
producer_level: c_int,
producer_window_size: usize,
transfer_sequence_count: usize,
transfer_src_size: usize,
transfer_result: usize,
}
unsafe fn external_sequence_producer_probe(
context: *mut c_void,
) -> &'static mut ExternalSequenceProducerProbe {
unsafe { &mut *context.cast::<ExternalSequenceProducerProbe>() }
}
unsafe extern "C" fn external_sequence_producer_test_producer(
context: *mut c_void,
out_seqs: *mut ZSTD_Sequence,
out_seqs_capacity: usize,
_src: *const c_void,
_src_size: usize,
dict: *const c_void,
dict_size: usize,
compression_level: c_int,
window_size: usize,
) -> usize {
let probe = unsafe { external_sequence_producer_probe(context) };
probe.events.push("produce");
probe.producer_dict = dict;
probe.producer_dict_size = dict_size;
probe.producer_level = compression_level;
probe.producer_window_size = window_size;
assert!(probe.sequences.len() <= out_seqs_capacity);
unsafe {
ptr::copy_nonoverlapping(probe.sequences.as_ptr(), out_seqs, probe.sequences.len());
}
probe.producer_result
}
unsafe extern "C" fn external_sequence_producer_test_transfer(
context: *mut c_void,
_seq_pos: *mut ZSTD_SequencePosition,
_in_seqs: *const ZSTD_Sequence,
in_seqs_size: usize,
_src: *const c_void,
block_size: usize,
) -> usize {
let probe = unsafe { external_sequence_producer_probe(context) };
probe.events.push("transfer");
probe.transfer_sequence_count = in_seqs_size;
probe.transfer_src_size = block_size;
probe.transfer_result
}
fn external_sequence_producer_test_state(
probe: &mut ExternalSequenceProducerProbe,
ext_seq_buf: &mut [ZSTD_Sequence],
ext_seq_buf_capacity: &usize,
src: &[u8],
src_size: &usize,
compression_level: &c_int,
window_size: &usize,
external_seq_count: &mut usize,
seq_store_complete: &mut c_int,
allow_fallback: &mut c_int,
) -> ZSTD_rust_externalSequenceProducerState {
let callback_context = (probe as *mut ExternalSequenceProducerProbe).cast();
ZSTD_rust_externalSequenceProducerState {
callback_context,
producer_state: callback_context,
producer: Some(external_sequence_producer_test_producer),
ext_seq_buf: ext_seq_buf.as_mut_ptr(),
ext_seq_buf_capacity,
src: src.as_ptr().cast(),
src_size,
compression_level,
window_size,
transfer: Some(external_sequence_producer_test_transfer),
external_seq_count,
seq_store_complete,
allow_fallback,
}
}
#[test]
fn external_sequence_producer_success_transfers_after_post_processing() {
let source = [1u8, 2, 3, 4, 5];
let mut probe = ExternalSequenceProducerProbe {
sequences: vec![ZSTD_Sequence {
offset: 0,
litLength: source.len() as u32,
matchLength: 0,
rep: 0,
}],
producer_result: 1,
..Default::default()
};
let mut ext_seq_buf = [ZSTD_Sequence {
offset: 0,
litLength: 0,
matchLength: 0,
rep: 0,
}; 2];
let ext_seq_buf_capacity = ext_seq_buf.len();
let src_size = source.len();
let compression_level = 7;
let window_size = 1 << 20;
let mut external_seq_count = 0;
let mut seq_store_complete = 0;
let mut allow_fallback = 0;
let state = external_sequence_producer_test_state(
&mut probe,
&mut ext_seq_buf,
&ext_seq_buf_capacity,
&source,
&src_size,
&compression_level,
&window_size,
&mut external_seq_count,
&mut seq_store_complete,
&mut allow_fallback,
);
let result = unsafe { ZSTD_rust_tryExternalSequenceProducer(&state) };
assert_eq!(result, 0);
assert_eq!(probe.events, ["produce", "transfer"]);
assert!(probe.producer_dict.is_null());
assert_eq!(probe.producer_dict_size, 0);
assert_eq!(probe.producer_level, compression_level);
assert_eq!(probe.producer_window_size, window_size);
assert_eq!(probe.transfer_sequence_count, 1);
assert_eq!(probe.transfer_src_size, source.len());
assert_eq!(external_seq_count, 1);
assert_eq!(seq_store_complete, 1);
assert_eq!(allow_fallback, 0);
}
#[test]
fn external_sequence_producer_errors_enable_fallback_before_transfer() {
let source = [1u8, 2, 3];
let mut probe = ExternalSequenceProducerProbe {
producer_result: ERROR(ZstdErrorCode::SequenceProducerFailed),
..Default::default()
};
let mut ext_seq_buf = [ZSTD_Sequence {
offset: 0,
litLength: 0,
matchLength: 0,
rep: 0,
}; 2];
let ext_seq_buf_capacity = ext_seq_buf.len();
let src_size = source.len();
let compression_level = 3;
let window_size = 1 << 20;
let mut external_seq_count = 0;
let mut seq_store_complete = 0;
let mut allow_fallback = 0;
let state = external_sequence_producer_test_state(
&mut probe,
&mut ext_seq_buf,
&ext_seq_buf_capacity,
&source,
&src_size,
&compression_level,
&window_size,
&mut external_seq_count,
&mut seq_store_complete,
&mut allow_fallback,
);
let result = unsafe { ZSTD_rust_tryExternalSequenceProducer(&state) };
assert_eq!(result, ERROR(ZstdErrorCode::SequenceProducerFailed));
assert_eq!(probe.events, ["produce"]);
assert_eq!(external_seq_count, result);
assert_eq!(seq_store_complete, 0);
assert_eq!(allow_fallback, 1);
}
#[test]
fn external_sequence_producer_invalid_length_does_not_enable_fallback() {
let source = [1u8, 2, 3];
let mut probe = ExternalSequenceProducerProbe {
sequences: vec![ZSTD_Sequence {
offset: 0,
litLength: (source.len() + 1) as u32,
matchLength: 0,
rep: 0,
}],
producer_result: 1,
..Default::default()
};
let mut ext_seq_buf = [ZSTD_Sequence {
offset: 0,
litLength: 0,
matchLength: 0,
rep: 0,
}; 2];
let ext_seq_buf_capacity = ext_seq_buf.len();
let src_size = source.len();
let compression_level = 3;
let window_size = 1 << 20;
let mut external_seq_count = 0;
let mut seq_store_complete = 0;
let mut allow_fallback = 0;
let state = external_sequence_producer_test_state(
&mut probe,
&mut ext_seq_buf,
&ext_seq_buf_capacity,
&source,
&src_size,
&compression_level,
&window_size,
&mut external_seq_count,
&mut seq_store_complete,
&mut allow_fallback,
);
let result = unsafe { ZSTD_rust_tryExternalSequenceProducer(&state) };
assert_eq!(result, ERROR(ZstdErrorCode::ExternalSequencesInvalid));
assert_eq!(probe.events, ["produce"]);
assert_eq!(seq_store_complete, 0);
assert_eq!(allow_fallback, 0);
}
#[test]
fn external_sequence_producer_transfer_errors_do_not_enable_fallback() {
let source = [1u8, 2, 3];
let mut probe = ExternalSequenceProducerProbe {
sequences: vec![ZSTD_Sequence {
offset: 0,
litLength: source.len() as u32,
matchLength: 0,
rep: 0,
}],
producer_result: 1,
transfer_result: ERROR(ZstdErrorCode::ExternalSequencesInvalid),
..Default::default()
};
let mut ext_seq_buf = [ZSTD_Sequence {
offset: 0,
litLength: 0,
matchLength: 0,
rep: 0,
}; 2];
let ext_seq_buf_capacity = ext_seq_buf.len();
let src_size = source.len();
let compression_level = 3;
let window_size = 1 << 20;
let mut external_seq_count = 0;
let mut seq_store_complete = 0;
let mut allow_fallback = 0;
let state = external_sequence_producer_test_state(
&mut probe,
&mut ext_seq_buf,
&ext_seq_buf_capacity,
&source,
&src_size,
&compression_level,
&window_size,
&mut external_seq_count,
&mut seq_store_complete,
&mut allow_fallback,
);
let result = unsafe { ZSTD_rust_tryExternalSequenceProducer(&state) };
assert_eq!(result, ERROR(ZstdErrorCode::ExternalSequencesInvalid));
assert_eq!(probe.events, ["produce", "transfer"]);
assert_eq!(seq_store_complete, 0);
assert_eq!(allow_fallback, 0);
}
#[derive(Default)]
struct Compress2TestContext {
events: Vec<&'static str>,