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 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 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 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(¬e.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(¤t, &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;