use std::borrow::Cow;
use anyhow::Result;
use storekey::{BorrowDecode, Encode};
use crate::catalog::{DatabaseId, NamespaceId};
use crate::cf::TableMutations;
use crate::key::category::{Categorise, Category};
use crate::kvs::impl_kv_key_storekey;
use crate::val::TableName;
#[derive(Clone, Debug, Eq, PartialEq, PartialOrd, Encode, BorrowDecode)]
#[storekey(format = "()")]
pub(crate) struct Cf<'a> {
__: u8,
_a: u8,
pub ns: NamespaceId,
_b: u8,
pub db: DatabaseId,
_d: u8,
pub ts: Cow<'a, [u8]>,
_c: u8,
pub tb: Cow<'a, TableName>,
}
impl_kv_key_storekey!(Cf<'_> => TableMutations);
impl Categorise for Cf<'_> {
fn categorise(&self) -> Category {
Category::ChangeFeed
}
}
impl<'a> Cf<'a> {
pub fn new(ns: NamespaceId, db: DatabaseId, ts: &'a [u8], tb: &'a TableName) -> Self {
Cf {
__: b'/',
_a: b'*',
ns,
_b: b'*',
db,
_d: b'#',
ts: Cow::Borrowed(ts),
_c: b'*',
tb: Cow::Borrowed(tb),
}
}
pub fn decode_key(k: &[u8]) -> Result<Cf<'_>> {
Ok(storekey::decode_borrow(k)?)
}
}
pub fn new<'a>(ns: NamespaceId, db: DatabaseId, ts: &'a [u8], tb: &'a TableName) -> Cf<'a> {
Cf::new(ns, db, ts, tb)
}
#[derive(Clone, Debug, Eq, PartialEq, PartialOrd, Encode, BorrowDecode)]
pub struct DatabaseChangeFeedRange {
__: u8,
_a: u8,
pub ns: NamespaceId,
_b: u8,
pub db: DatabaseId,
_c: u8,
_xx: u8,
}
impl DatabaseChangeFeedRange {
pub fn new_prefix(ns: NamespaceId, db: DatabaseId) -> Self {
Self {
__: b'/',
_a: b'*',
ns,
_b: b'*',
db,
_c: b'#',
_xx: 0x00,
}
}
pub fn new_suffix(ns: NamespaceId, db: DatabaseId) -> Self {
Self {
__: b'/',
_a: b'*',
ns,
_b: b'*',
db,
_c: b'#',
_xx: 0xff,
}
}
}
impl_kv_key_storekey!(DatabaseChangeFeedRange => Vec<u8>);
#[derive(Clone, Debug, Eq, PartialEq, PartialOrd, Encode, BorrowDecode)]
pub struct DatabaseChangeFeedTsRange<'a> {
__: u8,
_a: u8,
pub ns: NamespaceId,
_b: u8,
pub db: DatabaseId,
_c: u8,
pub ts: Cow<'a, [u8]>,
}
impl<'a> DatabaseChangeFeedTsRange<'a> {
pub fn new(ns: NamespaceId, db: DatabaseId, ts: &'a [u8]) -> Self {
Self {
__: b'/',
_a: b'*',
ns,
_b: b'*',
db,
_c: b'#',
ts: Cow::Borrowed(ts),
}
}
}
impl_kv_key_storekey!(DatabaseChangeFeedTsRange<'_> => TableMutations);
pub fn prefix_ts(ns: NamespaceId, db: DatabaseId, ts: &[u8]) -> DatabaseChangeFeedTsRange<'_> {
DatabaseChangeFeedTsRange::new(ns, db, ts)
}
#[expect(unused)]
pub fn prefix(ns: NamespaceId, db: DatabaseId) -> DatabaseChangeFeedRange {
DatabaseChangeFeedRange::new_prefix(ns, db)
}
pub fn suffix(ns: NamespaceId, db: DatabaseId) -> DatabaseChangeFeedRange {
DatabaseChangeFeedRange::new_suffix(ns, db)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::kvs::{HlcTimeStampImpl, KVKey, TimeStampImpl};
#[test]
fn cf_key() {
let ts_impl = HlcTimeStampImpl;
let buf = &mut [0u8; _];
let ts1 = ts_impl.create_from_versionstamp(12345).unwrap().encode(buf);
let tb = TableName::from("test");
let val = Cf::new(NamespaceId(1), DatabaseId(2), ts1, &tb);
let enc = Cf::encode_key(&val).unwrap();
assert_eq!(
enc,
&[
47, 42, 0, 0, 0, 1, 42, 0, 0, 0, 2, 35, 1, 0, 1, 0, 1, 0, 1, 0, 1, 0, 1, 0, 48, 57,
0, 42, 116, 101, 115, 116, 0
]
);
let buf = &mut [0; _];
let ts2 = ts_impl.create_from_versionstamp(12346).unwrap().encode(buf);
let val = Cf::new(NamespaceId(1), DatabaseId(2), ts2, &tb);
let enc = Cf::encode_key(&val).unwrap();
assert_eq!(
enc,
&[
47, 42, 0, 0, 0, 1, 42, 0, 0, 0, 2, 35, 1, 0, 1, 0, 1, 0, 1, 0, 1, 0, 1, 0, 48, 58,
0, 42, 116, 101, 115, 116, 0
]
);
}
#[test]
fn range_key() {
let val = DatabaseChangeFeedRange::new_prefix(NamespaceId(1), DatabaseId(2));
let enc = DatabaseChangeFeedRange::encode_key(&val).unwrap();
assert_eq!(enc, b"/*\x00\x00\x00\x01*\x00\x00\x00\x02#\x00");
let val = DatabaseChangeFeedRange::new_suffix(NamespaceId(1), DatabaseId(2));
let enc = DatabaseChangeFeedRange::encode_key(&val).unwrap();
assert_eq!(enc, b"/*\x00\x00\x00\x01*\x00\x00\x00\x02#\xff");
}
#[test]
fn ts_prefix_key() {
let ts_impl = HlcTimeStampImpl;
let buf = &mut [0u8; _];
let ts = ts_impl.create_from_versionstamp(12345).unwrap().encode(buf);
let val = DatabaseChangeFeedTsRange::new(NamespaceId(1), DatabaseId(2), ts);
let enc = DatabaseChangeFeedTsRange::encode_key(&val).unwrap();
assert_eq!(
enc,
&[
47, 42, 0, 0, 0, 1, 42, 0, 0, 0, 2, 35, 1, 0, 1, 0, 1, 0, 1, 0, 1, 0, 1, 0, 48, 57,
0
]
);
}
}