use crate::checkpoint::AckRef;
use crate::error::DeserError;
use crate::record::{Flow, RawPayload, Record};
pub trait RecFamily: 'static {
type Rec<'buf>: Send;
}
#[derive(Debug)]
pub struct Owned<T>(std::marker::PhantomData<fn() -> T>);
impl<T: Send + 'static> RecFamily for Owned<T> {
type Rec<'buf> = T;
}
pub trait EmitRecord<'buf, T> {
fn emit(&mut self, rec: Record<T>) -> Flow;
}
pub trait Deserializer<F: RecFamily>: Send {
fn deserialize<'buf>(
&mut self,
raw: &RawPayload<'buf>,
ack: &AckRef,
out: &mut dyn EmitRecord<'buf, F::Rec<'buf>>,
) -> Result<(), DeserError>;
}
#[derive(Clone, Copy, Debug, Default)]
pub struct BytesPassthrough;
impl Deserializer<Owned<Vec<u8>>> for BytesPassthrough {
fn deserialize<'buf>(
&mut self,
raw: &RawPayload<'buf>,
ack: &AckRef,
out: &mut dyn EmitRecord<'buf, Vec<u8>>,
) -> Result<(), DeserError> {
let _ = out.emit(Record {
payload: raw.bytes.to_vec(),
meta: raw.meta(),
ack: ack.clone(),
});
Ok(())
}
}
#[cfg(all(test, not(loom)))]
mod tests {
use super::*;
use crate::record::PartitionId;
struct Sink(Vec<Vec<u8>>);
impl EmitRecord<'_, Vec<u8>> for Sink {
fn emit(&mut self, rec: Record<Vec<u8>>) -> Flow {
self.0.push(rec.payload);
Flow::Continue
}
}
#[test]
fn passthrough_emits_one_owned_record() {
let (ack, _rx) = AckRef::test_pair();
let raw = RawPayload {
bytes: b"abc",
key: None,
partition: PartitionId(0),
offset: 1,
timestamp_ms: 2,
};
let mut sink = Sink(Vec::new());
BytesPassthrough.deserialize(&raw, &ack, &mut sink).unwrap();
assert_eq!(sink.0, vec![b"abc".to_vec()]);
}
}