use std::borrow::Cow;
use anyhow::Result;
use storekey::{BorrowDecode, Encode};
use crate::catalog::{DatabaseId, NamespaceId};
use crate::key::category::{Categorise, Category};
use crate::kvs::impl_kv_key_storekey;
use crate::lq::event::LiveEvents;
use crate::val::TableName;
#[derive(Clone, Debug, Eq, PartialEq, PartialOrd, Encode, BorrowDecode)]
#[storekey(format = "()")]
pub(crate) struct Lqe<'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!(Lqe<'_> => LiveEvents);
impl Categorise for Lqe<'_> {
fn categorise(&self) -> Category {
Category::LiveQueryEvent
}
}
impl<'a> Lqe<'a> {
pub fn new(ns: NamespaceId, db: DatabaseId, ts: &'a [u8], tb: &'a TableName) -> Self {
Lqe {
__: 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<Lqe<'_>> {
Ok(storekey::decode_borrow(k)?)
}
}
pub fn new<'a>(ns: NamespaceId, db: DatabaseId, ts: &'a [u8], tb: &'a TableName) -> Lqe<'a> {
Lqe::new(ns, db, ts, tb)
}
#[derive(Clone, Debug, Eq, PartialEq, PartialOrd, Encode, BorrowDecode)]
pub struct LqeTsRange<'a> {
__: u8,
_a: u8,
pub ns: NamespaceId,
_b: u8,
pub db: DatabaseId,
_c: u8,
pub ts: Cow<'a, [u8]>,
}
impl<'a> LqeTsRange<'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!(LqeTsRange<'_> => LiveEvents);
pub fn prefix_ts(ns: NamespaceId, db: DatabaseId, ts: &[u8]) -> LqeTsRange<'_> {
LqeTsRange::new(ns, db, ts)
}
#[derive(Clone, Debug, Eq, PartialEq, PartialOrd, Encode, BorrowDecode)]
pub struct LqeRange {
__: u8,
_a: u8,
pub ns: NamespaceId,
_b: u8,
pub db: DatabaseId,
_c: u8,
_xx: u8,
}
impl LqeRange {
pub fn new_suffix(ns: NamespaceId, db: DatabaseId) -> Self {
Self {
__: b'/',
_a: b'*',
ns,
_b: b'*',
db,
_c: b'%',
_xx: 0xff,
}
}
#[cfg(test)]
pub fn new_prefix(ns: NamespaceId, db: DatabaseId) -> Self {
Self {
__: b'/',
_a: b'*',
ns,
_b: b'*',
db,
_c: b'%',
_xx: 0x00,
}
}
}
impl_kv_key_storekey!(LqeRange => LiveEvents);
pub fn suffix(ns: NamespaceId, db: DatabaseId) -> LqeRange {
LqeRange::new_suffix(ns, db)
}
#[cfg(test)]
pub fn prefix(ns: NamespaceId, db: DatabaseId) -> LqeRange {
LqeRange::new_prefix(ns, db)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::kvs::{HlcTimeStampImpl, KVKey, TimeStampImpl};
#[test]
fn lqe_key_uses_percent_section_marker() {
let ts_impl = HlcTimeStampImpl;
let buf = &mut [0u8; _];
let ts = ts_impl.create_from_versionstamp(12345).unwrap().encode(buf);
let tb = TableName::from("test");
let enc = Lqe::new(NamespaceId(1), DatabaseId(2), ts, &tb).encode_key().unwrap();
assert_eq!(&enc[0..2], b"/*");
assert_eq!(enc[11], b'%', "section marker must be '%', distinct from changefeed '#'");
}
#[test]
fn lqe_range_is_disjoint_from_changefeed_range() {
let ts_impl = HlcTimeStampImpl;
let buf = &mut [0u8; _];
let ts = ts_impl.create_from_versionstamp(1).unwrap().encode(buf);
let cf_suffix =
crate::key::change::suffix(NamespaceId(1), DatabaseId(2)).encode_key().unwrap();
let lqe = LqeTsRange::new(NamespaceId(1), DatabaseId(2), ts).encode_key().unwrap();
assert!(lqe > cf_suffix, "lqe keys must sort after the changefeed suffix");
}
}