1use std::any::Any;
2
3use khive_storage::{note::Note, AtomicUnitOp, Entity, SqlStatement, SqlValue, TextDocument};
4use khive_types::SubstrateKind;
5use serde::Serialize;
6use uuid::Uuid;
7
8use crate::{KhiveRuntime, NamespaceToken, PostCommitDegradation, RuntimeError, RuntimeResult};
9
10#[derive(Clone, Debug, Serialize)]
13pub struct IndexRepairReport {
14 pub id: Uuid,
15 pub substrate: SubstrateKind,
16 pub namespace: String,
17 pub repaired: Vec<String>,
18 pub failures: Vec<PostCommitDegradation>,
19}
20
21enum Record {
22 Entity(Entity),
23 Note(Note),
24}
25
26impl Record {
27 fn document(&self) -> TextDocument {
28 match self {
29 Self::Entity(entity) => crate::entity_fts_document(entity),
30 Self::Note(note) => crate::note_fts_document(note),
31 }
32 }
33
34 fn version(&self) -> i64 {
35 match self {
36 Self::Entity(entity) => entity.version,
37 Self::Note(note) => note.version,
38 }
39 }
40
41 fn tables(&self) -> (&'static str, &'static str, &'static str) {
42 match self {
43 Self::Entity(_) => ("entities", "fts_entities", "entity.body"),
44 Self::Note(_) => ("notes", "fts_notes", "note.content"),
45 }
46 }
47
48 fn vector_statements(&self, table: &str, model: &str, vector: &[f32]) -> Vec<SqlStatement> {
49 match self {
50 Self::Entity(entity) => {
51 KhiveRuntime::entity_vector_insert_statements(table, entity, model, vector)
52 }
53 Self::Note(note) => crate::atomic_message::vector_insert_statements(
54 table,
55 ¬e.namespace,
56 note.id,
57 "note.content",
58 model,
59 vector,
60 "record-index-repair",
61 )
62 .into_iter()
63 .map(|plan| plan.statement)
64 .collect(),
65 }
66 }
67}
68
69#[derive(Clone, Copy)]
70enum Publication {
71 Repaired,
72 Healthy,
73 Changed,
74 Occupied,
75}
76
77fn statement(sql: String, params: Vec<SqlValue>) -> SqlStatement {
78 SqlStatement {
79 sql,
80 params,
81 label: Some("record-index-repair".into()),
82 }
83}
84
85fn vector_checks(record: &Record, doc: &TextDocument, model: &str) -> (SqlStatement, SqlStatement) {
86 let table = format!("vec_{}", crate::config::sanitize_key(model));
87 let healthy = statement(format!("SELECT 1 FROM {table} WHERE subject_id=?1 AND namespace=?2 AND kind=?3 AND field=?4 AND embedding_model=?5"), vec![
88 SqlValue::Text(doc.subject_id.to_string()), SqlValue::Text(doc.namespace.clone()), SqlValue::Text(doc.kind.to_string()), SqlValue::Text(record.tables().2.into()), SqlValue::Text(model.into()),
89 ]);
90 let occupied = statement(
93 format!("SELECT 1 FROM {table} WHERE subject_id=?1"),
94 vec![SqlValue::Text(doc.subject_id.to_string())],
95 );
96 (healthy, occupied)
97}
98
99fn same_document(a: &TextDocument, b: &TextDocument) -> bool {
100 a.subject_id == b.subject_id
101 && a.kind == b.kind
102 && a.record_kind == b.record_kind
103 && a.namespace == b.namespace
104 && a.title.as_deref().unwrap_or_default() == b.title.as_deref().unwrap_or_default()
105 && a.body == b.body
106 && a.tags == b.tags
107 && a.metadata == b.metadata
108 && a.updated_at == b.updated_at
109}
110
111impl IndexRepairReport {
112 fn failure(&mut self, stage: &'static str, error: impl ToString) {
113 self.failures.push(PostCommitDegradation {
114 stage,
115 error: error.to_string(),
116 });
117 }
118
119 fn publication(
120 &mut self,
121 result: RuntimeResult<Publication>,
122 stage: &'static str,
123 label: String,
124 ) -> bool {
125 match result {
126 Ok(Publication::Repaired) => self.repaired.push(label),
127 Ok(Publication::Healthy) => {},
128 Ok(Publication::Changed) => {
129 self.failure("source_revision", "record changed or was deleted before index publication; rerun against the current record");
130 return false;
131 }
132 Ok(Publication::Occupied) => self.failure(stage, format!("{label}: another vector identity occupies the subject; existing vector was preserved")),
133 Err(error) => self.failure(stage, format!("{label}: {error}")),
134 }
135 true
136 }
137}
138
139impl KhiveRuntime {
140 pub async fn repair_record_indexes(
145 &self,
146 token: &NamespaceToken,
147 id: Uuid,
148 ) -> RuntimeResult<IndexRepairReport> {
149 if self.is_read_only() {
150 return Err(RuntimeError::InvalidInput(
151 "record index repair requires a writable runtime".into(),
152 ));
153 }
154 let namespace = token.namespace().as_str();
155 let entity = self
156 .entities(token)?
157 .get_entity(id)
158 .await?
159 .filter(|entity| entity.namespace == namespace);
160 let note = self
161 .notes(token)?
162 .get_note(id)
163 .await?
164 .filter(|note| note.namespace == namespace);
165 let record = match (entity, note) {
166 (Some(entity), None) => Record::Entity(entity),
167 (None, Some(note)) => Record::Note(note),
168 (Some(_), Some(_)) => return Err(RuntimeError::InvalidInput(format!("record {id} is ambiguous: both an entity and a note exist in namespace {namespace}"))),
169 (None, None) => return Err(RuntimeError::NotFound(format!("live entity or note {id} in namespace {namespace}"))),
170 };
171 let doc = record.document();
172 let mut report = IndexRepairReport {
173 id,
174 substrate: doc.kind,
175 namespace: doc.namespace.clone(),
176 repaired: vec![],
177 failures: vec![],
178 };
179 let mut models = match &record {
180 Record::Entity(_) => self.registered_embedding_model_names(),
181 Record::Note(note) => self.embedding_models_for_note_kind(¬e.kind),
182 };
183 models.sort();
184 models.dedup();
185 let text = match &record {
186 Record::Entity(_) => self.text(token),
187 Record::Note(_) => self.text_for_notes(token),
188 };
189 match text {
190 Ok(text) => match text.get_document(&doc.namespace, id).await {
191 Ok(Some(current)) if same_document(¤t, &doc) => {}
192 Ok(_) => {
193 let result = self.repair_record_fts(&record, &doc).await;
194 if !report.publication(result, "fts", "fts".into()) {
195 return Ok(report);
196 }
197 }
198 Err(error) => report.failure("fts", error),
199 },
200 Err(error) => report.failure("fts", error),
201 }
202 let body = match &record {
203 Record::Entity(_) => doc.body.as_str(),
204 Record::Note(note) => crate::curation::note_embedding_text_ref(note),
205 };
206 if body.trim().is_empty() {
207 return Ok(report);
208 }
209 for model in models {
210 let (storage_model, dimensions) = match self.vector_model_metadata(&model) {
211 Ok(metadata) => metadata,
212 Err(error) => {
213 report.failure("vector_presence", format!("{model}: {error}"));
214 continue;
215 }
216 };
217 let store = match self.backend().vectors_for_namespace(
218 &crate::config::sanitize_key(&storage_model),
219 &storage_model,
220 dimensions,
221 &doc.namespace,
222 ) {
223 Ok(store) => store,
224 Err(error) => {
225 report.failure("vector_presence", format!("{model}: {error}"));
226 continue;
227 }
228 };
229 let (healthy, occupied) = vector_checks(&record, &doc, &storage_model);
230 let presence: RuntimeResult<Option<Publication>> = async {
231 let mut reader = self.sql().reader().await?;
232 if reader.query_scalar(healthy).await?.is_some() {
233 Ok(Some(Publication::Healthy))
234 } else if reader.query_scalar(occupied).await?.is_some() {
235 Ok(Some(Publication::Occupied))
236 } else {
237 Ok(None)
238 }
239 }
240 .await;
241 match presence {
242 Ok(Some(Publication::Healthy)) => {
243 #[cfg(test)]
244 tests::pause("vector_read").await;
245 match store
246 .get_vectors(&[id], &doc.namespace, record.tables().2)
247 .await
248 {
249 Ok(vectors) if vectors.contains_key(&id) => continue,
250 Ok(_) => report.failure(
251 "vector_presence",
252 format!(
253 "{model}: vector identity was present for {id}, \
254 but its vector was not returned"
255 ),
256 ),
257 Err(error) => {
258 report.failure("vector_presence", format!("{model}: {error}"))
259 }
260 }
261 continue;
262 }
263 Ok(Some(Publication::Occupied)) => {
264 report.failure("vector_presence", format!("{model}: another vector identity occupies the subject; existing vector was preserved"));
265 continue;
266 }
267 Ok(_) => {}
268 Err(error) => {
269 report.failure("vector_presence", format!("{model}: {error}"));
270 continue;
271 }
272 }
273 let outcome = match self
274 .embed_document_with_model_outcome_for_token(token, &model, body)
275 .await
276 {
277 Ok(outcome) => outcome,
278 Err(error) => {
279 report.failure("embedding", format!("{model}: {error}"));
280 continue;
281 }
282 };
283 if outcome.vector.len() != dimensions
284 || outcome.vector.iter().any(|value| !value.is_finite())
285 {
286 report.failure(
287 "vector_publication",
288 format!("{model}: embedding has invalid dimensions or non-finite values"),
289 );
290 continue;
291 }
292 let result = self
293 .repair_record_vector(&record, &doc, &storage_model, &outcome.vector)
294 .await;
295 if !report.publication(
296 result,
297 "vector_publication",
298 format!("vector:{storage_model}"),
299 ) {
300 break;
301 }
302 }
303 Ok(report)
304 }
305
306 async fn repair_record_fts(
307 &self,
308 record: &Record,
309 doc: &TextDocument,
310 ) -> RuntimeResult<Publication> {
311 let (_, table, _) = record.tables();
312 let canonical = khive_db::stores::text::insert_document_statement(table, doc);
313 let map = khive_db::stores::text::rowid_map_table(table);
314 let healthy = statement(
315 format!(
316 "SELECT 1 FROM {table} AS t JOIN {map} AS m ON m.rowid=t.rowid \
317 WHERE m.subject_id=?1 AND m.namespace=?6 AND t.subject_id=?1 AND t.kind=?2 \
318 AND t.title=?3 AND t.body=?4 AND t.tags=?5 AND t.namespace=?6 \
319 AND t.metadata IS ?7 AND t.updated_at=?8 AND t.record_kind IS ?9"
320 ),
321 canonical.params,
322 );
323 let statements = khive_db::stores::text::delete_document_statements(
324 table,
325 &doc.namespace,
326 doc.subject_id,
327 )
328 .into_iter()
329 .chain(khive_db::stores::text::insert_document_statements(
330 table, doc,
331 ))
332 .collect();
333 #[cfg(test)]
334 tests::pause("fts").await;
335 self.publish_record_repair(record, doc, healthy, None, statements)
336 .await
337 }
338
339 async fn repair_record_vector(
340 &self,
341 record: &Record,
342 doc: &TextDocument,
343 model: &str,
344 vector: &[f32],
345 ) -> RuntimeResult<Publication> {
346 let table = format!("vec_{}", crate::config::sanitize_key(model));
347 let (healthy, occupied) = vector_checks(record, doc, model);
348 let statements = record.vector_statements(&table, model, vector);
349 #[cfg(test)]
350 tests::pause("vector").await;
351 self.publish_record_repair(record, doc, healthy, Some(occupied), statements)
352 .await
353 }
354
355 async fn publish_record_repair(
356 &self,
357 record: &Record,
358 doc: &TextDocument,
359 healthy: SqlStatement,
360 occupied: Option<SqlStatement>,
361 statements: Vec<SqlStatement>,
362 ) -> RuntimeResult<Publication> {
363 let source = statement(
364 format!(
365 "SELECT version FROM {} WHERE id=?1 AND namespace=?2 AND deleted_at IS NULL",
366 record.tables().0
367 ),
368 vec![
369 SqlValue::Text(doc.subject_id.to_string()),
370 SqlValue::Text(doc.namespace.clone()),
371 ],
372 );
373 let version = record.version();
374 let op: AtomicUnitOp = Box::new(move |writer| {
375 Box::pin(async move {
376 let current = writer.query_scalar(source).await?;
377 let result = if !matches!(current, Some(SqlValue::Integer(current)) if current == version)
378 {
379 Publication::Changed
380 } else if writer.query_scalar(healthy).await?.is_some() {
381 Publication::Healthy
382 } else if let Some(occupied) = occupied {
383 if writer.query_scalar(occupied).await?.is_some() {
384 Publication::Occupied
385 } else {
386 for statement in statements {
387 writer.execute(statement).await?;
388 }
389 Publication::Repaired
390 }
391 } else {
392 for statement in statements {
393 writer.execute(statement).await?;
394 }
395 Publication::Repaired
396 };
397 Ok(Box::new(result) as Box<dyn Any + Send>)
398 })
399 });
400 self.sql()
401 .atomic_unit(op)
402 .await?
403 .downcast::<Publication>()
404 .map(|result| *result)
405 .map_err(|_| RuntimeError::Internal("invalid record index repair outcome".into()))
406 }
407}
408
409#[cfg(test)]
410#[path = "index_repair_tests.rs"]
411mod tests;