#![allow(clippy::unwrap_used)]
#![cfg(any(
feature = "protocol-ws",
feature = "kv-mem",
feature = "kv-rocksdb",
feature = "kv-tikv",
feature = "kv-surrealkv",
))]
use std::sync::Arc;
use std::time::Duration;
use futures::{Stream, StreamExt};
use surrealdb::method::QueryStream;
use surrealdb::opt::{Config, Resource};
use surrealdb::types::{Action, RecordId, SurrealValue, Value, object};
use surrealdb::{Notification, Result};
use tokio::sync::RwLock;
use tracing::info;
use ulid::Ulid;
use super::CreateDb;
use crate::api_integration::ApiRecordId;
const LQ_TIMEOUT: Duration = Duration::from_secs(2);
const MAX_NOTIFICATIONS: usize = 100;
pub async fn live_select_table(new_db: impl CreateDb) {
let config = Config::new();
let (permit, db) = new_db.create_db(config).await;
db.use_ns(Ulid::new().to_string()).use_db(Ulid::new().to_string()).await.unwrap();
{
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table}")).await.unwrap();
let mut users = db.select(&table).live().await.unwrap();
let created: Option<ApiRecordId> = db.create(table).await.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(created, Some(notification.data.clone()));
assert_eq!(notification.action, Action::Create);
let _: Option<ApiRecordId> = db
.update(¬ification.data.id)
.content(UpdateContent {
field: "bar".to_string(),
})
.await
.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.action, Action::Update);
let _: Option<ApiRecordId> = db.delete(¬ification.data.id).await.unwrap();
let notification: Notification<ApiRecordId> = users.next().await.unwrap().unwrap();
assert_eq!(notification.action, Action::Delete);
}
{
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table}")).await.unwrap();
let mut users = db.select(Resource::from(&table)).live().await.unwrap();
db.create(Resource::from(&table)).await.unwrap();
let notification =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert!(notification.data.is_object());
assert_eq!(notification.action, Action::Create);
}
drop(permit);
}
pub async fn live_select_record_id(new_db: impl CreateDb) {
let config = Config::new();
let (permit, db) = new_db.create_db(config).await;
db.use_ns(Ulid::new().to_string()).use_db(Ulid::new().to_string()).await.unwrap();
{
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table}")).await.unwrap();
let record_id = RecordId::new(table, "john");
let mut users = db.select(&record_id).live().await.unwrap();
let created: Option<ApiRecordId> = db.create(record_id).await.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(created, Some(notification.data.clone()));
assert_eq!(notification.action, Action::Create);
let _: Option<ApiRecordId> = db
.update(¬ification.data.id)
.content(UpdateContent {
field: "bar".to_string(),
})
.await
.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.action, Action::Update);
let _: Option<ApiRecordId> = db.delete(¬ification.data.id).await.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.action, Action::Delete);
}
{
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table}")).await.unwrap();
let record_id = RecordId::new(table, "john");
let mut users = db.select(Resource::from(&record_id)).live().await.unwrap();
db.create(Resource::from(record_id)).await.unwrap();
let notification: Notification<Value> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert!(notification.data.is_object());
assert_eq!(notification.action, Action::Create);
}
drop(permit);
}
pub async fn live_select_record_ranges(new_db: impl CreateDb) {
let config = Config::new();
let (permit, db) = new_db.create_db(config).await;
db.use_ns(Ulid::new().to_string()).use_db(Ulid::new().to_string()).await.unwrap();
{
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table}")).await.unwrap();
let mut users = db.select(&table).range("jane".."john").live().await.unwrap();
let created: Option<ApiRecordId> = db.create((table, "jane")).await.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(created, Some(notification.data.clone()));
assert_eq!(notification.action, Action::Create);
let _: Option<ApiRecordId> = db
.update(¬ification.data.id)
.content(UpdateContent {
field: "bar".to_string(),
})
.await
.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.action, Action::Update);
let _: Option<ApiRecordId> = db.delete(¬ification.data.id).await.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.action, Action::Delete);
}
{
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table}")).await.unwrap();
let mut users =
db.select(Resource::from(&table)).range("jane".."john").live().await.unwrap();
let created_value = db
.create(Resource::from((table, "job")))
.await
.unwrap()
.into_array()
.unwrap()
.remove(0)
.into_object()
.unwrap();
let notification: Notification<Value> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert!(notification.data.is_object());
assert_eq!(notification.action, Action::Create);
let thing = match created_value.get("id").unwrap() {
Value::RecordId(thing) => thing,
_ => panic!("Expected a thing"),
};
db.query("DELETE $item").bind(("item", thing.clone())).await.unwrap();
let notification: Notification<Value> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.action, Action::Delete);
let notification = match notification.data {
Value::Object(notification) => notification,
_ => panic!("Expected an object"),
};
assert_eq!(notification, created_value);
}
drop(permit);
}
pub async fn live_select_query(new_db: impl CreateDb) {
let config = Config::new();
let (permit, db) = new_db.create_db(config).await;
db.use_ns(Ulid::new().to_string()).use_db(Ulid::new().to_string()).await.unwrap();
{
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table}")).await.unwrap();
info!("Starting live query");
let users: QueryStream<Notification<ApiRecordId>> = db
.query(format!("LIVE SELECT * FROM {table}"))
.await
.unwrap()
.stream::<Notification<_>>(0)
.unwrap();
let users = Arc::new(RwLock::new(users));
info!("Creating record");
let created: Option<ApiRecordId> = db.create(table).await.unwrap();
let notifications = receive_all_pending_notifications(Arc::clone(&users), LQ_TIMEOUT).await;
assert_eq!(
notifications.iter().map(|n| n.action).collect::<Vec<_>>(),
vec![Action::Create],
"{:?}",
notifications
);
assert_eq!(created, Some(notifications[0].data.clone()));
info!("Updating record");
let _: Option<ApiRecordId> = db
.update(¬ifications[0].data.id)
.content(UpdateContent {
field: "bar".to_string(),
})
.await
.unwrap();
let notifications = receive_all_pending_notifications(Arc::clone(&users), LQ_TIMEOUT).await;
assert_eq!(
notifications.iter().map(|n| n.action).collect::<Vec<_>>(),
[Action::Update],
"{:?}",
notifications
);
info!("Deleting record");
let _: Option<ApiRecordId> = db.delete(¬ifications[0].data.id).await.unwrap();
let notifications = receive_all_pending_notifications(Arc::clone(&users), LQ_TIMEOUT).await;
assert_eq!(
notifications.iter().map(|n| n.action).collect::<Vec<_>>(),
[Action::Delete],
"{:?}",
notifications
);
}
{
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table} CHANGEFEED 10m INCLUDE ORIGINAL")).await.unwrap();
let mut users = db
.query(format!("LIVE SELECT * FROM {table}"))
.await
.unwrap()
.stream::<Value>(0)
.unwrap();
db.create(Resource::from(&table)).await.unwrap();
let notification =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert!(notification.data.is_object());
assert_eq!(notification.action, Action::Create);
}
{
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table} CHANGEFEED 10m INCLUDE ORIGINAL")).await.unwrap();
let mut users = db
.query(format!("LIVE SELECT * FROM {table}"))
.await
.unwrap()
.stream::<Notification<_>>(())
.unwrap();
let created: Option<ApiRecordId> = db.create(table).await.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(created, Some(notification.data.clone()));
assert_eq!(notification.action, Action::Create, "{:?}", notification);
let _: Option<ApiRecordId> = db
.update(¬ification.data.id)
.content(UpdateContent {
field: "bar".to_string(),
})
.await
.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.action, Action::Update, "{:?}", notification);
let _: Option<ApiRecordId> = db.delete(¬ification.data.id).await.unwrap();
let notification: Notification<ApiRecordId> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.action, Action::Delete, "{:?}", notification);
}
{
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table} CHANGEFEED 10m INCLUDE ORIGINAL")).await.unwrap();
let mut users = db
.query(format!("BEGIN; LIVE SELECT * FROM {table}; COMMIT"))
.await
.unwrap()
.stream::<Value>(())
.unwrap();
db.create(Resource::from(&table)).await.unwrap();
let notification =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert!(notification.data.is_object());
assert_eq!(notification.action, Action::Create);
}
drop(permit);
}
pub async fn live_query_delete_notifications(new_db: impl CreateDb) {
let config = Config::new();
let (permit, db) = new_db.create_db(config).await;
db.use_ns(Ulid::new().to_string()).use_db(Ulid::new().to_string()).await.unwrap();
db.query("DEFINE TABLE bar".to_string()).await.unwrap().check().unwrap();
let mut stream =
db.query("LIVE SELECT field FROM bar").await.unwrap().stream::<Value>(0).unwrap();
db.query("CREATE bar CONTENT { field: 'baz' }").await.unwrap().check().unwrap();
let notification =
tokio::time::timeout(LQ_TIMEOUT, stream.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.data, Value::Object(object! { field: "baz" }));
assert_eq!(notification.action, Action::Create);
db.query("UPDATE bar MERGE { data: 123 }").await.unwrap().check().unwrap();
let notification =
tokio::time::timeout(LQ_TIMEOUT, stream.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.data, Value::Object(object! { field: "baz" }));
assert_eq!(notification.action, Action::Update);
db.query("DELETE bar").await.unwrap().check().unwrap();
let notification =
tokio::time::timeout(LQ_TIMEOUT, stream.next()).await.unwrap().unwrap().unwrap();
assert_eq!(notification.data, Value::Object(object! { field: "baz" }));
assert_eq!(notification.action, Action::Delete);
drop(permit);
}
#[derive(Debug, Clone, SurrealValue, PartialEq, PartialOrd)]
struct ApiRecordIdWithFetchedLink {
id: RecordId,
link: Option<ApiRecordId>,
}
#[derive(Debug, Clone, SurrealValue, PartialEq, PartialOrd)]
struct ApiRecordIdWithUnfetchedLink {
id: RecordId,
link: RecordId,
}
#[derive(Debug, Clone, SurrealValue, PartialEq, PartialOrd)]
struct LinkContent {
link: RecordId,
}
#[derive(Debug, Clone, SurrealValue)]
struct UpdateContent {
field: String,
}
pub async fn live_select_with_fetch(new_db: impl CreateDb) {
let config = Config::new();
let (permit, db) = new_db.create_db(config).await;
db.use_ns(Ulid::new().to_string()).use_db(Ulid::new().to_string()).await.unwrap();
let table = format!("table_{}", Ulid::new());
let linktb = format!("link_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table}")).await.unwrap();
let mut users = db
.query(format!("LIVE SELECT * FROM {table} FETCH link"))
.await
.unwrap()
.stream::<Notification<_>>(())
.unwrap();
let link: Option<ApiRecordId> = db.create(&linktb).await.unwrap();
let linkone = link.unwrap().id;
let link: Option<ApiRecordId> = db.create(&linktb).await.unwrap();
let linktwo = link.unwrap().id;
let created: Option<ApiRecordIdWithUnfetchedLink> = db
.create(table)
.content(LinkContent {
link: linkone.clone(),
})
.await
.unwrap();
let notification: Notification<ApiRecordIdWithFetchedLink> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(
ApiRecordIdWithFetchedLink {
id: created.unwrap().id,
link: Some(ApiRecordId {
id: linkone,
}),
},
notification.data.clone()
);
assert_eq!(notification.action, Action::Create);
let updated: Option<ApiRecordIdWithUnfetchedLink> = db
.update(¬ification.data.id)
.content(LinkContent {
link: linktwo.clone(),
})
.await
.unwrap();
let notification: Notification<ApiRecordIdWithFetchedLink> =
tokio::time::timeout(LQ_TIMEOUT, users.next()).await.unwrap().unwrap().unwrap();
assert_eq!(
ApiRecordIdWithFetchedLink {
id: updated.unwrap().id,
link: Some(ApiRecordId {
id: linktwo,
}),
},
notification.data.clone()
);
assert_eq!(notification.action, Action::Update);
let _: Option<ApiRecordIdWithUnfetchedLink> = db.delete(¬ification.data.id).await.unwrap();
let notification: Notification<ApiRecordIdWithFetchedLink> =
users.next().await.unwrap().unwrap();
assert_eq!(notification.action, Action::Delete);
drop(permit);
}
async fn receive_all_pending_notifications<S: Stream<Item = Result<Notification<I>>> + Unpin, I>(
stream: Arc<RwLock<S>>,
timeout: Duration,
) -> Vec<Notification<I>> {
let mut results = Vec::new();
let we_expect_timeout = tokio::time::timeout(timeout, async {
while let Some(notification) = stream.write().await.next().await {
if results.len() >= MAX_NOTIFICATIONS {
panic!("too many notification!")
}
results.push(notification.unwrap())
}
})
.await;
assert!(we_expect_timeout.is_err());
results
}
pub async fn live_select_returns_uuid(new_db: impl CreateDb) {
let config = Config::new();
let (permit, db) = new_db.create_db(config).await;
db.use_ns(Ulid::new().to_string()).use_db(Ulid::new().to_string()).await.unwrap();
let table = format!("table_{}", Ulid::new());
db.query(format!("DEFINE TABLE {table}")).await.unwrap();
let mut response = db.query(format!("LIVE SELECT * FROM {table}")).await.unwrap();
assert_eq!(response.num_statements(), 1, "LIVE SELECT should return exactly one result");
let result: Value = response.take(0).unwrap();
assert!(result.is_uuid(), "LIVE SELECT should return a UUID, got: {:?}", result);
drop(permit);
}
define_include_tests!(live => {
#[test_log::test(tokio::test)]
live_select_table,
#[test_log::test(tokio::test)]
live_select_record_id,
#[test_log::test(tokio::test)]
live_select_record_ranges,
#[test_log::test(tokio::test)]
live_select_query,
#[test_log::test(tokio::test)]
live_select_with_fetch,
#[test_log::test(tokio::test)]
live_query_delete_notifications,
#[test_log::test(tokio::test)]
live_select_returns_uuid,
});