#![allow(clippy::bool_assert_comparison)]
#[macro_use]
extern crate serde_json;
extern crate serde;
use serde::{Deserialize, Serialize};
mod common;
use std::sync::{
atomic::{AtomicBool, Ordering},
Arc, Condvar, Mutex,
};
use dittolive_ditto::dql::*;
use futures::executor::block_on;
#[test]
fn exec_empty() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_empty_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_empty_impl(ditto);
}
fn exec_empty_impl(ditto: dittolive_ditto::prelude::Ditto) {
let documents = block_on(ditto.store().execute("SELECT * FROM cars")).unwrap();
assert_eq!(documents.item_count(), 0);
}
#[test]
fn exec_invalid() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_invalid_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_invalid_impl(ditto);
}
fn exec_invalid_impl(ditto: dittolive_ditto::prelude::Ditto) {
let res = block_on(ditto.store().execute("sudo make me a coffee."));
assert!(res.is_err());
}
#[test]
fn exec_empty_query() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_empty_query_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_empty_query_impl(ditto);
}
fn exec_empty_query_impl(ditto: dittolive_ditto::prelude::Ditto) {
let res = block_on(ditto.store().execute(""));
assert!(res.is_err());
}
#[test]
fn exec_iterate_result() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_iterate_result_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_iterate_result_impl(ditto);
}
fn exec_iterate_result_impl(ditto: dittolive_ditto::prelude::Ditto) {
#[derive(Serialize, Deserialize)]
struct Car {
color: String,
kms: u32,
broken_engine: bool,
}
futures::executor::block_on(ditto.store().execute((
"INSERT INTO cars DOCUMENTS (:car1), (:car2), (:car3), (:car4)",
json!({
"car1": Car {
color: "red".to_string(),
kms: 5000,
broken_engine: false,
},
"car2": Car {
color: "blue".to_string(),
kms: 45000,
broken_engine: false,
},
"car3": Car {
color: "red".to_string(),
kms: 445000,
broken_engine: true,
},
"car4": Car {
color: "black".to_string(),
kms: 52000,
broken_engine: true,
}
}),
)))
.unwrap();
let french_cars = block_on(
ditto
.store()
.execute("SELECT * FROM cars WHERE broken_engine = true"),
)
.unwrap();
for french_car in &french_cars {
let car: Car = french_car.deserialize_value().unwrap();
assert!(car.broken_engine);
}
}
#[test]
fn exec_manually_extract_result() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_manually_extract_result_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_manually_extract_result_impl(ditto);
}
fn exec_manually_extract_result_impl(ditto: dittolive_ditto::prelude::Ditto) {
#[derive(Serialize, Deserialize)]
struct Car {
color: String,
kms: u32,
broken_engine: bool,
}
futures::executor::block_on(ditto.store().execute((
"INSERT INTO cars DOCUMENTS (:car1), (:car2), (:car3), (:car4)",
json!({
"car1": Car {
color: "red".to_string(),
kms: 5000,
broken_engine: false,
},
"car2": Car {
color: "blue".to_string(),
kms: 45000,
broken_engine: false,
},
"car3": Car {
color: "red".to_string(),
kms: 45000,
broken_engine: true,
},
"car4": Car {
color: "black".to_string(),
kms: 52000,
broken_engine: true,
}
}),
)))
.unwrap();
let french_car = block_on(
ditto
.store()
.execute("SELECT * FROM cars WHERE broken_engine = true AND kms = 45000"),
)
.unwrap();
assert_eq!(french_car.item_count(), 1);
let mut first_item = french_car.get_item(0).unwrap();
assert_eq!(first_item.is_materialized(), false);
let data = first_item.value();
assert_eq!(first_item.is_materialized(), true);
first_item.dematerialize();
assert_eq!(first_item.is_materialized(), false);
assert!(matches!(&data["color"], ::serde_cbor::Value::Text(s) if s == "red"));
let car: Car = first_item.deserialize_value().unwrap();
assert_eq!(car.color, "red");
}
#[test]
fn exec_empty_args() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_empty_args_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_empty_args_impl(ditto);
}
fn exec_empty_args_impl(ditto: dittolive_ditto::prelude::Ditto) {
#[derive(Serialize)]
struct Args {
brand: String,
}
let args = Args {
brand: "Brandy".to_string(),
};
let documents = block_on(
ditto
.store()
.execute(("SELECT * FROM cars WHERE brand = :brand", args)),
)
.unwrap();
assert_eq!(documents.item_count(), 0);
}
#[test]
fn exec_null_query() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_null_query_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_null_query_impl(ditto);
}
fn exec_null_query_impl(ditto: dittolive_ditto::prelude::Ditto) {
futures::executor::block_on(ditto.store().execute((
"INSERT INTO cars DOCUMENTS (:doc1), (:doc2), (:doc3)",
json!({
"doc1": {"brand": "Intel"},
"doc2": {"brand": None::<String>},
"doc3": {"color": "red"}
}),
)))
.unwrap();
let documents = block_on(
ditto
.store()
.execute("SELECT * FROM cars WHERE brand IS NULL"),
)
.unwrap();
assert_eq!(documents.item_count(), 1);
}
#[test]
fn exec_missing_query() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_missing_query_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_missing_query_impl(ditto);
}
fn exec_missing_query_impl(ditto: dittolive_ditto::prelude::Ditto) {
futures::executor::block_on(ditto.store().execute((
"INSERT INTO cars DOCUMENTS (:doc1), (:doc2), (:doc3)",
json!({
"doc1": {"brand": "Intel"},
"doc2": {"brand": None::<String>},
"doc3": {"color": "red"}
}),
)))
.unwrap();
let documents = block_on(
ditto
.store()
.execute("SELECT * FROM cars WHERE brand IS MISSING"),
)
.unwrap();
assert_eq!(documents.item_count(), 1);
}
#[test]
fn exec_multiple() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_multiple_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_multiple_impl(ditto);
}
fn exec_multiple_impl(ditto: dittolive_ditto::prelude::Ditto) {
futures::executor::block_on(ditto.store().execute((
"INSERT INTO cars DOCUMENTS (:fiat), (:bmw)",
json!({
"fiat": {"brand": "Fiat"},
"bmw": {"brand": "BMW"}
}),
)))
.unwrap();
let documents = block_on(ditto.store().execute("SELECT * FROM cars")).unwrap();
assert_eq!(documents.item_count(), 2);
}
#[test]
fn exec_multiple_args() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_multiple_args_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
exec_multiple_args_impl(ditto);
}
fn exec_multiple_args_impl(ditto: dittolive_ditto::prelude::Ditto) {
#[derive(Serialize)]
struct Args {
color: String,
}
let args = Args {
color: "blue".to_string(),
};
futures::executor::block_on(ditto.store().execute((
"INSERT INTO cars DOCUMENTS (:car1), (:car2), (:car3), (:car4), (:car5)",
json!({
"car1": {"color": "red", "kms": 123},
"car2": {"color": "blue", "kms": 456},
"car3": {"color": "blue", "kms": 213456},
"car4": {"color": "red", "kms": 213456},
"car5": {"color": "black", "kms": 213456}
}),
)))
.unwrap();
let documents = block_on(
ditto
.store()
.execute(("SELECT * FROM cars WHERE color = :color", args)),
)
.unwrap();
assert_eq!(documents.mutated_document_ids().len(), 0);
assert_eq!(documents.item_count(), 2);
}
#[test]
fn replication_subscribe() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
replication_subscribe_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
replication_subscribe_impl(ditto);
}
fn replication_subscribe_impl(ditto: dittolive_ditto::prelude::Ditto) {
let sub = ditto
.sync()
.register_subscription("SELECT * from cars")
.unwrap();
sub.cancel();
}
#[test]
fn replication_subscribe_drop_after_ditto() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
replication_subscribe_drop_after_ditto_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
replication_subscribe_drop_after_ditto_impl(ditto);
}
fn replication_subscribe_drop_after_ditto_impl(ditto: dittolive_ditto::prelude::Ditto) {
let sub = ditto
.sync()
.register_subscription("SELECT * from cars")
.unwrap();
drop(ditto);
sub.cancel();
}
#[test]
fn change_observer() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
change_observer_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
change_observer_impl(ditto);
}
fn change_observer_impl(ditto: dittolive_ditto::prelude::Ditto) {
let finished = Arc::new((Mutex::new(false), Condvar::new()));
let finished_clone = Arc::clone(&finished);
let change_observer = ditto
.store()
.register_observer("SELECT * from cars", move |result: QueryResult| {
if result.item_count() == 3 {
let (lock, cvar) = &*finished_clone;
let mut done = lock.lock().unwrap();
*done = true;
cvar.notify_one();
}
})
.unwrap();
futures::executor::block_on(ditto.store().execute((
"INSERT INTO cars DOCUMENTS (:doc1), (:doc2), (:doc3)",
json!({
"doc1": {"brand": "Lada"},
"doc2": {"brand": "Mercedes"},
"doc3": {"brand": "Rezvani"}
}),
)))
.unwrap();
let (lock, cvar) = &*finished;
let mut done = lock.lock().unwrap();
while !*done {
done = cvar.wait(done).unwrap();
}
change_observer.cancel();
}
#[test]
fn change_observer_with_args() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
change_observer_with_args_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
change_observer_with_args_impl(ditto);
}
fn change_observer_with_args_impl(ditto: dittolive_ditto::prelude::Ditto) {
let finished = Arc::new(AtomicBool::new(false));
let finished_clone = Arc::clone(&finished);
#[derive(Serialize)]
struct Args {
brand: String,
}
let change_observer = ditto
.store()
.register_observer(
(
"SELECT * from cars WHERE brand = :brand",
Args {
brand: String::from("Lada"),
},
),
move |result: QueryResult| {
if result.item_count() == 2 {
finished_clone.store(true, Ordering::SeqCst);
}
},
)
.unwrap();
futures::executor::block_on(ditto.store().execute((
"INSERT INTO cars DOCUMENTS (:doc1), (:doc2), (:doc3), (:doc4)",
json!({
"doc1": {"brand": "Lada"},
"doc2": {"brand": "Mercedes"},
"doc3": {"brand": "Rezvani"},
"doc4": {"brand": "Lada"}
}),
)))
.unwrap();
while !finished.load(Ordering::SeqCst) {
std::thread::yield_now();
}
change_observer.cancel();
}
#[test]
fn change_observer_drop_after_ditto() {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
change_observer_drop_after_ditto_impl(ditto);
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
change_observer_drop_after_ditto_impl(ditto);
}
fn change_observer_drop_after_ditto_impl(ditto: dittolive_ditto::prelude::Ditto) {
let change_observer = ditto
.store()
.register_observer("SELECT * from cars", |_result| {})
.unwrap();
drop(ditto);
change_observer.cancel();
}
use dittolive_ditto::prelude::*;
#[tokio::test]
async fn change_observer_sequential_callbacks() -> anyhow::Result<()> {
fn setup_ditto(database_id: DatabaseId) -> anyhow::Result<(Arc<Ditto>, tempfile::TempDir)> {
let ditto_root_dir = tempfile::tempdir()?;
let config = DittoConfig::new(
database_id.to_string(),
DittoConfigConnect::SmallPeersOnly { private_key: None },
)
.with_persistence_directory(ditto_root_dir.path());
let ditto = Ditto::open_sync(config)?;
ditto.set_license_from_env("DITTO_LICENSE")?;
let ditto = Arc::new(ditto);
Ok((ditto, ditto_root_dir))
}
let database_id = DatabaseId::generate();
let (peer1, _peer1_db) = setup_ditto(database_id.clone())?;
let (peer2, _peer2_db) = setup_ditto(database_id)?;
let _ = peer1
.store()
.execute("ALTER SYSTEM SET DQL_RESTRICT_SUBSCRIPTIONS = false")
.await?;
{
let mut peer1_tc = TransportConfig::default();
peer1_tc.peer_to_peer.lan.enabled = false;
peer1_tc.peer_to_peer.bluetooth_le.enabled = false;
peer1_tc.listen.tcp.enabled = true;
peer1_tc.listen.tcp.interface_ip = "127.0.0.1".to_string();
peer1_tc.listen.tcp.port = 4040;
peer1.set_transport_config(peer1_tc);
peer1.sync().start()?;
println!("PEER1 START SYNC");
}
{
let mut peer2_tc = TransportConfig::default();
peer2_tc.peer_to_peer.lan.enabled = false;
peer2_tc.peer_to_peer.bluetooth_le.enabled = false;
peer2_tc
.connect
.tcp_servers
.insert("127.0.0.1:4040".to_string());
peer2.set_transport_config(peer2_tc);
peer2.sync().start()?;
println!("PEER2 START SYNC");
}
let peers = [&peer1, &peer2];
for (i, peer) in peers.iter().enumerate() {
while peer.presence().graph().remote_peers.len() != 1 {
tokio::time::sleep(std::time::Duration::from_millis(10)).await;
}
println!("PEER {i} CONNECTED");
}
println!("PEERS CONNECTED");
let assert_no_concurrent_callbacks = ::std::sync::Mutex::new(());
let assert_ordered_results = Arc::new(::std::sync::atomic::AtomicI64::new(0));
let peer1_inner = peer1.clone();
let ordered_inner = assert_ordered_results.clone();
let _peer1_observer = peer1.store().register_observer(
"SELECT * FROM cars WHERE purchased = false ORDER BY _id",
move |result: QueryResult| {
let _callback_lock = assert_no_concurrent_callbacks
.try_lock()
.expect("should not enter callback twice before exiting");
let items = result
.into_iter()
.map(|item| item.deserialize_value::<serde_json::Value>())
.collect::<Result<Vec<_>, _>>()
.expect("deserialize arbitrary json");
let peer1 = peer1_inner.clone();
let assert_ordered_results = ordered_inner.clone();
for value in &items {
let id = value["_id"].as_i64().expect("ID i64");
let expected_id =
assert_ordered_results.fetch_add(1, ::std::sync::atomic::Ordering::SeqCst);
assert_eq!(
id, expected_id,
"Saw ID that we've seen before, expected each ID to be processed once"
);
let _result = block_on(peer1.store().execute((
"UPDATE cars SET purchased = true WHERE _id = :id",
serde_json::json!({
"id": id,
}),
)))
.expect("SET purchased=true");
}
drop(_callback_lock);
},
)?;
let _peer1_subscription = peer1
.sync()
.register_subscription("SELECT * FROM cars WHERE purchased = false ORDER BY vin")?;
let car_count = 100_u32;
for i in 0..car_count {
peer2
.store()
.execute((
"INSERT INTO cars DOCUMENTS (:car)",
serde_json::json!({
"car": {
"_id": i,
"color": "blue",
"purchased": false
}
}),
))
.await?;
}
let last_car = car_count.saturating_sub(1);
let start = std::time::Instant::now();
loop {
if start.elapsed() > std::time::Duration::from_secs(20) {
panic!("Timed out while waiting for test to finish")
}
let results = peer1
.store()
.execute((
"SELECT * FROM cars WHERE purchased = true AND _id = :lastCar",
serde_json::json!({
"lastCar": last_car
}),
))
.await?;
if results.item_count() == 1 {
println!("CAR {last_car} HAS BEEN PURCHASED");
break;
}
}
let last_seen_id = assert_ordered_results.load(::std::sync::atomic::Ordering::SeqCst);
assert_eq!(
last_seen_id, car_count as i64,
"Expected to see {car_count} purchased cars"
);
Ok(())
}
#[test]
fn malformed_dql_error() -> anyhow::Result<()> {
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
malformed_dql_error_impl(ditto)?;
let (_root, ditto) = common::get_inactive_ditto(None).unwrap();
malformed_dql_error_impl(ditto)?;
Ok(())
}
fn malformed_dql_error_impl(ditto: dittolive_ditto::prelude::Ditto) -> anyhow::Result<()> {
let result = ditto.store().register_observer(
(
"SELECT * FROM my-collection WHERE _id = :id",
serde_json::json!({
"id": 4,
}),
),
|_query_result| {},
);
let err = result.err();
println!("REGISTER OBSERVER ERR: {err:?}");
Ok(())
}
#[test]
fn cancel_sync_subscription() -> anyhow::Result<()> {
let (_root, ditto) = common::get_inactive_ditto(None)?;
cancel_sync_subscription_impl(ditto)?;
let (_root, ditto) = common::get_inactive_ditto(None)?;
cancel_sync_subscription_impl(ditto)?;
Ok(())
}
fn cancel_sync_subscription_impl(ditto: dittolive_ditto::prelude::Ditto) -> anyhow::Result<()> {
let subs = ditto.sync().subscriptions();
assert_eq!(subs.len(), 0);
let sync_subscription = ditto.sync().register_subscription("SELECT * FROM cars")?;
let subs = ditto.sync().subscriptions();
let sub_count = subs.len();
assert_eq!(sub_count, 1);
sync_subscription.cancel();
let sub_count = ditto.sync().subscriptions().len();
assert_eq!(sub_count, 0);
let sync_subscription2 = ditto.sync().register_subscription("SELECT * FROM busses")?;
let sub_count = ditto.sync().subscriptions().len();
assert_eq!(sub_count, 1);
drop(sync_subscription2);
let sub_count = ditto.sync().subscriptions().len();
assert_eq!(sub_count, 1);
Ok(())
}
#[test]
fn commit_id_mutating_query() {
let (_root1, ditto1) = common::get_inactive_ditto(None).unwrap();
let (_root2, ditto2) = common::get_inactive_ditto(None).unwrap();
commit_id_mutating_query_impl(ditto1);
commit_id_mutating_query_impl(ditto2);
}
fn commit_id_mutating_query_impl(ditto: Ditto) {
let query_result = block_on(ditto.store().execute((
"INSERT INTO cats DOCUMENTS (:doc)",
json!({ "doc": { "_id": "test1", "value": 42 }}),
)))
.unwrap();
assert!(query_result.commit_id().is_some());
}
#[test]
fn commit_id_non_mutating_query() {
let (_root1, ditto1) = common::get_inactive_ditto(None).unwrap();
let (_root2, ditto2) = common::get_inactive_ditto(None).unwrap();
commit_id_non_mutating_query_impl(ditto1);
commit_id_non_mutating_query_impl(ditto2);
}
fn commit_id_non_mutating_query_impl(ditto: Ditto) {
block_on(ditto.store().execute((
"INSERT INTO cats DOCUMENTS (:doc)",
json!({ "doc": { "_id": "test2", "value": 42 }}),
)))
.unwrap();
let query_result = block_on(ditto.store().execute((
"SELECT * FROM cats WHERE _id = :id",
json!({ "id": "test2" }),
)))
.unwrap();
assert!(query_result.commit_id().is_none());
}
#[test]
fn commit_id_transaction_committed() {
let (_root1, ditto1) = common::get_inactive_ditto(None).unwrap();
let (_root2, ditto2) = common::get_inactive_ditto(None).unwrap();
commit_id_transaction_committed_impl(ditto1);
commit_id_transaction_committed_impl(ditto2);
}
fn commit_id_transaction_committed_impl(ditto: Ditto) {
let query_result = block_on(ditto.store().transaction(async |txn| {
let result = txn
.execute((
"INSERT INTO cats DOCUMENTS (:doc)",
json!({ "doc": { "_id": "test3", "value": 42 }}),
))
.await?;
Ok::<_, DittoError>((result, TransactionCompletionAction::Commit))
}))
.unwrap()
.0;
assert!(query_result.commit_id().is_some());
}
#[test]
fn commit_id_transaction_uncommitted() {
let (_root1, ditto1) = common::get_inactive_ditto(None).unwrap();
let (_root2, ditto2) = common::get_inactive_ditto(None).unwrap();
commit_id_transaction_uncommitted_impl(ditto1);
commit_id_transaction_uncommitted_impl(ditto2);
}
fn commit_id_transaction_uncommitted_impl(ditto: Ditto) {
let _ = block_on(ditto.store().transaction(async |txn| {
let query_result = txn
.execute((
"INSERT INTO cats DOCUMENTS (:doc)",
json!({ "doc": { "_id": "test4", "value": 42 }}),
))
.await?;
assert!(query_result.commit_id().is_none());
Ok::<_, DittoError>(TransactionCompletionAction::Commit)
}))
.unwrap();
}
#[test]
fn commit_id_transaction_non_mutating_after_mutation() {
let (_root1, ditto1) = common::get_inactive_ditto(None).unwrap();
let (_root2, ditto2) = common::get_inactive_ditto(None).unwrap();
commit_id_transaction_non_mutating_after_mutation_impl(ditto1);
commit_id_transaction_non_mutating_after_mutation_impl(ditto2);
}
fn commit_id_transaction_non_mutating_after_mutation_impl(ditto: Ditto) {
block_on(ditto.store().execute((
"INSERT INTO cats DOCUMENTS (:doc)",
json!({ "doc": { "_id": "test5", "value": 42 }}),
)))
.unwrap();
let query_result = block_on(ditto.store().transaction(async |txn| {
txn.execute((
"INSERT INTO cats DOCUMENTS (:doc)",
json!({ "doc": { "_id": "test6", "value": 43 }}),
))
.await?;
let result = txn
.execute((
"SELECT * FROM cats WHERE _id = :id",
json!({ "id": "test5" }),
))
.await?;
Ok::<_, DittoError>(result)
}))
.unwrap();
assert!(query_result.commit_id().is_some());
}
#[test]
fn commit_id_transaction_explicitly_rolled_back() {
let (_root1, ditto1) = common::get_inactive_ditto(None).unwrap();
let (_root2, ditto2) = common::get_inactive_ditto(None).unwrap();
commit_id_transaction_explicitly_rolled_back_impl(ditto1);
commit_id_transaction_explicitly_rolled_back_impl(ditto2);
}
fn commit_id_transaction_explicitly_rolled_back_impl(ditto: Ditto) {
let mut query_result = None;
let _ = block_on(ditto.store().transaction(async |txn| {
let result = txn
.execute((
"INSERT INTO cats DOCUMENTS (:doc)",
json!({ "doc": { "_id": "test7", "value": 42 }}),
))
.await?;
query_result = Some(result);
Ok::<_, DittoError>(TransactionCompletionAction::Rollback)
}))
.unwrap();
assert!(query_result.is_some());
assert!(query_result.as_ref().unwrap().commit_id().is_none());
}
#[test]
fn commit_id_transaction_implicitly_rolled_back() {
let (_root1, ditto1) = common::get_inactive_ditto(None).unwrap();
let (_root2, ditto2) = common::get_inactive_ditto(None).unwrap();
commit_id_transaction_implicitly_rolled_back_impl(ditto1);
commit_id_transaction_implicitly_rolled_back_impl(ditto2);
}
fn commit_id_transaction_implicitly_rolled_back_impl(ditto: Ditto) {
let mut query_result = None;
let transaction_result = block_on(ditto.store().transaction(async |txn| {
let result = txn
.execute((
"INSERT INTO cats DOCUMENTS (:doc)",
json!({ "doc": { "_id": "test8", "value": 42 }}),
))
.await?;
query_result = Some(result);
Err::<TransactionCompletionAction, DittoError>(
std::io::Error::other("Intentional error to trigger rollback").into(),
)
}));
assert!(transaction_result.is_err());
assert!(query_result.is_some());
assert!(query_result.as_ref().unwrap().commit_id().is_none());
}
#[test]
fn commit_id_store_observer() {
let (_root1, ditto1) =
common::get_inactive_ditto(None).expect("Failed to create test ditto instance");
let (_root2, ditto2) =
common::get_inactive_ditto(None).expect("Failed to create test ditto instance");
commit_id_store_observer_impl(ditto1);
commit_id_store_observer_impl(ditto2);
}
fn commit_id_store_observer_impl(ditto: Ditto) {
let create_test_doc = |id: &str, value: i32| json!({ "doc": { "_id": id, "value": value }});
let initial_doc = create_test_doc("commit_id_observer_initial", 100);
block_on(
ditto
.store()
.execute(("INSERT INTO cats DOCUMENTS (:doc)", initial_doc)),
)
.expect("Failed to insert initial test document");
let (tx, rx) = tokio::sync::oneshot::channel();
let mut maybe_tx = Some(tx);
let observer = ditto
.store()
.register_observer("SELECT * FROM cats", move |result: QueryResult| {
assert!(
result.commit_id().is_none(),
"Store observer should never return commit_id, got: {:?}",
result.commit_id()
);
if let Some(tx) = maybe_tx.take() {
let _ = tx.send(result); }
})
.expect("Failed to register store observer");
let trigger_doc = create_test_doc("commit_id_observer_trigger", 200);
block_on(
ditto
.store()
.execute(("INSERT INTO cats DOCUMENTS (:doc)", trigger_doc)),
)
.expect("Failed to insert trigger document");
let _result = rx
.blocking_recv()
.expect("Observer did not fire within timeout");
observer.cancel();
}