use std::borrow::Cow;
use anyhow::Result;
use reblessive::tree::Stk;
use crate::catalog::{DatabaseDefinition, IndexDefinition, TableDefinition};
use crate::ctx::FrozenContext;
use crate::dbs::{Force, Options};
use crate::doc::{CursorDoc, Document};
use crate::exe::FlowResultExt as _;
use crate::idx::docids::TableDocIds;
use crate::idx::index::IndexOperation;
use crate::key::schema::DocPendingKey;
use crate::kvs::index::{ConsumeResult, IndexMutation};
use crate::legacy::analyzer_function::LegacyAnalyzerFunction;
use crate::val::{RecordId, RecordIdentity, Value};
impl Document {
pub(super) async fn store_index_data(
&self,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
) -> Result<bool> {
let ixs = match &opt.force {
Force::All => self.doc_ctx.ix()?,
_ if self.is_modified() => self.doc_ctx.ix()?,
_ => return Ok(false),
};
let db = self.doc_ctx.db();
let tb = self.doc_ctx.tb()?;
if tb.drop {
return Ok(false);
}
let rid = self.id()?;
let mut doc_id_removal_deferred = false;
for ix in ixs.iter() {
if ix.prepare_remove {
continue;
}
let o = Self::build_opt_values(stk, ctx, opt, ix, &self.initial).await?;
let n = Self::build_opt_values(stk, ctx, opt, ix, &self.current).await?;
let count_cond_match = if let Some(cond) = &ix.count_cond {
let expr = &cond.0;
let old_matches = stk
.run(|stk| {
crate::legacy::expr_compute(expr, stk, ctx, opt, Some(&self.initial))
})
.await
.catch_return()?
.is_truthy();
let new_matches = stk
.run(|stk| {
crate::legacy::expr_compute(expr, stk, ctx, opt, Some(&self.current))
})
.await
.catch_return()?
.is_truthy();
Some((old_matches, new_matches))
} else {
None
};
let cond_changed = matches!(count_cond_match, Some((o, n)) if o != n);
if o != n || cond_changed {
doc_id_removal_deferred |=
Self::one_index(db, tb, stk, ctx, opt, ix, o, n, &rid, count_cond_match)
.await?;
}
}
Ok(doc_id_removal_deferred)
}
pub(super) async fn defer_doc_id_removal(&self, ctx: &FrozenContext) -> Result<()> {
let db = self.doc_ctx.db();
let rid = self.id()?;
let dp = DocPendingKey::new(
db.namespace_id,
db.database_id,
Cow::Borrowed(&rid.table),
RecordIdentity(rid.key.clone()),
);
ctx.tx().set_key(&dp, &()).await
}
pub(super) async fn remove_doc_id(&self, ctx: &FrozenContext) -> Result<()> {
let tb = self.doc_ctx.tb()?;
if tb.drop {
return Ok(());
}
let db = self.doc_ctx.db();
let ixs = self.doc_ctx.ix()?;
if !tb.graph_doc_ids && !ixs.iter().any(IndexDefinition::uses_shared_doc_ids) {
let settled =
ctx.tx().graph_doc_ids_settled(db.namespace_id, db.database_id, &tb.name).await?;
if !settled {
return Ok(());
}
}
let rid = self.id()?;
TableDocIds::new(db.namespace_id, db.database_id, rid.table.clone())
.remove(&ctx.tx(), &rid.key)
.await?;
Ok(())
}
#[allow(clippy::too_many_arguments)]
async fn one_index(
db: &DatabaseDefinition,
tb: &TableDefinition,
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
ix: &IndexDefinition,
o: Option<Vec<Value>>,
n: Option<Vec<Value>>,
rid: &RecordId,
count_cond_match: Option<(bool, bool)>,
) -> Result<bool> {
let doc_id_index = ix.uses_shared_doc_ids();
let (o, n) = if let Some(ib) = ctx.get_index_builder() {
let mutation = IndexMutation {
old_values: o,
new_values: n,
rid,
count_cond_match,
};
match ib.consume(db, ctx, ix, mutation).await? {
ConsumeResult::Enqueued => return Ok(doc_id_index),
ConsumeResult::Ignored(o, n) => (o, n),
ConsumeResult::Retired => return Ok(false),
}
} else {
(o, n)
};
let az_fn = LegacyAnalyzerFunction::new(ctx, opt);
let mut ic =
IndexOperation::new(ctx, db.namespace_id, db.database_id, tb.table_id, ix, o, n, rid);
if let Some((old_matches, new_matches)) = count_cond_match {
ic = ic.with_count_cond_match(old_matches, new_matches);
}
let mut require_compaction = false;
ic.compute(stk, &az_fn, &mut require_compaction).await?;
if require_compaction {
ic.trigger_compaction().await?;
}
Ok(false)
}
pub(crate) async fn build_opt_values(
stk: &mut Stk,
ctx: &FrozenContext,
opt: &Options,
ix: &IndexDefinition,
doc: &CursorDoc,
) -> Result<Option<Vec<Value>>> {
if doc.doc.as_ref().is_nullish() {
return Ok(None);
}
let mut o = Vec::with_capacity(ix.cols.len());
for idiom in ix.cols.iter() {
let v = crate::legacy::idiom_compute(idiom, stk, ctx, opt, Some(doc))
.await
.catch_return()?;
o.push(v);
}
Ok(Some(o))
}
}