This is an automated email from the ASF dual-hosted git repository.
numinnex 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 499d8aa46 perf(common): decode user headers without copying the blob
(#3708)
499d8aa46 is described below
commit 499d8aa4666d73764eaaa2d94ff9cfab6d747136
Author: Ryan Huang <[email protected]>
AuthorDate: Tue Jul 21 02:01:43 2026 +0800
perf(common): decode user headers without copying the blob (#3708)
---
core/common/src/types/message/message_view.rs | 78 ++++++++++++++++++----
core/common/src/wire_conversions.rs | 93 ++++++++++++++++++++++++++-
2 files changed, 155 insertions(+), 16 deletions(-)
diff --git a/core/common/src/types/message/message_view.rs
b/core/common/src/types/message/message_view.rs
index 2b8646d11..f37c83840 100644
--- a/core/common/src/types/message/message_view.rs
+++ b/core/common/src/types/message/message_view.rs
@@ -22,9 +22,9 @@ use crate::IggyByteSize;
use crate::Sizeable;
use crate::error::IggyError;
use crate::utils::checksum;
-use crate::wire_conversions::user_headers_from_wire;
+use crate::wire_conversions::user_headers_from_validated_slice;
use crate::{HeaderKey, IggyMessageHeaderView};
-use iggy_binary_protocol::WireUserHeaders;
+use iggy_binary_protocol::validate_user_headers;
use std::collections::BTreeMap;
use std::num::NonZeroUsize;
@@ -87,17 +87,14 @@ impl<'a> IggyMessageView<'a> {
/// Return instantiated user headers map
pub fn user_headers_map(&self) -> Result<Option<BTreeMap<HeaderKey,
HeaderValue>>, IggyError> {
if let Some(headers) = self.user_headers() {
- let wire = match WireUserHeaders::from_slice(headers) {
- Ok(w) => w,
- Err(e) => {
- tracing::warn!(
- "Failed to parse user headers: {e}, header_length={}",
- self.header().user_headers_length()
- );
- return Ok(None);
- }
- };
- let map = user_headers_from_wire(&wire)?;
+ if let Err(e) = validate_user_headers(headers) {
+ tracing::warn!(
+ "Failed to parse user headers: {e}, header_length={}",
+ self.header().user_headers_length()
+ );
+ return Ok(None);
+ }
+ let map = user_headers_from_validated_slice(headers)?;
Ok(Some(map))
} else {
Ok(None)
@@ -207,6 +204,7 @@ mod tests {
use super::*;
use crate::IggyMessage;
use bytes::Bytes;
+ use std::str::FromStr;
fn build_batch() -> crate::IggyMessagesBatch {
let messages = vec![
@@ -249,4 +247,58 @@ mod tests {
let last =
IggyMessageViewIterator::new(batch.buffer()).last().unwrap();
assert_eq!(last.payload(), b"three");
}
+
+ #[test]
+ fn
given_structurally_invalid_user_headers_when_mapped_should_return_none() {
+ // A TLV entry claiming a 4-byte key but carrying only 2 bytes.
+ let mut headers = Vec::new();
+ headers.push(6u8);
+ headers.extend_from_slice(&4u32.to_le_bytes());
+ headers.extend_from_slice(b"ab");
+
+ let message = IggyMessage::builder()
+ .payload(Bytes::from_static(b"payload"))
+ .build()
+ .unwrap();
+ let mut buffer = message.to_bytes().to_vec();
+ // Splice the malformed blob in and declare it in the fixed header.
+ let len = headers.len() as u32;
+
buffer[IGGY_MESSAGE_HEADERS_LENGTH_OFFSET_RANGE].copy_from_slice(&len.to_le_bytes());
+ buffer.extend_from_slice(&headers);
+
+ let view = IggyMessageView::new(&buffer).unwrap();
+ assert!(view.user_headers_map().unwrap().is_none());
+ }
+
+ #[test]
+ fn given_valid_user_headers_when_mapped_should_recover_every_entry() {
+ let user_headers = BTreeMap::from([
+ (
+ HeaderKey::from_str("content-type").unwrap(),
+ HeaderValue::from_str("text/plain").unwrap(),
+ ),
+ (HeaderKey::from_str("attempt").unwrap(), 3u32.into()),
+ ]);
+ let message = IggyMessage::builder()
+ .payload(Bytes::from_static(b"payload"))
+ .user_headers(user_headers.clone())
+ .build()
+ .unwrap();
+ let buffer = message.to_bytes().to_vec();
+
+ let view = IggyMessageView::new(&buffer).unwrap();
+ assert_eq!(view.user_headers_map().unwrap().unwrap(), user_headers);
+ }
+
+ #[test]
+ fn given_message_without_user_headers_when_mapped_should_return_none() {
+ let message = IggyMessage::builder()
+ .payload(Bytes::from_static(b"payload"))
+ .build()
+ .unwrap();
+ let buffer = message.to_bytes().to_vec();
+
+ let view = IggyMessageView::new(&buffer).unwrap();
+ assert!(view.user_headers_map().unwrap().is_none());
+ }
}
diff --git a/core/common/src/wire_conversions.rs
b/core/common/src/wire_conversions.rs
index b0e51841d..44975f159 100644
--- a/core/common/src/wire_conversions.rs
+++ b/core/common/src/wire_conversions.rs
@@ -30,7 +30,6 @@ use crate::{
RawPersonalAccessToken, Stats, Stream, StreamDetails, StreamPermissions,
Topic, TopicDetails,
TopicPermissions, TransportEndpoints, UserInfo, UserInfoDetails,
UserStatus,
};
-use iggy_binary_protocol::WireConsumer;
use iggy_binary_protocol::primitives::permissions::{
WireGlobalPermissions, WirePermissions, WireStreamPermissions,
WireTopicPermissions,
};
@@ -61,6 +60,7 @@ use
iggy_binary_protocol::responses::topics::get_topics::GetTopicsResponse;
use iggy_binary_protocol::responses::users::login_user::IdentityResponse;
use iggy_binary_protocol::responses::users::user_response::UserResponse;
use iggy_binary_protocol::responses::users::{GetUsersResponse,
UserDetailsResponse};
+use iggy_binary_protocol::{WireConsumer, WireUserHeaderIterator};
use std::collections::{BTreeMap, HashMap};
/// Sentinel value in the wire protocol indicating no authenticated user.
@@ -756,11 +756,27 @@ pub fn user_headers_to_wire(
pub fn user_headers_from_wire(
wire: &iggy_binary_protocol::WireUserHeaders,
) -> Result<BTreeMap<HeaderKey, HeaderValue>, IggyError> {
- if wire.is_empty() {
+ user_headers_from_validated_slice(wire.as_bytes())
+}
+
+/// Decode user headers from a slice that has already passed structural TLV
validation.
+///
+/// Callers holding only a borrowed buffer can use this to skip the copy that
+/// [`iggy_binary_protocol::WireUserHeaders::from_slice`] performs.
+///
+/// # Panics
+///
+/// Panics or yields garbage if `buf` has not been validated by
+/// [`iggy_binary_protocol::validate_user_headers`]. The underlying iterator
+/// slices each TLV field without bounds checks, relying on that validation.
+pub(crate) fn user_headers_from_validated_slice(
+ buf: &[u8],
+) -> Result<BTreeMap<HeaderKey, HeaderValue>, IggyError> {
+ if buf.is_empty() {
return Ok(BTreeMap::new());
}
let mut headers = BTreeMap::new();
- for entry in wire.iter() {
+ for entry in WireUserHeaderIterator::new(buf) {
let key_kind = HeaderKind::from_code(entry.key_kind.0)?;
if let Some(expected) = key_kind.expected_size()
&& entry.key.len() != expected
@@ -782,3 +798,74 @@ pub fn user_headers_from_wire(
}
Ok(headers)
}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use iggy_binary_protocol::WireUserHeaders;
+ use std::str::FromStr;
+
+ fn sample_headers() -> BTreeMap<HeaderKey, HeaderValue> {
+ BTreeMap::from([
+ (
+ HeaderKey::from_str("content-type").unwrap(),
+ HeaderValue::from_str("text/plain").unwrap(),
+ ),
+ (HeaderKey::from_str("retries").unwrap(), 7u32.into()),
+ ])
+ }
+
+ #[test]
+ fn
given_encoded_headers_when_decoded_from_slice_should_recover_the_originals() {
+ let wire = user_headers_to_wire(&sample_headers());
+
+ // The originals are the oracle here. Comparing against
`user_headers_from_wire`
+ // would be vacuous, since it delegates to the function under test.
+ let decoded =
user_headers_from_validated_slice(wire.as_bytes()).unwrap();
+
+ assert_eq!(decoded, sample_headers());
+ for (key, value) in &decoded {
+ let (expected_key, expected_value) = sample_headers()
+ .into_iter()
+ .find(|(k, _)| k == key)
+ .unwrap();
+ assert_eq!(key.kind(), expected_key.kind());
+ assert_eq!(value.kind(), expected_value.kind());
+ assert_eq!(value.as_bytes(), expected_value.as_bytes());
+ }
+ }
+
+ #[test]
+ fn
given_headers_copied_into_owned_wire_when_decoded_should_match_borrowed_decode()
{
+ // Guards the assumption that decoding a borrowed slice is equivalent
to
+ // decoding an independently copied buffer.
+ let wire = user_headers_to_wire(&sample_headers());
+ let copied = WireUserHeaders::from_slice(wire.as_bytes()).unwrap();
+
+ assert_eq!(
+ user_headers_from_validated_slice(wire.as_bytes()).unwrap(),
+ user_headers_from_validated_slice(copied.as_bytes()).unwrap()
+ );
+ }
+
+ #[test]
+ fn given_empty_buffer_when_decoded_from_slice_should_return_empty_map() {
+ assert!(user_headers_from_validated_slice(&[]).unwrap().is_empty());
+ }
+
+ #[test]
+ fn given_unknown_kind_code_when_decoded_from_slice_should_return_error() {
+ // Structurally valid TLV whose key kind code has no domain meaning.
+ let mut buf = Vec::new();
+ buf.push(0xFF);
+ buf.extend_from_slice(&3u32.to_le_bytes());
+ buf.extend_from_slice(b"key");
+ buf.push(0xFF);
+ buf.extend_from_slice(&3u32.to_le_bytes());
+ buf.extend_from_slice(b"val");
+
+ assert!(iggy_binary_protocol::validate_user_headers(&buf).is_ok());
+ assert!(user_headers_from_validated_slice(&buf).is_err());
+
assert!(user_headers_from_wire(&WireUserHeaders::from_slice(&buf).unwrap()).is_err());
+ }
+}