This is an automated email from the ASF dual-hosted git repository.

hubcio pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git


The following commit(s) were added to refs/heads/master by this push:
     new 60b878923 perf(server-ng): compute batch checksum in a single produce 
pass (#3746)
60b878923 is described below

commit 60b878923b8e4f67241be94e4008c04633de9fb6
Author: Hubert Gruszecki <[email protected]>
AuthorDate: Tue Aug 4 13:15:27 2026 +0200

    perf(server-ng): compute batch checksum in a single produce pass (#3746)
---
 core/partitions/src/iggy_partition.rs    |  58 +-
 core/partitions/src/iggy_partitions.rs   |   8 +-
 core/partitions/src/journal.rs           |  13 +-
 core/server-ng/src/responses.rs          |   2 +-
 core/server_common/src/send_messages2.rs | 987 +++++++++++++++++++++++++++----
 core/simulator/src/client.rs             |   4 +-
 6 files changed, 944 insertions(+), 128 deletions(-)

diff --git a/core/partitions/src/iggy_partition.rs 
b/core/partitions/src/iggy_partition.rs
index aded81684..5c6b14882 100644
--- a/core/partitions/src/iggy_partition.rs
+++ b/core/partitions/src/iggy_partition.rs
@@ -57,7 +57,8 @@ use server_common::{
     MESSAGE_ALIGN, Message, SegmentStorage,
     iobuf::{Frozen, Owned},
     send_messages2::{
-        convert_request_message, decode_prepare_slice, 
stamp_prepare_for_persistence,
+        ChecksumMode, convert_request_message, decode_prepare_slice, 
decode_prepare_slice_trusted,
+        stamp_prepare_for_persistence, verify_received_send_messages,
     },
     sharding::IggyNamespace,
 };
@@ -1149,7 +1150,13 @@ where
             );
 
             let message = if message.header().operation == 
Operation::SendMessages {
-                match convert_request_message(namespace, message) {
+                // Skip the batch-checksum pass: on the partition ingest path
+                // nothing reads it before `stamp_prepare_for_persistence`
+                // recomputes it over the stamped header. An already-canonical
+                // batch (native v2, or the plane's pre-encrypt convert output)
+                // returns early above, so Skip only affects the legacy
+                // transcode, whose output goes straight to project/stamp.
+                match convert_request_message(namespace, message, 
ChecksumMode::Skip) {
                     Ok(message) => message,
                     Err(error) => {
                         emit_partition_diag(
@@ -1532,6 +1539,32 @@ where
                 );
             }
         }
+        // First blob-integrity check on the replicated path. The consensus
+        // layer never validates the body (PrepareHeader integrity fields are
+        // inert zeros) and the batch checksum is recomputed locally at stamp,
+        // so a follower must verify each message's stamp-invariant per-message
+        // checksum before journaling transit bytes. Follower-only: the primary
+        // (and single-node self-replicate) produced these bytes and already
+        // checked the client batch at ingest, so they must not pay this pass.
+        // Fail closed on mismatch - drop without journaling, forwarding, or
+        // acking; the primary retransmits on prepare-timeout.
+        if is_backup
+            && header.operation == Operation::SendMessages
+            && let Err(error) = 
verify_received_send_messages(message.as_slice())
+        {
+            emit_partition_diag(
+                tracing::Level::WARN,
+                &PartitionDiagEvent::new(
+                    self.diag_ctx(),
+                    "rejecting replicated send_messages: per-message checksum 
mismatch",
+                )
+                .with_operation(header.operation)
+                .with_op(header.op)
+                .with_error(error.to_string()),
+            );
+            return;
+        }
+
         // Durability-before-ack: clone for chain-replicate, forward only
         // AFTER apply_replicated_operation persists. Forward-first would
         // give downstream an op whose WAL entry we never wrote, that violates
@@ -1941,11 +1974,15 @@ where
                         }
                         continue;
                     }
-                    // A resident committed SendMessages entry decoded once at 
append
-                    // (the offset index) with its checksum stamped over these 
exact
-                    // bytes, so it must decode again here. Guard the 
invariant for a
-                    // future disk read-back path that could make decode 
fallible.
-                    let Ok(batch) = decode_prepare_slice(entry.as_slice()) 
else {
+                    // Resident committed SendMessages entry: this node 
stamped it
+                    // in `append_messages` (recomputing the batch checksum 
over these
+                    // exact bytes), so a validating re-decode would only 
re-hash ~1
+                    // MiB to confirm our own write. Trust the structural 
decode; the
+                    // batch-checksum recompute belongs at network ingress 
(repair
+                    // validation + the follower receive gate), not on 
locally-stamped
+                    // bytes. Guard the invariant for a future disk read-back 
path that
+                    // could make decode fallible.
+                    let Ok(batch) = 
decode_prepare_slice_trusted(entry.as_slice()) else {
                         tracing::error!(
                             target: "iggy.partitions.diag",
                             namespace_raw = self.namespace().inner(),
@@ -2345,8 +2382,11 @@ where
         let Some(entry) = self.log.journal().inner.entry(prepare_header).await 
else {
             return Err(IggyError::InvalidCommand);
         };
-        let batch =
-            decode_prepare_slice(entry.as_slice()).map_err(|_| 
IggyError::InvalidCommand)?;
+        // Trusted (no batch-hash): the entry was read back from this replica's
+        // own journal, where it was stamped/validated at append; only header
+        // stats are needed, so re-hashing the ~1 MiB blob is redundant.
+        let batch = decode_prepare_slice_trusted(entry.as_slice())
+            .map_err(|_| IggyError::InvalidCommand)?;
         let message_count = batch.message_count();
         if message_count == 0 {
             return Ok(None);
diff --git a/core/partitions/src/iggy_partitions.rs 
b/core/partitions/src/iggy_partitions.rs
index c16b14843..c73e8a93c 100644
--- a/core/partitions/src/iggy_partitions.rs
+++ b/core/partitions/src/iggy_partitions.rs
@@ -26,7 +26,7 @@ use iggy_binary_protocol::{
     Command2, ConsensusHeader, Operation, PrepareHeader, PrepareOkHeader, 
RequestHeader,
 };
 use message_bus::MessageBus;
-use server_common::send_messages2::{convert_request_message, 
encrypt_batch_request};
+use server_common::send_messages2::{ChecksumMode, convert_request_message, 
encrypt_batch_request};
 use server_common::sharding::{IggyNamespace, LocalIdx, ShardId};
 #[cfg(debug_assertions)]
 use std::cell::Cell;
@@ -512,7 +512,11 @@ where
         let message = if message.header().operation == Operation::SendMessages
             && let Some(encryptor) = &self.config().encryptor
         {
-            let canonical = convert_request_message(namespace, message)
+            // Compute the batch checksum: this canonical output is validated 
by
+            // `encrypt_batch_request`'s decode before re-encryption, and the
+            // re-encrypted batch (checksum kept by `encrypt_batch_request`) 
then
+            // re-enters `convert` as the canonical-vs-legacy discriminator.
+            let canonical = convert_request_message(namespace, message, 
ChecksumMode::Compute)
                 .and_then(|message| encrypt_batch_request(message, encryptor));
             match canonical {
                 Ok(message) => message,
diff --git a/core/partitions/src/journal.rs b/core/partitions/src/journal.rs
index da648fcb6..0b2eaf8d0 100644
--- a/core/partitions/src/journal.rs
+++ b/core/partitions/src/journal.rs
@@ -19,7 +19,7 @@ use iggy_binary_protocol::{Operation, PrepareHeader};
 use journal::{Journal, Storage};
 use server_common::{
     iobuf::{Frozen, Owned},
-    send_messages2::{COMMAND_HEADER_SIZE, SendMessages2Ref, 
decode_prepare_slice},
+    send_messages2::{COMMAND_HEADER_SIZE, SendMessages2Ref, 
decode_prepare_slice_trusted},
 };
 use std::io;
 use std::{
@@ -538,8 +538,12 @@ impl PartitionJournal<PartitionJournalMemStorage> {
         // One decode feeds both the offset/timestamp index (keyed on
         // `origin_timestamp`) and the surfaced accounting meta 
(`base_timestamp`,
         // size, count); the two timestamps are distinct fields, do not 
conflate.
+        // Trusted (no batch-hash): every entry reaching append was just 
stamped
+        // by `stamp_prepare_for_persistence` (its checksum recomputed over 
this
+        // exact blob) or re-appended from an already-validated resident entry,
+        // so re-hashing the ~1 MiB blob here only to read the header is waste.
         let (index_offset_timestamp, meta) = if header.operation == 
Operation::SendMessages {
-            match decode_prepare_slice(entry.as_slice()) {
+            match decode_prepare_slice_trusted(entry.as_slice()) {
                 Ok(batch) if batch.message_count() != 0 => {
                     let message_count = batch.message_count();
                     let meta = RetainedBatchMeta {
@@ -885,7 +889,10 @@ fn try_push_resident_entry(
     if header.operation != Operation::SendMessages {
         return;
     }
-    let Ok(batch) = decode_prepare_slice(prepare.as_slice()) else {
+    // Resident entries were locally stamped in `append_messages` or validated
+    // at repair ingress, so a validating re-decode would only re-hash our own
+    // write. See the invariant note at the committed-prefix flush walk.
+    let Ok(batch) = decode_prepare_slice_trusted(prepare.as_slice()) else {
         return;
     };
     let Some(selection) = select_batch_slice(&batch, query, *matched_messages) 
else {
diff --git a/core/server-ng/src/responses.rs b/core/server-ng/src/responses.rs
index bfa459c68..6cc585230 100644
--- a/core/server-ng/src/responses.rs
+++ b/core/server-ng/src/responses.rs
@@ -1437,7 +1437,7 @@ where
 /// Size of the in-storage (`IggyMessage2`) per-message header inside a
 /// `SendMessages2` batch blob: `checksum`(8) + `id`(16) + `offset_delta`(4)
 /// + `timestamp_delta`(4) + `user_headers_length`(4) + `payload_length`(4)
-/// + reserved(8). See `server_common::send_messages2::from_legacy_request`.
+/// + reserved(8). See 
`server_common::send_messages2::SendMessages2Owned::from_messages`.
 const STORED_MESSAGE_HEADER_SIZE: usize = 48;
 
 /// Build the `PolledMessages` reply body from the owning shard's poll
diff --git a/core/server_common/src/send_messages2.rs 
b/core/server_common/src/send_messages2.rs
index a64cd0664..f0d9a7f88 100644
--- a/core/server_common/src/send_messages2.rs
+++ b/core/server_common/src/send_messages2.rs
@@ -173,84 +173,12 @@ impl SendMessages2Owned {
             header[32..36].copy_from_slice(&user_headers_length.to_le_bytes());
             header[36..40].copy_from_slice(&payload_length.to_le_bytes());
 
-            let checksum = calculate_checksum_parts(&header[8..], 
&message.payload, user_headers);
-            header[0..8].copy_from_slice(&checksum.to_le_bytes());
-
+            let msg_start = blob.len();
             blob.extend_from_slice(&header);
             blob.extend_from_slice(&message.payload);
             blob.extend_from_slice(user_headers);
-        }
-
-        let blob = blob.freeze();
-        let mut header = SendMessages2Header::new(
-            namespace.partition_id() as u64,
-            origin_timestamp,
-            u64::try_from(COMMAND_HEADER_SIZE + blob.len())
-                .map_err(|_| IggyError::InvalidCommand)?,
-            message_count,
-        );
-        header.batch_checksum = calculate_batch_checksum(&header, &blob);
-
-        Ok(Self { header, blob })
-    }
-
-    pub fn from_legacy_request(namespace: IggyNamespace, body: &[u8]) -> 
Result<Self, IggyError> {
-        let (message_count, messages) = legacy_messages_slice(body)?;
-        let mut parsed = Vec::with_capacity(message_count as usize);
-        let mut origin_timestamp = u64::MAX;
-        let mut cursor = 0usize;
-
-        while cursor < messages.len() && parsed.len() < message_count as usize 
{
-            let legacy = LegacyMessageRef::decode(&messages[cursor..])?;
-            origin_timestamp = origin_timestamp.min(legacy.origin_timestamp);
-            cursor += legacy.total_size;
-            parsed.push(legacy);
-        }
-
-        if parsed.len() != message_count as usize || cursor != messages.len() {
-            return Err(IggyError::InvalidCommand);
-        }
-
-        if origin_timestamp == u64::MAX {
-            origin_timestamp = 0;
-        }
-
-        let mut blob = BytesMut::with_capacity(messages.len());
-        for (index, legacy) in parsed.iter().enumerate() {
-            let id = if legacy.id == 0 {
-                random_id::get_uuid()
-            } else {
-                legacy.id
-            };
-            let offset_delta = u32::try_from(index).map_err(|_| 
IggyError::InvalidCommand)?;
-            let timestamp_delta = legacy
-                .origin_timestamp
-                .checked_sub(origin_timestamp)
-                .ok_or(IggyError::InvalidCommand)?;
-            if timestamp_delta > MAX_TIMESTAMP_DELTA_MICROS {
-                return 
Err(IggyError::InvalidMessageTimestampDelta(timestamp_delta));
-            }
-            let timestamp_delta =
-                u32::try_from(timestamp_delta).map_err(|_| 
IggyError::InvalidCommand)?;
-            let user_headers_length =
-                u32::try_from(legacy.user_headers.len()).map_err(|_| 
IggyError::InvalidCommand)?;
-            let payload_length =
-                u32::try_from(legacy.payload.len()).map_err(|_| 
IggyError::InvalidCommand)?;
-
-            let mut header = [0u8; MESSAGE_HEADER_SIZE];
-            header[8..24].copy_from_slice(&id.to_le_bytes());
-            header[24..28].copy_from_slice(&offset_delta.to_le_bytes());
-            header[28..32].copy_from_slice(&timestamp_delta.to_le_bytes());
-            header[32..36].copy_from_slice(&user_headers_length.to_le_bytes());
-            header[36..40].copy_from_slice(&payload_length.to_le_bytes());
-
-            let checksum =
-                calculate_checksum_parts(&header[8..], legacy.payload, 
legacy.user_headers);
-            header[0..8].copy_from_slice(&checksum.to_le_bytes());
-
-            blob.extend_from_slice(&header);
-            blob.extend_from_slice(legacy.payload);
-            blob.extend_from_slice(legacy.user_headers);
+            let checksum = XxHash3_64::oneshot(&blob[msg_start + 8..]);
+            blob[msg_start..msg_start + 
8].copy_from_slice(&checksum.to_le_bytes());
         }
 
         let blob = blob.freeze();
@@ -573,12 +501,12 @@ pub fn encrypt_batch_request(
         
header[28..32].copy_from_slice(&view.header.timestamp_delta.to_le_bytes());
         header[32..36].copy_from_slice(&user_headers_length.to_le_bytes());
         header[36..40].copy_from_slice(&payload_length.to_le_bytes());
-        let checksum = calculate_checksum_parts(&header[8..], 
&encrypted_payload, user_headers);
-        header[0..8].copy_from_slice(&checksum.to_le_bytes());
-
+        let msg_start = blob.len();
         blob.extend_from_slice(&header);
         blob.extend_from_slice(&encrypted_payload);
         blob.extend_from_slice(user_headers);
+        let checksum = XxHash3_64::oneshot(&blob[msg_start + 8..]);
+        blob[msg_start..msg_start + 
8].copy_from_slice(&checksum.to_le_bytes());
     }
 
     let blob = blob.freeze();
@@ -590,22 +518,165 @@ pub fn encrypt_batch_request(
     SendMessages2Owned { header, blob }.encode_request(request_header)
 }
 
+/// Whether the legacy transcode stamps a batch checksum onto its output.
+///
+/// The recompute is an `XxHash3` batch-checksum pass, needed only when a 
reader
+/// validates the transcoded batch before [`stamp_prepare_for_persistence`]
+/// recomputes it: the encrypt ingest path re-decodes the canonicalized batch
+/// (`encrypt_batch_request`'s validating decode, then the second `convert` its
+/// output re-enters as the canonical-vs-legacy discriminator). The partition
+/// ingest path has no such reader, so it skips the pass and the checksum stays
+/// zero until stamp.
+#[derive(Debug, Clone, Copy, PartialEq, Eq)]
+pub enum ChecksumMode {
+    /// Compute the batch checksum for the transcoded batch.
+    Compute,
+    /// Leave the batch checksum zero; `stamp_prepare_for_persistence` fills 
it.
+    Skip,
+}
+
 pub fn convert_request_message(
     namespace: IggyNamespace,
     message: Message<RequestHeader>,
+    checksum: ChecksumMode,
 ) -> Result<Message<RequestHeader>, IggyError> {
     let request_header = *message.header();
     let total_size = request_header.size as usize;
     let body = 
&message.as_slice()[std::mem::size_of::<RequestHeader>()..total_size];
-    if decode_batch_slice(body).is_ok() {
-        return Ok(message);
+    // A canonical body enters the pipeline verbatim, so it must end exactly at
+    // `batch_length`: `size` and `batch_length` are independent 
client-supplied
+    // fields and `decode_batch_slice` only lower-bounds the frame. A suffix 
past
+    // `batch_length` is covered by no checksum, still rides the buffer to 
disk,
+    // and desyncs the segment walk that advances by `batch_length`.
+    match decode_batch_slice(body).map(|batch| batch.header.total_size()) {
+        Ok(batch_length) if body.len() == batch_length => Ok(message),
+        Ok(_) => Err(IggyError::InvalidCommand),
+        Err(_) => transcode_legacy_request(namespace, body, request_header, 
checksum),
     }
-    SendMessages2Owned::from_legacy_request(namespace, 
body)?.encode_request(request_header)
 }
 
-/// Decode one batch slice (`[256B command header][blob]`), validating the
-/// batch checksum. The persisted segment-file record and the request slice
-/// share this layout, so both decode through here.
+/// Transcode a legacy `SendMessages` request body directly into the canonical
+/// `[RequestHeader][256B SendMessages2Header][blob]` form, writing each 
message
+/// record straight into the final aligned buffer.
+///
+/// Fused replacement for the `from_legacy_request(..).encode_request(..)`
+/// two-step: a size walk over the legacy input sizes the single output
+/// allocation, then a write walk lays down each canonical record in place. 
This
+/// drops the intermediate blob allocation and the full-blob copy the two-step
+/// paid. Output bytes are identical to that path.
+///
+/// `checksum` selects whether the output carries a batch checksum (see
+/// [`ChecksumMode`]); [`ChecksumMode::Skip`] leaves it zero for the partition
+/// ingest path, where stamp recomputes it.
+fn transcode_legacy_request(
+    namespace: IggyNamespace,
+    body: &[u8],
+    mut request_header: RequestHeader,
+    checksum: ChecksumMode,
+) -> Result<Message<RequestHeader>, IggyError> {
+    let (message_count, messages) = legacy_messages_slice(body)?;
+    let mut parsed = Vec::with_capacity(message_count as usize);
+    let mut origin_timestamp = u64::MAX;
+    let mut cursor = 0usize;
+    let mut blob_len = 0usize;
+
+    while cursor < messages.len() && parsed.len() < message_count as usize {
+        let legacy = LegacyMessageRef::decode(&messages[cursor..])?;
+        origin_timestamp = origin_timestamp.min(legacy.origin_timestamp);
+        cursor += legacy.total_size;
+        blob_len = blob_len
+            .checked_add(MESSAGE_HEADER_SIZE + legacy.payload.len() + 
legacy.user_headers.len())
+            .ok_or(IggyError::InvalidCommand)?;
+        parsed.push(legacy);
+    }
+
+    if parsed.len() != message_count as usize || cursor != messages.len() {
+        return Err(IggyError::InvalidCommand);
+    }
+
+    if origin_timestamp == u64::MAX {
+        origin_timestamp = 0;
+    }
+
+    let header_size = std::mem::size_of::<RequestHeader>();
+    let batch_length = COMMAND_HEADER_SIZE
+        .checked_add(blob_len)
+        .ok_or(IggyError::InvalidCommand)?;
+    let total_size = header_size
+        .checked_add(batch_length)
+        .ok_or(IggyError::InvalidCommand)?;
+    request_header.size = u32::try_from(total_size).map_err(|_| 
IggyError::InvalidCommand)?;
+
+    let mut buffer = Owned::<MESSAGE_ALIGN>::zeroed(total_size);
+    let bytes = buffer.as_mut_slice();
+    bytes[0..header_size].copy_from_slice(bytemuck::bytes_of(&request_header));
+
+    let mut write = PREPARE_SPLIT_POINT;
+    for (index, legacy) in parsed.iter().enumerate() {
+        let id = if legacy.id == 0 {
+            random_id::get_uuid()
+        } else {
+            legacy.id
+        };
+        let offset_delta = u32::try_from(index).map_err(|_| 
IggyError::InvalidCommand)?;
+        let timestamp_delta = legacy
+            .origin_timestamp
+            .checked_sub(origin_timestamp)
+            .ok_or(IggyError::InvalidCommand)?;
+        if timestamp_delta > MAX_TIMESTAMP_DELTA_MICROS {
+            return 
Err(IggyError::InvalidMessageTimestampDelta(timestamp_delta));
+        }
+        let timestamp_delta =
+            u32::try_from(timestamp_delta).map_err(|_| 
IggyError::InvalidCommand)?;
+        let user_headers_length =
+            u32::try_from(legacy.user_headers.len()).map_err(|_| 
IggyError::InvalidCommand)?;
+        let payload_length =
+            u32::try_from(legacy.payload.len()).map_err(|_| 
IggyError::InvalidCommand)?;
+
+        let mut header = [0u8; MESSAGE_HEADER_SIZE];
+        header[8..24].copy_from_slice(&id.to_le_bytes());
+        header[24..28].copy_from_slice(&offset_delta.to_le_bytes());
+        header[28..32].copy_from_slice(&timestamp_delta.to_le_bytes());
+        header[32..36].copy_from_slice(&user_headers_length.to_le_bytes());
+        header[36..40].copy_from_slice(&payload_length.to_le_bytes());
+        let msg_start = write;
+        bytes[write..write + MESSAGE_HEADER_SIZE].copy_from_slice(&header);
+        write += MESSAGE_HEADER_SIZE;
+        bytes[write..write + 
legacy.payload.len()].copy_from_slice(legacy.payload);
+        write += legacy.payload.len();
+        bytes[write..write + 
legacy.user_headers.len()].copy_from_slice(legacy.user_headers);
+        write += legacy.user_headers.len();
+        // The cover is [msg_start + 8 .. write], including the 8 reserved zero
+        // header bytes. This relies on the stack header being 
zero-initialized,
+        // not on the output buffer being pre-zeroed.
+        let checksum = XxHash3_64::oneshot(&bytes[msg_start + 8..write]);
+        bytes[msg_start..msg_start + 
8].copy_from_slice(&checksum.to_le_bytes());
+    }
+
+    let mut command = SendMessages2Header::new(
+        namespace.partition_id() as u64,
+        origin_timestamp,
+        batch_length as u64,
+        message_count,
+    );
+    if checksum == ChecksumMode::Compute {
+        command.batch_checksum = calculate_batch_checksum(
+            &command,
+            &bytes[PREPARE_SPLIT_POINT..PREPARE_SPLIT_POINT + blob_len],
+        );
+    }
+    command.encode_into(&mut bytes[header_size..header_size + 
COMMAND_HEADER_SIZE]);
+
+    Message::try_from(buffer).map_err(|_| IggyError::InvalidCommand)
+}
+
+/// Decode one batch slice (`[256B command header][blob]`), validating the 
batch
+/// checksum and every per-message checksum. The persisted segment-file record
+/// and the request slice share this layout, so both decode through here.
+///
+/// `body` may extend past the batch: the poll disk walk hands in the rest of 
the
+/// chunk and steps by `batch_length`. Callers whose buffer is meant to BE the
+/// batch must reject the surplus themselves - see [`convert_request_message`].
 pub fn decode_batch_slice(body: &[u8]) -> Result<SendMessages2Ref<'_>, 
IggyError> {
     if body.len() < COMMAND_HEADER_SIZE {
         return Err(IggyError::InvalidCommand);
@@ -618,7 +689,8 @@ pub fn decode_batch_slice(body: &[u8]) -> 
Result<SendMessages2Ref<'_>, IggyError
     }
 
     let blob = &body[COMMAND_HEADER_SIZE..COMMAND_HEADER_SIZE + blob_len];
-    let expected_checksum = calculate_batch_checksum(&header, blob);
+    let batch = SendMessages2Ref { header, blob };
+    let expected_checksum = verify_and_recompute_batch_checksum(&batch)?;
     if header.batch_checksum != expected_checksum {
         return Err(IggyError::InvalidBatchChecksum(
             header.batch_checksum,
@@ -627,10 +699,11 @@ pub fn decode_batch_slice(body: &[u8]) -> 
Result<SendMessages2Ref<'_>, IggyError
         ));
     }
 
-    Ok(SendMessages2Ref { header, blob })
+    Ok(batch)
 }
 
-/// Decode a `Prepare` message from a slice of bytes.
+/// Decode a `Prepare` message from a slice of bytes, validating the batch
+/// checksum and every per-message checksum.
 ///
 /// `bytes` must be 16-byte aligned (`PrepareHeader` has `u128` fields). Source
 /// from `Frozen<MESSAGE_ALIGN>` / `Owned<MESSAGE_ALIGN>` / `Message<H>`.
@@ -638,9 +711,43 @@ pub fn decode_batch_slice(body: &[u8]) -> 
Result<SendMessages2Ref<'_>, IggyError
 ///
 /// # Errors
 ///
-/// `IggyError::InvalidCommand` on: short buffer, bad bit pattern, `size`
-/// outside `[header_size, bytes.len()]`, short/checksum-mismatched body.
+/// `IggyError::InvalidCommand` on a short buffer, bad bit pattern, `size`
+/// outside `[header_size, bytes.len()]`, a `size` that does not describe the
+/// batch exactly, or frames that do not tile the batch;
+/// `InvalidBatchChecksum` / `InvalidMessageChecksum` on an integrity mismatch.
 pub fn decode_prepare_slice(bytes: &[u8]) -> Result<SendMessages2Ref<'_>, 
IggyError> {
+    decode_prepare_slice_inner(bytes, true)
+}
+
+/// Like [`decode_prepare_slice`] but skips the per-message checksum
+/// verification and batch-checksum recompute, extracting only the header meta.
+/// Every cheap structural check (length, 16-byte alignment, `size` bounds, and
+/// `size` describing the batch exactly) is still enforced.
+///
+/// INVARIANT: `bytes` MUST be node-local self-stamped -
+/// [`stamp_prepare_for_persistence`] recomputed the batch checksum over the
+/// exact blob on THIS node - or already integrity-checked at their network
+/// ingress. There is no consensus-layer blob validation: the `PrepareHeader`
+/// integrity fields are inert zeros. A replicated `SendMessages` prepare is
+/// gated per-message on receipt by [`verify_received_send_messages`], and a
+/// repaired prepare is validated via [`decode_prepare_slice`]; both run BEFORE
+/// the bytes reach any trusted decode. NEVER call this on unvalidated network
+/// bytes - it would let a corrupted blob pass undetected. The full-body
+/// per-message checksum pass dominates produce-path CPU, so trusted call sites
+/// that only read header meta skip it.
+///
+/// # Errors
+///
+/// Same structural errors as [`decode_prepare_slice`], minus
+/// `InvalidBatchChecksum` and `InvalidMessageChecksum`.
+pub fn decode_prepare_slice_trusted(bytes: &[u8]) -> 
Result<SendMessages2Ref<'_>, IggyError> {
+    decode_prepare_slice_inner(bytes, false)
+}
+
+fn decode_prepare_slice_inner(
+    bytes: &[u8],
+    validate_checksum: bool,
+) -> Result<SendMessages2Ref<'_>, IggyError> {
     let header_size = std::mem::size_of::<PrepareHeader>();
     if bytes.len() < header_size {
         return Err(IggyError::InvalidCommand);
@@ -670,23 +777,29 @@ pub fn decode_prepare_slice(bytes: &[u8]) -> 
Result<SendMessages2Ref<'_>, IggyEr
     }
 
     let header = SendMessages2Header::decode(&body[..COMMAND_HEADER_SIZE])?;
-    let blob = &body[COMMAND_HEADER_SIZE..];
     let blob_len = header.blob_len()?;
-    if body.len() < header.total_size() {
+    // Exact, not a lower bound: a prepare frame IS one batch, so bytes past
+    // `batch_length` belong to nobody - no checksum covers them, yet the flush
+    // writes them, desyncing the segment walk. Readers walking a multi-batch
+    // chunk use `decode_batch_slice`, which bounds the blob by design.
+    if body.len() != header.total_size() {
         return Err(IggyError::InvalidCommand);
     }
 
-    let blob = &blob[..blob_len];
-    let expected_checksum = calculate_batch_checksum(&header, blob);
-    if header.batch_checksum != expected_checksum {
-        return Err(IggyError::InvalidBatchChecksum(
-            header.batch_checksum,
-            expected_checksum,
-            header.base_offset,
-        ));
+    let blob = &body[COMMAND_HEADER_SIZE..COMMAND_HEADER_SIZE + blob_len];
+    let batch = SendMessages2Ref { header, blob };
+    if validate_checksum {
+        let expected_checksum = verify_and_recompute_batch_checksum(&batch)?;
+        if header.batch_checksum != expected_checksum {
+            return Err(IggyError::InvalidBatchChecksum(
+                header.batch_checksum,
+                expected_checksum,
+                header.base_offset,
+            ));
+        }
     }
 
-    Ok(SendMessages2Ref { header, blob })
+    Ok(batch)
 }
 
 pub fn stamp_prepare_for_persistence(
@@ -711,6 +824,35 @@ pub fn stamp_prepare_for_persistence(
     Ok((message, command, command.message_count))
 }
 
+/// Verify every per-message checksum in a received `SendMessages` prepare.
+///
+/// The FIRST blob-integrity check on the replicated path: the `PrepareHeader`
+/// integrity fields are inert zeros and the batch checksum is recomputed
+/// locally at stamp, so transit corruption of a message body would otherwise
+/// reach apply undetected. Backups call this before journaling a replicated
+/// prepare; on a mismatch the caller fails closed (drop, no `PrepareOk`) and 
the
+/// primary retransmits on prepare-timeout.
+///
+/// The stored `batch_checksum` is not consulted (a received prepare is
+/// pre-stamp, `base_offset` / `base_timestamp` zero): integrity rests on the
+/// per-message checksums, recomputed over the stamp-invariant cover
+/// (`header[8..48] || payload || user_headers`), which excludes the 256B 
command
+/// header, so it holds whether or not this node has stamped yet. Shares the
+/// frame walk with the validating decoders via
+/// [`verify_and_recompute_batch_checksum`], discarding its recomputed batch
+/// value.
+///
+/// # Errors
+///
+/// [`IggyError::InvalidCommand`] if the records do not tile `message_count`
+/// exactly (a length-field corruption desyncs the walk);
+/// [`IggyError::InvalidMessageChecksum`] on the first per-message mismatch.
+pub fn verify_received_send_messages(bytes: &[u8]) -> Result<(), IggyError> {
+    let batch = decode_prepare_slice_trusted(bytes)?;
+    verify_and_recompute_batch_checksum(&batch)?;
+    Ok(())
+}
+
 fn legacy_messages_slice(body: &[u8]) -> Result<(u32, &[u8]), IggyError> {
     if body.len() < 4 {
         return Err(IggyError::InvalidCommand);
@@ -778,26 +920,84 @@ impl<'a> LegacyMessageRef<'a> {
     }
 }
 
-// Hash in storage order: header tail, payload, user headers (the message
-// sections follow the legacy wire layout).
-fn calculate_checksum_parts(header_tail: &[u8], payload: &[u8], user_headers: 
&[u8]) -> u64 {
+/// Batch checksum v2: streaming `XxHash3_64` over the six batch header meta
+/// fields followed by each message's stored 8-byte checksum field in message
+/// order - NOT the message bodies.
+///
+/// Bodies are bound only transitively: each per-message checksum already 
covers
+/// `header[8..48] || payload || user_headers`, so hashing the checksum fields
+/// binds every body byte IFF a reader also re-verifies the per-message
+/// checksums. Stamp (produce) hashes `N * 8` bytes instead of the whole blob;
+/// validating decoders pay the one body pass as the per-message verify in
+/// [`verify_and_recompute_batch_checksum`], which hashes the checksum-field
+/// bytes in the same order so its recompute matches a compute here.
+///
+/// Assumes a well-formed blob whose frames tile exactly; every compute site
+/// builds the blob and satisfies this.
+fn calculate_batch_checksum(header: &SendMessages2Header, blob: &[u8]) -> u64 {
     let mut hasher = XxHash3_64::new();
-    hasher.write(header_tail);
-    hasher.write(payload);
-    hasher.write(user_headers);
+    write_batch_header_fields(&mut hasher, header);
+    let batch = SendMessages2Ref {
+        header: *header,
+        blob,
+    };
+    for framed in batch.iter_with_offsets() {
+        hasher.write(&blob[framed.start..framed.start + 8]);
+    }
     hasher.finish()
 }
 
-fn calculate_batch_checksum(header: &SendMessages2Header, blob: &[u8]) -> u64 {
-    let mut hasher = XxHash3_64::new();
+fn write_batch_header_fields(hasher: &mut XxHash3_64, header: 
&SendMessages2Header) {
     hasher.write(&header.partition_id.to_le_bytes());
     hasher.write(&header.base_offset.to_le_bytes());
     hasher.write(&header.base_timestamp.to_le_bytes());
     hasher.write(&header.origin_timestamp.to_le_bytes());
     hasher.write(&header.batch_length.to_le_bytes());
     hasher.write(&header.message_count.to_le_bytes());
-    hasher.write(blob);
-    hasher.finish()
+}
+
+/// Verify every per-message checksum in `batch` and return the recomputed v2
+/// batch checksum (see [`calculate_batch_checksum`]) from a single frame walk.
+///
+/// The per-message pass is the equal-integrity half of v2: the batch value
+/// binds bodies only through the checksum fields, so a validating decode must
+/// re-verify each message here or body corruption that leaves the checksum
+/// field intact would pass. This is the one full-body pass a validating decode
+/// pays; the caller then compares the returned value against the stored
+/// `batch_checksum`.
+///
+/// # Errors
+///
+/// [`IggyError::InvalidMessageChecksum`] on the first per-message mismatch;
+/// [`IggyError::InvalidCommand`] if the frames do not tile `message_count`
+/// exactly.
+fn verify_and_recompute_batch_checksum(batch: &SendMessages2Ref<'_>) -> 
Result<u64, IggyError> {
+    let blob = batch.blob();
+    let mut hasher = XxHash3_64::new();
+    write_batch_header_fields(&mut hasher, &batch.header);
+    let mut verified = 0u32;
+    let mut covered = 0usize;
+    for framed in batch.iter_with_offsets() {
+        // Cover (`header[8..48] || payload || user_headers`) hashed raw from 
the
+        // blob, byte-exact with the encoder's, so a flipped body byte fails 
even
+        // when the stored checksum field is left intact.
+        let stored = framed.message.header.checksum;
+        let expected = XxHash3_64::oneshot(&blob[framed.start + 
8..framed.end]);
+        if expected != stored {
+            return Err(IggyError::InvalidMessageChecksum(
+                stored,
+                expected,
+                batch.header.base_offset + 
u64::from(framed.message.header.offset_delta),
+            ));
+        }
+        hasher.write(&blob[framed.start..framed.start + 8]);
+        verified += 1;
+        covered = framed.end;
+    }
+    if verified != batch.message_count() || covered != blob.len() {
+        return Err(IggyError::InvalidCommand);
+    }
+    Ok(hasher.finish())
 }
 
 fn read_u32(bytes: &[u8], offset: usize) -> Result<u32, IggyError> {
@@ -827,7 +1027,8 @@ fn read_u128(bytes: &[u8], offset: usize) -> Result<u128, 
IggyError> {
 #[cfg(test)]
 mod tests {
     use super::*;
-    use iggy_binary_protocol::Command2;
+    use iggy_binary_protocol::{Command2, Operation};
+    use iggy_common::Aes256GcmEncryptor;
 
     fn aligned_prepare_bytes(size: u32) -> Owned<MESSAGE_ALIGN> {
         let mut owned = 
Owned::<MESSAGE_ALIGN>::zeroed(std::mem::size_of::<PrepareHeader>());
@@ -839,6 +1040,90 @@ mod tests {
         owned
     }
 
+    /// Assemble an already-stamped batch into a `Prepare`:
+    /// `[PrepareHeader][256B batch header][blob]`, copying `owned`'s header 
and
+    /// blob verbatim. Shared by every real-batch fixture.
+    fn prepare_from_owned(owned: &SendMessages2Owned) -> Owned<MESSAGE_ALIGN> {
+        let header_size = std::mem::size_of::<PrepareHeader>();
+        let total = header_size + owned.header.total_size();
+        let mut buffer = Owned::<MESSAGE_ALIGN>::zeroed(total);
+        {
+            let prepare: &mut PrepareHeader =
+                bytemuck::checked::try_from_bytes_mut(&mut 
buffer.as_mut_slice()[..header_size])
+                    .expect("zeroed bytes form a valid PrepareHeader");
+            prepare.command = Command2::Prepare;
+            prepare.size = u32::try_from(total).expect("prepare size fits 
u32");
+        }
+        let bytes = buffer.as_mut_slice();
+        owned
+            .header
+            .encode_into(&mut bytes[header_size..header_size + 
COMMAND_HEADER_SIZE]);
+        bytes[PREPARE_SPLIT_POINT..PREPARE_SPLIT_POINT + owned.blob.len()]
+            .copy_from_slice(&owned.blob);
+        buffer
+    }
+
+    /// A checksum-consistent STAMPED `Prepare` carrying real per-message 
records,
+    /// stamped at a non-zero `base_offset` / `base_timestamp` with a v2
+    /// `batch_checksum` over the final header fields + per-message checksum 
fields.
+    fn valid_prepare_bytes() -> Owned<MESSAGE_ALIGN> {
+        let namespace = IggyNamespace::new(1, 1, 7);
+        let mut owned = SendMessages2Owned::from_messages(namespace, 
&sample_messages())
+            .expect("build send batch");
+        owned.header.base_offset = 10;
+        owned.header.base_timestamp = 20;
+        owned.header.batch_checksum = 
owned.header.checksum_for_blob(&owned.blob);
+        prepare_from_owned(&owned)
+    }
+
+    #[test]
+    fn decode_prepare_slice_trusted_matches_validating_for_valid_batch() {
+        // The trusted variant must surface byte-identical header meta to the
+        // validating decode for a checksum-consistent batch; only the
+        // per-message and batch-checksum passes are skipped.
+        let owned = valid_prepare_bytes();
+
+        let validated = decode_prepare_slice(owned.as_slice()).expect("valid 
batch decodes");
+        let trusted =
+            decode_prepare_slice_trusted(owned.as_slice()).expect("valid batch 
decodes trusted");
+
+        assert_eq!(validated.header.base_offset, trusted.header.base_offset);
+        assert_eq!(
+            validated.header.base_timestamp,
+            trusted.header.base_timestamp
+        );
+        assert_eq!(
+            validated.header.origin_timestamp,
+            trusted.header.origin_timestamp
+        );
+        assert_eq!(validated.header.batch_length, trusted.header.batch_length);
+        assert_eq!(validated.message_count(), trusted.message_count());
+        assert_eq!(validated.header.total_size(), trusted.header.total_size());
+        assert_eq!(validated.blob(), trusted.blob());
+    }
+
+    #[test]
+    fn decode_prepare_slice_trusted_skips_batch_checksum() {
+        // A stored batch_checksum mutated after stamping fails the validating
+        // decode but passes the trusted one: exactly why the trusted variant 
is
+        // confined to locally-produced bytes (see its doc invariant).
+        let mut owned = valid_prepare_bytes();
+        let corrupt_index = std::mem::size_of::<PrepareHeader>() + 
BATCH_CHECKSUM_OFFSET;
+        owned.as_mut_slice()[corrupt_index] ^= 0xFF;
+
+        assert!(
+            matches!(
+                decode_prepare_slice(owned.as_slice()),
+                Err(IggyError::InvalidBatchChecksum(..))
+            ),
+            "validating decode must reject a mutated batch checksum",
+        );
+        assert!(
+            decode_prepare_slice_trusted(owned.as_slice()).is_ok(),
+            "trusted decode skips the batch-checksum recomputation",
+        );
+    }
+
     #[test]
     fn decode_prepare_slice_size_below_header_size_does_not_panic() {
         // Regression: without the `total_size < header_size` guard,
@@ -869,4 +1154,484 @@ mod tests {
         );
         let _ = decode_prepare_slice(misaligned);
     }
+
+    fn sample_messages() -> IggyMessages2 {
+        let mut messages = IggyMessages2::with_capacity(2);
+        messages.push(IggyMessage2 {
+            header: IggyMessage2Header {
+                id: 7,
+                origin_timestamp: 1_000,
+                ..Default::default()
+            },
+            payload: Bytes::from_static(b"first-payload"),
+            user_headers: None,
+        });
+        messages.push(IggyMessage2 {
+            header: IggyMessage2Header {
+                id: 8,
+                origin_timestamp: 1_050,
+                ..Default::default()
+            },
+            payload: Bytes::from_static(b"second-payload"),
+            user_headers: Some(Bytes::from_static(b"user-header-bytes")),
+        });
+        messages
+    }
+
+    /// `[PrepareHeader][256B batch header][blob]` carrying real per-message
+    /// records + checksums from the production encoder, left pre-stamp
+    /// (`base_offset` / `base_timestamp` zero) as a follower receives it.
+    fn prepare_with_messages(messages: &IggyMessages2) -> Owned<MESSAGE_ALIGN> 
{
+        let namespace = IggyNamespace::new(1, 1, 7);
+        let owned =
+            SendMessages2Owned::from_messages(namespace, 
messages).expect("build send batch");
+        prepare_from_owned(&owned)
+    }
+
+    #[test]
+    fn verify_received_send_messages_accepts_clean_batch() {
+        let owned = prepare_with_messages(&sample_messages());
+        verify_received_send_messages(owned.as_slice())
+            .expect("a clean batch passes the receive gate");
+    }
+
+    #[test]
+    fn verify_received_send_messages_rejects_flipped_payload_byte() {
+        let mut owned = prepare_with_messages(&sample_messages());
+        // First payload begins right after the first message's 48B header.
+        let payload_index = PREPARE_SPLIT_POINT + MESSAGE_HEADER_SIZE;
+        owned.as_mut_slice()[payload_index] ^= 0xFF;
+        assert!(
+            matches!(
+                verify_received_send_messages(owned.as_slice()),
+                Err(IggyError::InvalidMessageChecksum(..))
+            ),
+            "a flipped payload byte must fail the per-message checksum",
+        );
+    }
+
+    #[test]
+    fn verify_received_send_messages_rejects_flipped_stored_checksum() {
+        let mut owned = prepare_with_messages(&sample_messages());
+        // The first message's stored checksum is the first 8 bytes of the 
blob.
+        owned.as_mut_slice()[PREPARE_SPLIT_POINT] ^= 0xFF;
+        assert!(
+            matches!(
+                verify_received_send_messages(owned.as_slice()),
+                Err(IggyError::InvalidMessageChecksum(..))
+            ),
+            "a flipped stored checksum must fail the per-message check",
+        );
+    }
+
+    #[test]
+    fn checksum_oneshot_matches_streaming_reference() {
+        // Formula pin: the per-message checksum is XxHash3-64 (default seed)
+        // over `header[8..48] || payload || user_headers` as one byte stream.
+        // The encoders hash the concatenation in a single oneshot pass; this
+        // streaming reference feeds the same parts separately. Both must agree
+        // for every shape, or checksums at rest stop verifying.
+        fn streaming_reference(header_tail: &[u8], payload: &[u8], 
user_headers: &[u8]) -> u64 {
+            let mut hasher = XxHash3_64::new();
+            hasher.write(header_tail);
+            hasher.write(payload);
+            hasher.write(user_headers);
+            hasher.finish()
+        }
+
+        let header_tail: Vec<u8> = (0u8..40).collect();
+        let kilobyte: Vec<u8> = (0..1024u32).map(|index| (index % 251) as 
u8).collect();
+        let cases: &[(&[u8], &[u8])] = &[
+            (&[], &[]),
+            (b"payload-bytes", &[]),
+            (b"payload-bytes", b"user-header-bytes"),
+            (&kilobyte, &[]),
+            (&kilobyte, &kilobyte[..7]),
+            (&kilobyte[..1023], &kilobyte[..7]),
+        ];
+        for (payload, user_headers) in cases {
+            let mut concatenated =
+                Vec::with_capacity(header_tail.len() + payload.len() + 
user_headers.len());
+            concatenated.extend_from_slice(&header_tail);
+            concatenated.extend_from_slice(payload);
+            concatenated.extend_from_slice(user_headers);
+            assert_eq!(
+                XxHash3_64::oneshot(&concatenated),
+                streaming_reference(&header_tail, payload, user_headers),
+                "oneshot must match the streaming reference for payload {} B, 
user headers {} B",
+                payload.len(),
+                user_headers.len(),
+            );
+        }
+    }
+
+    #[test]
+    fn batch_checksum_v2_pins_header_fields_then_message_checksum_fields() {
+        // Formula pin for batch checksum v2: XxHash3-64 (default seed) 
streaming
+        // over the six batch header meta fields (LE, in field order) then each
+        // message's stored 8-byte checksum field in message order - never the
+        // bodies. This reference walks the blob by the KNOWN input message 
sizes,
+        // independent of the production frame decoder, and must equal what the
+        // encoder stamped, or a stamp will not verify against a read-back
+        // recompute.
+        let namespace = IggyNamespace::new(1, 1, 7);
+        let messages = sample_messages();
+        let mut owned =
+            SendMessages2Owned::from_messages(namespace, 
&messages).expect("build batch");
+        owned.header.base_offset = 100;
+        owned.header.base_timestamp = 200;
+        owned.header.batch_checksum = 
owned.header.checksum_for_blob(&owned.blob);
+
+        let mut hasher = XxHash3_64::new();
+        hasher.write(&owned.header.partition_id.to_le_bytes());
+        hasher.write(&owned.header.base_offset.to_le_bytes());
+        hasher.write(&owned.header.base_timestamp.to_le_bytes());
+        hasher.write(&owned.header.origin_timestamp.to_le_bytes());
+        hasher.write(&owned.header.batch_length.to_le_bytes());
+        hasher.write(&owned.header.message_count.to_le_bytes());
+        let mut frame_start = 0usize;
+        for message in messages.iter() {
+            hasher.write(&owned.blob[frame_start..frame_start + 8]);
+            let user_headers = 
message.user_headers.as_deref().unwrap_or_default();
+            frame_start += MESSAGE_HEADER_SIZE + message.payload.len() + 
user_headers.len();
+        }
+        let reference = hasher.finish();
+
+        assert_eq!(
+            frame_start,
+            owned.blob.len(),
+            "reference walk must consume the whole blob",
+        );
+        assert_eq!(
+            owned.header.batch_checksum, reference,
+            "v2 batch checksum must equal hash(6 header fields || per-message 
checksum fields)",
+        );
+    }
+
+    #[test]
+    fn decode_batch_slice_rejects_body_corruption_with_intact_checksum_field() 
{
+        // Equal-integrity: v2 binds bodies only through the per-message 
checksum
+        // fields, so a flipped body byte that leaves the 8-byte checksum field
+        // intact keeps the batch value matching. The validating decode must 
still
+        // reject it via the per-message verify - the sole at-rest read-back 
check
+        // (the poll disk walk) decodes through here.
+        let namespace = IggyNamespace::new(1, 1, 7);
+        let owned =
+            SendMessages2Owned::from_messages(namespace, 
&sample_messages()).expect("build batch");
+        let mut body = vec![0u8; COMMAND_HEADER_SIZE + owned.blob.len()];
+        owned.header.encode_into(&mut body[..COMMAND_HEADER_SIZE]);
+        body[COMMAND_HEADER_SIZE..].copy_from_slice(&owned.blob);
+
+        decode_batch_slice(&body).expect("the clean batch decodes");
+
+        // First payload byte sits right after the command header and the first
+        // message's 48B frame header, leaving that frame's checksum field 
intact.
+        let payload_index = COMMAND_HEADER_SIZE + MESSAGE_HEADER_SIZE;
+        body[payload_index] ^= 0xFF;
+        assert!(
+            matches!(
+                decode_batch_slice(&body),
+                Err(IggyError::InvalidMessageChecksum(..))
+            ),
+            "body corruption with an intact checksum field must fail the 
per-message verify",
+        );
+    }
+
+    #[test]
+    fn 
decode_prepare_slice_rejects_body_corruption_with_intact_checksum_field() {
+        // The same equal-integrity guarantee at the resident/repair validating
+        // decode, plus proof that the batch value alone is blind to it.
+        let mut owned = prepare_with_messages(&sample_messages());
+        decode_prepare_slice(owned.as_slice()).expect("the clean prepare 
decodes");
+
+        let payload_index = PREPARE_SPLIT_POINT + MESSAGE_HEADER_SIZE;
+        owned.as_mut_slice()[payload_index] ^= 0xFF;
+        assert!(
+            matches!(
+                decode_prepare_slice(owned.as_slice()),
+                Err(IggyError::InvalidMessageChecksum(..))
+            ),
+            "body corruption with an intact checksum field must fail the 
per-message verify",
+        );
+        assert!(
+            decode_prepare_slice_trusted(owned.as_slice()).is_ok(),
+            "the intact checksum field leaves the batch value matching, so 
trusted still passes",
+        );
+    }
+
+    /// Legacy `SendMessages` request body: `[metadata_len=4][message_count]`
+    /// then `count` skipped index slots, then the 64B-header legacy records.
+    fn legacy_send_messages_body(messages: &IggyMessages2) -> Vec<u8> {
+        let count = messages.count();
+        let mut body = Vec::new();
+        body.extend_from_slice(&4u32.to_le_bytes());
+        body.extend_from_slice(&count.to_le_bytes());
+        body.extend_from_slice(&vec![0u8; count as usize * INDEX_SIZE]);
+        for message in messages.iter() {
+            let user_headers = 
message.user_headers.as_deref().unwrap_or_default();
+            let mut header = [0u8; LEGACY_MESSAGE_HEADER_SIZE];
+            header[8..24].copy_from_slice(&message.header.id.to_le_bytes());
+            
header[40..48].copy_from_slice(&message.header.origin_timestamp.to_le_bytes());
+            header[48..52].copy_from_slice(&(user_headers.len() as 
u32).to_le_bytes());
+            header[52..56].copy_from_slice(&(message.payload.len() as 
u32).to_le_bytes());
+            body.extend_from_slice(&header);
+            body.extend_from_slice(&message.payload);
+            body.extend_from_slice(user_headers);
+        }
+        body
+    }
+
+    fn legacy_request_message(body: &[u8]) -> Message<RequestHeader> {
+        let header_size = std::mem::size_of::<RequestHeader>();
+        let total = header_size + body.len();
+        let mut buffer = Owned::<MESSAGE_ALIGN>::zeroed(total);
+        {
+            let header: &mut RequestHeader =
+                bytemuck::checked::try_from_bytes_mut(&mut 
buffer.as_mut_slice()[..header_size])
+                    .expect("zeroed bytes form a valid RequestHeader");
+            header.command = Command2::Request;
+            header.operation = Operation::SendMessages;
+            header.client = 1;
+            header.session = 1;
+            header.request = 1;
+            header.size = u32::try_from(total).expect("size fits u32");
+        }
+        buffer.as_mut_slice()[header_size..].copy_from_slice(body);
+        Message::try_from(buffer).expect("legacy request message is valid")
+    }
+
+    #[test]
+    fn convert_request_message_transcodes_legacy_to_canonical_bytes() {
+        // Golden: the fused legacy transcode must emit the exact canonical 
batch
+        // the native builder (`from_messages`) produces for the same messages 
-
+        // command header + blob, byte for byte. Explicit non-zero ids keep it
+        // deterministic (no `random_id` substitution).
+        let namespace = IggyNamespace::new(1, 1, 3);
+        let messages = sample_messages();
+
+        let owned =
+            SendMessages2Owned::from_messages(namespace, 
&messages).expect("build canonical batch");
+        let mut expected_body = vec![0u8; COMMAND_HEADER_SIZE + 
owned.blob.len()];
+        owned
+            .header
+            .encode_into(&mut expected_body[..COMMAND_HEADER_SIZE]);
+        expected_body[COMMAND_HEADER_SIZE..].copy_from_slice(&owned.blob);
+
+        let legacy = 
legacy_request_message(&legacy_send_messages_body(&messages));
+        let converted = convert_request_message(namespace, legacy, 
ChecksumMode::Compute)
+            .expect("legacy body transcodes");
+        let header_size = std::mem::size_of::<RequestHeader>();
+        let actual_body = 
&converted.as_slice()[header_size..converted.header().size as usize];
+
+        assert_eq!(
+            actual_body, expected_body,
+            "legacy transcode must be byte-identical to the canonical native 
batch",
+        );
+
+        // And the emitted batch is self-consistent: it validates through the
+        // batch-checksum decode and yields the original messages.
+        let decoded = decode_batch_slice(actual_body).expect("transcoded batch 
checksum is valid");
+        assert_eq!(decoded.message_count(), messages.count());
+        let payloads: Vec<&[u8]> = decoded.iter().map(|view| 
view.payload).collect();
+        assert_eq!(
+            payloads,
+            vec![&b"first-payload"[..], &b"second-payload"[..]]
+        );
+    }
+
+    #[test]
+    fn convert_request_message_skip_leaves_batch_checksum_zero_until_stamp() {
+        // The partition ingest path passes Skip: the transcoded batch must 
carry
+        // a zero checksum (stamp fills it) and be otherwise byte-identical to 
the
+        // Compute output - the flag toggles nothing but that one hash.
+        let namespace = IggyNamespace::new(1, 1, 3);
+        let messages = sample_messages();
+        let body = legacy_send_messages_body(&messages);
+        let header_size = std::mem::size_of::<RequestHeader>();
+
+        let computed = convert_request_message(
+            namespace,
+            legacy_request_message(&body),
+            ChecksumMode::Compute,
+        )
+        .expect("compute transcode");
+        let skipped =
+            convert_request_message(namespace, legacy_request_message(&body), 
ChecksumMode::Skip)
+                .expect("skip transcode");
+
+        let computed_body = 
&computed.as_slice()[header_size..computed.header().size as usize];
+        let skipped_body = 
&skipped.as_slice()[header_size..skipped.header().size as usize];
+
+        let skipped_header = 
SendMessages2Header::decode(&skipped_body[..COMMAND_HEADER_SIZE])
+            .expect("decode skipped header");
+        assert_eq!(
+            skipped_header.batch_checksum, 0,
+            "skip leaves the batch checksum zero until stamp",
+        );
+
+        // Patch only the 8-byte batch_checksum field into the skipped body; it
+        // must then equal the computed body, proving nothing else diverges.
+        let mut patched = skipped_body.to_vec();
+        patched[BATCH_CHECKSUM_OFFSET..BATCH_CHECKSUM_OFFSET + 8]
+            
.copy_from_slice(&computed_body[BATCH_CHECKSUM_OFFSET..BATCH_CHECKSUM_OFFSET + 
8]);
+        assert_eq!(
+            patched.as_slice(),
+            computed_body,
+            "skip and compute differ only in the batch_checksum field",
+        );
+    }
+
+    #[test]
+    fn encrypt_ingest_path_stays_canonical_through_flag_split() {
+        // Mirror the plane encrypt ingest sequence: convert(Compute) -> the
+        // validating decode encrypt performs on its input -> encrypt -> the
+        // validating decode the second convert performs as its discriminator 
->
+        // convert(Skip) (the partition convert), which sees an 
already-canonical
+        // batch and returns it unchanged. Every decode must succeed.
+        let namespace = IggyNamespace::new(1, 1, 3);
+        let messages = sample_messages();
+        let header_size = std::mem::size_of::<RequestHeader>();
+
+        let legacy = 
legacy_request_message(&legacy_send_messages_body(&messages));
+        let canonical = convert_request_message(namespace, legacy, 
ChecksumMode::Compute)
+            .expect("pre-encrypt transcode");
+        let canonical_body = 
&canonical.as_slice()[header_size..canonical.header().size as usize];
+        decode_batch_slice(canonical_body).expect("encrypt input decode 
validates the checksum");
+
+        let encryptor =
+            EncryptorKind::Aes256Gcm(Aes256GcmEncryptor::new(&[7u8; 
32]).expect("valid 32B key"));
+        let encrypted = encrypt_batch_request(canonical, 
&encryptor).expect("encrypt batch");
+        let encrypted_body: Vec<u8> =
+            encrypted.as_slice()[header_size..encrypted.header().size as 
usize].to_vec();
+        decode_batch_slice(&encrypted_body)
+            .expect("encrypt output drives the 2nd-convert discriminator");
+
+        let repassed = convert_request_message(namespace, encrypted, 
ChecksumMode::Skip)
+            .expect("second convert passes the canonical batch");
+        let repassed_body = 
&repassed.as_slice()[header_size..repassed.header().size as usize];
+        assert_eq!(
+            repassed_body,
+            encrypted_body.as_slice(),
+            "an already-canonical encrypted batch passes the partition convert 
untouched",
+        );
+    }
+
+    /// Junk suffixes that must be refused at both ingest boundaries: one 
below a
+    /// frame header (the frame walk stops on a short read) and one frame-sized
+    /// but undecodable (`reserved != 0`). Neither is covered by any checksum, 
so
+    /// a walk that stops at the last decodable frame cannot see them.
+    const TRAILING_JUNK_CASES: [&[u8]; 2] = [&[0xAA], &[0xFF; 64]];
+
+    /// Canonical `SendMessages` request carrying `junk` past `batch_length`, 
with
+    /// `RequestHeader.size` inflated to cover it. `size` and `batch_length` 
are
+    /// independent wire fields, so a non-conforming client can emit this.
+    fn canonical_request_with_trailing_bytes(junk: &[u8]) -> 
Message<RequestHeader> {
+        let namespace = IggyNamespace::new(1, 1, 3);
+        let owned =
+            SendMessages2Owned::from_messages(namespace, 
&sample_messages()).expect("build batch");
+        let header_size = std::mem::size_of::<RequestHeader>();
+        let total = header_size + owned.header.total_size() + junk.len();
+        let mut buffer = Owned::<MESSAGE_ALIGN>::zeroed(total);
+        {
+            let header: &mut RequestHeader =
+                bytemuck::checked::try_from_bytes_mut(&mut 
buffer.as_mut_slice()[..header_size])
+                    .expect("zeroed bytes form a valid RequestHeader");
+            header.command = Command2::Request;
+            header.operation = Operation::SendMessages;
+            header.client = 1;
+            header.session = 1;
+            header.request = 1;
+            header.size = u32::try_from(total).expect("size fits u32");
+        }
+        let bytes = buffer.as_mut_slice();
+        owned
+            .header
+            .encode_into(&mut bytes[header_size..header_size + 
COMMAND_HEADER_SIZE]);
+        let blob_end = PREPARE_SPLIT_POINT + owned.blob.len();
+        bytes[PREPARE_SPLIT_POINT..blob_end].copy_from_slice(&owned.blob);
+        bytes[blob_end..].copy_from_slice(junk);
+        Message::try_from(buffer).expect("request message is valid")
+    }
+
+    /// The replicated counterpart: a pre-stamp `Prepare` whose `size` covers
+    /// `junk` past `batch_length`.
+    fn prepare_with_trailing_bytes(junk: &[u8]) -> Owned<MESSAGE_ALIGN> {
+        let namespace = IggyNamespace::new(1, 1, 7);
+        let owned =
+            SendMessages2Owned::from_messages(namespace, 
&sample_messages()).expect("build batch");
+        let header_size = std::mem::size_of::<PrepareHeader>();
+        let total = header_size + owned.header.total_size() + junk.len();
+        let mut buffer = Owned::<MESSAGE_ALIGN>::zeroed(total);
+        {
+            let prepare: &mut PrepareHeader =
+                bytemuck::checked::try_from_bytes_mut(&mut 
buffer.as_mut_slice()[..header_size])
+                    .expect("zeroed bytes form a valid PrepareHeader");
+            prepare.command = Command2::Prepare;
+            prepare.size = u32::try_from(total).expect("prepare size fits 
u32");
+        }
+        let bytes = buffer.as_mut_slice();
+        owned
+            .header
+            .encode_into(&mut bytes[header_size..header_size + 
COMMAND_HEADER_SIZE]);
+        let blob_end = PREPARE_SPLIT_POINT + owned.blob.len();
+        bytes[PREPARE_SPLIT_POINT..blob_end].copy_from_slice(&owned.blob);
+        bytes[blob_end..].copy_from_slice(junk);
+        buffer
+    }
+
+    #[test]
+    fn convert_request_message_rejects_canonical_batch_with_trailing_bytes() {
+        // Client ingest boundary. Accepting the request would carry the suffix
+        // into the journal and onto disk: the flush writes the whole frame 
while
+        // every reader advances by `batch_length`, so the segment walk lands
+        // inside the junk and every later batch becomes unreadable.
+        let namespace = IggyNamespace::new(1, 1, 3);
+        for junk in TRAILING_JUNK_CASES {
+            for mode in [ChecksumMode::Compute, ChecksumMode::Skip] {
+                let message = canonical_request_with_trailing_bytes(junk);
+                let result = convert_request_message(namespace, message, mode);
+                assert!(
+                    matches!(result, Err(IggyError::InvalidCommand)),
+                    "{} trailing bytes ({mode:?}) must be rejected, got 
{result:?}",
+                    junk.len(),
+                );
+            }
+        }
+    }
+
+    #[test]
+    fn convert_request_message_accepts_exact_canonical_batch() {
+        // The same builder with no suffix must still pass untouched, so the
+        // rejection above is the suffix and not the fixture.
+        let namespace = IggyNamespace::new(1, 1, 3);
+        let message = canonical_request_with_trailing_bytes(&[]);
+        let expected = message.as_slice().to_vec();
+        let converted = convert_request_message(namespace, message, 
ChecksumMode::Skip)
+            .expect("an exact canonical batch passes untouched");
+        assert_eq!(converted.as_slice(), expected.as_slice());
+    }
+
+    #[test]
+    fn 
verify_received_send_messages_rejects_trailing_bytes_past_batch_length() {
+        // Replica ingest boundary. The gate clamps the blob to `batch_length`
+        // before verifying, so without an exact-frame check a primary could 
plant
+        // bytes that no per-message checksum covers on every backup.
+        for junk in TRAILING_JUNK_CASES {
+            let owned = prepare_with_trailing_bytes(junk);
+            let result = verify_received_send_messages(owned.as_slice());
+            assert!(
+                matches!(result, Err(IggyError::InvalidCommand)),
+                "{} trailing bytes must fail the receive gate, got {result:?}",
+                junk.len(),
+            );
+            assert!(
+                matches!(
+                    decode_prepare_slice(owned.as_slice()),
+                    Err(IggyError::InvalidCommand)
+                ),
+                "{} trailing bytes must fail the validating decode",
+                junk.len(),
+            );
+        }
+    }
 }
diff --git a/core/simulator/src/client.rs b/core/simulator/src/client.rs
index f511c1000..028adb38c 100644
--- a/core/simulator/src/client.rs
+++ b/core/simulator/src/client.rs
@@ -72,7 +72,7 @@ pub struct SimClient {
     partition_counter: Cell<u64>,
     /// Deterministic per-message id source for produced messages. The real SDK
     /// sends `id: 0` and lets the server mint a random UUID
-    /// (`SendMessages2::from_legacy_request` -> `random_id::get_uuid`); that
+    /// (`transcode_legacy_request` -> `random_id::get_uuid`); that
     /// mint is unseeded, so under the deterministic executor a produce's
     /// replicated body bytes (and their checksums) would differ run to run,
     /// silently breaking seeded replay. Stamping a deterministic id here keeps
@@ -516,7 +516,7 @@ impl SimClient {
     /// client). VSR clients resolve to an explicit partition before sending, 
so
     /// the sim always emits `WirePartitioning::PartitionId`: that is the shape
     /// the shell's `resolve_partition_request_namespace` decodes, and the raw
-    /// path converts it to `SendMessages2` via `from_legacy_request`.
+    /// path converts it to `SendMessages2` via `transcode_legacy_request`.
     ///
     /// # Panics
     /// Panics if a namespace id exceeds `u32` or the request buffer is 
invalid.

Reply via email to