Skip to main content

khive_runtime/curation/
entity_curation.rs

1#[cfg(test)]
2use super::race_seam;
3use super::{
4    entity_fts_document, map_merge_entity_storage_error, merge_entity_sql, merge_properties,
5    stale_entity_snapshot_error, Any, AtomicUnitOp, ContentMergeStrategy, EmbeddingModelPlan,
6    Entity, EntityDedupMergePolicy, EntityMergeValidation, EntityPatch, EventAttribution,
7    EventKind, HashMap, KhiveRuntime, MergeEventContext, MergeSqlError, MergeSummary,
8    MergeTxLimits, NamespaceToken, RuntimeError, RuntimeResult, RuntimeWriteOperation,
9    SqlStatement, SqlValue, SqliteError, SubstrateKind, Uuid, Value,
10};
11
12impl KhiveRuntime {
13    /// Patch-style entity update.
14    ///
15    /// Only fields set to `Some(_)` are changed. Re-indexes FTS5 (and vectors if configured)
16    /// when `name`, `description`, or `entity_type` changes; skips re-indexing for
17    /// property/tag-only patches.
18    ///
19    /// Returns `RuntimeError::NotFound` if the entity does not exist or belongs to a different
20    /// namespace. Namespace isolation is enforced at the runtime layer.
21    /// Computes the patched `Entity`, `reindex_required`, and `changed_fields` without
22    /// writing anything, so both the normal write path and the atomic-prepare path
23    /// share one source of truth for what a patched entity looks like.
24    pub(crate) async fn prepare_update_entity(
25        &self,
26        token: &NamespaceToken,
27        id: Uuid,
28        patch: EntityPatch,
29    ) -> RuntimeResult<(Entity, bool, Vec<&'static str>, i64, Option<i64>)> {
30        self.prepare_guarded_entity_update(token, id, patch, None, &[])
31            .await
32    }
33
34    async fn prepare_guarded_entity_update(
35        &self,
36        token: &NamespaceToken,
37        id: Uuid,
38        patch: EntityPatch,
39        expected: Option<&Entity>,
40        remove_properties: &[&str],
41    ) -> RuntimeResult<(Entity, bool, Vec<&'static str>, i64, Option<i64>)> {
42        crate::secret_gate::reject_reserved_secret_gate_property(patch.properties.as_ref())?;
43        if !remove_properties.is_empty() {
44            let removals = Value::Object(
45                remove_properties
46                    .iter()
47                    .map(|key| ((*key).to_string(), Value::Null))
48                    .collect(),
49            );
50            crate::secret_gate::reject_reserved_secret_gate_property(Some(&removals))?;
51        }
52        if let Some(ref name) = patch.name {
53            crate::secret_gate::check_at(name, "entity", "name")?;
54        }
55        if let Some(Some(ref desc)) = patch.description {
56            crate::secret_gate::check_at(desc, "entity", "description")?;
57        }
58        if let Some(ref props) = patch.properties {
59            crate::secret_gate::check_json_at(props, "entity", "properties")?;
60        }
61        if let Some(ref tags) = patch.tags {
62            crate::secret_gate::check_tags_at(tags, "entity", "tags")?;
63        }
64        let store = self.entities(token)?;
65        let mut entity = store.get_entity(id).await?.ok_or_else(|| {
66            if expected.is_some() {
67                stale_entity_snapshot_error(id)
68            } else {
69                RuntimeError::NotFound(format!("entity {id}"))
70            }
71        })?;
72        if let Some(expected) = expected {
73            let actual = serde_json::to_value(&entity)
74                .map_err(|error| RuntimeError::Internal(error.to_string()))?;
75            let expected = serde_json::to_value(expected)
76                .map_err(|error| RuntimeError::Internal(error.to_string()))?;
77            if actual != expected {
78                return Err(stale_entity_snapshot_error(id));
79            }
80        }
81        let expected_updated_at = entity.updated_at;
82        let expected_deleted_at = entity.deleted_at;
83        #[cfg(test)]
84        race_seam::pause_after_read().await;
85
86        // ADR-014 tri-state: outer `None` = unchanged; `Some(None)` = explicit
87        // clear (no vocabulary validation — there is no value to validate);
88        // `Some(Some(raw))` = set, validated and normalized.
89        let validated_entity_type = match &patch.entity_type {
90            Some(None) => Some(None),
91            Some(Some(raw)) => Some(Some(
92                self.validate_entity_type_for_kind(&entity.kind, Some(raw))?
93                    .expect("set branch always yields a normalized value"),
94            )),
95            None => None,
96        };
97
98        let mut reindex_required = false;
99        let mut changed_fields: Vec<&'static str> = Vec::new();
100
101        if let Some(name) = patch.name {
102            reindex_required |= entity.name != name;
103            entity.name = name;
104            changed_fields.push("name");
105        }
106        if let Some(desc_patch) = patch.description {
107            reindex_required |= entity.description != desc_patch;
108            entity.description = desc_patch;
109            changed_fields.push("description");
110        }
111        if let Some(props) = patch.properties {
112            let (merged, _) = merge_properties(
113                &entity.properties,
114                &Some(props),
115                EntityDedupMergePolicy::PreferFrom,
116            );
117            entity.properties = merged;
118            changed_fields.push("properties");
119        }
120        if let Some(Value::Object(properties)) = entity.properties.as_mut() {
121            let mut removed = false;
122            for key in remove_properties {
123                removed |= properties.remove(*key).is_some();
124            }
125            if removed && !changed_fields.contains(&"properties") {
126                changed_fields.push("properties");
127            }
128        }
129        if let Some(tags) = patch.tags {
130            entity.tags = tags;
131            changed_fields.push("tags");
132        }
133        if let Some(entity_type) = validated_entity_type {
134            reindex_required |= entity.entity_type != entity_type;
135            entity.entity_type = entity_type;
136            changed_fields.push("entity_type");
137        }
138
139        // A patch may carry properties from the stored row into the full
140        // replacement. Validate the final object, including that carry.
141        crate::secret_gate::reject_reserved_secret_gate_property(entity.properties.as_ref())?;
142
143        if expected.is_some() && changed_fields.is_empty() {
144            return Ok((
145                entity,
146                reindex_required,
147                changed_fields,
148                expected_updated_at,
149                expected_deleted_at,
150            ));
151        }
152
153        // #2943: `entity.properties`, `entity.entity_type`, and `entity.tags`
154        // are all final here — the owning pack's KindHook, if any, validates
155        // the resulting record, mirroring `prepare_note_update_hook` on the
156        // note side. `Ok(())` when no pack registered a hook for this entity
157        // kind (the trait default, or the runtime-layer aggregate was never
158        // installed). Placed after the no-op early return above so a
159        // genuinely unchanged guarded update never re-runs the hook for
160        // nothing; only `updated_at` remains to be bumped after this point.
161        if let Some(hook) = self.entity_kind_hook(&entity.kind) {
162            hook.validate_entity_update(self, token, &entity, entity.properties.as_ref())
163                .await?;
164        }
165
166        // `updated_at` is also the optimistic-concurrency revision for
167        // full-entity replacement. Make it strictly advance even when two
168        // operations land inside one clock microsecond. Saturation is not a
169        // valid fallback: reusing i64::MAX would make the CAS accept a write
170        // without advancing its revision.
171        let minimum_updated_at = expected_updated_at.checked_add(1).ok_or_else(|| {
172            RuntimeError::Internal(format!(
173                "entity {id} updated_at is already at i64::MAX and cannot advance"
174            ))
175        })?;
176        entity.updated_at = chrono::Utc::now()
177            .timestamp_micros()
178            .max(minimum_updated_at);
179        Ok((
180            entity,
181            reindex_required,
182            changed_fields,
183            expected_updated_at,
184            expected_deleted_at,
185        ))
186    }
187
188    #[cfg(test)]
189    pub(crate) async fn update_entity(
190        &self,
191        token: &NamespaceToken,
192        id: Uuid,
193        patch: EntityPatch,
194    ) -> RuntimeResult<Entity> {
195        Ok(self
196            .update_entity_with_embedding_report(token, id, patch)
197            .await?
198            .0)
199    }
200
201    pub async fn update_entity_with_embedding_report(
202        &self,
203        token: &NamespaceToken,
204        id: Uuid,
205        patch: EntityPatch,
206    ) -> RuntimeResult<(Entity, crate::retrieval::EmbeddingTruncationReport)> {
207        self.update_entity_with_expected_version_and_embedding_report(token, id, patch, None)
208            .await
209    }
210
211    /// Entity update with an optional caller revision, checked inside the writer transaction.
212    pub async fn update_entity_with_expected_version_and_embedding_report(
213        &self,
214        token: &NamespaceToken,
215        id: Uuid,
216        patch: EntityPatch,
217        expected_version: Option<i64>,
218    ) -> RuntimeResult<(Entity, crate::retrieval::EmbeddingTruncationReport)> {
219        crate::entity_write::validate_expected_version(expected_version)?;
220        let (entity, reindex_required, changed_fields, expected_updated_at, expected_deleted_at) =
221            self.prepare_update_entity(token, id, patch).await?;
222
223        self.persist_prepared_entity_update(
224            token,
225            entity,
226            reindex_required,
227            changed_fields,
228            expected_updated_at,
229            expected_deleted_at,
230            expected_version,
231        )
232        .await
233    }
234
235    /// Apply an admin patch only if the entity still matches the full read snapshot.
236    /// Property removals apply after the normal merge and preserve all other keys.
237    /// Missing keys alone are a no-op; reserved runtime-owned keys cannot be removed.
238    /// A changed, deleted, or missing entity returns a conflict without writing.
239    /// A bounded embedding returns a non-retryable error with the committed ID
240    /// and truncation report; use the report-aware variant to retain the record.
241    pub async fn update_entity_if_unchanged(
242        &self,
243        token: &NamespaceToken,
244        expected: &Entity,
245        patch: EntityPatch,
246        remove_properties: &[&str],
247    ) -> RuntimeResult<Entity> {
248        let (entity, embedding) = self
249            .update_entity_if_unchanged_with_embedding_report(
250                token,
251                expected,
252                patch,
253                remove_properties,
254            )
255            .await?;
256        crate::operations::legacy_post_commit_result_with_embedding(
257            "update_entity_if_unchanged",
258            entity.id,
259            entity,
260            embedding,
261            Vec::new(),
262        )
263    }
264
265    /// Apply a guarded admin patch and retain embedding truncation accounting.
266    pub async fn update_entity_if_unchanged_with_embedding_report(
267        &self,
268        token: &NamespaceToken,
269        expected: &Entity,
270        patch: EntityPatch,
271        remove_properties: &[&str],
272    ) -> RuntimeResult<(Entity, crate::retrieval::EmbeddingTruncationReport)> {
273        let (entity, reindex_required, changed_fields, expected_updated_at, expected_deleted_at) =
274            self.prepare_guarded_entity_update(
275                token,
276                expected.id,
277                patch,
278                Some(expected),
279                remove_properties,
280            )
281            .await?;
282        if changed_fields.is_empty() {
283            return Ok((
284                entity,
285                crate::retrieval::EmbeddingTruncationReport::default(),
286            ));
287        }
288        self.persist_prepared_entity_update(
289            token,
290            entity,
291            reindex_required,
292            changed_fields,
293            expected_updated_at,
294            expected_deleted_at,
295            None,
296        )
297        .await
298    }
299
300    #[allow(clippy::too_many_arguments)]
301    pub(super) async fn persist_prepared_entity_update(
302        &self,
303        token: &NamespaceToken,
304        mut entity: Entity,
305        reindex_required: bool,
306        changed_fields: Vec<&'static str>,
307        expected_updated_at: i64,
308        expected_deleted_at: Option<i64>,
309        expected_version: Option<i64>,
310    ) -> RuntimeResult<(Entity, crate::retrieval::EmbeddingTruncationReport)> {
311        // This final whole-object replacement must reserve the complete
312        // candidate, even if a future caller bypasses the patch preparer.
313        crate::secret_gate::reject_reserved_secret_gate_property(entity.properties.as_ref())?;
314        let id = entity.id;
315        let _ = self.entities(token)?;
316        let next_version = entity
317            .version
318            .checked_add(1)
319            .ok_or_else(|| RuntimeError::InvalidInput("entity version overflow".into()))?;
320        use crate::atomic_plan::{AffectedRowGuard, PlanStatement, PostCommitEffect, UpdatePlan};
321        use crate::atomic_runner::{
322            run_atomic_unit, AtomicOpFailure, AtomicOpPlan, AtomicRunOutcome,
323        };
324        let plan = UpdatePlan {
325            target_id: id,
326            statements: vec![PlanStatement {
327                statement: khive_db::stores::entity::entity_replace_if_unchanged_statement(
328                    &entity,
329                    expected_updated_at,
330                    expected_deleted_at,
331                ),
332                guard: Some(AffectedRowGuard::exactly(1)),
333            }],
334            post_commit: PostCommitEffect::None,
335            edge_natural_key: None,
336            idempotent_noop: false,
337            entity_guard: expected_version.map(|expected_version| {
338                crate::entity_write::EntityWriteGuard {
339                    id,
340                    expected_version,
341                }
342            }),
343            note_guard: None,
344            note_vector_purge: None,
345            note_embedding_inheritance: None,
346            graph_effects: Vec::new(),
347        };
348        match run_atomic_unit(
349            self.sql().as_ref(),
350            vec![AtomicOpPlan::Update(Box::new(plan))],
351        )
352        .await
353        {
354            Ok(AtomicRunOutcome::Committed { .. }) => entity.version = next_version,
355            Ok(AtomicRunOutcome::RolledBack {
356                failure: AtomicOpFailure::EntityConflict(conflict),
357                ..
358            }) => return Err(conflict.into_error().into()),
359            Ok(AtomicRunOutcome::RolledBack {
360                failure: AtomicOpFailure::GuardFailed { .. },
361                ..
362            }) => return Err(stale_entity_snapshot_error(id)),
363            Ok(AtomicRunOutcome::RolledBack { failure, .. }) => {
364                return Err(RuntimeError::Internal(format!(
365                    "entity update rolled back: {failure:?}"
366                )))
367            }
368            Err(error) => return Err(RuntimeError::Storage(error.0)),
369        }
370
371        let event_token =
372            token.with_namespace(crate::Namespace::parse(&entity.namespace).map_err(|error| {
373                RuntimeError::Internal(format!("entity namespace invalid: {error}"))
374            })?);
375        let event_store = self.events(&event_token)?;
376        let event = khive_storage::event::Event::new(
377            entity.namespace.clone(),
378            "update",
379            EventKind::EntityUpdated,
380            SubstrateKind::Entity,
381            "",
382        )
383        .with_target(entity.id)
384        .with_payload(serde_json::json!({
385            "id": entity.id,
386            "namespace": entity.namespace,
387            "changed_fields": changed_fields,
388        }));
389        let event_result = event_store.append_event(event).await.map_err(|e| {
390            RuntimeError::Internal(format!("update_entity: event store write failed: {e}"))
391        });
392
393        let embedding_report = if reindex_required {
394            self.reindex_entity(token, &entity).await?
395        } else {
396            crate::retrieval::EmbeddingTruncationReport::default()
397        };
398
399        event_result?;
400
401        Ok((entity, embedding_report))
402    }
403
404    /// Merge `from_id` into `into_id`.
405    ///
406    /// All edges incident to `from_id` are rewired to `into_id`. Self-loops that would
407    /// result from the rewire are dropped. Properties and tags are merged per `strategy`.
408    /// `from_id` is tombstoned with merge provenance and removed from indexes. Returns a summary.
409    ///
410    /// If `dry_run` is true, computes and returns the planned summary without mutating any rows.
411    ///
412    /// Atomic: all SQL (entity reads/writes, edge rewires, FTS updates, vec-index
413    /// delete, merge event with destructive edge preimages) runs on one pool
414    /// connection inside one `BEGIN IMMEDIATE` transaction via
415    /// `merge_entity_sql`. If embedding vectors are configured, the vector re-insert for
416    /// `into_id` is performed after the transaction (requires async embedding computation).
417    pub async fn merge_entity(
418        &self,
419        token: &NamespaceToken,
420        into_id: Uuid,
421        from_id: Uuid,
422        strategy: EntityDedupMergePolicy,
423        content_strategy: ContentMergeStrategy,
424        dry_run: bool,
425    ) -> RuntimeResult<MergeSummary> {
426        self.merge_entity_with_reason(
427            token,
428            into_id,
429            from_id,
430            strategy,
431            content_strategy,
432            dry_run,
433            None,
434        )
435        .await
436    }
437
438    /// Merge `from_id` into `into_id` and include an optional reason in the audit event.
439    // REASON: these arguments mirror the merge verb's policy, content strategy,
440    // dry-run, and audit-reason fields; a builder would only move that surface.
441    #[allow(clippy::too_many_arguments)]
442    pub async fn merge_entity_with_reason(
443        &self,
444        token: &NamespaceToken,
445        into_id: Uuid,
446        from_id: Uuid,
447        strategy: EntityDedupMergePolicy,
448        content_strategy: ContentMergeStrategy,
449        dry_run: bool,
450        reason: Option<String>,
451    ) -> RuntimeResult<MergeSummary> {
452        self.merge_entity_with_validation(
453            token,
454            into_id,
455            from_id,
456            strategy,
457            content_strategy,
458            dry_run,
459            reason,
460            EntityMergeValidation::LegacyKind,
461        )
462        .await
463    }
464
465    /// Merge two entities with an explicit override for the entity safety floor.
466    ///
467    /// Non-forced calls enforce entity kind, name similarity, and project compatibility
468    /// against the rows reread inside the merge transaction. Legacy merge methods retain
469    /// their historical same-kind-only policy.
470    /// A non-dry-run override is recorded as `force: true` in the merge event.
471    // REASON: these arguments mirror the merge verb's policy, content strategy,
472    // dry-run, audit-reason, and force fields; a builder would only move that surface.
473    #[allow(clippy::too_many_arguments)]
474    pub async fn merge_entity_with_reason_and_force(
475        &self,
476        token: &NamespaceToken,
477        into_id: Uuid,
478        from_id: Uuid,
479        strategy: EntityDedupMergePolicy,
480        content_strategy: ContentMergeStrategy,
481        dry_run: bool,
482        reason: Option<String>,
483        force: bool,
484    ) -> RuntimeResult<MergeSummary> {
485        let validation = if force {
486            EntityMergeValidation::Forced
487        } else {
488            EntityMergeValidation::SafetyFloor
489        };
490        self.merge_entity_with_validation(
491            token,
492            into_id,
493            from_id,
494            strategy,
495            content_strategy,
496            dry_run,
497            reason,
498            validation,
499        )
500        .await
501    }
502
503    #[allow(clippy::too_many_arguments)]
504    async fn merge_entity_with_validation(
505        &self,
506        token: &NamespaceToken,
507        into_id: Uuid,
508        from_id: Uuid,
509        strategy: EntityDedupMergePolicy,
510        content_strategy: ContentMergeStrategy,
511        dry_run: bool,
512        reason: Option<String>,
513        validation: EntityMergeValidation,
514    ) -> RuntimeResult<MergeSummary> {
515        if let Some(reason) = reason.as_deref() {
516            crate::secret_gate::check_at(reason, "merge", "reason")?;
517        }
518        if into_id == from_id {
519            return Err(RuntimeError::InvalidInput(
520                "cannot merge an entity into itself".into(),
521            ));
522        }
523        let ns = token.namespace().as_str().to_owned();
524        let fts_table = "fts_entities".to_string();
525        // One immutable registry view governs transactional deletion, table
526        // preparation, and survivor reindex. A late model belongs to a later
527        // write/backfill rather than only one leg of this merge.
528        let embedding_plan = EmbeddingModelPlan::capture(self);
529        let vec_tables = embedding_plan.vector_tables();
530        // Loaded once here (sync, cheap) so the rewire loop can evaluate the
531        // endpoint contract without an async round-trip per edge (khive#1216).
532        let pack_rules = self.pack_edge_rules();
533
534        // Ensure all required tables exist (idempotent DDL) before the transaction.
535        let _ = self.entities(token)?;
536        let _ = self.graph(token)?;
537        let _ = self.text(token)?;
538        let _ = self.events(token)?;
539        // vectors_for_model (not the default-model-only self.vectors()) so
540        // custom-only runtimes (no default embedding_model) still get DDL primed.
541        for model_name in embedding_plan.model_names() {
542            let _ = self.vectors_for_model(token, model_name)?;
543        }
544
545        let pool = self.backend().pool_arc();
546        let writer_task = pool
547            .writer_task_for_runtime_write(RuntimeWriteOperation::MergeEntity)
548            .map_err(RuntimeError::Storage)?;
549        // Minted before the transaction so the tombstone and its in-transaction
550        // EntityMerged event carry the same id.
551        let merge_event_id = Uuid::new_v4();
552        let event_context = MergeEventContext {
553            attribution: EventAttribution::from_token(token),
554            reason,
555            force: validation == EntityMergeValidation::Forced,
556            strategy,
557            content_strategy,
558            kind: EventKind::EntityMerged,
559            substrate: SubstrateKind::Entity,
560            event_id: Some(merge_event_id),
561        };
562
563        let (mut summary, updated_entity) = if let Some(writer_task) = writer_task {
564            writer_task
565                .send(move |conn| {
566                    merge_entity_sql(
567                        conn,
568                        ns,
569                        fts_table,
570                        vec_tables,
571                        into_id,
572                        from_id,
573                        strategy,
574                        content_strategy,
575                        dry_run,
576                        pack_rules,
577                        validation,
578                        MergeTxLimits::default(),
579                        merge_event_id,
580                        Some(event_context),
581                    )
582                    .map_err(|e| {
583                        khive_storage::StorageError::driver(
584                            khive_storage::StorageCapability::Entities,
585                            "merge_entity",
586                            e,
587                        )
588                    })
589                })
590                .await
591                .inspect_err(|error| khive_storage::usage::account_event_write(Err(error)))
592                .map_err(map_merge_entity_storage_error)?
593        } else {
594            tokio::task::spawn_blocking(move || {
595                let guard = pool.writer()?;
596                let mut refusal = None;
597                let result = guard.transaction(|conn| {
598                    merge_entity_sql(
599                        conn,
600                        ns,
601                        fts_table,
602                        vec_tables,
603                        into_id,
604                        from_id,
605                        strategy,
606                        content_strategy,
607                        dry_run,
608                        pack_rules,
609                        validation,
610                        MergeTxLimits::default(),
611                        merge_event_id,
612                        Some(event_context),
613                    )
614                    .map_err(|error| match error {
615                        MergeSqlError::Sqlite(error) => error,
616                        MergeSqlError::Refusal(error) => {
617                            refusal = Some(error);
618                            SqliteError::InvalidData(
619                                "entity merge refused by transactional policy".to_string(),
620                            )
621                        }
622                    })
623                });
624                match refusal {
625                    Some(error) => Err(error),
626                    None => result.map_err(RuntimeError::from),
627                }
628            })
629            .await
630            .map_err(|e| RuntimeError::Internal(e.to_string()))??
631        };
632
633        // Count only committed event rows; dry-run never inserts an event.
634        if !dry_run {
635            khive_storage::usage::account_event_write(Ok(1));
636            tracing::info!(
637                into_id = %summary.kept_id,
638                from_id = %summary.removed_id,
639                budget_rows = summary.tx_budget.rows_charged,
640                budget_bytes = summary.tx_budget.bytes_charged,
641                budget_max_rows = summary.tx_budget.max_rows,
642                budget_max_bytes = summary.tx_budget.max_bytes,
643                "merge_entity: transaction materialization budget"
644            );
645        }
646
647        // FTS and vec-deletes already committed inside the transaction above;
648        // only the embedding re-insert needs an async step outside it.
649        if !dry_run && !embedding_plan.is_empty() {
650            match self
651                .reindex_entity_with_plan(token, &updated_entity, &embedding_plan, None)
652                .await
653            {
654                Ok(report) => summary.embedding_truncation = report,
655                Err(error) => {
656                    tracing::warn!(
657                        into_id = %summary.kept_id,
658                        from_id = %summary.removed_id,
659                        error = %error,
660                        "merge_entity: committed merge but survivor reindex failed"
661                    );
662                    summary.post_commit_reindex_error = Some(error.to_string());
663                }
664            }
665        }
666
667        Ok(summary)
668    }
669
670    // ---- Internal helpers ----
671
672    async fn apply_entity_index_revision(
673        &self,
674        entity: &Entity,
675        statements: Vec<SqlStatement>,
676    ) -> RuntimeResult<bool> {
677        let namespace = entity.namespace.clone();
678        let id = entity.id.to_string();
679        let version = entity.version;
680        let op: AtomicUnitOp = Box::new(move |writer| {
681            Box::pin(async move {
682                let current = writer
683                    .query_scalar(SqlStatement {
684                        sql: "SELECT version FROM entities \
685                              WHERE namespace=?1 AND id=?2 AND deleted_at IS NULL"
686                            .into(),
687                        params: vec![SqlValue::Text(namespace), SqlValue::Text(id)],
688                        label: Some("entity-index-revision".into()),
689                    })
690                    .await?;
691                if !matches!(current, Some(SqlValue::Integer(current)) if current == version) {
692                    return Ok(Box::new(false) as Box<dyn Any + Send>);
693                }
694                for statement in statements {
695                    writer.execute(statement).await?;
696                }
697                Ok(Box::new(true) as Box<dyn Any + Send>)
698            })
699        });
700        self.sql()
701            .atomic_unit(op)
702            .await?
703            .downcast::<bool>()
704            .map(|result| *result)
705            .map_err(|_| RuntimeError::Internal("invalid entity index outcome".into()))
706    }
707
708    pub(crate) fn entity_vector_insert_statements(
709        table: &str,
710        entity: &Entity,
711        model_name: &str,
712        vector: &[f32],
713    ) -> Vec<SqlStatement> {
714        let subject = entity.id.to_string();
715        let model_key = table
716            .strip_prefix("vec_")
717            .expect("runtime vector tables use the vec_ prefix");
718        let kind = SubstrateKind::Entity.to_string();
719        let field = "entity.body";
720        let blob = khive_storage::encode_f32_native(vector);
721        vec![
722            SqlStatement {
723                sql: format!(
724                    "INSERT INTO ann_write_log \
725                     (namespace, embedding_model, kind, field, subject_id, op) \
726                     SELECT namespace, embedding_model, kind, field, subject_id, 'delete' \
727                     FROM {table} WHERE subject_id=?1 AND NOT \
728                     (namespace=?2 AND embedding_model=?3 AND kind=?4 AND field=?5)"
729                ),
730                params: vec![
731                    SqlValue::Text(subject.clone()),
732                    SqlValue::Text(entity.namespace.clone()),
733                    SqlValue::Text(model_name.to_string()),
734                    SqlValue::Text(kind.clone()),
735                    SqlValue::Text(field.into()),
736                ],
737                label: Some("entity-reindex-log-delete".into()),
738            },
739            SqlStatement {
740                sql: format!("DELETE FROM {table} WHERE subject_id=?1"),
741                params: vec![SqlValue::Text(subject.clone())],
742                label: Some("entity-reindex-vector-delete".into()),
743            },
744            // This raw replacement cannot attest the embedded input. Clear the old
745            // sidecar in the same atomic index revision even when the new BLOB is
746            // byte-identical to the old one.
747            SqlStatement {
748                sql: "DELETE FROM vector_provenance \
749                      WHERE model_key = ?1 AND subject_id = ?2"
750                    .into(),
751                params: vec![
752                    SqlValue::Text(model_key.to_string()),
753                    SqlValue::Text(subject.clone()),
754                ],
755                label: Some("entity-reindex-provenance-clear".into()),
756            },
757            SqlStatement {
758                sql: format!(
759                    "INSERT INTO {table} \
760                     (subject_id, namespace, kind, field, embedding_model, embedding) \
761                     VALUES (?1, ?2, ?3, ?4, ?5, ?6)"
762                ),
763                params: vec![
764                    SqlValue::Text(subject.clone()),
765                    SqlValue::Text(entity.namespace.clone()),
766                    SqlValue::Text(kind.clone()),
767                    SqlValue::Text(field.into()),
768                    SqlValue::Text(model_name.to_string()),
769                    SqlValue::Blob(blob),
770                ],
771                label: Some("entity-reindex-vector-insert".into()),
772            },
773            SqlStatement {
774                sql: "INSERT INTO ann_write_log \
775                      (namespace, embedding_model, kind, field, subject_id, op) \
776                      VALUES (?1, ?2, ?3, ?4, ?5, 'upsert')"
777                    .into(),
778                params: vec![
779                    SqlValue::Text(entity.namespace.clone()),
780                    SqlValue::Text(model_name.to_string()),
781                    SqlValue::Text(kind),
782                    SqlValue::Text(field.into()),
783                    SqlValue::Text(subject),
784                ],
785                label: Some("entity-reindex-log-upsert".into()),
786            },
787        ]
788    }
789
790    pub(crate) async fn publish_entity_vector_revision(
791        &self,
792        token: &NamespaceToken,
793        entity: &Entity,
794        model_name: &str,
795        vector: &[f32],
796    ) -> RuntimeResult<bool> {
797        self.vectors_for_model(token, model_name)?;
798        let (storage_model, dimensions) = self.vector_model_metadata(model_name)?;
799        if let Some(index) = vector.iter().position(|value| !value.is_finite()) {
800            return Err(RuntimeError::InvalidInput(format!(
801                "non-finite entity vector at index {index}"
802            )));
803        }
804        if vector.len() != dimensions {
805            return Err(RuntimeError::InvalidInput(format!(
806                "entity vector has {} dimensions; expected {dimensions}",
807                vector.len()
808            )));
809        }
810        let table = format!("vec_{}", crate::config::sanitize_key(&storage_model));
811        let statements =
812            Self::entity_vector_insert_statements(&table, entity, &storage_model, vector);
813        #[cfg(test)]
814        race_seam::pause_before_entity_vector_publish().await;
815        self.apply_entity_index_revision(entity, statements).await
816    }
817
818    /// Re-upsert FTS5 document and vector(s) for the entity across all registered models.
819    ///
820    /// Uses `entity.namespace` — the authoritative namespace stored on the record — rather
821    /// than the caller-supplied `namespace` parameter. This prevents a cross-namespace
822    /// reindex from writing the search document into the wrong namespace's FTS index.
823    ///
824    /// Best-effort for vectors: if embedding or inserting for a particular model fails,
825    /// logs a warning and continues to the next model. The FTS step is fail-closed
826    /// (propagates error). Callers (update_entity, merge_entity) have already committed
827    /// the entity row, so a partial embed miss leaves a stale vector rather than
828    /// rolling back the update. Each failed guarded vector replacement rolls back
829    /// its own index writes, keeping the prior row searchable.
830    pub(crate) async fn reindex_entity(
831        &self,
832        token: &NamespaceToken,
833        entity: &Entity,
834    ) -> RuntimeResult<crate::retrieval::EmbeddingTruncationReport> {
835        let embedding_plan = EmbeddingModelPlan::capture(self);
836        self.reindex_entity_with_plan(token, entity, &embedding_plan, None)
837            .await
838    }
839
840    pub(crate) async fn reindex_entity_with_precomputed(
841        &self,
842        token: &NamespaceToken,
843        entity: &Entity,
844        mut precomputed: HashMap<String, crate::retrieval::DocumentEmbeddingOutcome>,
845    ) -> RuntimeResult<crate::retrieval::EmbeddingTruncationReport> {
846        let embedding_plan = EmbeddingModelPlan::capture(self);
847        self.reindex_entity_with_plan(token, entity, &embedding_plan, Some(&mut precomputed))
848            .await
849    }
850
851    pub(super) async fn reindex_entity_with_plan(
852        &self,
853        token: &NamespaceToken,
854        entity: &Entity,
855        embedding_plan: &EmbeddingModelPlan,
856        mut precomputed: Option<&mut HashMap<String, crate::retrieval::DocumentEmbeddingOutcome>>,
857    ) -> RuntimeResult<crate::retrieval::EmbeddingTruncationReport> {
858        // Test-only fault seam: force the post-commit FTS leg to fail after a
859        // merge or update has already persisted its entity row.
860        #[cfg(test)]
861        if crate::operations::consume_fts_fail_fault(&entity.namespace) {
862            return Err(RuntimeError::Internal("injected FTS failure".to_string()));
863        }
864        // Use entity.namespace (authoritative) rather than token.namespace().as_str() (caller claim).
865        let doc = entity_fts_document(entity);
866        let embed_body = doc.body.clone();
867        let _ = self.text(token)?;
868        #[cfg(test)]
869        race_seam::pause_before_entity_index_publish().await;
870        let statements = khive_db::stores::text::delete_document_statements(
871            "fts_entities",
872            &entity.namespace,
873            entity.id,
874        )
875        .into_iter()
876        .chain(khive_db::stores::text::insert_document_statements(
877            "fts_entities",
878            &doc,
879        ))
880        .collect();
881        if !self.apply_entity_index_revision(entity, statements).await? {
882            return Ok(crate::retrieval::EmbeddingTruncationReport::default());
883        }
884
885        let mut report = crate::retrieval::EmbeddingTruncationReport::default();
886        for model_name in embedding_plan.model_names() {
887            let embedding = match precomputed
888                .as_mut()
889                .and_then(|outcomes| outcomes.remove(model_name))
890            {
891                Some(outcome) => Ok(outcome),
892                None => {
893                    self.embed_document_with_model_outcome_for_token(token, model_name, &embed_body)
894                        .await
895                }
896            };
897            match embedding {
898                Ok(outcome) => {
899                    report.observe(&outcome);
900                    match self
901                        .publish_entity_vector_revision(token, entity, model_name, &outcome.vector)
902                        .await
903                    {
904                        Ok(true) => {}
905                        Ok(false) => break,
906                        Err(error) => tracing::warn!(
907                            model = model_name,
908                            id = %entity.id,
909                            "reindex_entity: vector insert failed, skipping model: {error}"
910                        ),
911                    }
912                }
913                Err(e) => {
914                    tracing::warn!(
915                        model = model_name,
916                        id = %entity.id,
917                        "reindex_entity: embed failed for model, skipping: {e}"
918                    );
919                }
920            }
921        }
922
923        Ok(report)
924    }
925}