use std::thread;
use reifydb_codec::row::bytes::EncodedBytes;
use reifydb_core::{
common::CommitVersion,
delta::Delta,
interface::{
catalog::id::QueueId,
store::{MultiVersionCommit, MultiVersionGet},
},
key::{any::TaggedKey, queue::QueueDeduplicationKey},
};
use reifydb_store_multi::store::StandardMultiStore;
use reifydb_value::util::cowvec::CowVec;
fn encoded_bytes(bytes: &[u8]) -> EncodedBytes {
EncodedBytes(CowVec::new(bytes.to_vec()))
}
#[test]
fn concurrent_reads_during_writes_no_deadlock() {
let (store, _guard) = StandardMultiStore::testing_memory_with_persistent_sqlite();
let key: TaggedKey =
QueueDeduplicationKey::new(QueueId(1), b"k".iter().map(|b| !b).collect::<Vec<u8>>()).into();
MultiVersionCommit::commit(
&store,
CowVec::new(vec![Delta::Set {
key: key.clone(),
bytes: encoded_bytes(b"v0"),
}]),
CommitVersion(1),
)
.unwrap();
let last: u64 = 200;
let readers: Vec<_> = (0..4)
.map(|_| {
let store = store.clone();
let key = key.clone();
thread::spawn(move || {
for _ in 0..500 {
let got = store.get(&key, CommitVersion(u64::MAX)).unwrap();
assert!(got.is_some());
}
})
})
.collect();
for v in 2..=last {
MultiVersionCommit::commit(
&store,
CowVec::new(vec![Delta::Set {
key: key.clone(),
bytes: encoded_bytes(format!("v{v}").as_bytes()),
}]),
CommitVersion(v),
)
.unwrap();
}
for reader in readers {
reader.join().expect("reader thread panicked (deadlock or read error)");
}
let final_value = store.get(&key, CommitVersion(u64::MAX)).unwrap().unwrap();
assert_eq!(final_value.bytes.as_slice(), format!("v{last}").as_bytes());
}