#![cfg(feature = "kv-tikv")]
use std::collections::HashSet;
use std::sync::Arc;
use surrealdb_catalog::{DatabaseId, NamespaceId};
use surrealdb_strand::TableName;
use uuid::Uuid;
use crate::CommunityComposer;
use crate::key::schema::{RootRoot, VersionKey};
use crate::kvs::sequences::Sequences;
use crate::kvs::{Datastore, TransactionFactory, TransactionType};
async fn fresh_tikv_tf() -> TransactionFactory {
let ds = Datastore::builder()
.with_id(Uuid::new_v4())
.build_with_factory_path("tikv:127.0.0.1:2379", CommunityComposer())
.await
.unwrap();
let tx = ds.transaction(TransactionType::Write).await.unwrap();
tx.delr(RootRoot {}.range_subtree().unwrap()).await.unwrap();
tx.del_key(&VersionKey {}).await.unwrap();
tx.commit().await.unwrap();
ds.transaction_factory().clone()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn table_doc_ids_unique_across_nodes() {
use tokio::sync::Barrier;
const NODES: usize = 4;
const PER_NODE: usize = 100;
const BATCH: u32 = 1;
let tf = fresh_tikv_tf().await;
let ns = NamespaceId(1);
let db = DatabaseId(1);
let tb: TableName = "t".into();
let barrier = Arc::new(Barrier::new(NODES));
let mut handles = Vec::with_capacity(NODES);
for _ in 0..NODES {
let seqs = Sequences::new(tf.clone(), Uuid::new_v4());
let tb = tb.clone();
let barrier = Arc::clone(&barrier);
handles.push(tokio::spawn(async move {
barrier.wait().await;
let mut ids = Vec::with_capacity(PER_NODE);
for _ in 0..PER_NODE {
let id = seqs.next_table_doc_id(None, ns, db, tb.clone(), BATCH).await.unwrap();
ids.push(id);
}
ids
}));
}
let mut all = Vec::with_capacity(NODES * PER_NODE);
for h in handles {
all.extend(h.await.unwrap());
}
let unique: HashSet<_> = all.iter().copied().collect();
assert_eq!(
unique.len(),
all.len(),
"table doc-IDs must be globally unique across nodes; {} duplicate(s) out of {}",
all.len() - unique.len(),
all.len(),
);
}
async fn define_sequence(
tf: &TransactionFactory,
sqs: &Sequences,
ns: NamespaceId,
db: DatabaseId,
name: &str,
start: i64,
batch: u32,
) -> anyhow::Result<()> {
use std::borrow::Cow;
use crate::key::schema::{SeqBatchPrefix, SeqStatePrefix, SequenceKey};
let tx = tf.transaction(TransactionType::Write, sqs.clone()).await?;
let res = async {
tx.set_key(
&SequenceKey {
ns,
db,
sq: Cow::Borrowed(name),
},
&surrealdb_catalog::SequenceDefinition {
name: name.into(),
batch,
start,
timeout: None,
},
)
.await?;
tx.delr(
SeqBatchPrefix {
ns,
db,
sq: Cow::Borrowed(name),
}
.range()?,
)
.await?;
tx.delr(
SeqStatePrefix {
ns,
db,
sq: Cow::Borrowed(name),
}
.range()?,
)
.await?;
Ok::<(), anyhow::Error>(())
}
.await;
match res {
Ok(()) => tx.commit().await,
Err(e) => {
tx.cancel().await?;
Err(e)
}
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn allocation_never_reissues_across_repositions_on_tikv() {
use std::time::Duration;
let tf = fresh_tikv_tf().await;
let ns = NamespaceId(1);
let db = DatabaseId(1);
let name = "sq";
const BATCH: u32 = 10;
const NODES: usize = 3;
const TASKS_PER_NODE: usize = 2;
const ALLOCATIONS: usize = 500;
const STRIDE: i64 = 10_000;
const MAX_EPOCHS: i64 = 40;
let seed = Sequences::new(tf.clone(), Uuid::new_v4());
define_sequence(&tf, &seed, ns, db, name, 1, BATCH).await.unwrap();
let done = Arc::new(std::sync::atomic::AtomicBool::new(false));
let mut workers = Vec::new();
for _ in 0..NODES {
let seqs = Sequences::new(tf.clone(), Uuid::new_v4());
for _ in 0..TASKS_PER_NODE {
let seqs = seqs.clone();
let tf = tf.clone();
workers.push(tokio::spawn(async move {
let mut got = Vec::new();
for _ in 0..ALLOCATIONS {
let tx = tf.transaction(TransactionType::Read, seqs.clone()).await.unwrap();
let res = seqs.next_user_sequence_id(None, &tx, ns, db, name).await;
let _ = tx.cancel().await;
if let Ok(v) = res {
got.push(v);
}
}
got
}));
}
}
let repositioner = {
let tf = tf.clone();
let done = Arc::clone(&done);
tokio::spawn(async move {
let sqs = Sequences::new(tf.clone(), Uuid::new_v4());
let mut epoch = 0i64;
while !done.load(std::sync::atomic::Ordering::Relaxed) && epoch < MAX_EPOCHS {
epoch += 1;
while define_sequence(&tf, &sqs, ns, db, name, epoch * STRIDE, BATCH).await.is_err()
{
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
epoch
})
};
let mut issued = Vec::new();
for w in workers {
issued.extend(
tokio::time::timeout(Duration::from_secs(120), w)
.await
.expect("allocation stalled")
.expect("a worker panicked"),
);
}
done.store(true, std::sync::atomic::Ordering::Relaxed);
let repositions = tokio::time::timeout(Duration::from_secs(120), repositioner)
.await
.expect("repositioning stalled")
.expect("the repositioner panicked");
assert!(
repositions > 1,
"allocation finished before repositioning got going, so nothing was exercised: \
{repositions} reposition(s)"
);
assert!(!issued.is_empty(), "the fixture issued nothing, so it asserts nothing");
let unique: HashSet<i64> = issued.iter().copied().collect();
assert_eq!(
unique.len(),
issued.len(),
"an id was issued twice across a reposition on TiKV ({} of {} distinct)",
unique.len(),
issued.len(),
);
}