Skip to main content

khive_runtime/
index_repair.rs

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/// Per-record index work that actually committed. Failed stages can coexist
11/// with repairs; rerunning repairs only the gaps that still remain.
12#[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                &note.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 vector_checks(record: &Record, doc: &TextDocument, model: &str) -> (SqlStatement, SqlStatement) {
78    let table = format!("vec_{}", crate::config::sanitize_key(model));
79    let healthy = SqlStatement::new(format!("SELECT 1 FROM {table} WHERE subject_id=?1 AND namespace=?2 AND kind=?3 AND field=?4 AND embedding_model=?5"), vec![
80        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()),
81    ]).labelled("record-index-repair");
82    // Even a mismatched identity may occupy the vec0 subject primary key.
83    // Repair preserves it instead of silently replacing an existing vector.
84    let occupied = SqlStatement::new(
85        format!("SELECT 1 FROM {table} WHERE subject_id=?1"),
86        vec![SqlValue::Text(doc.subject_id.to_string())],
87    )
88    .labelled("record-index-repair");
89    (healthy, occupied)
90}
91
92fn same_document(a: &TextDocument, b: &TextDocument) -> bool {
93    a.subject_id == b.subject_id
94        && a.kind == b.kind
95        && a.record_kind == b.record_kind
96        && a.namespace == b.namespace
97        && a.title.as_deref().unwrap_or_default() == b.title.as_deref().unwrap_or_default()
98        && a.body == b.body
99        && a.tags == b.tags
100        && a.metadata == b.metadata
101        && a.updated_at == b.updated_at
102}
103
104impl IndexRepairReport {
105    fn failure(&mut self, stage: &'static str, error: impl ToString) {
106        self.failures.push(PostCommitDegradation {
107            stage,
108            error: error.to_string(),
109        });
110    }
111
112    fn publication(
113        &mut self,
114        result: RuntimeResult<Publication>,
115        stage: &'static str,
116        label: String,
117    ) -> bool {
118        match result {
119            Ok(Publication::Repaired) => self.repaired.push(label),
120            Ok(Publication::Healthy) => {},
121            Ok(Publication::Changed) => {
122                self.failure("source_revision", "record changed or was deleted before index publication; rerun against the current record");
123                return false;
124            }
125            Ok(Publication::Occupied) => self.failure(stage, format!("{label}: another vector identity occupies the subject; existing vector was preserved")),
126            Err(error) => self.failure(stage, format!("{label}: {error}")),
127        }
128        true
129    }
130}
131
132impl KhiveRuntime {
133    /// Repair indexes for one unambiguous live ID in the token namespace.
134    /// Only missing/stale FTS and missing kind-selected vectors are repaired.
135    /// Existing vectors, substrate rows and other records are never rewritten.
136    /// A later failure does not undo earlier reported repairs.
137    pub async fn repair_record_indexes(
138        &self,
139        token: &NamespaceToken,
140        id: Uuid,
141    ) -> RuntimeResult<IndexRepairReport> {
142        if self.is_read_only() {
143            return Err(RuntimeError::InvalidInput(
144                "record index repair requires a writable runtime".into(),
145            ));
146        }
147        let namespace = token.namespace().as_str();
148        let entity = self
149            .entities(token)?
150            .get_entity(id)
151            .await?
152            .filter(|entity| entity.namespace == namespace);
153        let note = self
154            .notes(token)?
155            .get_note(id)
156            .await?
157            .filter(|note| note.namespace == namespace);
158        let record = match (entity, note) {
159            (Some(entity), None) => Record::Entity(entity),
160            (None, Some(note)) => Record::Note(note),
161            (Some(_), Some(_)) => return Err(RuntimeError::InvalidInput(format!("record {id} is ambiguous: both an entity and a note exist in namespace {namespace}"))),
162            (None, None) => return Err(RuntimeError::NotFound(format!("live entity or note {id} in namespace {namespace}"))),
163        };
164        let doc = record.document();
165        let mut report = IndexRepairReport {
166            id,
167            substrate: doc.kind,
168            namespace: doc.namespace.clone(),
169            repaired: vec![],
170            failures: vec![],
171        };
172        let mut models = match &record {
173            Record::Entity(_) => self.registered_embedding_model_names(),
174            Record::Note(note) => self.embedding_models_for_note_kind(&note.kind),
175        };
176        models.sort();
177        models.dedup();
178        let text = match &record {
179            Record::Entity(_) => self.text(token),
180            Record::Note(_) => self.text_for_notes(token),
181        };
182        match text {
183            Ok(text) => match text.get_document(&doc.namespace, id).await {
184                Ok(Some(current)) if same_document(&current, &doc) => {}
185                Ok(_) => {
186                    let result = self.repair_record_fts(&record, &doc).await;
187                    if !report.publication(result, "fts", "fts".into()) {
188                        return Ok(report);
189                    }
190                }
191                Err(error) => report.failure("fts", error),
192            },
193            Err(error) => report.failure("fts", error),
194        }
195        let body = match &record {
196            Record::Entity(_) => doc.body.as_str(),
197            Record::Note(note) => crate::curation::note_embedding_text_ref(note),
198        };
199        if body.trim().is_empty() {
200            return Ok(report);
201        }
202        for model in models {
203            let (storage_model, dimensions) = match self.vector_model_metadata(&model) {
204                Ok(metadata) => metadata,
205                Err(error) => {
206                    report.failure("vector_presence", format!("{model}: {error}"));
207                    continue;
208                }
209            };
210            let store = match self.backend().vectors_for_namespace(
211                &crate::config::sanitize_key(&storage_model),
212                &storage_model,
213                dimensions,
214                &doc.namespace,
215            ) {
216                Ok(store) => store,
217                Err(error) => {
218                    report.failure("vector_presence", format!("{model}: {error}"));
219                    continue;
220                }
221            };
222            let (healthy, occupied) = vector_checks(&record, &doc, &storage_model);
223            let presence: RuntimeResult<Option<Publication>> = async {
224                let mut reader = self.sql().reader().await?;
225                if reader.query_scalar(healthy).await?.is_some() {
226                    Ok(Some(Publication::Healthy))
227                } else if reader.query_scalar(occupied).await?.is_some() {
228                    Ok(Some(Publication::Occupied))
229                } else {
230                    Ok(None)
231                }
232            }
233            .await;
234            match presence {
235                Ok(Some(Publication::Healthy)) => {
236                    #[cfg(test)]
237                    tests::pause("vector_read").await;
238                    match store
239                        .get_vectors(&[id], &doc.namespace, record.tables().2)
240                        .await
241                    {
242                        Ok(vectors) if vectors.contains_key(&id) => continue,
243                        Ok(_) => report.failure(
244                            "vector_presence",
245                            format!(
246                                "{model}: vector identity was present for {id}, \
247                                 but its vector was not returned"
248                            ),
249                        ),
250                        Err(error) => {
251                            report.failure("vector_presence", format!("{model}: {error}"))
252                        }
253                    }
254                    continue;
255                }
256                Ok(Some(Publication::Occupied)) => {
257                    report.failure("vector_presence", format!("{model}: another vector identity occupies the subject; existing vector was preserved"));
258                    continue;
259                }
260                Ok(_) => {}
261                Err(error) => {
262                    report.failure("vector_presence", format!("{model}: {error}"));
263                    continue;
264                }
265            }
266            let outcome = match self
267                .embed_document_with_model_outcome_for_token(token, &model, body)
268                .await
269            {
270                Ok(outcome) => outcome,
271                Err(error) => {
272                    report.failure("embedding", format!("{model}: {error}"));
273                    continue;
274                }
275            };
276            if outcome.vector.len() != dimensions
277                || outcome.vector.iter().any(|value| !value.is_finite())
278            {
279                report.failure(
280                    "vector_publication",
281                    format!("{model}: embedding has invalid dimensions or non-finite values"),
282                );
283                continue;
284            }
285            let result = self
286                .repair_record_vector(&record, &doc, &storage_model, &outcome.vector)
287                .await;
288            if !report.publication(
289                result,
290                "vector_publication",
291                format!("vector:{storage_model}"),
292            ) {
293                break;
294            }
295        }
296        Ok(report)
297    }
298
299    async fn repair_record_fts(
300        &self,
301        record: &Record,
302        doc: &TextDocument,
303    ) -> RuntimeResult<Publication> {
304        let (_, table, _) = record.tables();
305        let canonical = khive_db::stores::text::insert_document_statement(table, doc);
306        let map = khive_db::stores::text::rowid_map_table(table);
307        let healthy = SqlStatement::new(
308            format!(
309                "SELECT 1 FROM {table} AS t JOIN {map} AS m ON m.rowid=t.rowid \
310            WHERE m.subject_id=?1 AND m.namespace=?6 AND t.subject_id=?1 AND t.kind=?2 \
311            AND t.title=?3 AND t.body=?4 AND t.tags=?5 AND t.namespace=?6 \
312            AND t.metadata IS ?7 AND t.updated_at=?8 AND t.record_kind IS ?9"
313            ),
314            canonical.params,
315        )
316        .labelled("record-index-repair");
317        let statements = khive_db::stores::text::delete_document_statements(
318            table,
319            &doc.namespace,
320            doc.subject_id,
321        )
322        .into_iter()
323        .chain(khive_db::stores::text::insert_document_statements(
324            table, doc,
325        ))
326        .collect();
327        #[cfg(test)]
328        tests::pause("fts").await;
329        self.publish_record_repair(record, doc, healthy, None, statements)
330            .await
331    }
332
333    async fn repair_record_vector(
334        &self,
335        record: &Record,
336        doc: &TextDocument,
337        model: &str,
338        vector: &[f32],
339    ) -> RuntimeResult<Publication> {
340        let table = format!("vec_{}", crate::config::sanitize_key(model));
341        let (healthy, occupied) = vector_checks(record, doc, model);
342        let statements = record.vector_statements(&table, model, vector);
343        #[cfg(test)]
344        tests::pause("vector").await;
345        self.publish_record_repair(record, doc, healthy, Some(occupied), statements)
346            .await
347    }
348
349    async fn publish_record_repair(
350        &self,
351        record: &Record,
352        doc: &TextDocument,
353        healthy: SqlStatement,
354        occupied: Option<SqlStatement>,
355        statements: Vec<SqlStatement>,
356    ) -> RuntimeResult<Publication> {
357        let source = SqlStatement::new(
358            format!(
359                "SELECT version FROM {} WHERE id=?1 AND namespace=?2 AND deleted_at IS NULL",
360                record.tables().0
361            ),
362            vec![
363                SqlValue::Text(doc.subject_id.to_string()),
364                SqlValue::Text(doc.namespace.clone()),
365            ],
366        )
367        .labelled("record-index-repair");
368        let version = record.version();
369        let op: AtomicUnitOp = Box::new(move |writer| {
370            Box::pin(async move {
371                let current = writer.query_scalar(source).await?;
372                let result = if !matches!(current, Some(SqlValue::Integer(current)) if current == version)
373                {
374                    Publication::Changed
375                } else if writer.query_scalar(healthy).await?.is_some() {
376                    Publication::Healthy
377                } else if let Some(occupied) = occupied {
378                    if writer.query_scalar(occupied).await?.is_some() {
379                        Publication::Occupied
380                    } else {
381                        for statement in statements {
382                            writer.execute(statement).await?;
383                        }
384                        Publication::Repaired
385                    }
386                } else {
387                    for statement in statements {
388                        writer.execute(statement).await?;
389                    }
390                    Publication::Repaired
391                };
392                Ok(Box::new(result) as Box<dyn Any + Send>)
393            })
394        });
395        self.sql()
396            .atomic_unit(op)
397            .await?
398            .downcast::<Publication>()
399            .map(|result| *result)
400            .map_err(|_| RuntimeError::Internal("invalid record index repair outcome".into()))
401    }
402}
403
404#[cfg(test)]
405#[path = "index_repair_tests.rs"]
406mod tests;