use std::borrow::Cow;
use std::collections::HashMap;
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU8, AtomicU64, Ordering};
use std::time::Duration;
use anyhow::{Result, ensure};
use chrono::Utc;
use common::time::sleep;
use futures::channel::oneshot::{Receiver, Sender, channel};
use surrealdb_datastore::close::{CommitAction, RollbackAction};
#[cfg(not(target_family = "wasm"))]
use tokio::spawn;
use tokio::sync::{Notify, RwLock};
use uuid::Uuid;
#[cfg(target_family = "wasm")]
use wasm_bindgen_futures::spawn_local as spawn;
use web_time::Instant;
use super::replay::CountPrimaryProgress;
use super::state::{
build_owner_expired, delete_durable_build_artifacts, 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, build_abort_deadline,
};
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::idx::trees::hnsw::index::HnswCompactionOutcome;
use crate::key::schema::{IdxRoot, RecordKey, RecordPrefix};
use crate::key::{KVKey, KVKeyDecode, Resumable};
use crate::kvs::sequences::Sequences;
#[cfg(test)]
use crate::kvs::testing::{
NonRetryableErrorSite, RetryableConflictSite, maybe_inject_non_retryable_error,
maybe_inject_retryable_conflict,
};
use crate::kvs::{
DatastoreError, Direction, INDEXING_BATCH_MAX_BYTES, INDEXING_BATCH_SIZE,
INDEXING_PROBE_BATCH_SIZE, Transaction, TransactionFactory, TransactionType,
is_retryable_transaction_conflict, is_shutdown_error,
};
const INDEX_BUILD_CLEANUP_RETRY_SLEEP: Duration = Duration::from_millis(100);
use crate::mem::ALLOC;
use crate::val::{RecordId, RecordIdKey, TableName, Value};
pub(super) type SharedIndexKey = Arc<IndexKey>;
const RELEASE_GATE_BUILDER: u8 = 1 << 0;
const RELEASE_GATE_STATEMENT: u8 = 1 << 1;
const RELEASE_GATE_OPEN: u8 = RELEASE_GATE_BUILDER | RELEASE_GATE_STATEMENT;
fn is_memory_threshold_error(err: &anyhow::Error) -> bool {
matches!(err.downcast_ref::<DatastoreError>(), Some(DatastoreError::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),
};
building.defer_ownership_release(ctx).await;
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(),
DatastoreError::IndexAlreadyBuilding {
name: building.ix.name.to_string(),
}
);
}
indexes.insert(Arc::clone(&building.ix_key), Arc::clone(&building));
}
let b = Arc::clone(&building);
let guard = BuildingFinishGuard(Arc::clone(&building));
spawn(async move {
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(DatastoreError::IndexingBuildingCancelled {
reason: format!("Index {} build state no longer exists", building.ix.name),
}
.into());
};
match state.phase {
IndexBuildPhase::Online => return Ok(()),
IndexBuildPhase::Error => {
return Err(DatastoreError::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();
building.defer_ownership_release(ctx).await;
self.start_acquired_building(Arc::clone(&building), acquired, Some(s))
.await?;
return r.await.map_err(|_| DatastoreError::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<()>>>> {
expect_not_prepare_remove(&ix)?;
let (ns, db) = ctx.expect_ns_db_ids(&opt).await?;
let key = Arc::new(IndexKey::new(ns, db, &ix.table_name.clone(), 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(),
DatastoreError::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.clone(), 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,
release_gate: AtomicU8,
finished_notify: Notify,
}
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),
release_gate: AtomicU8::new(RELEASE_GATE_STATEMENT),
finished_notify: Notify::new(),
})
}
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_key(&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.put_compare_key(&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_key(&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,
};
if let Some(previous) = existing.as_ref().map(|s| s.generation) {
let previous_bt = self.ikb.new_bt_key(previous);
if let Some(current) = tx.get_key(&previous_bt, None).await? {
tx.del_compare_key(&previous_bt, Some(¤t)).await?;
}
}
tx.set_key(&self.ikb.new_bt_key(generation), &0).await?;
let res = tx.put_compare_key(&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_key(&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_key(&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.put_compare_key(&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,
update: F,
) -> Result<IndexBuildState>
where
F: FnMut(&mut IndexBuildState),
{
self.update_owned_build_state_inner(generation, update, false).await
}
async fn update_owned_build_state_inner<F>(
&self,
generation: BuildGeneration,
mut update: F,
fence_ticket_allocation: bool,
) -> 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_key(&state_key, None).await? else {
tx.cancel().await?;
return Err(DatastoreError::CorruptedIndex(
"Index build state is missing during state update",
)
.into());
};
if current.generation != generation || current.owner != Some(self.owner) {
tx.cancel().await?;
return Err(DatastoreError::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
};
if fence_ticket_allocation {
let bt = self.ikb.new_bt_key(generation);
if let Some(ticket) = tx.get_key(&bt, None).await?
&& let Err(err) =
tx.put_compare_key(&bt, &ticket.saturating_add(1), Some(&ticket)).await
{
let _ = tx.cancel().await;
return Err(err);
}
}
let res = tx.put_compare_key(&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_key(&state_key, None).await? else {
return Err(DatastoreError::CorruptedIndex(
"Index build state is missing during build-state update",
)
.into());
};
if current.generation != generation || current.owner != Some(self.owner) {
return Err(DatastoreError::IndexingBuildingCancelled {
reason: format!("Index build ownership was lost for {}", self.ix.name),
}
.into());
}
let mut next = current.clone();
update(&mut next);
tx.put_compare_key(&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_inner(
generation,
|state| {
if state.phase == IndexBuildPhase::Building {
state.phase = IndexBuildPhase::Closing;
state.error = None;
state.report_status = Some(IndexBuildReportStatus::Indexing);
}
},
true,
)
.await?;
Ok(())
}
async fn ensure_single_ticket_allocator(&self, generation: BuildGeneration) -> Result<()> {
let tx = self.new_read_tx().await?;
let counter = catch!(tx, tx.get_key(&self.ikb.new_bt_key(generation), None).await);
let state = catch!(tx, tx.get_key(&self.ikb.new_bs_key(), None).await);
tx.cancel().await?;
let Some(state) = state.filter(|state| state.generation == generation) else {
return Ok(());
};
if counter.is_some() && state.next_ticket != 0 {
return Err(DatastoreError::IndexingBuildingCancelled {
reason: format!(
"Index {} was built while nodes of different versions allocated writer \
tickets for build generation {generation}, so queued writes may have been \
overwritten. Run `REBUILD INDEX {} ON {}` to rebuild it",
self.ix.name, self.ix.name, self.ix.table_name
),
}
.into());
}
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.initial_complete = true;
state.error = None;
Self::set_report(
state,
IndexBuildReportStatus::Ready,
Some(initial),
Some(0),
Some(updated),
);
})
.await?;
Ok(())
}
pub(super) async fn defer_ownership_release(self: &Arc<Self>, ctx: &FrozenContext) {
let Some(txn) = ctx.try_tx() else {
return;
};
self.release_gate.store(0, Ordering::Release);
txn.on_commit(Box::new(ReleaseBuildOwnership(Arc::clone(self)))).await;
txn.on_rollback(Box::new(ReleaseBuildOwnership(Arc::clone(self)))).await;
}
fn arrive_at_release_gate(&self, arrival: u8) -> bool {
(self.release_gate.fetch_or(arrival, Ordering::AcqRel) | arrival) == RELEASE_GATE_OPEN
}
pub(super) async fn release_build_ownership(&self, generation: BuildGeneration) {
match self.read_durable_build_state().await {
Ok(Some(state))
if state.generation == generation
&& state.phase == IndexBuildPhase::Online
&& state.owner == Some(self.owner) => {}
Ok(_) => return,
Err(err) => {
warn!(
index = %self.ix.name,
table = %self.ix.table_name,
error = %err,
"could not read durable build state to release build ownership"
);
return;
}
}
match self
.update_owned_build_state(generation, |state| {
if state.phase == IndexBuildPhase::Online {
state.owner = None;
}
})
.await
{
Ok(_) => self.requeue_compaction_after_release().await,
Err(err) => warn!(
index = %self.ix.name,
table = %self.ix.table_name,
error = %err,
"could not release durable build ownership; background compaction \
will skip this index until its next build"
),
}
}
async fn requeue_compaction_after_release(&self) {
let fenced = match &self.ix.index {
Index::Hnsw(_) => true,
#[cfg(diskann)]
Index::DiskAnn(_) => true,
_ => false,
};
if !fenced {
return;
}
if let Err(err) = self.enqueue_compaction_request().await {
warn!(
index = %self.ix.name,
table = %self.ix.table_name,
error = %err,
"could not queue index compaction after releasing build ownership; \
pendings queued during the build wait for the next write to the index"
);
}
}
async fn enqueue_compaction_request(&self) -> Result<()> {
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
tx.set_key(&self.ikb.new_ic_key(self.ctx.node_id()), &()).await?;
tx.trigger_index_compaction();
tx.commit().await
}
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_key(&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.put_compare_key(&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_key(&state_key, None).await? else {
return Err(DatastoreError::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(DatastoreError::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.put_compare_key(&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
}
pub(super) 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, 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, 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, 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(), self.ix.format_version)
.await
{
warn!("Failed to evict HNSW index after index-builder compaction apply error: {err}");
}
}
async fn discard_failed_hnsw_apply(&self, outcome: HnswCompactionOutcome) {
if outcome == HnswCompactionOutcome::InPlace {
self.evict_cached_hnsw_index().await;
}
}
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.clone(),
&self.ix.name,
None,
)
.await?
{
expect_not_prepare_remove(&ix)?;
}
*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.getu_key(&self.ikb.new_bs_key()).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 mut resume_cursor = if scanning_initial {
acquired.initial_cursor.clone()
} else {
None
};
if matches!(self.ix.index, Index::Count(_))
&& let Some(cursor) = &resume_cursor
&& self.resume_strands_an_older_marker(generation, cursor).await?
{
resume_cursor = 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 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) = crate::idx::wipe_index_data(&tx, &self.ikb, &self.ix.index).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 {
self.rebuild_primary_appendings().await?;
let mut range = RecordPrefix {
ns: self.ix_key.ns,
db: self.ix_key.db,
tb: Cow::Borrowed(self.ikb.table()),
}
.range()?;
if let Some(cursor) = &resume_cursor {
let checkpoint = RecordKey {
ns: self.ix_key.ns,
db: self.ix_key.db,
tb: Cow::Borrowed(self.ikb.table()),
id: Cow::Borrowed(cursor),
}
.encode_key()?;
range = range.resume_after(&checkpoint, Direction::Forward);
}
let mut next = Some(range);
let mut v1_appending_sentinel = false;
let mut count_progress = matches!(self.ix.index, Index::Count(_))
.then(|| CountPrimaryProgress::resuming_at(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_raw(rng.clone(), scan_batch_size, None).await
);
tx.cancel().await?;
res
};
next = batch
.next
.and(batch.result.last())
.map(|(last, _)| rng.resume_after(last, Direction::Forward));
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 = RecordKey::decode_key(last_key)?.id.into_owned();
let indexed = loop {
if self.is_aborted().await {
return Ok(());
}
let ctx = self.new_write_tx_ctx().await?;
let tx = ctx.tx();
let mut attempt = count_progress.clone();
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);
}
let indexed = match self
.index_initial_batch(
&ctx,
&tx,
&values,
initial_count,
&mut v1_appending_sentinel,
&mut attempt,
)
.await
{
Ok(indexed) => indexed,
Err(err) => {
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
{
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(),
) {
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())
{
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(()) => {
count_progress = attempt;
break indexed;
}
Err(err) => {
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_progress.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 mut attempt = count_progress.clone();
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);
}
let indexed = match self
.index_remaining_count_primary_appendings(
&ctx,
&tx,
&mut attempt,
initial_count,
)
.await
{
Ok(indexed) => indexed,
Err(err) => {
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
{
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) => {
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.ensure_single_ticket_allocator(generation).await?;
self.mark_durable_online(generation, initial_count, updates_count).await?;
if let Err(err) = self.reclaim_deferred_doc_ids().await {
warn!(
index = %self.ix.name,
table = %self.ix.table_name,
error = %err,
"deferred doc-ID reclaim sweep failed; leftover markers will be \
reclaimed by the next doc-ID index build on the table"
);
}
let drained = self.compact_hnsw_pendings(&mut last_prepare_remove_check).await;
#[cfg(diskann)]
let drained = match drained {
Ok(()) => self.compact_diskann_pendings(&mut last_prepare_remove_check).await,
Err(err) => Err(err),
};
if self.arrive_at_release_gate(RELEASE_GATE_BUILDER) {
self.release_build_ownership(generation).await;
}
drained
}
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 started = Instant::now();
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 prepared = started.elapsed();
let started = Instant::now();
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,
p,
self.ix.format_version,
plan,
)
.await;
match res {
Ok(outcome) if outcome.applied() => {
let applied = started.elapsed();
let started = Instant::now();
#[cfg(test)]
if let Err(err) = maybe_inject_retryable_conflict(
RetryableConflictSite::ConcurrentIndexHnswPendingCompaction,
self.ctx.node_id(),
) {
self.discard_failed_hnsw_apply(outcome).await;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict compacting HNSW pendings, retrying",
)
.await
{
continue;
}
return Err(err);
}
if let Err(err) = tx.commit().await {
self.discard_failed_hnsw_apply(outcome).await;
if self
.cancel_and_retryable_conflict(
&tx,
&err,
"transient conflict compacting HNSW pendings, retrying",
)
.await
{
continue;
}
return Err(err);
}
debug!(
target: "surrealdb::core::kvs::index",
index = %self.ix.name,
?outcome,
prepare_ms = prepared.as_millis() as u64,
apply_ms = applied.as_millis() as u64,
commit_ms = started.elapsed().as_millis() as u64,
"committed HNSW compaction batch"
);
}
Ok(_) => {
tx.cancel().await?;
return Ok(());
}
Err(err) => {
let _ = tx.cancel().await;
self.evict_cached_hnsw_index().await;
if self
.retryable_conflict(
&err,
"transient conflict compacting HNSW pendings, retrying",
)
.await
{
continue;
}
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(());
}
#[cfg(test)]
if let Err(err) = maybe_inject_retryable_conflict(
RetryableConflictSite::ConcurrentIndexDiskAnnPendingCompaction,
self.ctx.node_id(),
) {
let _ = tx.cancel().await;
if self
.retryable_conflict(
&err,
"transient conflict compacting DiskANN pendings, retrying",
)
.await
{
continue;
}
return Err(err);
}
let res = IndexOperation::apply_diskann_compaction(
&ctx,
ctx.get_index_stores(),
&self.ikb,
p,
self.ix.format_version,
plan,
)
.await;
if !tx.closed() {
let _ = tx.cancel().await;
}
match res {
Ok(true) => {}
Ok(false) => return Ok(()),
Err(err) => {
if self
.retryable_conflict(
&err,
"transient conflict compacting DiskANN pendings, retrying",
)
.await
{
continue;
}
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(DatastoreError::QueryBeyondMemoryThreshold))
} else {
Ok(())
}
}
pub(super) fn is_finished(&self) -> bool {
self.finished.load(Ordering::Acquire)
}
pub(super) async fn wait_finished(&self) {
loop {
let notified = self.finished_notify.notified();
tokio::pin!(notified);
notified.as_mut().enable();
if self.is_finished() {
return;
}
notified.await;
}
}
}
struct ReleaseBuildOwnership(IndexBuilding);
impl ReleaseBuildOwnership {
async fn release(self) {
if self.0.arrive_at_release_gate(RELEASE_GATE_STATEMENT) {
let generation = self.0.build_generation.load(Ordering::Acquire);
self.0.release_build_ownership(generation).await;
}
}
}
impl CommitAction for ReleaseBuildOwnership {
fn run(self: Box<Self>) -> Pin<Box<dyn Future<Output = ()> + Send>> {
Box::pin((*self).release())
}
}
impl RollbackAction for ReleaseBuildOwnership {
fn run(
self: Box<Self>,
_drain_started_at: Instant,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
Box::pin(async move {
(*self).release().await;
Ok(())
})
}
}
struct BuildingFinishGuard(IndexBuilding);
impl Drop for BuildingFinishGuard {
fn drop(&mut self) {
self.0.finished.store(true, Ordering::Release);
self.0.finished_notify.notify_waiters();
}
}
fn expect_not_prepare_remove(ix: &IndexDefinition) -> anyhow::Result<()> {
if ix.prepare_remove {
Err(anyhow::Error::new(crate::kvs::DatastoreError::IndexingBuildingCancelled {
reason: "Prepare remove.".to_string(),
}))
} else {
Ok(())
}
}
pub(crate) struct AbortLocalBuild {
builder: IndexBuilder,
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
ix: IndexId,
}
impl AbortLocalBuild {
pub(crate) fn boxed(
builder: IndexBuilder,
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
ix: IndexId,
) -> Box<dyn CommitAction> {
Box::new(Self {
builder,
ns,
db,
tb,
ix,
})
}
}
impl CommitAction for AbortLocalBuild {
fn run(self: Box<Self>) -> Pin<Box<dyn Future<Output = ()> + Send>> {
Box::pin(async move {
self.builder.remove_index(self.ns, self.db, &self.tb, self.ix).await;
})
}
}
struct IndexBuildCleanup {
builder: IndexBuilder,
tf: TransactionFactory,
sequences: Sequences,
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
ix: IndexId,
}
impl IndexBuildCleanup {
fn new(
builder: IndexBuilder,
tf: TransactionFactory,
sequences: Sequences,
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
ix: IndexId,
) -> Self {
Self {
builder,
tf,
sequences,
ns,
db,
tb,
ix,
}
}
async fn stop_builder(&self, abort_deadline: Instant) {
self.builder
.remove_index_and_wait(self.ns, self.db, &self.tb, self.ix, abort_deadline)
.await;
}
async fn delete_once(&self) -> Result<()> {
let tx = self.tf.transaction(TransactionType::Write, self.sequences.clone()).await?;
let ikb = IndexKeyBase::new(self.ns, self.db, self.tb.clone(), self.ix);
let result: Result<()> = async {
delete_durable_build_artifacts(&tx, &ikb, false).await?;
let index_prefix = IdxRoot {
ns: self.ns,
db: self.db,
tb: Cow::Borrowed(&self.tb),
ix: self.ix,
}
.range()?;
tx.delr(index_prefix).await?;
tx.commit_bare().await?;
Ok(())
}
.await;
if let Err(err) = result {
let _ = tx.cancel_bare().await;
return Err(err);
}
Ok(())
}
}
pub(crate) struct CleanUncommittedBuild {
cleanup: IndexBuildCleanup,
}
impl CleanUncommittedBuild {
pub(crate) fn boxed(
builder: IndexBuilder,
tf: TransactionFactory,
sequences: Sequences,
ns: NamespaceId,
db: DatabaseId,
tb: TableName,
ix: IndexId,
) -> Box<dyn RollbackAction> {
Box::new(Self {
cleanup: IndexBuildCleanup::new(builder, tf, sequences, ns, db, tb, ix),
})
}
}
impl RollbackAction for CleanUncommittedBuild {
fn run(
self: Box<Self>,
drain_started_at: Instant,
) -> Pin<Box<dyn Future<Output = Result<()>> + Send>> {
let abort_deadline = build_abort_deadline(drain_started_at);
Box::pin(async move {
self.cleanup.stop_builder(abort_deadline).await;
loop {
match self.cleanup.delete_once().await {
Ok(()) => return Ok(()),
Err(err) if is_retryable_transaction_conflict(&err) => {
debug!(
error = %err,
"retryable conflict while cleaning uncommitted index build, retrying"
);
sleep(INDEX_BUILD_CLEANUP_RETRY_SLEEP).await;
}
Err(err) => return Err(err),
}
}
})
}
}