#![cfg(feature = "kv-mem")]
use std::time::Duration;
use futures::StreamExt;
use surrealdb::Surreal;
use surrealdb::engine::local::Mem;
use surrealdb::method::StreamItem;
async fn db() -> Surreal<surrealdb::engine::local::Db> {
let db = Surreal::new::<Mem>(()).await.expect("an embedded connection");
db.use_ns("test").use_db("test").await.expect("use");
db
}
fn alive_tasks() -> usize {
tokio::runtime::Handle::current().metrics().num_alive_tasks()
}
async fn wait_for_tasks(target: usize) -> usize {
let mut alive = alive_tasks();
for _ in 0..200 {
if alive <= target {
break;
}
tokio::time::sleep(Duration::from_millis(25)).await;
alive = alive_tasks();
}
alive
}
async fn settled_tasks() -> usize {
let mut last = alive_tasks();
for _ in 0..200 {
tokio::time::sleep(Duration::from_millis(25)).await;
let now = alive_tasks();
if now == last {
return now;
}
last = now;
}
last
}
#[tokio::test]
async fn rows_and_statement_ends_arrive_in_order() {
let db = db().await;
db.query("FOR $i IN array::range(0, 100) { CREATE type::record('row', $i) SET n = $i }")
.await
.expect("send")
.check()
.expect("seed");
let mut items =
db.query("SELECT * FROM row").query("RETURN 'done'").stream_items().expect("a stream");
let mut rows = 0usize;
let mut ends = Vec::new();
while let Some(item) = items.next().await {
match item.expect("no stream-level failure") {
StreamItem::Row {
statement,
..
} => {
assert!(
!ends.contains(&statement),
"statement {statement} produced a row after it finished"
);
rows += 1;
}
StreamItem::StatementEnd {
statement,
result,
..
} => {
result.expect("both statements succeed");
ends.push(statement);
}
}
}
assert_eq!(rows, 101, "every row, plus the second statement's one value");
assert_eq!(ends, vec![0, 1]);
}
#[tokio::test]
async fn streaming_and_buffered_agree() {
let db = db().await;
db.query("FOR $i IN array::range(0, 50) { CREATE type::record('row', $i) SET n = $i }")
.await
.expect("send")
.check()
.expect("seed");
let mut buffered = db.query("SELECT n FROM row ORDER BY n").await.expect("send");
let buffered: Vec<surrealdb::types::Value> = buffered.take(0).expect("rows");
let mut items = db.query("SELECT n FROM row ORDER BY n").stream_items().expect("a stream");
let mut streamed = Vec::new();
while let Some(item) = items.next().await {
if let StreamItem::Row {
value,
..
} = item.expect("no stream-level failure")
{
streamed.push(value);
}
}
assert_eq!(streamed, buffered, "the same rows, in the same order");
}
#[tokio::test]
async fn a_failed_statement_ends_with_its_error() {
let db = db().await;
let mut items = db
.query("RETURN 1")
.query("THROW 'nope'")
.query("RETURN 3")
.stream_items()
.expect("a stream");
let mut outcomes = Vec::new();
while let Some(item) = items.next().await {
if let StreamItem::StatementEnd {
statement,
result,
..
} = item.expect("no stream-level failure")
{
outcomes.push((statement, result.is_ok()));
}
}
assert_eq!(outcomes, vec![(0, true), (1, false), (2, true)]);
}
#[tokio::test]
async fn a_dropped_stream_releases_the_connection() {
let db = db().await;
db.query("FOR $i IN array::range(0, 500) { CREATE type::record('row', $i) SET n = $i }")
.await
.expect("send")
.check()
.expect("seed");
let mut items = db.query("SELECT * FROM row").stream_items().expect("a stream");
items.next().await.expect("at least one item").expect("no stream-level failure");
drop(items);
let mut response =
tokio::time::timeout(std::time::Duration::from_secs(10), db.query("RETURN 'still here'"))
.await
.expect("the connection is usable after abandoning a stream")
.expect("send");
let answer: Option<String> = response.take(0).expect("a result");
assert_eq!(answer.as_deref(), Some("still here"));
}
#[tokio::test]
async fn a_dropped_stream_stops_the_execution_behind_it() {
let db = db().await;
db.query("FOR $i IN array::range(0, 10000) { CREATE type::record('row', $i) SET n = $i }")
.await
.expect("send")
.check()
.expect("seed");
let baseline = settled_tasks().await;
let mut items = db.query("SELECT * FROM row").stream_items().expect("a stream");
items.next().await.expect("at least one item").expect("no stream-level failure");
assert!(
alive_tasks() > baseline,
"the execution is running while the stream is being read, or the count below proves nothing",
);
drop(items);
let alive = wait_for_tasks(baseline).await;
assert_eq!(
alive, baseline,
"the execution is still alive after its stream was dropped, holding a transaction that \
will never be committed or cancelled",
);
}