use crate::checkpoint::AckRef;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Flow {
Continue,
Blocked,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct PartitionId(pub u32);
#[derive(Clone, Copy, Debug)]
pub struct RecordMeta {
pub partition: PartitionId,
pub offset: i64,
pub event_time_ms: i64,
pub key_hash: Option<u64>,
}
#[derive(Debug)]
pub struct Record<T> {
pub payload: T,
pub meta: RecordMeta,
pub ack: AckRef,
}
impl<T> Record<T> {
#[inline(always)]
pub fn map<U>(self, f: impl FnOnce(T) -> U) -> Record<U> {
Record {
payload: f(self.payload),
meta: self.meta,
ack: self.ack,
}
}
}
#[derive(Clone, Copy, Debug)]
pub struct RawPayload<'buf> {
pub bytes: &'buf [u8],
pub key: Option<&'buf [u8]>,
pub partition: PartitionId,
pub offset: i64,
pub timestamp_ms: i64,
}
impl RawPayload<'_> {
#[inline]
pub fn meta(&self) -> RecordMeta {
RecordMeta {
partition: self.partition,
offset: self.offset,
event_time_ms: self.timestamp_ms,
key_hash: self.key.map(stable_key_hash),
}
}
}
#[inline]
#[must_use]
pub fn stable_key_hash(key: &[u8]) -> u64 {
const OFFSET: u64 = 0xcbf2_9ce4_8422_2325;
const PRIME: u64 = 0x0000_0100_0000_01b3;
let mut h = OFFSET;
for &b in key {
h ^= u64::from(b);
h = h.wrapping_mul(PRIME);
}
h
}
#[cfg(test)]
mod tests {
use super::*;
#[cfg(not(loom))]
#[test]
fn record_map_preserves_meta_and_ack() {
let (ack, _rx) = crate::checkpoint::AckRef::test_pair();
let rec = Record {
payload: 21u32,
meta: RecordMeta {
partition: PartitionId(3),
offset: 42,
event_time_ms: 1_000,
key_hash: Some(7),
},
ack,
};
let mapped = rec.map(|v| u64::from(v) * 2);
assert_eq!(mapped.payload, 42);
assert_eq!(mapped.meta.partition, PartitionId(3));
assert_eq!(mapped.meta.offset, 42);
assert_eq!(mapped.meta.key_hash, Some(7));
}
#[test]
fn key_hash_is_stable() {
assert_eq!(stable_key_hash(b""), 0xcbf2_9ce4_8422_2325);
assert_eq!(stable_key_hash(b"a"), 0xaf63_dc4c_8601_ec8c);
assert_ne!(stable_key_hash(b"user-1"), stable_key_hash(b"user-2"));
}
#[test]
fn raw_payload_meta_hashes_key() {
let payload = RawPayload {
bytes: b"v",
key: Some(b"k"),
partition: PartitionId(0),
offset: 9,
timestamp_ms: 5,
};
assert_eq!(payload.meta().key_hash, Some(stable_key_hash(b"k")));
let keyless = RawPayload {
key: None,
..payload
};
assert_eq!(keyless.meta().key_hash, None);
}
#[test]
fn meta_is_copy_and_small() {
assert!(size_of::<RecordMeta>() <= 40);
fn assert_copy<T: Copy>() {}
assert_copy::<RecordMeta>();
assert_copy::<Flow>();
assert_copy::<PartitionId>();
}
}