surrealdb-core 3.3.1

A scalable, distributed, collaborative, document-graph database, for the realtime web
//! Index operation helpers for constructing and maintaining KV index entries.
//!
//! This module orchestrates index mutations for a single document: removing old
//! entries and inserting new ones across different index types (UNIQUE, regular,
//! search, fulltext, etc.). Index keys are built with key::index and field
//! values are encoded via key::value::Array.
//!
//! Numeric normalization in keys:
//! - Array normalizes Number values (Int/Float/Decimal) using a lexicographic numeric encoding so
//!   that byte-wise order matches numeric order. As a result, numerically equal values (e.g., 0,
//!   0.0, 0dec) map to the same key bytes.
//! - UNIQUE index behavior leverages this: equal numerics across variants will collide on the same
//!   index key and cause a uniqueness violation.
//!
//! Range scans and lookups benefit because a single probe/range can be used for
//! numeric predicates without fanning out per numeric variant.

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> {
		// Collect indexes or skip
		let ixs = match &opt.force {
			Force::All => self.doc_ctx.ix()?,
			_ if self.is_modified() => self.doc_ctx.ix()?,
			_ => return Ok(false),
		};
		// Get the current database
		let db = self.doc_ctx.db();
		// Get the document table
		let tb = self.doc_ctx.tb()?;
		// Check if the table is DROP
		if tb.drop {
			return Ok(false);
		}
		// Get the record id
		let rid = self.id()?;
		// Whether a doc-ID index deferred this record's removal to an in-progress
		// build (see `one_index`). When true the caller must not drop the shared
		// doc-ID mapping yet — the builder's replay still needs it.
		let mut doc_id_removal_deferred = false;
		// Loop through all index statements
		for ix in ixs.iter() {
			// Decommissioned indexes are ignored
			if ix.prepare_remove {
				continue;
			}
			// Calculate old values
			let o = Self::build_opt_values(stk, ctx, opt, ix, &self.initial).await?;
			// Calculate new values
			let n = Self::build_opt_values(stk, ctx, opt, ix, &self.current).await?;
			// For COUNT indexes with a condition, evaluate against the full document.
			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
			};
			// Update the index entries. Regular indexes update when the
			// stored values change. COUNT WHERE indexes have no indexed
			// values (so `o == n` is trivially true) and need to react
			// to the predicate-match flag flipping.
			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?;
			}
		}
		// Carry on
		Ok(doc_id_removal_deferred)
	}

	/// Durably marks the record's shared doc-ID mapping for a later reclaim,
	/// because an in-progress index build enqueued this delete for replay and
	/// may still need the mapping (see [`store_index_data`](Self::store_index_data)).
	///
	/// The `!dp` marker is written in this same (delete) transaction —
	/// atomically with the enqueued mutation — so the reclaim obligation exists
	/// from the moment the central removal is skipped. If the build replays the
	/// delete, the replay completes (or re-defers) the reclaim and consumes the
	/// marker; if the build never replays it (error, abort, `REMOVE INDEX`),
	/// the marker survives and the next doc-ID index build's sweep completes
	/// the reclaim instead (see `kvs::index::replay`).
	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
	}

	/// Removes the record's entry from the table's shared doc-ID space.
	///
	/// The doc-ID mapping (`!di`/`!dd`) is shared by every index on the table, so
	/// it must be dropped exactly once — after all per-index maintenance has
	/// released the record — rather than by any individual index. Called on record
	/// deletion, after [`store_index_data`](Self::store_index_data).
	///
	/// No-op for tables without a doc-ID consumer, so consumer-free tables
	/// never allocate or write a mapping. The consumers are the doc-ID
	/// indexes (full-text / HNSW / DiskAnn / doc-ID-format b-tree) and —
	/// once the table's `graph_doc_ids` marker is set — the numeric
	/// adjacency blocks holding its record ids: a mapping surviving its
	/// record would hand a re-created record its predecessor's id, which a
	/// block entry naming that id would then resolve to.
	pub(super) async fn remove_doc_id(&self, ctx: &FrozenContext) -> Result<()> {
		// A table being dropped has its whole prefix reclaimed separately.
		let tb = self.doc_ctx.tb()?;
		if tb.drop {
			return Ok(());
		}
		let db = self.doc_ctx.db();
		// Only tables with a consumer of the shared space maintain it.
		let ixs = self.doc_ctx.ix()?;
		if !tb.graph_doc_ids && !ixs.iter().any(IndexDefinition::uses_shared_doc_ids) {
			// `tb` is a plain snapshot read that registers no read-conflict, so a
			// concurrent fold committing the table's first numeric block between
			// this snapshot and this transaction's commit would otherwise go
			// undetected: the two transactions touch disjoint write-sets, so no
			// write-write check catches it either, even though this read
			// happened-before that write. Re-read the flag through the
			// transaction's settled read — locked, so commit-time OCC
			// validates it, and memoized, so a multi-record delete issues it
			// once per table — if a fold's flag-set is racing this
			// transaction, one of the two is forced to retry and the retry
			// observes the settled value.
			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> {
		// Does this index consume the table's shared doc-ID space?
		let doc_id_index = ix.uses_shared_doc_ids();
		// Get the index builder
		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? {
				// The index builder consumed the value, which means it is currently building the
				// index asynchronously, we don't index the document and let the index builder
				// do it later. For a doc-ID index this defers the record's removal — and its use
				// of the shared `!di`/`!dd` mapping — to the builder's replay, so report the
				// deferral to the caller so it does not drop the mapping before the replay runs.
				ConsumeResult::Enqueued => return Ok(doc_id_index),
				// The index builder is done, the index has been built; we can proceed normally
				ConsumeResult::Ignored(o, n) => (o, n),
				// The definition was retired after it was read from a cache.
				ConsumeResult::Retired => return Ok(false),
			}
		} else {
			(o, n)
		};
		// Store all the variables and parameters required by the index operation
		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);
		}
		// Keep track of compaction requests, we need to trigger them after the index operation
		let mut require_compaction = false;
		// Execute the index operation
		ic.compute(stk, &az_fn, &mut require_compaction).await?;
		// Did any compaction request have to be triggered?
		if require_compaction {
			ic.trigger_compaction().await?;
		}
		Ok(false)
	}

	/// Extract from the given document, the values required by the index and put then in an array.
	/// Eg. If the index is composed of the columns `name` and `instrument`
	/// Given this doc: { "id": 1, "instrument": "piano", "name": "Tobie" }
	/// It will return: ["Tobie", "piano"]
	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))
	}
}