From 7b054ab2ad628b2ca2f43e5fb88af759d85977bf Mon Sep 17 00:00:00 2001 From: SickMcNugget Date: Thu, 11 Dec 2025 12:20:43 +0800 Subject: [PATCH] Insert headers to &mut OwnedHeaders inplace This commit adds a new function insert_inplace(), that takes a mutable reference to an OwnedHeaders, and can insert headers to that reference. The body of insert() has been deduplicated into insert_internal(), and the new functionality is also tested inside of test_headers(). --- src/message.rs | 47 +++++++++++++++++++++++++++++++++++++++++------ 1 file changed, 41 insertions(+), 6 deletions(-) diff --git a/src/message.rs b/src/message.rs index e72741219..a96a09ced 100644 --- a/src/message.rs +++ b/src/message.rs @@ -513,6 +513,12 @@ unsafe impl KafkaDrop for RDKafkaHeaders { unsafe impl Send for OwnedHeaders {} unsafe impl Sync for OwnedHeaders {} +impl AsMut for OwnedHeaders { + fn as_mut(&mut self) -> &mut OwnedHeaders { + self + } +} + impl OwnedHeaders { /// Creates a new `OwnedHeaders` struct with initial capacity 5. pub fn new() -> OwnedHeaders { @@ -529,11 +535,12 @@ impl OwnedHeaders { } } - /// Inserts a new header. - pub fn insert(self, header: Header<'_, &V>) -> OwnedHeaders + fn insert_internal(mut owned_headers: R, header: Header<'_, &V>) -> R where V: ToBytes + ?Sized, + R: AsMut, { + let headers = owned_headers.as_mut(); let (value_ptr, value_len) = match header.value { None => (ptr::null_mut(), 0), Some(value) => { @@ -546,16 +553,32 @@ impl OwnedHeaders { }; let err = unsafe { rdsys::rd_kafka_header_add( - self.ptr(), + headers.ptr(), header.key.as_ptr() as *const c_char, - header.key.len() as isize, + header.key.as_bytes().len() as isize, value_ptr, value_len, ) }; // OwnedHeaders should always represent writable instances of RDKafkaHeaders assert!(!err.is_error()); - self + owned_headers + } + + /// Inserts a new header. + pub fn insert(self, header: Header<'_, &V>) -> OwnedHeaders + where + V: ToBytes + ?Sized, + { + Self::insert_internal(self, header) + } + + /// Inserts a header in-place. + pub fn insert_inplace(&mut self, header: Header<'_, &V>) + where + V: ToBytes + ?Sized, + { + Self::insert_internal(self, header); } pub(crate) fn ptr(&self) -> *mut RDKafkaHeaders { @@ -856,7 +879,7 @@ mod test { #[test] fn test_headers() { - let owned = OwnedHeaders::new() + let mut owned = OwnedHeaders::new() .insert(Header { key: "key1", value: Some("value1"), @@ -879,5 +902,17 @@ mod test { value: Some("value2") }) ); + + owned.insert_inplace(Header { + key: "key3", + value: Some("value3"), + }); + assert_eq!( + owned.get_as::(2), + Ok(Header { + key: "key3", + value: Some("value3") + }) + ); } }