use anyhow::Result;
use crate::catalog::{DatabaseId, NamespaceId};
use crate::cf::{ChangeSet, DatabaseMutation, TableMutations};
use crate::err::Error;
use crate::expr::statements::show::ShowSince;
use crate::key::change;
#[cfg(debug_assertions)]
use crate::key::debug::Sprintable;
use crate::kvs::{KVKey, KVValue, Transaction};
use crate::val::TableName;
pub async fn read(
tx: &Transaction,
ns: NamespaceId,
db: DatabaseId,
tb: Option<&TableName>,
start: ShowSince,
limit: Option<u32>,
) -> Result<Vec<ChangeSet>> {
let ts_impl = tx.timestamp_impl();
let ts = match start {
ShowSince::Versionstamp(x) => {
ts_impl.create_from_versionstamp(x as u128).ok_or_else(|| Error::Query {
message: format!(
"Invalid versionstamp `{x}`, outside of range for kv-store timestamps"
),
})?
}
ShowSince::Timestamp(x) => {
ts_impl.create_from_datetime(x.0).ok_or_else(|| Error::Query {
message: format!(
"Invalid versionstamp `{x}`, outside of range for kv-store timestamps"
),
})?
}
};
let buf = &mut [0u8; _];
let ts_bytes = ts.encode(buf);
let beg = change::prefix_ts(ns, db, ts_bytes).encode_key()?;
let end = change::suffix(ns, db).encode_key()?;
let limit = limit.unwrap_or(100).min(1000);
let mut current_ts: Option<Vec<u8>> = None;
let mut buf: Vec<TableMutations> = Vec::new();
let mut res = Vec::<ChangeSet>::new();
#[cfg(debug_assertions)]
let mut prev_ts: Option<Vec<u8>> = None;
for (k, v) in tx.scan(beg..end, limit, 0, None).await? {
#[cfg(debug_assertions)]
trace!("Reading change feed entry: {}", k.sprint());
let key = crate::key::change::Cf::decode_key(&k)?;
#[cfg(debug_assertions)]
{
if let Some(p) = &prev_ts {
assert!(
key.ts.as_ref() >= p.as_slice(),
"changefeed scan returned out-of-order versionstamps"
);
}
prev_ts = Some(key.ts.to_vec());
}
if tb.is_some_and(|tb| *tb != *key.tb) {
continue;
}
let tb_muts = TableMutations::kv_decode_value(&v, ())?;
match current_ts {
Some(ref x) => {
if key.ts != x.as_slice() {
let db_mut = DatabaseMutation(buf);
let version = ts_impl.decode(x)?.as_versionstamp();
res.push(ChangeSet(version, db_mut));
buf = Vec::new();
current_ts = Some(key.ts.into_owned())
}
}
None => {
current_ts = Some(key.ts.into_owned());
}
}
buf.push(tb_muts);
}
if !buf.is_empty() {
let db_mut = DatabaseMutation(buf);
let ts_bytes = current_ts.expect("timestamp should be set when mutations exist");
let version = ts_impl.decode(ts_bytes.as_slice())?.as_versionstamp();
res.push(ChangeSet(version, db_mut));
}
Ok(res)
}