use std::collections::HashMap;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::time::Duration;
use anyhow::{Result, ensure};
use chrono::Utc;
use futures::channel::oneshot::{Receiver, Sender, channel};
#[cfg(not(target_family = "wasm"))]
use tokio::spawn;
use tokio::sync::RwLock;
use tokio::time::sleep;
use uuid::Uuid;
#[cfg(target_family = "wasm")]
use wasm_bindgen_futures::spawn_local as spawn;
use web_time::Instant;
use super::state::{
build_owner_expired, delete_stale_build_queues, durable_index_error_reason,
durable_report_count, is_condition_not_met, report_status_from_phase,
};
use super::{
AcquiredBuild, BUILD_CLOSING_SLEEP, BuildGeneration, IndexBuildPhase, IndexBuildReportStatus,
IndexBuildState, IndexBuilding,
};
use crate::catalog::providers::TableProvider;
use crate::catalog::{DatabaseId, Index, IndexDefinition, IndexId, NamespaceId, TableId};
use crate::ctx::{Context, FrozenContext};
use crate::dbs::Options;
use crate::err::Error;
use crate::idx::IndexKeyBase;
use crate::idx::index::IndexOperation;
use crate::key::index::all as index_all;
use crate::key::record;
use crate::kvs::LockType::Optimistic;
use crate::kvs::ds::TransactionFactory;
#[cfg(test)]
use crate::kvs::testing::{
NonRetryableErrorSite, RetryableConflictSite, maybe_inject_non_retryable_error,
maybe_inject_retryable_conflict,
};
use crate::kvs::util::advance_key;
use crate::kvs::{
INDEXING_BATCH_MAX_BYTES, INDEXING_BATCH_SIZE, INDEXING_PROBE_BATCH_SIZE, KVKey, Transaction,
TransactionType, is_retryable_transaction_conflict, is_shutdown_error,
};
use crate::mem::ALLOC;
use crate::val::{RecordId, RecordIdKey, TableName, Value};
pub(super) type SharedIndexKey = Arc<IndexKey>;
fn is_memory_threshold_error(err: &anyhow::Error) -> bool {
matches!(err.downcast_ref::<Error>(), Some(Error::QueryBeyondMemoryThreshold))
}
#[cfg(test)]
fn maybe_inject_initial_batch_commit_error(node_id: Uuid) -> Result<()> {
maybe_inject_non_retryable_error(
NonRetryableErrorSite::ConcurrentIndexInitialBatchCommit,
node_id,
)?;
maybe_inject_non_retryable_error(
NonRetryableErrorSite::ConcurrentIndexInitialBatchShutdown,
node_id,
)?;
maybe_inject_non_retryable_error(
NonRetryableErrorSite::ConcurrentIndexInitialBatchMemoryThreshold,
node_id,
)
}
#[derive(Hash, PartialEq, Eq)]
pub(super) struct IndexKey {
pub(super) ns: NamespaceId,
pub(super) db: DatabaseId,
pub(super) tb: TableName,
pub(super) ix: IndexId,
}
impl IndexKey {
pub(super) fn new(ns: NamespaceId, db: DatabaseId, tb: &TableName, ix: IndexId) -> Self {
Self {
ns,
db,
tb: tb.to_owned(),
ix,
}
}
}
#[derive(Clone)]
pub(crate) struct IndexBuilder {
pub(super) tf: TransactionFactory,
pub(super) indexes: Arc<RwLock<HashMap<SharedIndexKey, IndexBuilding>>>,
}
enum BuildStart {
Started,
RemoteOwner(IndexBuilding),
}
pub(crate) struct IndexMutation<'a> {
pub(crate) old_values: Option<Vec<Value>>,
pub(crate) new_values: Option<Vec<Value>>,
pub(crate) rid: &'a RecordId,
pub(crate) count_cond_match: Option<(bool, bool)>,
}
impl IndexBuilder {
pub(in crate::kvs) fn new(tf: TransactionFactory) -> Self {
Self {
tf,
indexes: Default::default(),
}
}
pub(crate) fn transaction_factory(&self) -> TransactionFactory {
self.tf.clone()
}
pub(crate) async fn has_unfinished_build(&self) -> bool {
self.indexes.read().await.values().any(|building| !building.is_finished())
}
#[allow(clippy::too_many_arguments)]
async fn start_building(
&self,
ctx: &FrozenContext,
opt: Options,
tb: TableId,
ix: Arc<IndexDefinition>,
ix_key: SharedIndexKey,
sdr: Option<Sender<Result<()>>>,
) -> Result<BuildStart> {
let building = Arc::new(Building::new(ctx, self.tf.clone(), opt, tb, ix, ix_key)?);
let acquired = match building.acquire_build_state().await {
Ok(Some(acquired)) => acquired,
Ok(None) => return Ok(BuildStart::RemoteOwner(building)),
Err(err) => return Err(err),
};
self.start_acquired_building(Arc::clone(&building), acquired, sdr).await?;
Ok(BuildStart::Started)
}
async fn start_acquired_building(
&self,
building: IndexBuilding,
acquired: AcquiredBuild,
sdr: Option<Sender<Result<()>>>,
) -> Result<()> {
{
let mut indexes = self.indexes.write().await;
if let Some(existing) = indexes.get(&building.ix_key) {
ensure!(
existing.is_finished(),
Error::IndexAlreadyBuilding {
name: building.ix.name.to_string(),
}
);
}
indexes.insert(Arc::clone(&building.ix_key), Arc::clone(&building));
}
let b = Arc::clone(&building);
spawn(async move {
let guard = BuildingFinishGuard(Arc::clone(&b));
let r = b.run_acquired(acquired).await;
let generation = b.build_generation.load(Ordering::Acquire);
if let Err(err) = &r {
if is_shutdown_error(err) {
info!(
index = %b.ix.name,
table = %b.ix.table_name,
"index build interrupted by datastore shutdown; \
it will resume after restart"
);
} else if is_memory_threshold_error(err) {
warn!(
index = %b.ix.name,
table = %b.ix.table_name,
"index build interrupted by the memory threshold; \
it will resume from its checkpoint once the owner \
lease expires"
);
} else if generation != 0 {
let _ = b.mark_durable_error(generation, err.to_string()).await;
}
} else if b.aborted.load(Ordering::Acquire) && generation != 0 {
let _ = b.mark_durable_aborted(generation).await;
}
drop(guard);
if let Some(s) = sdr
&& s.send(r).is_err()
{
warn!("Failed to send index building result to the consumer");
}
});
Ok(())
}
async fn wait_for_remote_building(
&self,
ctx: &FrozenContext,
building: IndexBuilding,
) -> Result<()> {
loop {
if let Some(reason) = ctx.done(true)? {
return Err(Error::from(reason).into());
}
let Some(state) = building.read_durable_build_state().await? else {
return Err(Error::IndexingBuildingCancelled {
reason: format!("Index {} build state no longer exists", building.ix.name),
}
.into());
};
match state.phase {
IndexBuildPhase::Online => return Ok(()),
IndexBuildPhase::Error => {
return Err(Error::IndexingBuildingCancelled {
reason: format!(
"{}. Run `REBUILD INDEX {} ON {}` to retry the build",
durable_index_error_reason(&building.ix, &state),
building.ix.name,
building.ix.table_name
),
}
.into());
}
IndexBuildPhase::Building | IndexBuildPhase::Closing => {
if build_owner_expired(&state, Utc::now())
&& let Some(acquired) = building.takeover_expired_build_state().await?
{
let (s, r) = channel();
self.start_acquired_building(Arc::clone(&building), acquired, Some(s))
.await?;
return r.await.map_err(|_| Error::IndexingBuildingCancelled {
reason: "Channel shutdown".to_string(),
})?;
}
sleep(BUILD_CLOSING_SLEEP).await;
}
}
}
}
pub(crate) async fn build(
&self,
ctx: &FrozenContext,
opt: Options,
tb: TableId,
ix: Arc<IndexDefinition>,
blocking: bool,
) -> Result<Option<Receiver<Result<()>>>> {
ix.expect_not_prepare_remove()?;
let (ns, db) = ctx.expect_ns_db_ids(&opt).await?;
let key = Arc::new(IndexKey::new(ns, db, &ix.table_name, ix.index_id));
let (rcv, sdr) = if blocking {
let (s, r) = channel();
(Some(r), Some(s))
} else {
(None, None)
};
if let Some(existing) = self.indexes.read().await.get(&key) {
ensure!(
existing.is_finished(),
Error::IndexAlreadyBuilding {
name: ix.name.to_string(),
}
);
}
match self.start_building(ctx, opt, tb, ix, key, sdr).await? {
BuildStart::Started => Ok(rcv),
BuildStart::RemoteOwner(building) if blocking => {
self.wait_for_remote_building(ctx, building).await?;
Ok(None)
}
BuildStart::RemoteOwner(_) => Ok(None),
}
}
pub(crate) async fn resume_stalled(
&self,
ctx: &FrozenContext,
opt: Options,
ns: NamespaceId,
db: DatabaseId,
tb: TableId,
ix: Arc<IndexDefinition>,
) -> Result<bool> {
let key = Arc::new(IndexKey::new(ns, db, &ix.table_name, ix.index_id));
if let Some(existing) = self.indexes.read().await.get(&key)
&& !existing.is_finished()
{
return Ok(false);
}
let building = Arc::new(Building::new(ctx, self.tf.clone(), opt, tb, ix, key)?);
let Some(state) = building.read_durable_build_state().await? else {
return Ok(false);
};
if !matches!(state.phase, IndexBuildPhase::Building | IndexBuildPhase::Closing)
|| !build_owner_expired(&state, Utc::now())
{
return Ok(false);
}
match building.takeover_expired_build_state().await? {
Some(acquired) => {
self.start_acquired_building(building, acquired, None).await?;
Ok(true)
}
None => Ok(false),
}
}
}
pub(super) struct Building {
pub(super) ctx: FrozenContext,
pub(super) owner: Uuid,
pub(super) opt: Options,
pub(super) tb: TableId,
pub(super) ikb: IndexKeyBase,
pub(super) tf: TransactionFactory,
pub(super) ix: Arc<IndexDefinition>,
pub(super) ix_key: SharedIndexKey,
pub(super) build_generation: AtomicU64,
pub(super) aborted: AtomicBool,
pub(super) finished: AtomicBool,
}
impl Building {
pub(super) fn new(
ctx: &FrozenContext,
tf: TransactionFactory,
opt: Options,
tb: TableId,
ix: Arc<IndexDefinition>,
ix_key: SharedIndexKey,
) -> Result<Self> {
let ikb = IndexKeyBase::new(ix_key.ns, ix_key.db, ix.table_name.clone(), ix.index_id);
Ok(Self {
ctx: Context::new_concurrent(ctx).freeze(),
owner: Uuid::now_v7(),
opt,
tb,
ikb,
tf,
ix,
ix_key,
build_generation: AtomicU64::new(0),
aborted: AtomicBool::new(false),
finished: AtomicBool::new(false),
})
}
pub(super) async fn acquire_build_state(&self) -> Result<Option<AcquiredBuild>> {
loop {
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let state_key = self.ikb.new_bs_key();
let existing = tx.get(&state_key, None).await?;
if let Some(current) = existing.as_ref()
&& matches!(current.phase, IndexBuildPhase::Building | IndexBuildPhase::Closing)
{
if !build_owner_expired(current, Utc::now()) {
tx.cancel().await?;
return Ok(None);
}
let mut next = current.clone();
next.owner = Some(self.owner);
next.error = None;
next.report_status =
next.report_status.or_else(|| Some(report_status_from_phase(next.phase)));
let now = Utc::now();
next.updated_at = now;
next.owner_heartbeat_at = Some(now);
let res = tx.putc(&state_key, &next, Some(current)).await;
match res {
Ok(()) => {
if self
.commit_and_retryable_conflict(
&tx,
"transient conflict acquiring build ownership, retrying",
)
.await?
{
continue;
}
self.build_generation.store(current.generation, Ordering::Release);
return Ok(Some(AcquiredBuild {
generation: current.generation,
phase: current.phase,
initial_complete: current.initial_complete,
initial_count: durable_report_count(current.initial),
updates_count: durable_report_count(current.updated),
initial_cursor: current.initial_cursor.clone(),
}));
}
Err(err) if is_condition_not_met(&err) => {
let _ = tx.cancel().await;
continue;
}
Err(err) => {
let _ = tx.cancel().await;
return Err(err);
}
}
}
tx.cancel().await?;
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let existing = tx.get(&state_key, None).await?;
if let Some(current) = existing.as_ref()
&& matches!(current.phase, IndexBuildPhase::Building | IndexBuildPhase::Closing)
{
tx.cancel().await?;
continue;
}
let generation = existing.as_ref().map(|s| s.generation.saturating_add(1)).unwrap_or(1);
let now = Utc::now();
let state = IndexBuildState {
generation,
phase: IndexBuildPhase::Building,
owner: Some(self.owner),
next_ticket: 0,
initial_complete: false,
updated_at: now,
owner_heartbeat_at: Some(now),
error: None,
report_status: Some(IndexBuildReportStatus::Started),
initial: None,
updated: None,
pending: None,
initial_cursor: None,
};
let res = tx.putc(&state_key, &state, existing.as_ref()).await;
match res {
Ok(()) => {
if self
.commit_and_retryable_conflict(
&tx,
"transient conflict acquiring build ownership, retrying",
)
.await?
{
continue;
}
}
Err(err) if is_condition_not_met(&err) => {
let _ = tx.cancel().await;
continue;
}
Err(err) => {
let _ = tx.cancel().await;
return Err(err);
}
}
self.build_generation.store(generation, Ordering::Release);
if generation > 1 {
self.wait_for_prior_generation_reservations(generation).await?;
self.wipe_stale_build_queues(generation).await?;
}
return Ok(Some(AcquiredBuild {
generation,
phase: IndexBuildPhase::Building,
initial_complete: false,
initial_count: 0,
updates_count: 0,
initial_cursor: None,
}));
}
}
async fn wipe_stale_build_queues(&self, below: BuildGeneration) -> Result<()> {
loop {
if self.is_aborted().await {
return Ok(());
}
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
if let Err(err) = delete_stale_build_queues(&tx, &self.ikb, below).await {
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict wiping stale build queues, retrying",
)
.await
{
continue;
}
return Err(err);
}
match tx.commit().await {
Ok(()) => return Ok(()),
Err(err) => {
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict wiping stale build queues, retrying",
)
.await
{
continue;
}
return Err(err);
}
}
}
}
async fn read_durable_build_state(&self) -> Result<Option<IndexBuildState>> {
let tx = self.new_read_tx().await?;
let state = catch!(tx, tx.get(&self.ikb.new_bs_key(), None).await);
tx.cancel().await?;
Ok(state)
}
async fn takeover_expired_build_state(&self) -> Result<Option<AcquiredBuild>> {
loop {
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let state_key = self.ikb.new_bs_key();
let Some(current) = tx.get(&state_key, None).await? else {
tx.cancel().await?;
return Ok(None);
};
if !matches!(current.phase, IndexBuildPhase::Building | IndexBuildPhase::Closing)
|| !build_owner_expired(¤t, Utc::now())
{
tx.cancel().await?;
return Ok(None);
}
let mut next = current.clone();
next.owner = Some(self.owner);
next.error = None;
next.report_status =
next.report_status.or_else(|| Some(report_status_from_phase(next.phase)));
let now = Utc::now();
next.updated_at = now;
next.owner_heartbeat_at = Some(now);
let res = tx.putc(&state_key, &next, Some(¤t)).await;
match res {
Ok(()) => {
if self
.commit_and_retryable_conflict(
&tx,
"transient conflict taking over build ownership, retrying",
)
.await?
{
continue;
}
self.build_generation.store(current.generation, Ordering::Release);
return Ok(Some(AcquiredBuild {
generation: current.generation,
phase: current.phase,
initial_complete: current.initial_complete,
initial_count: durable_report_count(current.initial),
updates_count: durable_report_count(current.updated),
initial_cursor: current.initial_cursor.clone(),
}));
}
Err(err) if is_condition_not_met(&err) => {
let _ = tx.cancel().await;
continue;
}
Err(err) => {
let _ = tx.cancel().await;
return Err(err);
}
}
}
}
async fn update_owned_build_state<F>(
&self,
generation: BuildGeneration,
mut update: F,
) -> Result<IndexBuildState>
where
F: FnMut(&mut IndexBuildState),
{
loop {
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let state_key = self.ikb.new_bs_key();
let Some(current) = tx.get(&state_key, None).await? else {
tx.cancel().await?;
return Err(Error::CorruptedIndex(
"Index build state is missing during state update",
)
.into());
};
if current.generation != generation || current.owner != Some(self.owner) {
tx.cancel().await?;
return Err(Error::IndexingBuildingCancelled {
reason: format!("Index build ownership was lost for {}", self.ix.name),
}
.into());
}
let mut next = current.clone();
update(&mut next);
let now = Utc::now();
next.updated_at = now;
next.owner_heartbeat_at = if next.owner == Some(self.owner) {
Some(now)
} else {
None
};
let res = tx.putc(&state_key, &next, Some(¤t)).await;
match res {
Ok(()) => {
if self
.commit_and_retryable_conflict(
&tx,
"transient conflict updating build state, retrying",
)
.await?
{
continue;
}
return Ok(next);
}
Err(err) if is_condition_not_met(&err) => {
let _ = tx.cancel().await;
continue;
}
Err(err) => {
let _ = tx.cancel().await;
return Err(err);
}
}
}
}
fn set_report(
state: &mut IndexBuildState,
status: IndexBuildReportStatus,
initial: Option<usize>,
pending: Option<usize>,
updated: Option<usize>,
) {
state.report_status = Some(status);
state.initial = initial.map(|v| v as u64);
state.pending = pending.map(|v| v as u64);
state.updated = updated.map(|v| v as u64);
if status != IndexBuildReportStatus::Error {
state.error = None;
}
}
pub(super) async fn mark_durable_report(
&self,
generation: BuildGeneration,
status: IndexBuildReportStatus,
initial: Option<usize>,
pending: Option<usize>,
updated: Option<usize>,
) -> Result<()> {
self.update_owned_build_state(generation, |state| {
Self::set_report(state, status, initial, pending, updated);
})
.await?;
Ok(())
}
pub(super) async fn mark_durable_initial_complete(
&self,
generation: BuildGeneration,
) -> Result<()> {
self.update_owned_build_state(generation, |state| {
if state.phase == IndexBuildPhase::Building {
state.initial_complete = true;
state.initial_cursor = None;
state.error = None;
}
})
.await?;
Ok(())
}
async fn update_build_state_in_tx<F>(
&self,
tx: &Transaction,
generation: BuildGeneration,
update: F,
) -> Result<()>
where
F: FnOnce(&mut IndexBuildState),
{
let state_key = self.ikb.new_bs_key();
let Some(current) = tx.get(&state_key, None).await? else {
return Err(Error::CorruptedIndex(
"Index build state is missing during build-state update",
)
.into());
};
if current.generation != generation || current.owner != Some(self.owner) {
return Err(Error::IndexingBuildingCancelled {
reason: format!("Index build ownership was lost for {}", self.ix.name),
}
.into());
}
let mut next = current.clone();
update(&mut next);
tx.putc(&state_key, &next, Some(¤t)).await?;
Ok(())
}
async fn checkpoint_initial_scan(
&self,
tx: &Transaction,
generation: BuildGeneration,
cursor: &RecordIdKey,
initial_count: usize,
) -> Result<()> {
self.update_build_state_in_tx(tx, generation, |state| {
state.initial_cursor = Some(cursor.clone());
state.initial = Some(initial_count as u64);
})
.await
}
async fn complete_initial_scan(
&self,
tx: &Transaction,
generation: BuildGeneration,
initial_count: usize,
) -> Result<()> {
self.update_build_state_in_tx(tx, generation, |state| {
state.initial_complete = true;
state.initial_cursor = None;
state.initial = Some(initial_count as u64);
state.error = None;
})
.await
}
pub(super) async fn mark_durable_closing(&self, generation: BuildGeneration) -> Result<()> {
self.update_owned_build_state(generation, |state| {
if state.phase == IndexBuildPhase::Building {
state.phase = IndexBuildPhase::Closing;
state.error = None;
state.report_status = Some(IndexBuildReportStatus::Indexing);
}
})
.await?;
Ok(())
}
pub(super) async fn mark_durable_online(
&self,
generation: BuildGeneration,
initial: usize,
updated: usize,
) -> Result<()> {
self.update_owned_build_state(generation, |state| {
state.phase = IndexBuildPhase::Online;
state.owner = None;
state.initial_complete = true;
state.error = None;
Self::set_report(
state,
IndexBuildReportStatus::Ready,
Some(initial),
Some(0),
Some(updated),
);
})
.await?;
Ok(())
}
async fn mark_durable_error(&self, generation: BuildGeneration, error: String) -> Result<()> {
loop {
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let state_key = self.ikb.new_bs_key();
let Some(current) = tx.get(&state_key, None).await? else {
tx.cancel().await?;
return Ok(());
};
if current.generation != generation || current.owner != Some(self.owner) {
tx.cancel().await?;
return Ok(());
}
let mut next = current.clone();
next.phase = IndexBuildPhase::Error;
next.owner = None;
next.owner_heartbeat_at = None;
next.error = Some(error.clone());
next.report_status = Some(IndexBuildReportStatus::Error);
next.updated_at = Utc::now();
let res = tx.putc(&state_key, &next, Some(¤t)).await;
match res {
Ok(()) => {
if self
.commit_and_retryable_conflict(
&tx,
"transient conflict marking build error, retrying",
)
.await?
{
continue;
}
return Ok(());
}
Err(err) if is_condition_not_met(&err) => {
let _ = tx.cancel().await;
continue;
}
Err(err) => {
let _ = tx.cancel().await;
return Err(err);
}
}
}
}
async fn mark_durable_aborted(&self, generation: BuildGeneration) -> Result<()> {
self.update_owned_build_state(generation, |state| {
state.phase = IndexBuildPhase::Error;
state.owner = None;
state.report_status = Some(IndexBuildReportStatus::Aborted);
state.error = None;
})
.await?;
Ok(())
}
pub(super) async fn maintain_build_ownership(
&self,
tx: &Transaction,
generation: BuildGeneration,
allowed: &[IndexBuildPhase],
) -> Result<()> {
let state_key = self.ikb.new_bs_key();
let Some(current) = tx.get(&state_key, None).await? else {
return Err(Error::CorruptedIndex(
"Index build state is missing during ownership heartbeat",
)
.into());
};
if current.generation != generation
|| current.owner != Some(self.owner)
|| !allowed.contains(¤t.phase)
{
return Err(Error::IndexingBuildingCancelled {
reason: format!("Index build ownership was lost for {}", self.ix.name),
}
.into());
}
let mut next = current.clone();
let now = Utc::now();
next.updated_at = now;
next.owner_heartbeat_at = Some(now);
tx.putc(&state_key, &next, Some(¤t)).await?;
Ok(())
}
async fn retryable_conflict(&self, err: &anyhow::Error, action: &str) -> bool {
if is_retryable_transaction_conflict(err) {
debug!(
target: "surrealdb::core::kvs::index",
index = %self.ix.name,
table = %self.ix.table_name,
action,
error = %err,
"retryable conflict during concurrent index build, retrying"
);
sleep(Duration::from_millis(100)).await;
true
} else {
false
}
}
pub(super) async fn cancel_and_retryable_conflict(
&self,
tx: &Transaction,
err: &anyhow::Error,
action: &str,
) -> bool {
let _ = tx.cancel().await;
self.retryable_conflict(err, action).await
}
async fn commit_and_retryable_conflict(&self, tx: &Transaction, action: &str) -> Result<bool> {
match tx.commit().await {
Ok(()) => Ok(false),
Err(err) => {
if self.cancel_and_retryable_conflict(tx, &err, action).await {
Ok(true)
} else {
Err(err)
}
}
}
}
pub(super) async fn new_read_tx(&self) -> Result<Transaction> {
self.tf
.transaction(TransactionType::Read, Optimistic, self.ctx.try_get_sequences()?.clone())
.await
}
pub(super) async fn new_write_tx_ctx(&self) -> Result<FrozenContext> {
let tx = self
.tf
.transaction(TransactionType::Write, Optimistic, self.ctx.try_get_sequences()?.clone())
.await?
.into();
let mut ctx = Context::new_child(&self.ctx);
ctx.set_transaction(tx);
Ok(ctx.freeze())
}
pub(super) async fn new_read_tx_ctx(&self) -> Result<FrozenContext> {
let tx = self
.tf
.transaction(TransactionType::Read, Optimistic, self.ctx.try_get_sequences()?.clone())
.await?
.into();
let mut ctx = Context::new_child(&self.ctx);
ctx.set_transaction(tx);
Ok(ctx.freeze())
}
async fn evict_cached_hnsw_index(&self) {
if let Err(err) =
self.ctx.get_index_stores().remove_hnsw_index(self.tb, self.ikb.clone()).await
{
warn!("Failed to evict HNSW index after index-builder compaction error: {err}");
}
}
pub(super) async fn check_prepare_remove_with_tx(
&self,
last_prepare_remove_check: &mut Instant,
tx: &Transaction,
) -> Result<()> {
if last_prepare_remove_check.elapsed() < Duration::from_secs(5) {
return Ok(());
};
if let Some(ix) = tx
.get_tb_index(self.ix_key.ns, self.ix_key.db, &self.ix.table_name, &self.ix.name, None)
.await?
{
ix.expect_not_prepare_remove()?;
}
*last_prepare_remove_check = Instant::now();
Ok(())
}
pub(super) async fn check_prepare_remove(
&self,
last_prepare_remove_check: &mut Instant,
) -> Result<()> {
let tx = self.new_read_tx().await?;
catch!(tx, self.check_prepare_remove_with_tx(last_prepare_remove_check, &tx).await);
tx.cancel().await?;
Ok(())
}
pub(super) async fn compaction_write_still_owns_index(
&self,
tx: &Transaction,
generation: BuildGeneration,
) -> Result<bool> {
if generation == 0 {
return Ok(false);
}
let Some(state) = tx.get(&self.ikb.new_bs_key(), None).await? else {
return Ok(false);
};
if state.generation != generation || state.phase != IndexBuildPhase::Online {
return Ok(false);
}
if let Some(ix) = tx
.get_tb_index_by_id(
self.ix_key.ns,
self.ix_key.db,
&self.ix_key.tb,
self.ix_key.ix,
None,
)
.await?
{
return Ok(ix.index_id == self.ix.index_id
&& ix.name == self.ix.name
&& !ix.prepare_remove);
}
Ok(true)
}
#[cfg(test)]
#[cfg_attr(not(feature = "kv-mem"), allow(dead_code))]
pub(super) async fn run(&self) -> Result<()> {
let Some(acquired) = self.acquire_build_state().await? else {
return Ok(());
};
let generation = acquired.generation;
let res = self.run_acquired(acquired).await;
if res.is_ok() && self.aborted.load(Ordering::Acquire) {
let _ = self.mark_durable_aborted(generation).await;
}
res
}
pub(super) async fn run_acquired(&self, acquired: AcquiredBuild) -> Result<()> {
let mut last_prepare_remove_check = Instant::now();
let generation = acquired.generation;
let scanning_initial =
acquired.phase == IndexBuildPhase::Building && !acquired.initial_complete;
let resume_cursor = if scanning_initial {
acquired.initial_cursor.clone()
} else {
None
};
let restarting_initial_scan = scanning_initial && resume_cursor.is_none();
let mut initial_count = if restarting_initial_scan {
0
} else {
acquired.initial_count
};
let mut updates_count = if restarting_initial_scan {
0
} else {
acquired.updates_count
};
if restarting_initial_scan {
self.mark_durable_report(
generation,
IndexBuildReportStatus::Cleaning,
None,
None,
None,
)
.await?;
self.wait_for_prior_generation_reservations(generation).await?;
self.wipe_stale_build_queues(generation).await?;
loop {
if self.is_aborted().await {
return Ok(());
}
let ctx = self.new_write_tx_ctx().await?;
let key =
index_all::new(self.ix_key.ns, self.ix_key.db, &self.ix_key.tb, self.ix_key.ix);
let tx = ctx.tx();
if let Err(err) = self
.maintain_build_ownership(&tx, generation, &[IndexBuildPhase::Building])
.await
{
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict maintaining build ownership, retrying",
)
.await
{
continue;
}
return Err(err);
}
if let Err(err) = tx.delp(&key).await {
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict while cleaning existing index data, retrying",
)
.await
{
continue;
}
return Err(err);
}
#[cfg(test)]
if let Err(err) = maybe_inject_retryable_conflict(
RetryableConflictSite::ConcurrentIndexInitialCleanup,
self.ctx.node_id(),
) {
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict while cleaning existing index data, retrying",
)
.await
{
continue;
}
return Err(err);
}
match tx.commit().await {
Ok(()) => break,
Err(err) => {
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict while cleaning existing index data, retrying",
)
.await
{
continue;
}
return Err(err);
}
}
}
}
if scanning_initial {
let beg = if let Some(cursor) = &resume_cursor {
let mut key = record::new(self.ix_key.ns, self.ix_key.db, self.ikb.table(), cursor)
.encode_key()?;
advance_key(&mut key);
key
} else {
record::prefix(self.ix_key.ns, self.ix_key.db, self.ikb.table())?
};
let end = record::suffix(self.ix_key.ns, self.ix_key.db, self.ikb.table())?;
let mut next = Some(beg..end);
let mut v1_appending_sentinel = false;
let mut count_primary_cursor =
matches!(self.ix.index, Index::Count(_)).then(|| resume_cursor.clone());
let mut scan_batch_size = INDEXING_PROBE_BATCH_SIZE;
self.mark_durable_report(
generation,
IndexBuildReportStatus::Indexing,
Some(initial_count),
Some(0),
None,
)
.await?;
while let Some(rng) = next {
if self.is_aborted().await {
return Ok(());
}
self.is_beyond_threshold(None)?;
let batch = {
let tx = self.new_read_tx().await?;
catch!(
tx,
self.check_prepare_remove_with_tx(&mut last_prepare_remove_check, &tx)
.await
);
let res = catch!(tx, tx.batch_keys_vals(rng, scan_batch_size, None).await);
tx.cancel().await?;
res
};
next = batch.next;
if batch.result.is_empty() {
break;
}
{
let bytes: usize = batch.result.iter().map(|(k, v)| k.len() + v.len()).sum();
let avg = (bytes / batch.result.len()).max(1);
scan_batch_size = (INDEXING_BATCH_MAX_BYTES / avg)
.clamp(1, INDEXING_BATCH_SIZE as usize) as u32;
}
{
let values = batch.result;
let Some((last_key, _)) = values.last() else {
break;
};
let batch_cursor = record::RecordKey::decode_key(last_key)?.id;
let indexed = loop {
if self.is_aborted().await {
return Ok(());
}
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let saved_count_primary_cursor = count_primary_cursor.clone();
if let Err(err) = self
.maintain_build_ownership(&tx, generation, &[IndexBuildPhase::Building])
.await
{
count_primary_cursor = saved_count_primary_cursor;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict maintaining build ownership, retrying",
)
.await
{
continue;
}
return Err(err);
}
let indexed = match self
.index_initial_batch(
&ctx,
&tx,
&values,
initial_count,
&mut v1_appending_sentinel,
&mut count_primary_cursor,
)
.await
{
Ok(indexed) => indexed,
Err(err) => {
count_primary_cursor = saved_count_primary_cursor;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict in initial index batch, retrying",
)
.await
{
continue;
}
return Err(err);
}
};
if self.is_aborted().await {
tx.cancel().await?;
return Ok(());
}
if let Err(err) = self
.checkpoint_initial_scan(
&tx,
generation,
&batch_cursor,
initial_count + indexed,
)
.await
{
count_primary_cursor = saved_count_primary_cursor;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict checkpointing the initial scan, retrying",
)
.await
{
continue;
}
return Err(err);
}
#[cfg(test)]
if let Err(err) = maybe_inject_retryable_conflict(
RetryableConflictSite::ConcurrentIndexInitialBatch,
self.ctx.node_id(),
) {
count_primary_cursor = saved_count_primary_cursor;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict on initial index batch commit, retrying",
)
.await
{
continue;
}
return Err(err);
}
#[cfg(test)]
if let Err(err) =
maybe_inject_initial_batch_commit_error(self.ctx.node_id())
{
count_primary_cursor = saved_count_primary_cursor;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict on initial index batch commit, retrying",
)
.await
{
continue;
}
return Err(err);
}
match tx.commit().await {
Ok(()) => break indexed,
Err(err) => {
count_primary_cursor = saved_count_primary_cursor;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict on initial index batch commit, retrying",
)
.await
{
continue;
}
return Err(err);
}
}
};
initial_count += indexed;
if !self.is_aborted().await {
self.mark_durable_report(
generation,
IndexBuildReportStatus::Indexing,
Some(initial_count),
Some(0),
None,
)
.await?;
}
}
}
if count_primary_cursor.is_some() {
let indexed = loop {
if self.is_aborted().await {
return Ok(());
}
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let saved_count_primary_cursor = count_primary_cursor.clone();
if let Err(err) = self
.maintain_build_ownership(&tx, generation, &[IndexBuildPhase::Building])
.await
{
count_primary_cursor = saved_count_primary_cursor;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict maintaining build ownership, retrying",
)
.await
{
continue;
}
return Err(err);
}
let indexed = match self
.index_remaining_count_primary_appendings(
&ctx,
&tx,
&mut count_primary_cursor,
initial_count,
)
.await
{
Ok(indexed) => indexed,
Err(err) => {
count_primary_cursor = saved_count_primary_cursor;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict in initial count appending range, retrying",
)
.await
{
continue;
}
return Err(err);
}
};
if self.is_aborted().await {
tx.cancel().await?;
return Ok(());
}
if let Err(err) =
self.complete_initial_scan(&tx, generation, initial_count + indexed).await
{
count_primary_cursor = saved_count_primary_cursor;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict completing the initial scan, retrying",
)
.await
{
continue;
}
return Err(err);
}
match tx.commit().await {
Ok(()) => {
#[cfg(test)]
maybe_inject_non_retryable_error(
NonRetryableErrorSite::ConcurrentIndexCountTailCommitted,
self.ctx.node_id(),
)?;
break indexed;
}
Err(err) => {
count_primary_cursor = saved_count_primary_cursor;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict on initial count appending commit, retrying",
)
.await
{
continue;
}
return Err(err);
}
}
};
initial_count += indexed;
} else {
self.mark_durable_initial_complete(generation).await?;
}
}
self.mark_durable_report(
generation,
IndexBuildReportStatus::Indexing,
Some(initial_count),
Some(0),
Some(updates_count),
)
.await?;
self.index_appending_loop(
initial_count,
&mut updates_count,
&mut last_prepare_remove_check,
)
.await?;
if acquired.phase == IndexBuildPhase::Building {
self.mark_durable_closing(generation).await?;
}
self.index_appending_loop(
initial_count,
&mut updates_count,
&mut last_prepare_remove_check,
)
.await?;
self.wait_for_durable_reservations(generation, &mut last_prepare_remove_check).await?;
self.index_appending_loop(
initial_count,
&mut updates_count,
&mut last_prepare_remove_check,
)
.await?;
self.mark_durable_online(generation, initial_count, updates_count).await?;
self.compact_hnsw_pendings(&mut last_prepare_remove_check).await?;
#[cfg(diskann)]
self.compact_diskann_pendings(&mut last_prepare_remove_check).await?;
Ok(())
}
async fn compact_hnsw_pendings(&self, last_prepare_remove_check: &mut Instant) -> Result<()> {
let Index::Hnsw(p) = &self.ix.index else {
return Ok(());
};
loop {
if self.is_aborted().await {
return Ok(());
}
self.is_beyond_threshold(None)?;
self.check_prepare_remove(last_prepare_remove_check).await?;
let plan = {
let ctx = self.new_read_tx_ctx().await?;
let tx = ctx.tx();
let res = IndexOperation::prepare_hnsw_compaction(&ctx, &self.ikb).await;
let cancel = tx.cancel().await;
match res {
Ok(plan) => {
cancel?;
plan
}
Err(err) => {
let _ = cancel;
return Err(err);
}
}
};
if !plan.has_work() {
return Ok(());
}
let has_more = plan.has_more();
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let generation = self.build_generation.load(Ordering::Acquire);
if !self.compaction_write_still_owns_index(&tx, generation).await? {
tx.cancel().await?;
return Ok(());
}
let res = IndexOperation::apply_hnsw_compaction(
&ctx,
ctx.get_index_stores(),
&self.ikb,
&self.ix,
p,
plan,
)
.await;
match res {
Ok(true) => {
if let Err(err) = tx.commit().await {
self.evict_cached_hnsw_index().await;
return Err(err);
}
}
Ok(false) => {
tx.cancel().await?;
return Ok(());
}
Err(err) => {
let _ = tx.cancel().await;
self.evict_cached_hnsw_index().await;
return Err(err);
}
}
if !has_more {
return Ok(());
}
}
}
#[cfg(diskann)]
async fn compact_diskann_pendings(
&self,
last_prepare_remove_check: &mut Instant,
) -> Result<()> {
let Index::DiskAnn(p) = &self.ix.index else {
return Ok(());
};
loop {
if self.is_aborted().await {
return Ok(());
}
self.is_beyond_threshold(None)?;
self.check_prepare_remove(last_prepare_remove_check).await?;
let plan = {
let ctx = self.new_read_tx_ctx().await?;
let tx = ctx.tx();
let res = IndexOperation::prepare_diskann_compaction(&ctx, &self.ikb).await;
let cancel = tx.cancel().await;
match res {
Ok(plan) => {
cancel?;
plan
}
Err(err) => {
let _ = cancel;
return Err(err);
}
}
};
if !plan.requires_apply() {
return Ok(());
}
let has_more = plan.has_more();
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let generation = self.build_generation.load(Ordering::Acquire);
if !self.compaction_write_still_owns_index(&tx, generation).await? {
tx.cancel().await?;
return Ok(());
}
let res = IndexOperation::apply_diskann_compaction(
&ctx,
ctx.get_index_stores(),
&self.ikb,
&self.ix,
p,
plan,
)
.await;
if !tx.closed() {
let _ = tx.cancel().await;
}
match res {
Ok(true) => {}
Ok(false) => return Ok(()),
Err(err) => return Err(err),
}
if !has_more {
return Ok(());
}
}
}
pub(super) fn abort(&self) {
self.aborted.store(true, Ordering::Relaxed);
}
pub(super) async fn is_aborted(&self) -> bool {
self.aborted.load(Ordering::Relaxed)
}
pub(super) fn is_beyond_threshold(&self, count: Option<usize>) -> Result<()> {
if let Some(count) = count
&& count % 100 != 0
{
return Ok(());
}
if ALLOC.is_beyond_threshold() {
Err(anyhow::Error::new(Error::QueryBeyondMemoryThreshold))
} else {
Ok(())
}
}
pub(super) fn is_finished(&self) -> bool {
self.finished.load(Ordering::Relaxed)
}
}
struct BuildingFinishGuard(IndexBuilding);
impl Drop for BuildingFinishGuard {
fn drop(&mut self) {
self.0.finished.store(true, Ordering::Relaxed);
}
}