use anyhow::Result;
use reblessive::tree::Stk;
use surrealdb_strand::TableName;
use uuid::Uuid;
use crate::catalog::providers::TableProvider;
use crate::catalog::{Error, TableDefinition};
use crate::ctx::FrozenContext;
use crate::dbs::Options;
use crate::doc::CursorDoc;
use crate::expr::Base;
use crate::expr::statements::remove::index::RemoveIndexStatement;
use crate::iam::{Action, ResourceKind};
use crate::kvs::index::{AbortLocalBuild, retire_durable_index};
use crate::legacy::expr_to_ident;
use crate::val::Value;
pub(crate) async fn remove_index_statement_compute(
this: &RemoveIndexStatement,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
doc: Option<&CursorDoc>,
) -> Result<Value> {
ctx.is_allowed(opt, Action::Edit, ResourceKind::Index, Base::Db)?;
let name = expr_to_ident(stk, ctx, opt, doc, &this.name, "index name").await?;
let table_name = TableName::new(expr_to_ident(stk, ctx, opt, doc, &this.what, "what").await?);
let (ns_name, db_name) = opt.ns_db()?;
let (ns, db) = ctx.expect_ns_db_ids(opt).await?;
let txn = ctx.tx();
let res = txn.expect_tb_index(ns, db, &table_name, &name).await;
let ix = match res {
Err(e) => {
if this.if_exists && matches!(e.downcast_ref(), Some(Error::IxNotFound { .. })) {
return Ok(Value::None);
}
return Err(e);
}
Ok(ix) => ix,
};
let tb = txn.expect_tb(ns, db, &table_name).await?;
let removed_last_doc_id_index = ix.uses_shared_doc_ids()
&& !tb.graph_doc_ids
&& !txn
.all_tb_indexes(ns, db, &table_name, None)
.await?
.iter()
.any(|other| other.index_id != ix.index_id && other.uses_shared_doc_ids());
ctx.get_index_stores().index_removed(ns, db, &tb, &ix).await?;
if let Some(index_builder) = ctx.get_index_builder() {
txn.on_commit(AbortLocalBuild::boxed(
index_builder.clone(),
ns,
db,
table_name.clone(),
ix.index_id,
))
.await;
}
retire_durable_index(&txn, ns, db, &table_name, ix.index_id, false).await?;
txn.del_tb_index_deferred(ns, db, &table_name, &name).await?;
if removed_last_doc_id_index {
txn.del_tb_doc_ids_deferred(ns, db, &table_name).await?;
}
txn.replace_tb(
ns_name,
db_name,
&TableDefinition {
cache_indexes_ts: Uuid::now_v7(),
..(*tb).clone()
},
)
.await?;
txn.clear_cache();
Ok(Value::None)
}
#[cfg(test)]
#[cfg(feature = "kv-tikv")]
mod tikv_concurrency {
use std::borrow::Cow;
use std::sync::Arc;
use surrealdb_strand::TableName;
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use web_time::Duration;
use crate::CommunityComposer;
use crate::catalog::providers::{DatabaseProvider, NamespaceProvider};
use crate::dbs::Session;
use crate::key::schema::{DocKeyPrefix, DocLookupPrefix, RootRoot, VersionKey};
use crate::kvs::{Datastore, TransactionType};
async fn fresh_tikv_ds() -> Arc<Datastore> {
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
}
async fn drain_reclaim_queue(ds: &Arc<Datastore>) {
const MAX_PASSES: usize = 64;
for _ in 0..MAX_PASSES {
let (batches, _) = Datastore::reclaim_tombstones(
Arc::clone(ds),
Duration::from_secs(1),
Duration::ZERO,
CancellationToken::new(),
)
.await
.unwrap();
if batches == 0 {
return;
}
}
}
async fn remove_index(ds: Arc<Datastore>, sql: String) {
let ses = Session::owner().with_ns("test").with_db("test");
let mut last = String::new();
for _ in 0..200 {
match ds.execute(&sql, &ses, None).await {
Ok(responses) => {
if responses.into_iter().all(|r| r.result.is_ok()) {
return;
}
}
Err(e) => last = e.to_string(),
}
}
panic!("`{sql}` did not succeed after retries: {last}");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn concurrent_last_consumer_removal_purges_mappings() {
const ITERS: usize = 6;
let ds = fresh_tikv_ds().await;
let ses = Session::owner().with_ns("test").with_db("test");
for k in 0..ITERS {
let tb = format!("t{k}");
let setup = format!(
"DEFINE ANALYZER a{k} TOKENIZERS blank;
DEFINE INDEX i1 ON {tb} FIELDS a1 FULLTEXT ANALYZER a{k} BM25;
DEFINE INDEX i2 ON {tb} FIELDS a2 FULLTEXT ANALYZER a{k} BM25;
CREATE {tb}:1 SET a1 = 'hello', a2 = 'world';
CREATE {tb}:2 SET a1 = 'foo', a2 = 'bar';"
);
for r in ds.execute(&setup, &ses, None).await.unwrap() {
r.result.unwrap();
}
let a = tokio::spawn(remove_index(
Arc::clone(&ds),
format!("REMOVE INDEX IF EXISTS i1 ON {tb}"),
));
let b = tokio::spawn(remove_index(
Arc::clone(&ds),
format!("REMOVE INDEX IF EXISTS i2 ON {tb}"),
));
a.await.unwrap();
b.await.unwrap();
drain_reclaim_queue(&ds).await;
let tx = ds.transaction(TransactionType::Read).await.unwrap();
let ns = tx.get_ns_by_name("test", None).await.unwrap().unwrap().namespace_id;
let db = tx.get_db_by_name("test", "test", None).await.unwrap().unwrap().database_id;
let tb_name: TableName = tb.as_str().into();
let forward = tx
.getr(DocLookupPrefix::new(ns, db, Cow::Borrowed(&tb_name)).range().unwrap(), None)
.await
.unwrap();
let reverse = tx
.getr(DocKeyPrefix::new(ns, db, Cow::Borrowed(&tb_name)).range().unwrap(), None)
.await
.unwrap();
tx.cancel().await.unwrap();
assert!(
forward.is_empty() && reverse.is_empty(),
"iteration {k}: shared doc-ID mappings leaked after removing the last consumer: \
{} forward, {} reverse",
forward.len(),
reverse.len(),
);
}
}
async fn execute_retrying(ds: Arc<Datastore>, sql: String) {
let ses = Session::owner().with_ns("test").with_db("test");
let mut last = String::new();
for _ in 0..200 {
match ds.execute(&sql, &ses, None).await {
Ok(responses) => {
if responses.into_iter().all(|r| r.result.is_ok()) {
return;
}
}
Err(e) => last = e.to_string(),
}
}
panic!("`{sql}` did not succeed after retries: {last}");
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
#[serial_test::serial]
async fn concurrent_last_consumer_overwrite_purges_mappings() {
const ITERS: usize = 6;
let ds = fresh_tikv_ds().await;
let ses = Session::owner().with_ns("test").with_db("test");
for k in 0..ITERS {
let tb = format!("o{k}");
let setup = format!(
"DEFINE ANALYZER oa{k} TOKENIZERS blank;
DEFINE INDEX i1 ON {tb} FIELDS a1 FULLTEXT ANALYZER oa{k} BM25;
DEFINE INDEX i2 ON {tb} FIELDS a2 FULLTEXT ANALYZER oa{k} BM25;
CREATE {tb}:1 SET a1 = 'hello', a2 = 'world';
CREATE {tb}:2 SET a1 = 'foo', a2 = 'bar';"
);
for r in ds.execute(&setup, &ses, None).await.unwrap() {
r.result.unwrap();
}
let a = tokio::spawn(execute_retrying(
Arc::clone(&ds),
format!("DEFINE INDEX OVERWRITE i1 ON {tb} FIELDS a1"),
));
let b = tokio::spawn(execute_retrying(
Arc::clone(&ds),
format!("DEFINE INDEX OVERWRITE i2 ON {tb} FIELDS a2"),
));
a.await.unwrap();
b.await.unwrap();
drain_reclaim_queue(&ds).await;
let tx = ds.transaction(TransactionType::Read).await.unwrap();
let ns = tx.get_ns_by_name("test", None).await.unwrap().unwrap().namespace_id;
let db = tx.get_db_by_name("test", "test", None).await.unwrap().unwrap().database_id;
let tb_name: TableName = tb.as_str().into();
let forward = tx
.getr(DocLookupPrefix::new(ns, db, Cow::Borrowed(&tb_name)).range().unwrap(), None)
.await
.unwrap();
let reverse = tx
.getr(DocKeyPrefix::new(ns, db, Cow::Borrowed(&tb_name)).range().unwrap(), None)
.await
.unwrap();
tx.cancel().await.unwrap();
assert!(
forward.is_empty() && reverse.is_empty(),
"iteration {k}: shared doc-ID mappings leaked after overwriting the last \
consumer: {} forward, {} reverse",
forward.len(),
reverse.len(),
);
}
}
}