use std::borrow::Cow;
use reifydb_codec::key::{deserializer::KeyDeserializer, encoded::EncodedKey, serializer::KeySerializer};
use reifydb_macro::KeyCodec;
use reifydb_value::value::{datetime::DateTime, row_number::RowNumber};
use smallvec::{SmallVec, smallvec};
use super::KeyTag;
use crate::{
interface::catalog::id::QueueId,
key::{
any::{ByteEncoding, Field, KeyFields, Width},
bound::TaggedKeyBoundRange,
},
};
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = Queue)]
pub struct QueueKey {
pub queue: QueueId,
}
impl QueueKey {
pub fn new(queue: QueueId) -> Self {
Self {
queue,
}
}
pub fn encoded(queue: impl Into<QueueId>) -> EncodedKey {
Self::new(queue.into()).encode()
}
pub fn full_scan() -> TaggedKeyBoundRange {
TaggedKeyBoundRange::kind(Self::TAG)
}
}
#[cfg(test)]
mod queue_key_tests {
use std::ops::Bound;
use super::*;
#[test]
fn test_encode_decode_roundtrip() {
let encoded = QueueKey::encoded(QueueId(42));
let decoded = QueueKey::decode(&encoded).unwrap();
assert_eq!(decoded.queue, QueueId(42));
}
#[test]
fn test_decode_rejects_foreign_kind() {
let mut serializer = KeySerializer::with_capacity(9);
serializer.extend_u8(KeyTag::NamespaceQueue as u8).extend_u64(7u64);
assert!(QueueKey::decode(&serializer.to_encoded_key()).is_none());
}
#[test]
fn test_full_scan_brackets_every_queue_key() {
let range = QueueKey::full_scan().encode();
let Bound::Included(start) = &range.start else {
panic!("expected an included start bound")
};
let Bound::Included(end) = &range.end else {
panic!("expected an included end bound")
};
assert_eq!(start.as_slice(), &[!(KeyTag::Queue as u8)]);
assert_eq!(end.as_slice(), &[!(KeyTag::Queue as u8 - 1)]);
assert!(start.as_slice() < end.as_slice(), "the range must be non-empty under byte order");
for id in [QueueId(1), QueueId(u64::MAX)] {
let key = QueueKey::encoded(id);
assert!(
key.as_slice() >= start.as_slice() && key.as_slice() <= end.as_slice(),
"queue {id:?} must fall inside the scan range"
);
}
}
#[test]
fn test_full_scan_excludes_the_neighbouring_kind() {
let range = QueueKey::full_scan().encode();
let Bound::Included(start) = &range.start else {
panic!("expected an included start bound")
};
let Bound::Included(end) = &range.end else {
panic!("expected an included end bound")
};
let mut serializer = KeySerializer::with_capacity(9);
serializer.extend_u8(KeyTag::NamespaceQueue as u8).extend_u64(1u64);
let foreign = serializer.to_encoded_key();
assert!(
foreign.as_slice() < start.as_slice() || foreign.as_slice() > end.as_slice(),
"a NamespaceQueue key must fall outside the QueueKey scan range"
);
}
}
#[cfg(test)]
mod byte_identical_check_queue_key {
use reifydb_codec::key::serializer::KeySerializer;
use super::*;
fn legacy_encode(key: &QueueKey) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(9);
serializer.extend_u8(KeyTag::Queue as u8).extend_u64(key.queue);
serializer.to_encoded_key()
}
#[test]
fn matches_legacy_byte_layout() {
for id in [QueueId(0), QueueId(1), QueueId(u64::MAX)] {
let key = QueueKey {
queue: id,
};
assert_eq!(legacy_encode(&key).as_slice(), key.encode().as_slice());
}
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = QueueAttempt)]
pub struct QueueAttemptKey {
pub queue: QueueId,
pub row: RowNumber,
pub attempt: u32,
}
impl QueueAttemptKey {
pub fn new(queue: impl Into<QueueId>, row: impl Into<RowNumber>, attempt: u32) -> Self {
Self {
queue: queue.into(),
row: row.into(),
attempt,
}
}
pub fn encoded(queue: impl Into<QueueId>, row: impl Into<RowNumber>, attempt: u32) -> EncodedKey {
Self::new(queue, row, attempt).encode()
}
pub fn item_scan(queue: QueueId, row: RowNumber) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(
Self::TAG,
[Field::UDesc(Width::U64, queue.0 as u128), Field::UDesc(Width::U64, row.0 as u128)],
)
}
pub fn queue_scan(queue: QueueId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, [Field::UDesc(Width::U64, queue.0 as u128)])
}
pub fn full_scan() -> TaggedKeyBoundRange {
TaggedKeyBoundRange::kind(Self::TAG)
}
}
#[cfg(test)]
mod queue_item_state_key_tests {
use std::ops::Bound;
use reifydb_codec::key::encoded::EncodedKeyRange;
use super::*;
fn contains(range: &EncodedKeyRange, key: &EncodedKey) -> bool {
let after_start = match &range.start {
Bound::Included(start) => key.as_slice() >= start.as_slice(),
Bound::Excluded(start) => key.as_slice() > start.as_slice(),
Bound::Unbounded => true,
};
let before_end = match &range.end {
Bound::Included(end) => key.as_slice() <= end.as_slice(),
Bound::Excluded(end) => key.as_slice() < end.as_slice(),
Bound::Unbounded => true,
};
after_start && before_end
}
#[test]
fn test_attempt_key_roundtrips() {
let key = QueueAttemptKey {
queue: QueueId(7),
row: RowNumber(42),
attempt: u32::MAX,
};
assert_eq!(QueueAttemptKey::decode(&key.encode()), Some(key));
}
#[test]
fn test_attempt_zero_roundtrips() {
let key = QueueAttemptKey {
queue: QueueId(0),
row: RowNumber(0),
attempt: 0,
};
assert_eq!(QueueAttemptKey::decode(&key.encode()), Some(key));
}
#[test]
fn test_item_scan_excludes_neighbouring_items_and_queues() {
let range = QueueAttemptKey::item_scan(QueueId(3), RowNumber(5)).encode();
assert!(contains(&range, &QueueAttemptKey::encoded(QueueId(3), RowNumber(5), 0)));
assert!(contains(&range, &QueueAttemptKey::encoded(QueueId(3), RowNumber(5), u32::MAX)));
assert!(!contains(&range, &QueueAttemptKey::encoded(QueueId(3), RowNumber(6), 0)));
assert!(!contains(&range, &QueueAttemptKey::encoded(QueueId(3), RowNumber(4), 0)));
assert!(!contains(&range, &QueueAttemptKey::encoded(QueueId(4), RowNumber(5), 0)));
}
#[test]
fn test_queue_scan_covers_every_item_of_one_queue_only() {
let range = QueueAttemptKey::queue_scan(QueueId(3)).encode();
assert!(contains(&range, &QueueAttemptKey::encoded(QueueId(3), RowNumber(0), 0)));
assert!(contains(&range, &QueueAttemptKey::encoded(QueueId(3), RowNumber(u64::MAX), 9)));
assert!(!contains(&range, &QueueAttemptKey::encoded(QueueId(2), RowNumber(1), 0)));
assert!(!contains(&range, &QueueAttemptKey::encoded(QueueId(4), RowNumber(1), 0)));
}
#[test]
fn test_a_foreign_kind_does_not_decode() {
let foreign = QueueItemStateKey::encoded(QueueId(1), 0, RowNumber(1));
assert_eq!(QueueAttemptKey::decode(&foreign), None);
}
}
#[derive(Debug, Clone, PartialEq, Hash)]
pub struct QueueDeduplicationKey {
pub queue: QueueId,
pub tail: EncodedKey,
}
impl QueueDeduplicationKey {
pub fn new(queue: impl Into<QueueId>, tail: impl AsRef<[u8]>) -> Self {
Self {
queue: queue.into(),
tail: EncodedKey::new(tail),
}
}
pub fn encoded(queue: impl Into<QueueId>, tail: impl AsRef<[u8]>) -> EncodedKey {
Self::new(queue, tail).encode()
}
pub fn full_scan(queue: QueueId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, [Field::UDesc(Width::U64, queue.0 as u128)])
}
}
impl QueueDeduplicationKey {
pub const TAG: KeyTag = KeyTag::QueueDeduplication;
pub fn encode(&self) -> EncodedKey {
let mut serializer = KeySerializer::with_capacity(9 + self.tail.len() + 1);
serializer.extend_u8(Self::TAG as u8).extend_u64(self.queue).extend_bytes(&self.tail);
serializer.to_encoded_key()
}
pub fn decode(key: &EncodedKey) -> Option<Self> {
let mut de = KeyDeserializer::from_bytes(key.as_slice());
let kind: KeyTag = de.read_u8().ok()?.try_into().ok()?;
if kind != Self::TAG {
return None;
}
let queue = de.read_u64().ok()?;
let tail = de.read_bytes().ok()?;
Some(Self {
queue: QueueId(queue),
tail: EncodedKey::new(tail),
})
}
}
#[cfg(test)]
mod queue_deduplication_key_tests {
use std::ops::Bound;
use super::*;
#[test]
fn test_encode_decode_roundtrip() {
let encoded = QueueDeduplicationKey::encoded(QueueId(3), b"invoice-42".to_vec());
let decoded = QueueDeduplicationKey::decode(&encoded).unwrap();
assert_eq!(decoded.queue, QueueId(3));
assert_eq!(decoded.tail.as_slice(), b"invoice-42");
}
#[test]
fn test_arbitrary_bytes_survive_the_tail_encoding() {
for key in [
vec![],
vec![0x00],
vec![0xff],
vec![0xff, 0x00, 0xff],
"order/\u{00e9}\u{4e2d}".as_bytes().to_vec(),
] {
let encoded = QueueDeduplicationKey::encoded(QueueId(1), key.clone());
let decoded = QueueDeduplicationKey::decode(&encoded).unwrap();
assert_eq!(decoded.tail.as_slice(), key.as_slice(), "tail {key:?} must round-trip unchanged");
}
}
#[test]
fn test_the_same_key_in_two_queues_encodes_differently() {
let a = QueueDeduplicationKey::encoded(QueueId(1), b"same".to_vec());
let b = QueueDeduplicationKey::encoded(QueueId(2), b"same".to_vec());
assert_ne!(a, b);
}
#[test]
fn test_full_scan_contains_only_the_target_queue() {
let range = QueueDeduplicationKey::full_scan(QueueId(3)).encode();
let Bound::Included(start) = &range.start else {
panic!("expected an included start bound")
};
let Bound::Excluded(end) = &range.end else {
panic!("expected an excluded end bound")
};
assert!(start.as_slice() < end.as_slice(), "the range must be non-empty under byte order");
for key in [vec![], b"a".to_vec(), vec![0xff; 64]] {
let inside = QueueDeduplicationKey::encoded(QueueId(3), key.clone());
assert!(
inside.as_slice() >= start.as_slice() && inside.as_slice() < end.as_slice(),
"key {key:?} in queue 3 must fall inside the scan range"
);
}
for queue in [QueueId(2), QueueId(4)] {
let neighbour = QueueDeduplicationKey::encoded(queue, b"a".to_vec());
assert!(
neighbour.as_slice() < start.as_slice() || neighbour.as_slice() >= end.as_slice(),
"queue {queue:?} must fall outside queue 3's scan range"
);
}
}
#[test]
fn test_a_foreign_or_truncated_key_does_not_decode() {
let encoded = QueueDeduplicationKey::encoded(QueueId(3), b"invoice-42".to_vec());
let mut wrong_kind = encoded.as_slice().to_vec();
wrong_kind[0] = KeyTag::Queue as u8;
assert_eq!(QueueDeduplicationKey::decode(&EncodedKey::new(wrong_kind)), None);
let truncated = encoded.as_slice()[..5].to_vec();
assert_eq!(QueueDeduplicationKey::decode(&EncodedKey::new(truncated)), None);
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = QueuePartition)]
pub struct QueuePartitionKey {
pub queue: QueueId,
pub partition: u16,
}
impl QueuePartitionKey {
pub fn new(queue: impl Into<QueueId>, partition: u16) -> Self {
Self {
queue: queue.into(),
partition,
}
}
pub fn encoded(queue: impl Into<QueueId>, partition: u16) -> EncodedKey {
Self {
queue: queue.into(),
partition,
}
.encode()
}
pub fn queue_scan(queue: QueueId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, [Field::UDesc(Width::U64, queue.0 as u128)])
}
pub fn full_scan() -> TaggedKeyBoundRange {
TaggedKeyBoundRange::kind(Self::TAG)
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = QueueItemState)]
pub struct QueueItemStateKey {
pub queue: QueueId,
pub partition: u16,
pub row: RowNumber,
}
impl QueueItemStateKey {
pub fn new(queue: impl Into<QueueId>, partition: u16, row: impl Into<RowNumber>) -> Self {
Self {
queue: queue.into(),
partition,
row: row.into(),
}
}
pub fn encoded(queue: impl Into<QueueId>, partition: u16, row: impl Into<RowNumber>) -> EncodedKey {
Self {
queue: queue.into(),
partition,
row: row.into(),
}
.encode()
}
pub fn partition_scan(queue: QueueId, partition: u16) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(
Self::TAG,
[Field::UDesc(Width::U64, queue.0 as u128), Field::UDesc(Width::U16, partition as u128)],
)
}
pub fn queue_scan(queue: QueueId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, [Field::UDesc(Width::U64, queue.0 as u128)])
}
pub fn full_scan() -> TaggedKeyBoundRange {
TaggedKeyBoundRange::kind(Self::TAG)
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = QueueDue)]
pub struct QueueDueKey {
pub queue: QueueId,
pub partition: u16,
pub due: DateTime,
pub row: RowNumber,
}
impl QueueDueKey {
pub fn new(queue: impl Into<QueueId>, partition: u16, due: DateTime, row: impl Into<RowNumber>) -> Self {
Self {
queue: queue.into(),
partition,
due,
row: row.into(),
}
}
pub fn encoded(
queue: impl Into<QueueId>,
partition: u16,
due: DateTime,
row: impl Into<RowNumber>,
) -> EncodedKey {
Self::new(queue, partition, due, row).encode()
}
pub fn partition_scan(queue: QueueId, partition: u16) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(
Self::TAG,
[Field::UDesc(Width::U64, queue.0 as u128), Field::UDesc(Width::U16, partition as u128)],
)
}
pub fn queue_scan(queue: QueueId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, [Field::UDesc(Width::U64, queue.0 as u128)])
}
pub fn full_scan() -> TaggedKeyBoundRange {
TaggedKeyBoundRange::kind(Self::TAG)
}
}
#[derive(Debug, Clone, PartialEq, KeyCodec, Hash)]
#[key(tag = QueueKeyActive)]
pub struct QueueKeyActiveKey {
pub queue: QueueId,
pub partition: u16,
pub key_hash: u64,
pub row: RowNumber,
}
impl QueueKeyActiveKey {
pub fn new(queue: impl Into<QueueId>, partition: u16, key_hash: u64, row: impl Into<RowNumber>) -> Self {
Self {
queue: queue.into(),
partition,
key_hash,
row: row.into(),
}
}
pub fn encoded(
queue: impl Into<QueueId>,
partition: u16,
key_hash: u64,
row: impl Into<RowNumber>,
) -> EncodedKey {
Self {
queue: queue.into(),
partition,
key_hash,
row: row.into(),
}
.encode()
}
pub fn key_scan(queue: QueueId, partition: u16, key_hash: u64) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(
Self::TAG,
[
Field::UDesc(Width::U64, queue.0 as u128),
Field::UDesc(Width::U16, partition as u128),
Field::UDesc(Width::U64, key_hash as u128),
],
)
}
pub fn partition_scan(queue: QueueId, partition: u16) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(
Self::TAG,
[Field::UDesc(Width::U64, queue.0 as u128), Field::UDesc(Width::U16, partition as u128)],
)
}
pub fn queue_scan(queue: QueueId) -> TaggedKeyBoundRange {
TaggedKeyBoundRange::prefix(Self::TAG, [Field::UDesc(Width::U64, queue.0 as u128)])
}
pub fn full_scan() -> TaggedKeyBoundRange {
TaggedKeyBoundRange::kind(Self::TAG)
}
}
#[cfg(test)]
mod queue_partition_key_tests {
use std::ops::Bound;
use reifydb_codec::key::encoded::EncodedKeyRange;
use super::*;
fn contains(range: &EncodedKeyRange, key: &EncodedKey) -> bool {
let after_start = match &range.start {
Bound::Included(start) => key.as_slice() >= start.as_slice(),
Bound::Excluded(start) => key.as_slice() > start.as_slice(),
Bound::Unbounded => true,
};
let before_end = match &range.end {
Bound::Included(end) => key.as_slice() <= end.as_slice(),
Bound::Excluded(end) => key.as_slice() < end.as_slice(),
Bound::Unbounded => true,
};
after_start && before_end
}
#[test]
fn test_partition_key_roundtrips_at_both_partition_bounds() {
for partition in [0u16, 1, 1023] {
let encoded = QueuePartitionKey::encoded(QueueId(7), partition);
let decoded = QueuePartitionKey::decode(&encoded).unwrap();
assert_eq!(decoded.queue, QueueId(7));
assert_eq!(decoded.partition, partition);
}
}
#[test]
fn test_item_state_key_roundtrips() {
let encoded = QueueItemStateKey::encoded(QueueId(3), 5, RowNumber(42));
let decoded = QueueItemStateKey::decode(&encoded).unwrap();
assert_eq!(decoded.queue, QueueId(3));
assert_eq!(decoded.partition, 5);
assert_eq!(decoded.row, RowNumber(42));
}
#[test]
fn test_due_key_roundtrips_at_epoch_and_far_future() {
for nanos in [0u64, 1, 4_102_444_800_000_000_000, u64::MAX] {
let due = DateTime::from_nanos(nanos);
let encoded = QueueDueKey::encoded(QueueId(1), 2, due, RowNumber(9));
let decoded = QueueDueKey::decode(&encoded).unwrap();
assert_eq!(decoded.due.to_nanos(), nanos);
assert_eq!(decoded.row, RowNumber(9));
assert_eq!(decoded.partition, 2);
}
}
#[test]
fn test_partition_scan_excludes_neighbouring_partitions() {
for partition in [0u16, 1, 1023] {
let range = QueueItemStateKey::partition_scan(QueueId(4), partition).encode();
assert!(contains(&range, &QueueItemStateKey::encoded(QueueId(4), partition, RowNumber(0))));
assert!(contains(
&range,
&QueueItemStateKey::encoded(QueueId(4), partition, RowNumber(u64::MAX))
));
for other in [partition.wrapping_sub(1), partition + 1] {
if other == partition {
continue;
}
let neighbour = QueueItemStateKey::encoded(QueueId(4), other, RowNumber(1));
assert!(
!contains(&range, &neighbour),
"partition {other} must fall outside {partition}"
);
}
let other_queue = QueueItemStateKey::encoded(QueueId(5), partition, RowNumber(1));
assert!(!contains(&range, &other_queue), "queue 5 must fall outside queue 4's partition scan");
}
}
#[test]
fn test_queue_scan_covers_every_partition_of_one_queue_only() {
let range = QueueDueKey::queue_scan(QueueId(4)).encode();
for partition in [0u16, 1, 1023] {
let inside = QueueDueKey::encoded(QueueId(4), partition, DateTime::from_nanos(7), RowNumber(1));
assert!(contains(&range, &inside), "partition {partition} must fall inside the queue scan");
}
for queue in [QueueId(3), QueueId(5)] {
let outside = QueueDueKey::encoded(queue, 0, DateTime::from_nanos(7), RowNumber(1));
assert!(!contains(&range, &outside), "queue {queue:?} must fall outside queue 4's scan");
}
}
#[test]
fn test_due_keys_sort_latest_due_first() {
let earlier = QueueDueKey::encoded(QueueId(1), 0, DateTime::from_nanos(1_000), RowNumber(1));
let later = QueueDueKey::encoded(QueueId(1), 0, DateTime::from_nanos(2_000), RowNumber(1));
assert!(later.as_slice() < earlier.as_slice(), "the later due time must encode to the smaller key");
}
#[test]
fn test_a_foreign_kind_does_not_decode() {
let encoded = QueueItemStateKey::encoded(QueueId(1), 0, RowNumber(1));
assert_eq!(QueuePartitionKey::decode(&encoded), None);
assert_eq!(QueueDueKey::decode(&encoded), None);
assert_eq!(QueueKeyActiveKey::decode(&EncodedKey::new(encoded.as_slice()[..3].to_vec())), None);
}
#[test]
fn test_key_active_key_roundtrips() {
let encoded = QueueKeyActiveKey::encoded(QueueId(3), 5, 0xDEAD_BEEF_CAFE_F00D, RowNumber(42));
let decoded = QueueKeyActiveKey::decode(&encoded).unwrap();
assert_eq!(decoded.queue, QueueId(3));
assert_eq!(decoded.partition, 5);
assert_eq!(decoded.key_hash, 0xDEAD_BEEF_CAFE_F00D);
assert_eq!(decoded.row, RowNumber(42));
}
#[test]
fn test_key_active_keys_sort_largest_row_first() {
let first = QueueKeyActiveKey::encoded(QueueId(1), 0, 77, RowNumber(1));
let middle = QueueKeyActiveKey::encoded(QueueId(1), 0, 77, RowNumber(5));
let last = QueueKeyActiveKey::encoded(QueueId(1), 0, 77, RowNumber(9));
assert!(last.as_slice() < middle.as_slice());
assert!(middle.as_slice() < first.as_slice());
}
#[test]
fn test_key_scan_excludes_neighbouring_keys_and_partitions() {
let range = QueueKeyActiveKey::key_scan(QueueId(4), 2, 77).encode();
assert!(contains(&range, &QueueKeyActiveKey::encoded(QueueId(4), 2, 77, RowNumber(0))));
assert!(contains(&range, &QueueKeyActiveKey::encoded(QueueId(4), 2, 77, RowNumber(u64::MAX))));
for other_hash in [76u64, 78, 0, u64::MAX] {
let neighbour = QueueKeyActiveKey::encoded(QueueId(4), 2, other_hash, RowNumber(1));
assert!(!contains(&range, &neighbour), "key hash {other_hash} must fall outside key 77");
}
let other_partition = QueueKeyActiveKey::encoded(QueueId(4), 3, 77, RowNumber(1));
assert!(!contains(&range, &other_partition), "partition 3 must fall outside partition 2");
let other_queue = QueueKeyActiveKey::encoded(QueueId(5), 2, 77, RowNumber(1));
assert!(!contains(&range, &other_queue), "queue 5 must fall outside queue 4");
}
#[test]
fn test_partition_scan_covers_every_key_of_one_partition_only() {
let range = QueueKeyActiveKey::partition_scan(QueueId(4), 2).encode();
for key_hash in [0u64, 77, u64::MAX] {
let inside = QueueKeyActiveKey::encoded(QueueId(4), 2, key_hash, RowNumber(1));
assert!(contains(&range, &inside), "key hash {key_hash} must fall inside the partition scan");
}
let other_partition = QueueKeyActiveKey::encoded(QueueId(4), 3, 77, RowNumber(1));
assert!(!contains(&range, &other_partition));
}
#[test]
fn test_schedule_keys_match_legacy_byte_layout() {
let queue = QueueId(7);
let partition = 3u16;
let row = RowNumber(42);
let mut legacy = KeySerializer::with_capacity(11);
legacy.extend_u8(KeyTag::QueuePartition as u8).extend_u64(queue).extend_u16(partition);
assert_eq!(legacy.to_encoded_key().as_slice(), QueuePartitionKey::encoded(queue, partition).as_slice());
let mut legacy = KeySerializer::with_capacity(19);
legacy.extend_u8(KeyTag::QueueItemState as u8)
.extend_u64(queue)
.extend_u16(partition)
.extend_u64(row.0);
assert_eq!(
legacy.to_encoded_key().as_slice(),
QueueItemStateKey::encoded(queue, partition, row).as_slice()
);
let due = DateTime::from_nanos(1_000);
let mut legacy = KeySerializer::with_capacity(27);
legacy.extend_u8(KeyTag::QueueDue as u8)
.extend_u64(queue)
.extend_u16(partition)
.extend_datetime(&due)
.extend_u64(row.0);
assert_eq!(
legacy.to_encoded_key().as_slice(),
QueueDueKey::encoded(queue, partition, due, row).as_slice()
);
let key_hash = 0xDEAD_BEEFu64;
let mut legacy = KeySerializer::with_capacity(28);
legacy.extend_u8(KeyTag::QueueKeyActive as u8)
.extend_u64(queue)
.extend_u16(partition)
.extend_u64(key_hash)
.extend_u64(row.0);
assert_eq!(
legacy.to_encoded_key().as_slice(),
QueueKeyActiveKey::encoded(queue, partition, key_hash, row).as_slice()
);
}
}
impl KeyFields for QueueDeduplicationKey {
fn fields(&self) -> SmallVec<[Field<'_>; 6]> {
smallvec![
Field::UDesc(Width::U64, self.queue.0 as u128),
Field::BytesDesc(ByteEncoding::Escaped, Cow::Borrowed(self.tail.as_slice())),
]
}
}