use std::sync::Arc;
use anyhow::{Result, bail};
use reblessive::tree::Stk;
use surrealdb_types::ToSql;
use uuid::Uuid;
use crate::catalog::providers::TableProvider;
use crate::catalog::{
Error as CatalogError, INDEX_FORMAT_VERSION, IndexDefinition, TableDefinition, TableId,
};
use crate::ctx::FrozenContext;
use crate::dbs::Options;
use crate::doc::CursorDoc;
use crate::err::EngineError;
use crate::exe::FlowResultExt;
use crate::exec::Error as ExecError;
use crate::expr::index_kind::Index;
use crate::expr::statements::define::DefineKind;
use crate::expr::statements::define::index::DefineIndexStatement;
use crate::expr::{Base, Idiom, Part};
use crate::iam::{Action, ResourceKind};
use crate::kvs::index::{
AbortLocalBuild, CleanUncommittedBuild, IndexBuilder, retire_durable_index,
};
use crate::kvs::tx::DocIdsReclaimClaim;
use crate::kvs::{DatastoreError, Transaction};
use crate::legacy::{expr_to_ident, exprs_to_fields};
use crate::val::{TableName, Value};
#[instrument(level = "trace", name = "DefineIndexStatement::compute", skip_all)]
pub(crate) async fn define_index_statement_compute(
this: &DefineIndexStatement,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
doc: Option<&CursorDoc>,
) -> Result<Value> {
ctx.is_allowed(opt, Action::Edit, ResourceKind::Index, Base::Db)?;
let txn = ctx.tx();
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, "index table").await?);
let (ns, db) = opt.ns_db()?;
let tb = txn.get_or_add_tb(Some(ctx), ns, db, &table_name, None).await?;
let tb_name = tb.name.clone();
anyhow::ensure!(
crate::kvs::lightweight::lightweight_relation(&tb.table_type).is_none(),
crate::exec::Error::Thrown(
"a LIGHTWEIGHT relation cannot carry an index: its edges store no records".to_owned(),
)
);
let existing = txn.get_tb_index(tb.namespace_id, tb.database_id, &tb_name, &name, None).await?;
if existing.is_some() {
match this.kind {
DefineKind::Default => {
if !opt.import {
bail!(CatalogError::IxAlreadyExists {
name: this.name.to_sql(),
});
}
}
DefineKind::Overwrite => {}
DefineKind::IfNotExists => return Ok(Value::None),
}
}
let cols = exprs_to_fields(stk, ctx, opt, doc, this.cols.as_slice()).await?;
for idiom in cols.iter() {
let fd = idiom.to_raw_string();
if let Some(f) =
txn.get_tb_field(tb.namespace_id, tb.database_id, &tb_name, &fd, None).await?
{
if f.computed.is_some() {
bail!(ExecError::ComputedFieldCannotBeIndexed {
field: fd,
index: name
});
}
continue;
}
if let Some(Part::Field(first)) = idiom.0.first() &&
let Some(f) =
txn.get_tb_field(tb.namespace_id, tb.database_id, &tb_name, first, None).await?
&& f.field_kind.as_ref().is_none_or(|k| k.allows_sub_fields())
{
if f.computed.is_some() {
bail!(ExecError::ComputedFieldCannotBeIndexed {
field: first.as_str().to_owned(),
index: name
});
}
continue;
}
if tb.schemafull {
bail!(CatalogError::FdNotFound {
name: idiom.to_raw_string(),
});
}
}
let comment = stk
.run(|stk| crate::legacy::expr_compute(&this.comment, stk, ctx, opt, doc))
.await
.catch_return()?
.cast_to()?;
if let Some(ix) = existing.as_ref()
&& this.kind == DefineKind::Default
&& opt.import
&& import_replay_can_reuse_index(ix, &table_name, &cols, &this.index)
{
let index_def = IndexDefinition {
index_id: ix.index_id,
name: name.into(),
table_name,
cols: cols.clone(),
index: this.index.clone().into(),
count_cond: match &this.index {
Index::Count(Some(cond)) => Some(cond.clone()),
_ => None,
},
comment,
prepare_remove: false,
format_version: ix.format_version,
};
txn.put_tb_index(tb.namespace_id, tb.database_id, &tb_name, &index_def).await?;
refresh_table_index_cache(ctx, &txn, ns, db, &tb).await?;
return Ok(Value::None);
}
if !matches!(this.index, Index::Count(_))
&& txn.claim_tb_doc_ids_reclaim(tb.namespace_id, tb.database_id, &tb_name).await?
== DocIdsReclaimClaim::InProgress
{
bail!(DatastoreError::QueryNotExecuted {
message: format!(
"The shared document-ID space for table `{tb_name}` is still being reclaimed; retry DEFINE INDEX after cleanup completes"
),
});
}
if let Some(ix) = existing.as_ref() {
let replacement_uses_doc_ids = !matches!(this.index, Index::Count(_));
let purge_table_doc_ids = ix.uses_shared_doc_ids()
&& !replacement_uses_doc_ids
&& !tb.graph_doc_ids
&& !txn
.all_tb_indexes(tb.namespace_id, tb.database_id, &tb_name, None)
.await?
.iter()
.any(|other| other.index_id != ix.index_id && other.uses_shared_doc_ids());
ctx.get_index_stores().index_removed(tb.namespace_id, tb.database_id, &tb, ix).await?;
if let Some(index_builder) = ctx.get_index_builder() {
txn.on_commit(AbortLocalBuild::boxed(
index_builder.clone(),
tb.namespace_id,
tb.database_id,
tb_name.clone(),
ix.index_id,
))
.await;
}
retire_durable_index(&txn, tb.namespace_id, tb.database_id, &tb_name, ix.index_id, false)
.await?;
txn.del_tb_index_deferred(tb.namespace_id, tb.database_id, &tb_name, &name).await?;
if purge_table_doc_ids {
txn.del_tb_doc_ids_deferred(tb.namespace_id, tb.database_id, &tb_name).await?;
}
}
let index_id = ctx
.try_get_sequences()?
.next_index_id(Some(ctx), tb.namespace_id, tb.database_id, tb_name.clone())
.await?;
let index_def = IndexDefinition {
index_id,
name: name.clone().into(),
table_name,
cols: cols.clone(),
index: this.index.clone().into(),
count_cond: match &this.index {
Index::Count(Some(cond)) => Some(cond.clone()),
_ => None,
},
comment,
prepare_remove: false,
format_version: INDEX_FORMAT_VERSION,
};
txn.put_tb_index(tb.namespace_id, tb.database_id, &tb_name, &index_def).await?;
refresh_table_index_cache(ctx, &txn, ns, db, &tb).await?;
let index_builder =
ctx.get_index_builder().ok_or_else(|| EngineError::unreachable("No Index Builder"))?;
txn.on_rollback(CleanUncommittedBuild::boxed(
index_builder.clone(),
index_builder.transaction_factory(),
txn.sequences(),
tb.namespace_id,
tb.database_id,
tb_name.clone(),
index_id,
))
.await;
run_indexing_with_builder(
index_builder,
ctx,
opt,
tb.table_id,
Arc::new(index_def),
!this.concurrently,
)
.await?;
Ok(Value::None)
}
pub(crate) async fn refresh_table_index_cache(
_ctx: &FrozenContext,
txn: &Transaction,
ns: &str,
db: &str,
tb: &TableDefinition,
) -> Result<()> {
txn.replace_tb(
ns,
db,
&TableDefinition {
cache_indexes_ts: Uuid::now_v7(),
..tb.clone()
},
)
.await?;
txn.clear_cache();
Ok(())
}
pub(crate) async fn run_indexing(
ctx: &FrozenContext,
opt: &Options,
tb: TableId,
ix: Arc<IndexDefinition>,
blocking: bool,
) -> Result<()> {
let index_builder =
ctx.get_index_builder().ok_or_else(|| EngineError::unreachable("No Index Builder"))?;
run_indexing_with_builder(index_builder, ctx, opt, tb, ix, blocking).await
}
pub(crate) async fn run_indexing_with_builder(
index_builder: &IndexBuilder,
ctx: &FrozenContext,
opt: &Options,
tb: TableId,
ix: Arc<IndexDefinition>,
blocking: bool,
) -> Result<()> {
let rcv = index_builder.build(ctx, opt.clone(), tb, ix, blocking).await?;
if let Some(rcv) = rcv {
rcv.await.map_err(|_| DatastoreError::IndexingBuildingCancelled {
reason: "Channel shutdown".to_string(),
})?
} else {
Ok(())
}
}
pub(crate) fn import_replay_can_reuse_index(
ix: &IndexDefinition,
table_name: &TableName,
cols: &[Idiom],
index: &Index,
) -> bool {
!ix.prepare_remove
&& ix.table_name.as_str() == table_name.as_str()
&& ix.cols.as_slice() == cols
&& ix.index == crate::catalog::Index::from(index.clone())
}