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 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    // Even a mismatched identity may occupy the vec0 subject primary key.
91    // Repair preserves it instead of silently replacing an existing vector.
92    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    /// Repair indexes for one unambiguous live ID in the token namespace.
141    /// Only missing/stale FTS and missing kind-selected vectors are repaired.
142    /// Existing vectors, substrate rows and other records are never rewritten.
143    /// A later failure does not undo earlier reported repairs.
144    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(&note.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(&current, &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;