use std::collections::HashMap;
use anyhow::Result;
use parking_lot::Mutex;
use crate::catalog::{DatabaseId, NamespaceId, TableDefinition};
use crate::cf::TableMutations;
use crate::doc::CursorRecord;
use crate::kvs::KVValue;
use crate::val::{RecordId, TableName};
type PreparedWrite = (NamespaceId, DatabaseId, TableName, crate::kvs::Val);
#[derive(Hash, Eq, PartialEq, Debug)]
pub struct ChangeKey {
pub ns: NamespaceId,
pub db: DatabaseId,
pub tb: TableName,
}
pub struct Changefeed {
buffer: Mutex<HashMap<ChangeKey, TableMutations>>,
}
impl Changefeed {
pub(crate) fn new() -> Self {
Self {
buffer: Mutex::new(HashMap::new()),
}
}
pub(crate) fn buffer_table_change(
&self,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
dt: &TableDefinition,
) {
let mut buffer = self.buffer.lock();
buffer
.entry(ChangeKey {
ns,
db,
tb: tb.clone(),
})
.or_insert_with(|| TableMutations::new(tb.clone()))
.push_table_change(dt.to_owned());
}
#[expect(clippy::too_many_arguments)]
pub(crate) fn buffer_record_change(
&self,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
id: RecordId,
previous: CursorRecord,
current: CursorRecord,
store_difference: bool,
) {
let mut buffer = self.buffer.lock();
buffer
.entry(ChangeKey {
ns,
db,
tb: tb.clone(),
})
.or_insert_with(|| TableMutations::new(tb.clone()))
.push_record_change(id, previous, current, store_difference);
}
pub(crate) fn changes(&self) -> Result<Vec<PreparedWrite>> {
let buffer = self.buffer.lock();
if buffer.is_empty() {
return Ok(Vec::new());
}
let mut res = Vec::with_capacity(buffer.len());
for (key, mutations) in buffer.iter() {
let value = mutations.kv_encode_value()?;
res.push((key.ns, key.db, key.tb.clone(), value));
}
Ok(res)
}
pub(crate) fn clear(&self) {
self.buffer.lock().clear();
}
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use surrealdb_strand::Strand;
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider, TableProvider};
use crate::catalog::{
DatabaseDefinition, DatabaseId, NamespaceDefinition, NamespaceId, TableDefinition, TableId,
};
use crate::cf::ChangeSet;
use crate::expr::changefeed::ChangeFeed;
use crate::expr::statements::show::ShowSince;
use crate::kvs::LockType::*;
use crate::kvs::TransactionType::*;
use crate::kvs::{Datastore, Transaction};
use crate::val::{RecordId, RecordIdKey, TableName, Value};
const DONT_STORE_PREVIOUS: bool = false;
const NS: &str = "myns";
const DB: &str = "mydb";
const TB: &str = "mytb";
#[tokio::test]
async fn changefeed_read_write() {
let ds = init(false).await;
let tx = ds.transaction(Write, Optimistic).await.unwrap();
let tb_name = TableName::new(TB.to_owned());
let tb = tx.expect_tb_by_name(NS, DB, &tb_name).await.unwrap();
tx.commit().await.unwrap();
let tx1 = ds.transaction(Write, Optimistic).await.unwrap();
let record_a = RecordId {
table: tb_name.clone(),
key: RecordIdKey::String(Strand::new_static("A")),
};
let value_a: Value = "a".into();
let previous = Value::None;
tx1.changefeed_buffer_record_change(
tb.namespace_id,
tb.database_id,
&tb.name,
&record_a,
previous.clone().into(),
value_a.into(),
DONT_STORE_PREVIOUS,
);
tx1.commit().await.unwrap();
let tx2 = ds.transaction(Write, Optimistic).await.unwrap();
let record_c = RecordId {
table: tb_name.clone(),
key: RecordIdKey::String(Strand::new_static("C")),
};
let value_c: Value = "c".into();
tx2.changefeed_buffer_record_change(
tb.namespace_id,
tb.database_id,
&tb.name,
&record_c,
previous.clone().into(),
value_c.into(),
DONT_STORE_PREVIOUS,
);
tx2.commit().await.unwrap();
let tx3 = ds.transaction(Write, Optimistic).await.unwrap();
let record_b = RecordId {
table: tb_name.clone(),
key: RecordIdKey::String(Strand::new_static("B")),
};
let value_b: Value = "b".into();
tx3.changefeed_buffer_record_change(
tb.namespace_id,
tb.database_id,
&tb.name,
&record_b,
previous.clone().into(),
value_b.into(),
DONT_STORE_PREVIOUS,
);
let record_c2 = RecordId {
table: tb_name.clone(),
key: RecordIdKey::String(Strand::new_static("C")),
};
let value_c2: Value = "c2".into();
tx3.changefeed_buffer_record_change(
tb.namespace_id,
tb.database_id,
&tb.name,
&record_c2,
previous.clone().into(),
value_c2.into(),
DONT_STORE_PREVIOUS,
);
tx3.commit().await.unwrap();
let start: u64 = 0;
let tx4 = ds.transaction(Write, Optimistic).await.unwrap();
let r = crate::cf::read(
&tx4,
tb.namespace_id,
tb.database_id,
Some(&tb.name),
ShowSince::Versionstamp(start),
Some(10),
)
.await
.unwrap();
tx4.commit().await.unwrap();
assert_eq!(r.len(), 3);
assert_eq!(r[0].1.0.len(), 1); assert_eq!(r[0].1.0[0].1.len(), 1);
assert_eq!(r[1].1.0.len(), 1); assert_eq!(r[1].1.0[0].1.len(), 1);
assert_eq!(r[2].1.0.len(), 1); assert_eq!(r[2].1.0[0].1.len(), 2);
assert!(r[0].0 < r[1].0, "Versionstamps should be monotonically increasing");
assert!(r[1].0 < r[2].0, "Versionstamps should be monotonically increasing");
}
#[test_log::test(tokio::test)]
async fn scan_picks_up_from_offset() {
let ds = init(false).await;
let tx = ds.transaction(Write, Optimistic).await.unwrap();
let tb_name = TableName::new(TB.to_owned());
let tb = tx.expect_tb_by_name(NS, DB, &tb_name).await.unwrap();
tx.commit().await.unwrap();
let _id1 = record_change_feed_entry(
ds.transaction(Write, Optimistic).await.unwrap(),
&tb,
"First".to_string(),
)
.await;
let _id2 = record_change_feed_entry(
ds.transaction(Write, Optimistic).await.unwrap(),
&tb,
"Second".to_string(),
)
.await;
let r = change_feed_ts(ds.transaction(Write, Optimistic).await.unwrap(), &tb, 0).await;
assert_eq!(r.len(), 2);
let r = change_feed_ts(
ds.transaction(Write, Optimistic).await.unwrap(),
&tb,
r[0].0 as u64 + 1,
)
.await;
assert_eq!(r.len(), 1);
}
async fn change_feed_ts(tx: Transaction, tb: &TableDefinition, ts: u64) -> Vec<ChangeSet> {
let r = crate::cf::read(
&tx,
tb.namespace_id,
tb.database_id,
Some(&tb.name),
ShowSince::Versionstamp(ts),
Some(10),
)
.await
.unwrap();
tx.cancel().await.unwrap();
r
}
async fn record_change_feed_entry(
tx: Transaction,
tb: &TableDefinition,
id: String,
) -> RecordId {
let record_id = RecordId {
table: tb.name.clone(),
key: RecordIdKey::String(id.into()),
};
let value_a: Value = "a".into();
let previous = Value::None.into();
tx.changefeed_buffer_record_change(
tb.namespace_id,
tb.database_id,
&tb.name,
&record_id,
previous,
value_a.into(),
DONT_STORE_PREVIOUS,
);
tx.commit().await.unwrap();
record_id
}
async fn init(store_diff: bool) -> Datastore {
let namespace_id = NamespaceId(1);
let database_id = DatabaseId(2);
let table_id = TableId(3);
let ns_def = NamespaceDefinition {
namespace_id,
name: NS.into(),
comment: None,
};
let db_def = DatabaseDefinition {
namespace_id,
database_id,
name: DB.into(),
changefeed: Some(ChangeFeed {
expiry: Duration::from_secs(10),
store_diff,
}),
comment: None,
strict: false,
};
let mut tb_def = TableDefinition::new(
namespace_id,
database_id,
table_id,
TableName::new(TB.to_owned()),
);
tb_def.changefeed = Some(ChangeFeed {
expiry: Duration::from_secs(10 * 60),
store_diff,
});
let ds = Datastore::new("memory").await.unwrap();
let tx = ds.transaction(Write, Optimistic).await.unwrap();
tx.put_ns(ns_def).await.unwrap();
tx.put_db(NS, db_def).await.unwrap();
tx.put_tb(NS, DB, &tb_def).await.unwrap();
tx.commit().await.unwrap();
ds
}
}