use crate::cf::{TableMutation, TableMutations};
use crate::kvs::Key;
use crate::statements::DefineTableStatement;
use crate::thing::Thing;
use crate::value::Value;
use std::borrow::Cow;
use std::collections::HashMap;
type PreparedWrite = (Vec<u8>, Vec<u8>, Vec<u8>, crate::kvs::Val);
pub struct Writer {
buf: Buffer,
}
pub struct Buffer {
pub b: HashMap<ChangeKey, TableMutations>,
}
#[derive(Hash, Eq, PartialEq, Debug)]
pub struct ChangeKey {
pub ns: String,
pub db: String,
pub tb: String,
}
impl Buffer {
pub fn new() -> Self {
Self {
b: HashMap::new(),
}
}
pub fn push(&mut self, ns: String, db: String, tb: String, m: TableMutation) {
let tb2 = tb.clone();
let ms = self
.b
.entry(ChangeKey {
ns,
db,
tb,
})
.or_insert(TableMutations::new(tb2));
ms.1.push(m);
}
}
impl Writer {
pub(crate) fn new() -> Self {
Self {
buf: Buffer::new(),
}
}
pub(crate) fn update(&mut self, ns: &str, db: &str, tb: &str, id: Thing, v: Cow<'_, Value>) {
if v.is_some() {
self.buf.push(
ns.to_string(),
db.to_string(),
tb.to_string(),
TableMutation::Set(id, v.into_owned()),
);
} else {
self.buf.push(ns.to_string(), db.to_string(), tb.to_string(), TableMutation::Del(id));
}
}
pub(crate) fn define_table(&mut self, ns: &str, db: &str, tb: &str, dt: &DefineTableStatement) {
self.buf.push(
ns.to_string(),
db.to_string(),
tb.to_string(),
TableMutation::Def(dt.to_owned()),
)
}
pub(crate) fn get(&self) -> Vec<PreparedWrite> {
let mut r = Vec::<(Vec<u8>, Vec<u8>, Vec<u8>, crate::kvs::Val)>::new();
for (
ChangeKey {
ns,
db,
tb,
},
mutations,
) in self.buf.b.iter()
{
let ts_key: Key = crate::key::database::vs::new(ns, db).into();
let tc_key_prefix: Key = crate::key::change::versionstamped_key_prefix(ns, db);
let tc_key_suffix: Key = crate::key::change::versionstamped_key_suffix(tb.as_str());
r.push((ts_key, tc_key_prefix, tc_key_suffix, mutations.into()))
}
r
}
}
#[cfg(test)]
mod tests {
use std::borrow::Cow;
use std::time::Duration;
use crate::cf::{ChangeSet, DatabaseMutation, TableMutation, TableMutations};
use crate::changefeed::ChangeFeed;
use crate::id::Id;
use crate::key::key_req::KeyRequirements;
use crate::kvs::{Datastore, LockType::*, TransactionType::*};
use crate::statements::show::ShowSince;
use crate::statements::{
DefineDatabaseStatement, DefineNamespaceStatement, DefineTableStatement,
};
use crate::thing::Thing;
use crate::value::Value;
use crate::vs;
#[tokio::test]
async fn test_changefeed_read_write() {
let ts = crate::Datetime::default();
let ns = "myns";
let db = "mydb";
let tb = "mytb";
let dns = DefineNamespaceStatement {
name: crate::Ident(ns.to_string()),
..Default::default()
};
let ddb = DefineDatabaseStatement {
name: crate::Ident(db.to_string()),
changefeed: Some(ChangeFeed {
expiry: Duration::from_secs(10),
}),
..Default::default()
};
let dtb = DefineTableStatement {
name: tb.into(),
changefeed: Some(ChangeFeed {
expiry: Duration::from_secs(10),
}),
..Default::default()
};
let ds = Datastore::new("memory").await.unwrap();
let mut tx0 = ds.transaction(Write, Optimistic).await.unwrap();
let ns_root = crate::key::root::ns::new(ns);
tx0.put(ns_root.key_category(), &ns_root, dns).await.unwrap();
let db_root = crate::key::namespace::db::new(ns, db);
tx0.put(db_root.key_category(), &db_root, ddb).await.unwrap();
let tb_root = crate::key::database::tb::new(ns, db, tb);
tx0.put(tb_root.key_category(), &tb_root, dtb.clone()).await.unwrap();
tx0.commit().await.unwrap();
ds.tick_at(ts.0.timestamp().try_into().unwrap()).await.unwrap();
let mut tx1 = ds.transaction(Write, Optimistic).await.unwrap();
let thing_a = Thing {
tb: tb.to_owned(),
id: Id::String("A".to_string()),
};
let value_a: super::Value = "a".into();
tx1.record_change(ns, db, tb, &thing_a, Cow::Borrowed(&value_a));
tx1.complete_changes(true).await.unwrap();
tx1.commit().await.unwrap();
let mut tx2 = ds.transaction(Write, Optimistic).await.unwrap();
let thing_c = Thing {
tb: tb.to_owned(),
id: Id::String("C".to_string()),
};
let value_c: Value = "c".into();
tx2.record_change(ns, db, tb, &thing_c, Cow::Borrowed(&value_c));
tx2.complete_changes(true).await.unwrap();
tx2.commit().await.unwrap();
let x = ds.transaction(Write, Optimistic).await;
let mut tx3 = x.unwrap();
let thing_b = Thing {
tb: tb.to_owned(),
id: Id::String("B".to_string()),
};
let value_b: Value = "b".into();
tx3.record_change(ns, db, tb, &thing_b, Cow::Borrowed(&value_b));
let thing_c2 = Thing {
tb: tb.to_owned(),
id: Id::String("C".to_string()),
};
let value_c2: Value = "c2".into();
tx3.record_change(ns, db, tb, &thing_c2, Cow::Borrowed(&value_c2));
tx3.complete_changes(true).await.unwrap();
tx3.commit().await.unwrap();
let start: u64 = 0;
let mut tx4 = ds.transaction(Write, Optimistic).await.unwrap();
let r =
crate::cf::read(&mut tx4, ns, db, Some(tb), ShowSince::Versionstamp(start), Some(10))
.await
.unwrap();
tx4.commit().await.unwrap();
let want: Vec<ChangeSet> = vec![
ChangeSet(
vs::u64_to_versionstamp(2),
DatabaseMutation(vec![TableMutations(
"mytb".to_string(),
vec![TableMutation::Set(
Thing::from(("mytb".to_string(), "A".to_string())),
Value::from("a"),
)],
)]),
),
ChangeSet(
vs::u64_to_versionstamp(3),
DatabaseMutation(vec![TableMutations(
"mytb".to_string(),
vec![TableMutation::Set(
Thing::from(("mytb".to_string(), "C".to_string())),
Value::from("c"),
)],
)]),
),
ChangeSet(
vs::u64_to_versionstamp(4),
DatabaseMutation(vec![TableMutations(
"mytb".to_string(),
vec![
TableMutation::Set(
Thing::from(("mytb".to_string(), "B".to_string())),
Value::from("b"),
),
TableMutation::Set(
Thing::from(("mytb".to_string(), "C".to_string())),
Value::from("c2"),
),
],
)]),
),
];
assert_eq!(r, want);
let mut tx5 = ds.transaction(Write, Optimistic).await.unwrap();
crate::cf::gc_db(&mut tx5, ns, db, vs::u64_to_versionstamp(4), Some(10)).await.unwrap();
tx5.commit().await.unwrap();
let mut tx6 = ds.transaction(Write, Optimistic).await.unwrap();
let r =
crate::cf::read(&mut tx6, ns, db, Some(tb), ShowSince::Versionstamp(start), Some(10))
.await
.unwrap();
tx6.commit().await.unwrap();
let want: Vec<ChangeSet> = vec![ChangeSet(
vs::u64_to_versionstamp(4),
DatabaseMutation(vec![TableMutations(
"mytb".to_string(),
vec![
TableMutation::Set(
Thing::from(("mytb".to_string(), "B".to_string())),
Value::from("b"),
),
TableMutation::Set(
Thing::from(("mytb".to_string(), "C".to_string())),
Value::from("c2"),
),
],
)]),
)];
assert_eq!(r, want);
ds.tick_at((ts.0.timestamp() + 5).try_into().unwrap()).await.unwrap();
let mut tx7 = ds.transaction(Write, Optimistic).await.unwrap();
let r = crate::cf::read(&mut tx7, ns, db, Some(tb), ShowSince::Timestamp(ts), Some(10))
.await
.unwrap();
tx7.commit().await.unwrap();
assert_eq!(r, want);
}
}