use super::{
delete_document_statements, event_insert_statements, EventKind, KhiveRuntime, NamespaceToken,
PlanStatement, RuntimeError, RuntimeResult, SqlStatement, SqlValue, SubstrateKind, Uuid, Value,
};
#[cfg(doc)]
use crate::atomic_runner::apply_plan;
fn vector_table_names(runtime: &KhiveRuntime) -> Vec<String> {
runtime
.registered_embedding_model_names()
.iter()
.map(|name| format!("vec_{}", crate::config::sanitize_key(name)))
.collect()
}
fn purge_index_row_statement(
table: &str,
namespace: &str,
subject_id: Uuid,
label: &str,
) -> PlanStatement {
PlanStatement {
statement: SqlStatement {
sql: format!("DELETE FROM {table} WHERE namespace = ?1 AND subject_id = ?2"),
params: vec![
SqlValue::Text(namespace.to_string()),
SqlValue::Text(subject_id.to_string()),
],
label: Some(label.to_string()),
},
guard: None,
}
}
fn purge_vector_provenance_statement(
table: &str,
namespace: &str,
subject_id: Uuid,
label: &str,
) -> PlanStatement {
let model_key = table
.strip_prefix("vec_")
.expect("runtime vector tables use the vec_ prefix");
PlanStatement {
statement: SqlStatement {
sql: "DELETE FROM vector_provenance \
WHERE model_key = ?1 AND subject_id = ?2 AND namespace = ?3"
.to_string(),
params: vec![
SqlValue::Text(model_key.to_string()),
SqlValue::Text(subject_id.to_string()),
SqlValue::Text(namespace.to_string()),
],
label: Some(label.to_string()),
},
guard: None,
}
}
fn purge_fts_document_statements(
fts_table: &str,
namespace: &str,
subject_id: Uuid,
label_prefix: &str,
) -> [PlanStatement; 2] {
let [mut fts_stmt, mut map_stmt] = delete_document_statements(fts_table, namespace, subject_id);
fts_stmt.label = Some(label_prefix.to_string());
map_stmt.label = Some(format!("{label_prefix}-map"));
[
PlanStatement {
statement: fts_stmt,
guard: None,
},
PlanStatement {
statement: map_stmt,
guard: None,
},
]
}
fn log_vector_row_delete_statement(
table: &str,
namespace: &str,
subject_id: Uuid,
label: &str,
) -> PlanStatement {
PlanStatement {
statement: SqlStatement {
sql: format!(
"INSERT INTO ann_write_log \
(namespace, embedding_model, kind, field, subject_id, op) \
SELECT namespace, embedding_model, kind, field, subject_id, 'delete' \
FROM {table} WHERE namespace = ?1 AND subject_id = ?2"
),
params: vec![
SqlValue::Text(namespace.to_string()),
SqlValue::Text(subject_id.to_string()),
],
label: Some(label.to_string()),
},
guard: None,
}
}
async fn vector_table_exists(runtime: &KhiveRuntime, table: &str) -> RuntimeResult<bool> {
let mut reader = runtime
.sql()
.reader()
.await
.map_err(RuntimeError::Storage)?;
let row = reader
.query_scalar(SqlStatement {
sql: "SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?1".to_string(),
params: vec![SqlValue::Text(table.to_string())],
label: Some("atomic-delete-vec-table-exists".to_string()),
})
.await
.map_err(RuntimeError::Storage)?;
Ok(row.is_some())
}
pub(super) async fn push_index_purge_statements(
runtime: &KhiveRuntime,
statements: &mut Vec<PlanStatement>,
fts_table: &str,
namespace: &str,
subject_id: Uuid,
label_prefix: &str,
) -> RuntimeResult<()> {
statements.extend(purge_fts_document_statements(
fts_table,
namespace,
subject_id,
&format!("{label_prefix}-purge-fts"),
));
for vec_table in vector_table_names(runtime) {
if vector_table_exists(runtime, &vec_table).await? {
statements.push(log_vector_row_delete_statement(
&vec_table,
namespace,
subject_id,
&format!("{label_prefix}-log-delete-vec-{vec_table}"),
));
statements.push(purge_index_row_statement(
&vec_table,
namespace,
subject_id,
&format!("{label_prefix}-purge-vec-{vec_table}"),
));
statements.push(purge_vector_provenance_statement(
&vec_table,
namespace,
subject_id,
&format!("{label_prefix}-purge-vec-provenance-{vec_table}"),
));
}
}
Ok(())
}
pub(crate) fn event_append_statements(
token: &NamespaceToken,
namespace: &str,
verb: &str,
kind: EventKind,
substrate: SubstrateKind,
target_id: Uuid,
payload: Value,
) -> RuntimeResult<Vec<PlanStatement>> {
let record_token = token
.with_namespace(crate::Namespace::parse(namespace).map_err(|error| {
RuntimeError::Internal(format!("event namespace invalid: {error}"))
})?);
let event = crate::EventAttribution::from_token(&record_token).stamp(
khive_storage::event::Event::new(namespace.to_string(), verb, kind, substrate, "")
.with_target(target_id)
.with_payload(payload),
);
let statements = event_insert_statements(&event)
.map_err(|e| RuntimeError::Internal(format!("event_insert_statements: {e}")))?;
Ok(statements
.into_iter()
.map(|statement| PlanStatement {
statement,
guard: None,
})
.collect())
}