feat(compress): project stream continue states into Rust

Pass the shared continue and end projections directly through the buffered
stream ABI so Rust owns stream dispatch into the migrated compression paths.
Remove the redundant C continue/end forwarding callbacks while retaining the
C reset callback for private CCtx session state.

Test Plan:
- ulimit -v 41943040; CARGO_BUILD_JOBS=1 cargo fmt --all -- --check
- ulimit -v 41943040; CARGO_BUILD_JOBS=1 cargo test --lib zstd_compress::tests::compress_stream -- --nocapture
- ulimit -v 41943040; CARGO_BUILD_JOBS=1 cargo clippy --all-targets -- -D warnings
- ulimit -v 41943040; CARGO_BUILD_JOBS=1 cargo test
- ulimit -v 41943040; make -j1
- ulimit -v 41943040; make -j1 -C tests test-zstream ZSTREAM_TESTTIME=-T2s
- ulimit -v 41943040; make -j1 -C tests test-fuzzer FUZZERTEST=-T3s FUZZER_FLAGS=--no-big-tests
This commit is contained in:
2026-07-20 01:06:14 +02:00
parent 0a69df20a3
commit 6da0296af6
2 changed files with 239 additions and 42 deletions
+216 -17
View File
@@ -2914,16 +2914,15 @@ pub unsafe extern "C" fn ZSTD_rust_setParams(
0
}
type CompressStreamBlockFn =
unsafe extern "C" fn(*mut c_void, *mut c_void, usize, *const c_void, usize) -> usize;
type CompressStreamResetFn = unsafe extern "C" fn(*mut c_void) -> usize;
/// Explicit projection of the single-threaded `ZSTD_compressStream_generic`
/// state machine.
///
/// Rust owns buffering, output-drain policy, directive handling, and stream
/// stage transitions. The opaque callback context stays in C; callbacks
/// retain access to `ZSTD_CCtx`, frame construction, and session reset.
/// stage transitions. Continue/end state projections retain the migrated
/// frame construction path; the opaque callback context remains only for
/// session reset.
#[repr(C)]
pub struct ZSTD_rust_compressStreamState {
callback_context: *mut c_void,
@@ -2942,8 +2941,8 @@ pub struct ZSTD_rust_compressStreamState {
out_buff_content_size: *mut usize,
out_buff_flushed_size: *mut usize,
frame_ended: *mut c_uint,
compress_continue: CompressStreamBlockFn,
compress_end: CompressStreamBlockFn,
compress_continue_state: *const ZSTD_rust_compressContinueState,
compress_end_state: *const ZSTD_rust_compressEndState,
reset_session: CompressStreamResetFn,
}
@@ -3007,11 +3006,11 @@ const _: () = {
== 11 * size_of::<usize>() + 2 * size_of::<c_int>() + 2 * size_of::<usize>()
);
assert!(
offset_of!(ZSTD_rust_compressStreamState, compress_continue)
offset_of!(ZSTD_rust_compressStreamState, compress_continue_state)
== 12 * size_of::<usize>() + 2 * size_of::<c_int>() + 2 * size_of::<usize>()
);
assert!(
offset_of!(ZSTD_rust_compressStreamState, compress_end)
offset_of!(ZSTD_rust_compressStreamState, compress_end_state)
== 13 * size_of::<usize>() + 2 * size_of::<c_int>() + 2 * size_of::<usize>()
);
assert!(
@@ -3441,6 +3440,8 @@ unsafe fn compress_stream_generic_body_with(
|| state.out_buff_content_size.is_null()
|| state.out_buff_flushed_size.is_null()
|| state.frame_ended.is_null()
|| state.compress_continue_state.is_null()
|| state.compress_end_state.is_null()
{
return ERROR(ZstdErrorCode::Generic);
}
@@ -3488,8 +3489,8 @@ unsafe fn compress_stream_generic_body_with(
unsafe { input.src.cast::<u8>().add(input.pos).cast() }
};
let c_size = unsafe {
(state.compress_end)(
state.callback_context,
ZSTD_rust_compressEnd(
state.compress_end_state,
dst,
output_remaining,
src,
@@ -3591,20 +3592,22 @@ unsafe fn compress_stream_generic_body_with(
};
let c_size = unsafe {
if last_block {
(state.compress_end)(
state.callback_context,
ZSTD_rust_compressEnd(
state.compress_end_state,
output_dst,
output_size,
source,
input_size,
)
} else {
(state.compress_continue)(
state.callback_context,
ZSTD_rust_compressContinue(
state.compress_continue_state,
output_dst,
output_size,
source,
input_size,
1,
0,
)
}
};
@@ -11762,6 +11765,22 @@ mod tests {
continue_result: usize,
end_result: usize,
reset_result: usize,
window: ZSTD_rust_windowUpdateState,
force_non_contiguous: c_int,
next_to_update: c_uint,
window_projection: ZSTD_rust_compressContinueWindowProjection,
frame_chunk_clamp_state: Option<ZSTD_rust_frameChunkClampState>,
frame_chunk_prepare_state: Option<ZSTD_rust_frameChunkPrepareState>,
continue_frame_chunk_state: Option<ZSTD_rust_frameChunkState>,
end_frame_chunk_state: Option<ZSTD_rust_frameChunkState>,
continue_state: Option<ZSTD_rust_compressContinueState>,
end_continue_state: Option<ZSTD_rust_compressContinueState>,
end_state: Option<ZSTD_rust_compressEndState>,
compression_stage: c_int,
is_first_block: c_int,
end_is_first_block: c_int,
consumed_src_size: u64,
produced_c_size: u64,
}
unsafe fn compress_stream_test_context(
@@ -11776,6 +11795,7 @@ mod tests {
dst_capacity: usize,
_src: *const c_void,
_src_size: usize,
_last_block: c_uint,
) -> usize {
let context = unsafe { compress_stream_test_context(context) };
context.continue_calls += 1;
@@ -11792,6 +11812,7 @@ mod tests {
dst_capacity: usize,
_src: *const c_void,
_src_size: usize,
_last_block: c_uint,
) -> usize {
let context = unsafe { compress_stream_test_context(context) };
context.end_calls += 1;
@@ -11808,6 +11829,23 @@ mod tests {
context.reset_result
}
unsafe extern "C" fn compress_stream_test_frame_prepare_overflow(
_context: *mut c_void,
_src: *const c_void,
_block_size: usize,
) {
}
unsafe extern "C" fn compress_stream_test_frame_prepare_window(
_context: *mut c_void,
_src: *const c_void,
_block_size: usize,
_max_dist: c_uint,
) {
}
unsafe extern "C" fn compress_stream_test_trace(_context: *mut c_void, _extra_c_size: usize) {}
#[allow(clippy::too_many_arguments)]
fn compress_stream_test_state(
context: &mut CompressStreamTestContext,
@@ -11827,8 +11865,169 @@ mod tests {
out_buff_flushed_size: &mut usize,
frame_ended: &mut c_uint,
) -> ZSTD_rust_compressStreamState {
let base = WINDOW_INIT_SENTINEL.as_ptr().cast::<c_void>();
context.window = ZSTD_rust_windowUpdateState {
nextSrc: base,
base,
dictBase: base,
dictLimit: ZSTD_WINDOW_START_INDEX,
lowLimit: ZSTD_WINDOW_START_INDEX,
};
context.force_non_contiguous = 0;
context.next_to_update = ZSTD_WINDOW_START_INDEX;
context.compression_stage = ZSTD_COMPRESSION_STAGE_ONGOING;
context.is_first_block = 1;
context.end_is_first_block = 1;
context.consumed_src_size = 0;
context.produced_c_size = 0;
context.window_projection = ZSTD_rust_compressContinueWindowProjection {
next_src: ptr::addr_of_mut!(context.window.nextSrc),
base: ptr::addr_of_mut!(context.window.base),
dict_base: ptr::addr_of_mut!(context.window.dictBase),
dict_limit: ptr::addr_of_mut!(context.window.dictLimit),
low_limit: ptr::addr_of_mut!(context.window.lowLimit),
force_non_contiguous: ptr::addr_of_mut!(context.force_non_contiguous),
next_to_update: ptr::addr_of_mut!(context.next_to_update),
};
let callback_context = (context as *mut CompressStreamTestContext).cast();
context.frame_chunk_clamp_state = Some(ZSTD_rust_frameChunkClampState {
next_to_update: ptr::addr_of_mut!(context.next_to_update),
low_limit: ptr::addr_of!(context.window.lowLimit),
});
let frame_chunk_clamp_state = context
.frame_chunk_clamp_state
.as_ref()
.map_or(ptr::null(), |state| state as *const _);
context.frame_chunk_prepare_state = Some(ZSTD_rust_frameChunkPrepareState {
callback_context,
max_dist: 64,
correct_overflow: compress_stream_test_frame_prepare_overflow,
check_dict_validity: compress_stream_test_frame_prepare_window,
enforce_max_dist: compress_stream_test_frame_prepare_window,
clamp_state: frame_chunk_clamp_state,
});
let frame_chunk_prepare_state = context
.frame_chunk_prepare_state
.as_ref()
.map_or(ptr::null(), |state| state as *const _);
context.continue_frame_chunk_state = Some(ZSTD_rust_frameChunkState {
callback_context,
tmp_workspace: ptr::null_mut(),
checksum_state: ptr::null_mut(),
is_first_block: ptr::addr_of_mut!(context.is_first_block),
stage: ptr::addr_of_mut!(context.compression_stage),
tmp_wksp_size: 0,
block_size_max,
savings: 0,
pre_block_splitter_level: 1,
strategy: ZSTD_FAST,
use_target_c_block_size: 1,
block_splitter_enabled: 0,
checksum_flag: 0,
ending_stage: ZSTD_COMPRESSION_STAGE_ENDING,
prepare_state: frame_chunk_prepare_state,
compress_target: compress_stream_test_block,
compress_split: compress_stream_test_block,
compress_internal: compress_stream_test_block,
});
let continue_frame_chunk_state = context
.continue_frame_chunk_state
.as_ref()
.map_or(ptr::null(), |state| state as *const _);
context.end_frame_chunk_state = Some(ZSTD_rust_frameChunkState {
callback_context,
tmp_workspace: ptr::null_mut(),
checksum_state: ptr::null_mut(),
is_first_block: ptr::addr_of_mut!(context.end_is_first_block),
stage: ptr::addr_of_mut!(context.compression_stage),
tmp_wksp_size: 0,
block_size_max,
savings: 0,
pre_block_splitter_level: 1,
strategy: ZSTD_FAST,
use_target_c_block_size: 1,
block_splitter_enabled: 0,
checksum_flag: 0,
ending_stage: ZSTD_COMPRESSION_STAGE_ENDING,
prepare_state: frame_chunk_prepare_state,
compress_target: compress_stream_test_end,
compress_split: compress_stream_test_end,
compress_internal: compress_stream_test_end,
});
let end_frame_chunk_state = context
.end_frame_chunk_state
.as_ref()
.map_or(ptr::null(), |state| state as *const _);
context.continue_state = Some(ZSTD_rust_compressContinueState {
callback_context,
window_state: &context.window_projection,
ldm_window_state: ptr::null(),
overflow_state: ptr::null(),
frame_chunk_state: continue_frame_chunk_state,
compress_block: compress_stream_test_block,
stage: ptr::addr_of_mut!(context.compression_stage),
consumed_src_size: ptr::addr_of_mut!(context.consumed_src_size),
produced_c_size: ptr::addr_of_mut!(context.produced_c_size),
pledged_src_size_plus_one: 0,
block_size_max,
check_block_size: 0,
no_dict_id_flag: 0,
checksum_flag: 0,
content_size_flag: 0,
format: 0,
window_log: 20,
dict_id: 0,
ldm_enabled: 0,
});
let continue_state = context
.continue_state
.as_ref()
.map_or(ptr::null(), |state| state as *const _);
context.end_continue_state = Some(ZSTD_rust_compressContinueState {
callback_context,
window_state: &context.window_projection,
ldm_window_state: ptr::null(),
overflow_state: ptr::null(),
frame_chunk_state: end_frame_chunk_state,
compress_block: compress_stream_test_end,
stage: ptr::addr_of_mut!(context.compression_stage),
consumed_src_size: ptr::addr_of_mut!(context.consumed_src_size),
produced_c_size: ptr::addr_of_mut!(context.produced_c_size),
pledged_src_size_plus_one: 0,
block_size_max,
check_block_size: 0,
no_dict_id_flag: 0,
checksum_flag: 0,
content_size_flag: 0,
format: 0,
window_log: 20,
dict_id: 0,
ldm_enabled: 0,
});
let end_continue_state = context
.end_continue_state
.as_ref()
.map_or(ptr::null(), |state| state as *const _);
context.end_state = Some(ZSTD_rust_compressEndState {
callback_context,
compress_continue_state: end_continue_state,
trace: compress_stream_test_trace,
consumed_src_size: ptr::addr_of!(context.consumed_src_size),
pledged_src_size_plus_one: 0,
content_size_flag: 0,
stage: ptr::addr_of_mut!(context.compression_stage),
no_dict_id_flag: 0,
checksum_flag: 0,
format: 0,
window_log: 20,
checksum_state: &COMPRESS_END_TEST_CHECKSUM_STATE,
});
let end_state = context
.end_state
.as_ref()
.map_or(ptr::null(), |state| state as *const _);
ZSTD_rust_compressStreamState {
callback_context: (context as *mut CompressStreamTestContext).cast(),
callback_context,
in_buffer_mode,
out_buffer_mode,
stream_stage: stage,
@@ -11844,8 +12043,8 @@ mod tests {
out_buff_content_size,
out_buff_flushed_size,
frame_ended,
compress_continue: compress_stream_test_block,
compress_end: compress_stream_test_end,
compress_continue_state: continue_state,
compress_end_state: end_state,
reset_session: compress_stream_test_reset,
}
}