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 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 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 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 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 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 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 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 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 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 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 #[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 #[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 let embedding_plan = EmbeddingModelPlan::capture(self);
529 let vec_tables = embedding_plan.vector_tables();
530 let pack_rules = self.pack_edge_rules();
533
534 let _ = self.entities(token)?;
536 let _ = self.graph(token)?;
537 let _ = self.text(token)?;
538 let _ = self.events(token)?;
539 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 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 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 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 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 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 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 #[cfg(test)]
861 if crate::operations::consume_fts_fail_fault(&entity.namespace) {
862 return Err(RuntimeError::Internal("injected FTS failure".to_string()));
863 }
864 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}