Skip to content
Open
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 41 additions & 6 deletions src/message.rs
Original file line number Diff line number Diff line change
Expand Up @@ -513,6 +513,12 @@ unsafe impl KafkaDrop for RDKafkaHeaders {
unsafe impl Send for OwnedHeaders {}
unsafe impl Sync for OwnedHeaders {}

impl AsMut<OwnedHeaders> 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 {
Expand All @@ -529,11 +535,12 @@ impl OwnedHeaders {
}
}

/// Inserts a new header.
pub fn insert<V>(self, header: Header<'_, &V>) -> OwnedHeaders
fn insert_internal<V, R>(mut owned_headers: R, header: Header<'_, &V>) -> R
where
V: ToBytes + ?Sized,
R: AsMut<OwnedHeaders>,
{
let headers = owned_headers.as_mut();
let (value_ptr, value_len) = match header.value {
None => (ptr::null_mut(), 0),
Some(value) => {
Expand All @@ -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<V>(self, header: Header<'_, &V>) -> OwnedHeaders
where
V: ToBytes + ?Sized,
{
Self::insert_internal(self, header)
}

/// Inserts a header in-place.
pub fn insert_inplace<V>(&mut self, header: Header<'_, &V>)
where
V: ToBytes + ?Sized,
{
Self::insert_internal(self, header);
}

pub(crate) fn ptr(&self) -> *mut RDKafkaHeaders {
Expand Down Expand Up @@ -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"),
Expand All @@ -879,5 +902,17 @@ mod test {
value: Some("value2")
})
);

owned.insert_inplace(Header {
key: "key3",
value: Some("value3"),
});
assert_eq!(
owned.get_as::<str>(2),
Ok(Header {
key: "key3",
value: Some("value3")
})
);
}
}