use std::sync::Arc;
use bytes::Bytes;
use crate::Headers;
use crate::error::{KrafkaError, Result};
pub trait Serializer<T: ?Sized>: Send + Sync {
fn serialize(&self, topic: &str, headers: &mut Headers, value: &T) -> Result<Bytes>;
}
#[derive(Debug, Clone, Copy, Default)]
pub struct BytesSerializer;
impl Serializer<Bytes> for BytesSerializer {
fn serialize(&self, _topic: &str, _headers: &mut Headers, value: &Bytes) -> Result<Bytes> {
Ok(value.clone())
}
}
impl Serializer<Vec<u8>> for BytesSerializer {
fn serialize(&self, _topic: &str, _headers: &mut Headers, value: &Vec<u8>) -> Result<Bytes> {
Ok(Bytes::copy_from_slice(value))
}
}
impl Serializer<[u8]> for BytesSerializer {
fn serialize(&self, _topic: &str, _headers: &mut Headers, value: &[u8]) -> Result<Bytes> {
Ok(Bytes::copy_from_slice(value))
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct StringSerializer;
impl Serializer<String> for StringSerializer {
fn serialize(&self, _topic: &str, _headers: &mut Headers, value: &String) -> Result<Bytes> {
Ok(Bytes::copy_from_slice(value.as_bytes()))
}
}
impl Serializer<str> for StringSerializer {
fn serialize(&self, _topic: &str, _headers: &mut Headers, value: &str) -> Result<Bytes> {
Ok(Bytes::copy_from_slice(value.as_bytes()))
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct NoKey;
impl<T: ?Sized> Serializer<T> for NoKey {
fn serialize(&self, topic: &str, _headers: &mut Headers, _value: &T) -> Result<Bytes> {
Err(KrafkaError::serialization(format!(
"this producer sends no keys, but a key was given for topic {topic}"
)))
}
}
pub trait Deserializer: Send + Sync {
fn deserialize(
&self,
topic: &str,
headers: &Headers,
payload: Bytes,
is_key: bool,
) -> Result<Bytes>;
}
impl<T: Deserializer + ?Sized> Deserializer for Arc<T> {
fn deserialize(
&self,
topic: &str,
headers: &Headers,
payload: Bytes,
is_key: bool,
) -> Result<Bytes> {
(**self).deserialize(topic, headers, payload, is_key)
}
}
#[cfg(test)]
#[allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
mod tests {
use super::*;
#[derive(Debug)]
struct Prefix(&'static [u8]);
impl Deserializer for Prefix {
fn deserialize(
&self,
_topic: &str,
_headers: &Headers,
payload: Bytes,
_is_key: bool,
) -> Result<Bytes> {
Ok(payload.slice(self.0.len()..))
}
}
#[test]
fn deserializer_is_object_safe() {
let de: Arc<dyn Deserializer> = Arc::new(Prefix(b"\x00\x01"));
let plain = de
.deserialize(
"orders",
&Headers::new(),
Bytes::from_static(b"\x00\x01payload"),
false,
)
.unwrap();
assert_eq!(&plain[..], b"payload");
}
#[test]
fn deserialize_can_be_zero_copy() {
let de = Prefix(b"\x00\x01");
let input = Bytes::from_static(b"\x00\x01payload");
let out = de
.deserialize("orders", &Headers::new(), input.clone(), false)
.unwrap();
assert_eq!(out.as_ptr(), input[2..].as_ptr(), "expected a sub-slice");
}
#[test]
fn byte_and_string_serializers_pass_values_through() {
let mut headers = Vec::new();
let bytes = Bytes::from_static(b"raw");
assert_eq!(
BytesSerializer
.serialize("t", &mut headers, &bytes)
.unwrap(),
bytes
);
assert_eq!(
BytesSerializer
.serialize("t", &mut headers, &b"raw".to_vec())
.unwrap(),
bytes
);
assert_eq!(
BytesSerializer
.serialize("t", &mut headers, &b"raw"[..])
.unwrap(),
bytes
);
assert_eq!(
StringSerializer
.serialize("t", &mut headers, "héllo")
.unwrap(),
Bytes::from("héllo")
);
assert_eq!(
StringSerializer
.serialize("t", &mut headers, &"héllo".to_string())
.unwrap(),
Bytes::from("héllo")
);
assert!(
headers.is_empty(),
"pass-through serializers add no headers"
);
}
#[test]
fn no_key_refuses_a_key() {
let mut headers = Vec::new();
assert!(NoKey.serialize("t", &mut headers, "k").is_err());
}
}