use std::sync::Arc;
use anyhow::Result;
use chrono::{DateTime, Utc};
use revision::revisioned;
use serde::{Deserialize, Serialize};
use uuid::Uuid;
use super::{BUILD_OWNER_LEASE_SECS, BuildGeneration, BuildTicket};
use crate::catalog::providers::TableProvider;
use crate::catalog::{DatabaseId, IndexDefinition, IndexId, NamespaceId};
use crate::err::Error;
use crate::idx::IndexKeyBase;
use crate::key::table::bs::Bs;
use crate::kvs::{Error as KvsError, Transaction, impl_kv_value_revisioned};
use crate::val::{Object, TableName, Value};
#[revisioned(revision = 1)]
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub(crate) enum IndexBuildReportStatus {
Started,
Cleaning,
Indexing,
Ready,
Aborted,
Error,
}
impl IndexBuildReportStatus {
fn as_str(self) -> &'static str {
match self {
Self::Started => "started",
Self::Cleaning => "cleaning",
Self::Indexing => "indexing",
Self::Ready => "ready",
Self::Aborted => "aborted",
Self::Error => "error",
}
}
}
#[revisioned(revision = 1)]
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub(crate) enum IndexBuildPhase {
Building,
Closing,
Online,
Error,
}
#[revisioned(revision = 3)]
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub(crate) struct IndexBuildState {
pub(crate) generation: BuildGeneration,
pub(crate) phase: IndexBuildPhase,
pub(crate) owner: Option<Uuid>,
pub(crate) next_ticket: BuildTicket,
pub(crate) initial_complete: bool,
pub(crate) updated_at: DateTime<Utc>,
#[revision(start = 3)]
pub(crate) owner_heartbeat_at: Option<DateTime<Utc>>,
#[revision(start = 2)]
pub(crate) error: Option<String>,
#[revision(start = 3)]
pub(crate) report_status: Option<IndexBuildReportStatus>,
#[revision(start = 3)]
pub(crate) initial: Option<u64>,
#[revision(start = 3)]
pub(crate) updated: Option<u64>,
#[revision(start = 3)]
pub(crate) pending: Option<u64>,
}
impl_kv_value_revisioned!(IndexBuildState);
#[revisioned(revision = 1)]
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
pub(crate) struct IndexBuildReservation {
pub(crate) node: Uuid,
pub(crate) expires_at: DateTime<Utc>,
}
impl_kv_value_revisioned!(IndexBuildReservation);
pub(super) fn report_status_from_phase(phase: IndexBuildPhase) -> IndexBuildReportStatus {
match phase {
IndexBuildPhase::Building | IndexBuildPhase::Closing => IndexBuildReportStatus::Indexing,
IndexBuildPhase::Online => IndexBuildReportStatus::Ready,
IndexBuildPhase::Error => IndexBuildReportStatus::Error,
}
}
pub(super) fn durable_index_error_reason(ix: &IndexDefinition, state: &IndexBuildState) -> String {
state.error.clone().unwrap_or_else(|| format!("Index {} is in an error state", ix.name))
}
fn index_building_status_value(ix: &IndexDefinition, state: Option<IndexBuildState>) -> Value {
let Some(state) = state else {
let mut out = Object::default();
out.insert("status", IndexBuildReportStatus::Ready.as_str().into());
return out.into();
};
let status = state.report_status.unwrap_or_else(|| report_status_from_phase(state.phase));
let mut out = Object::default();
if let Some(initial) = state.initial {
out.insert("initial", initial.into());
}
if let Some(pending) = state.pending {
out.insert("pending", pending.into());
}
if let Some(updated) = state.updated {
out.insert("updated", updated.into());
}
if status == IndexBuildReportStatus::Error {
out.insert("error", durable_index_error_reason(ix, &state).into());
}
out.insert("status", status.as_str().into());
out.into()
}
pub(crate) async fn index_building_info(
tx: &Transaction,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
) -> Result<Value> {
let ikb = IndexKeyBase::new(ns, db, ix.table_name.clone(), ix.index_id);
let status = tx.get(&ikb.new_bs_key(), None).await?;
let mut out = Object::default();
out.insert("building", index_building_status_value(ix, status));
Ok(out.into())
}
pub(crate) async fn retire_durable_index(
tx: &Transaction,
ns: NamespaceId,
db: DatabaseId,
tb: &TableName,
ix: IndexId,
) -> Result<()> {
let ikb = IndexKeyBase::new(ns, db, tb.clone(), ix);
tx.del(&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?;
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 {
if matches!(err.downcast_ref::<KvsError>(), Some(KvsError::TransactionConditionNotMet)) {
return true;
}
matches!(err.downcast_ref::<Error>(), Some(Error::Kvs(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(super) async fn catalog_still_references_index(
tx: &Transaction,
ns: NamespaceId,
db: DatabaseId,
ix: &IndexDefinition,
) -> Result<bool> {
let Some(current) = tx.get_tb_index(ns, db, &ix.table_name, &ix.name, None).await? else {
return Ok(false);
};
Ok(!current.prepare_remove && current.index_id == ix.index_id)
}
pub(crate) async fn filter_online_indexes(
tx: &Transaction,
ns: NamespaceId,
db: DatabaseId,
indexes: Arc<[IndexDefinition]>,
) -> Result<Arc<[IndexDefinition]>> {
if indexes.is_empty() {
return Ok(indexes);
}
let state_keys: Vec<_> =
indexes.iter().map(|ix| Bs::new(ns, db, &ix.table_name, ix.index_id)).collect();
let states = tx.getm(state_keys, None).await?;
let mut filtered = Vec::new();
let mut filtered_any = false;
for (ix, state) in indexes.iter().zip(states) {
let online = match state {
Some(state) => state.phase == IndexBuildPhase::Online,
None => catalog_still_references_index(tx, ns, db, ix).await?,
};
if online {
filtered.push(ix.clone());
} else {
filtered_any = true;
}
}
if filtered_any {
Ok(Arc::from(filtered))
} else {
Ok(indexes)
}
}