use anyhow::Result;
use chrono::{DateTime, Utc};
pub(super) use surrealdb_datastore::index_state::{
CatalogIndexState, catalog_index_state, durable_index_error_reason, report_status_from_phase,
};
pub(crate) use surrealdb_datastore::index_state::{filter_online_indexes, index_building_info};
use super::{BUILD_OWNER_LEASE_SECS, IndexBuildPhase, IndexBuildState};
use crate::catalog::{DatabaseId, IndexId, NamespaceId};
use crate::idx::IndexKeyBase;
use crate::kvs::{Error as KvsError, Transaction, storage_error};
use crate::val::TableName;
pub(crate) async fn retire_durable_index(
tx: &Transaction,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
ix: IndexId,
expunge: bool,
) -> Result<()> {
let ikb = IndexKeyBase::new(ns, db, tb.clone(), ix);
delete_durable_build_artifacts(tx, &ikb, expunge).await?;
Ok(())
}
pub(super) async fn delete_durable_build_artifacts(
tx: &Transaction,
ikb: &IndexKeyBase,
expunge: bool,
) -> Result<()> {
if expunge {
tx.clr_key(&ikb.new_bs_key()).await?;
tx.clrr(ikb.new_bg_all_generations_range()?).await?;
tx.clrr(ikb.new_bp_all_generations_range()?).await?;
tx.clrr(ikb.new_br_all_generations_range()?).await?;
tx.clrr(ikb.new_bt_all_generations_range()?).await?;
} else {
tx.del_key(&ikb.new_bs_key()).await?;
delete_durable_build_queues(tx, ikb).await?;
}
Ok(())
}
pub(super) async fn delete_durable_build_queues(
tx: &Transaction,
ikb: &IndexKeyBase,
) -> Result<()> {
tx.delr(ikb.new_bg_all_generations_range()?).await?;
tx.delr(ikb.new_bp_all_generations_range()?).await?;
tx.delr(ikb.new_br_all_generations_range()?).await?;
tx.delr(ikb.new_bt_all_generations_range()?).await?;
Ok(())
}
pub(super) async fn delete_stale_build_queues(
tx: &Transaction,
ikb: &IndexKeyBase,
below: super::BuildGeneration,
) -> Result<()> {
tx.delr(ikb.new_bg_range_below(below)?).await?;
tx.delr(ikb.new_bp_range_below(below)?).await?;
tx.delr(ikb.new_br_range_below(below)?).await?;
tx.delr(ikb.new_bt_range_below(below)?).await?;
Ok(())
}
pub(super) fn durable_report_count(count: Option<u64>) -> usize {
match count {
Some(count) => usize::try_from(count).unwrap_or(usize::MAX),
None => 0,
}
}
pub(super) fn is_condition_not_met(err: &anyhow::Error) -> bool {
matches!(storage_error(err), Some(KvsError::TransactionConditionNotMet))
}
pub(super) fn build_owner_expired(state: &IndexBuildState, now: DateTime<Utc>) -> bool {
state.owner_heartbeat_at.unwrap_or(state.updated_at)
+ chrono::Duration::seconds(BUILD_OWNER_LEASE_SECS)
<= now
}
pub(crate) async fn build_owner_is_live(tx: &Transaction, ikb: &IndexKeyBase) -> Result<bool> {
Ok(build_owner_holds(tx.get_key(&ikb.new_bs_key(), None).await?))
}
pub(crate) async fn build_owner_is_live_locked(
tx: &Transaction,
ikb: &IndexKeyBase,
) -> Result<bool> {
Ok(build_owner_holds(tx.getu_key(&ikb.new_bs_key()).await?))
}
fn build_owner_holds(state: Option<IndexBuildState>) -> bool {
let Some(state) = state else {
return false;
};
match state.phase {
IndexBuildPhase::Building | IndexBuildPhase::Closing => true,
IndexBuildPhase::Online => state.owner.is_some(),
IndexBuildPhase::Error => false,
}
}