1use serde_json::Value;
20use uuid::Uuid;
21
22use khive_storage::types::SqlValue;
23use khive_storage::{AttachmentSubstrate, EdgeRelation, EdgeUpsertDisposition, SqlStatement};
24use khive_types::{EventKind, SubstrateKind};
25
26use crate::atomic_plan::{
27 AddEntityPlan, AddNotePlan, AffectedRowGuard, DeletePlan, EdgeNaturalKey, LinkPlan, MergePlan,
28 PlanStatement, PostCommitEffect, UpdatePlan,
29};
30use crate::atomic_runner::AtomicOpPlan;
31use crate::atomic_runner::CommittedPostCommitEffects;
32use crate::curation::{entity_fts_document, note_fts_document};
33use crate::error::{RuntimeError, RuntimeResult};
34use crate::operations::{
35 canonical_edge_endpoints, merge_dependency_kind, validate_edge_metadata, validate_edge_weight,
36 Resolved,
37};
38use crate::runtime::{KhiveRuntime, NamespaceToken};
39
40use khive_db::stores::attachment::delete_record_attachments_statement;
41use khive_db::stores::entity::{
42 entity_hard_delete_statement, entity_replace_if_unchanged_statement,
43 entity_soft_delete_statement, entity_upsert_statement,
44};
45use khive_db::stores::event::event_insert_statements;
46use khive_db::stores::event::hard_delete_lineage_warning_statements;
47use khive_db::stores::graph::{
48 edge_hard_delete_statement, edge_insert_new_guarded_by_endpoints_statement,
49 edge_link_replace_if_unchanged_and_endpoints_exist_statement,
50 edge_replace_if_unchanged_statement, edge_soft_delete_statement,
51 edge_symmetric_absorb_or_update_inplace_statement, edge_symmetric_delete_if_conflict_statement,
52 purge_incident_edges_statement,
53};
54use khive_db::stores::note::{
55 note_hard_delete_statement, note_soft_delete_statement, note_upsert_statement,
56};
57use khive_db::stores::text::{delete_document_statements, insert_document_statements};
58
59fn obj(args: &Value) -> RuntimeResult<&serde_json::Map<String, Value>> {
64 args.as_object()
65 .ok_or_else(|| RuntimeError::InvalidInput("op args must be a JSON object".into()))
66}
67
68fn require_str<'a>(args: &'a Value, key: &str) -> RuntimeResult<&'a str> {
69 obj(args)?
70 .get(key)
71 .and_then(|v| v.as_str())
72 .ok_or_else(|| RuntimeError::InvalidInput(format!("missing required field {key:?}")))
73}
74
75fn require_uuid(args: &Value, key: &str) -> RuntimeResult<Uuid> {
88 let raw = require_str(args, key)?;
89 Uuid::parse_str(raw).map_err(|_| {
90 RuntimeError::InvalidInput(format!(
91 "{key} must be a full UUID; got {raw:?}. This atomic-plan stage consumes an \
92 already-resolved record and performs no namespace-scoped search of its own, so a \
93 short hex prefix cannot be resolved here — resolve it to a full UUID first (e.g. \
94 via `get`) and pass that."
95 ))
96 })
97}
98
99fn optional_str<'a>(args: &'a Value, key: &str) -> Option<&'a str> {
100 obj(args).ok()?.get(key).and_then(|v| v.as_str())
101}
102
103fn optional_create_string(args: &Value, key: &str) -> RuntimeResult<Option<String>> {
104 match obj(args)?.get(key) {
105 None | Some(Value::Null) => Ok(None),
106 Some(Value::String(value)) => Ok(Some(value.clone())),
107 Some(other) => Err(RuntimeError::InvalidInput(format!(
108 "{key} must be a string or null, got: {other}"
109 ))),
110 }
111}
112
113fn optional_entity_type_patch(args: &Value, key: &str) -> RuntimeResult<Option<Option<String>>> {
120 match obj(args)?.get(key) {
121 None => Ok(None),
122 Some(Value::Null) => Ok(Some(None)),
123 Some(Value::String(value)) => Ok(Some(Some(value.clone()))),
124 Some(other) => Err(RuntimeError::InvalidInput(format!(
125 "{key} must be a string or null, got: {other}"
126 ))),
127 }
128}
129
130fn optional_string_patch(args: &Value, key: &str) -> RuntimeResult<Option<Option<String>>> {
143 match obj(args)?.get(key) {
144 None | Some(Value::Null) => Ok(None),
145 Some(Value::String(s)) => Ok(Some(Some(s.clone()))),
146 Some(other) => Err(RuntimeError::InvalidInput(format!(
147 "{key} must be a string or null, got: {other}"
148 ))),
149 }
150}
151
152fn entity_name_patch(args: &Value) -> RuntimeResult<Option<String>> {
163 match obj(args)?.get("name") {
164 None | Some(Value::Null) => Ok(None),
165 Some(Value::String(s)) => Ok(Some(s.clone())),
166 Some(other) => Err(RuntimeError::InvalidInput(format!(
167 "name must be a string, got: {other}"
168 ))),
169 }
170}
171
172fn optional_properties(args: &Value, key: &str) -> RuntimeResult<Option<Value>> {
181 match obj(args)?.get(key) {
182 None | Some(Value::Null) => Ok(None),
183 Some(v) => Ok(Some(v.clone())),
184 }
185}
186
187fn optional_tags(args: &Value) -> RuntimeResult<Option<Vec<String>>> {
194 match obj(args)?.get("tags") {
195 None | Some(Value::Null) => Ok(None),
196 Some(Value::Array(items)) => {
197 let mut tags = Vec::with_capacity(items.len());
198 for item in items {
199 let s = item.as_str().ok_or_else(|| {
200 RuntimeError::InvalidInput("tags must be an array of strings".into())
201 })?;
202 tags.push(s.to_string());
203 }
204 Ok(Some(tags))
205 }
206 Some(_) => Err(RuntimeError::InvalidInput(
207 "tags must be an array of strings".into(),
208 )),
209 }
210}
211
212fn optional_f64(args: &Value, key: &str) -> RuntimeResult<Option<f64>> {
213 match obj(args)?.get(key) {
214 None => Ok(None),
215 Some(Value::Null) => Ok(None),
216 Some(v) => v
217 .as_f64()
218 .map(Some)
219 .ok_or_else(|| RuntimeError::InvalidInput(format!("{key} must be a number"))),
220 }
221}
222
223fn optional_f64_patch(args: &Value, key: &str) -> RuntimeResult<Option<Option<f64>>> {
229 match obj(args)?.get(key) {
230 None => Ok(None),
231 Some(Value::Null) => Ok(Some(None)),
232 Some(v) => v
233 .as_f64()
234 .map(|f| Some(Some(f)))
235 .ok_or_else(|| RuntimeError::InvalidInput(format!("{key} must be a number"))),
236 }
237}
238
239fn vector_table_names(runtime: &KhiveRuntime) -> Vec<String> {
244 runtime
245 .registered_embedding_model_names()
246 .iter()
247 .map(|name| format!("vec_{}", crate::config::sanitize_key(name)))
248 .collect()
249}
250
251fn purge_index_row_statement(
261 table: &str,
262 namespace: &str,
263 subject_id: Uuid,
264 label: &str,
265) -> PlanStatement {
266 PlanStatement {
267 statement: SqlStatement {
268 sql: format!("DELETE FROM {table} WHERE namespace = ?1 AND subject_id = ?2"),
269 params: vec![
270 SqlValue::Text(namespace.to_string()),
271 SqlValue::Text(subject_id.to_string()),
272 ],
273 label: Some(label.to_string()),
274 },
275 guard: None,
276 }
277}
278
279fn purge_fts_document_statements(
286 fts_table: &str,
287 namespace: &str,
288 subject_id: Uuid,
289 label_prefix: &str,
290) -> [PlanStatement; 2] {
291 let [mut fts_stmt, mut map_stmt] = delete_document_statements(fts_table, namespace, subject_id);
292 fts_stmt.label = Some(label_prefix.to_string());
293 map_stmt.label = Some(format!("{label_prefix}-map"));
294 [
295 PlanStatement {
296 statement: fts_stmt,
297 guard: None,
298 },
299 PlanStatement {
300 statement: map_stmt,
301 guard: None,
302 },
303 ]
304}
305
306fn log_vector_row_delete_statement(
307 table: &str,
308 namespace: &str,
309 subject_id: Uuid,
310 label: &str,
311) -> PlanStatement {
312 PlanStatement {
313 statement: SqlStatement {
314 sql: format!(
315 "INSERT INTO ann_write_log \
316 (namespace, embedding_model, kind, field, subject_id, op) \
317 SELECT namespace, embedding_model, kind, field, subject_id, 'delete' \
318 FROM {table} WHERE namespace = ?1 AND subject_id = ?2"
319 ),
320 params: vec![
321 SqlValue::Text(namespace.to_string()),
322 SqlValue::Text(subject_id.to_string()),
323 ],
324 label: Some(label.to_string()),
325 },
326 guard: None,
327 }
328}
329
330async fn vector_table_exists(runtime: &KhiveRuntime, table: &str) -> RuntimeResult<bool> {
335 let mut reader = runtime
336 .sql()
337 .reader()
338 .await
339 .map_err(RuntimeError::Storage)?;
340 let row = reader
341 .query_scalar(SqlStatement {
342 sql: "SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = ?1".to_string(),
343 params: vec![SqlValue::Text(table.to_string())],
344 label: Some("atomic-delete-vec-table-exists".to_string()),
345 })
346 .await
347 .map_err(RuntimeError::Storage)?;
348 Ok(row.is_some())
349}
350
351async fn push_index_purge_statements(
367 runtime: &KhiveRuntime,
368 statements: &mut Vec<PlanStatement>,
369 fts_table: &str,
370 namespace: &str,
371 subject_id: Uuid,
372 label_prefix: &str,
373) -> RuntimeResult<()> {
374 statements.extend(purge_fts_document_statements(
375 fts_table,
376 namespace,
377 subject_id,
378 &format!("{label_prefix}-purge-fts"),
379 ));
380 for vec_table in vector_table_names(runtime) {
381 if vector_table_exists(runtime, &vec_table).await? {
382 statements.push(log_vector_row_delete_statement(
383 &vec_table,
384 namespace,
385 subject_id,
386 &format!("{label_prefix}-log-delete-vec-{vec_table}"),
387 ));
388 statements.push(purge_index_row_statement(
389 &vec_table,
390 namespace,
391 subject_id,
392 &format!("{label_prefix}-purge-vec-{vec_table}"),
393 ));
394 }
395 }
396 Ok(())
397}
398
399pub(crate) fn event_append_statements(
416 token: &NamespaceToken,
417 namespace: &str,
418 verb: &str,
419 kind: EventKind,
420 substrate: SubstrateKind,
421 target_id: Uuid,
422 payload: Value,
423) -> RuntimeResult<Vec<PlanStatement>> {
424 let record_token = token
425 .with_namespace(crate::Namespace::parse(namespace).map_err(|error| {
426 RuntimeError::Internal(format!("event namespace invalid: {error}"))
427 })?);
428 let event = crate::EventAttribution::from_token(&record_token).stamp(
429 khive_storage::event::Event::new(namespace.to_string(), verb, kind, substrate, "")
430 .with_target(target_id)
431 .with_payload(payload),
432 );
433 let statements = event_insert_statements(&event)
434 .map_err(|e| RuntimeError::Internal(format!("event_insert_statements: {e}")))?;
435 Ok(statements
436 .into_iter()
437 .map(|statement| PlanStatement {
438 statement,
439 guard: None,
440 })
441 .collect())
442}
443
444pub async fn prepare_op(
454 runtime: &KhiveRuntime,
455 token: &NamespaceToken,
456 tool: &str,
457 args: &Value,
458) -> RuntimeResult<AtomicOpPlan> {
459 match tool {
460 "update" => prepare_update(runtime, token, args, None).await,
469 "delete" => prepare_delete(runtime, token, args, None).await,
478 "link" => prepare_link(runtime, token, args).await,
479 "merge" => prepare_merge(runtime, token, args).await,
480 "propose" | "review" | "withdraw" => prepare_governance_unimplemented(tool),
481 other => Err(RuntimeError::InvalidInput(format!(
482 "{other:?} has no atomic_prepare::prepare_op implementation; the CLI \
483 admissibility check should have rejected this before prepare"
484 ))),
485 }
486}
487
488fn prepare_governance_unimplemented(tool: &str) -> RuntimeResult<AtomicOpPlan> {
489 Err(RuntimeError::InvalidInput(format!(
490 "{tool:?} is on the ADR-099 v1 admissible verb list but has no --atomic \
491 prepare/apply implementation yet: its lifecycle (ADR-046) is an \
492 event-sourced changeset-interpreter over a dedicated `proposals_open` \
493 table, not a small guarded-DML plan — a faithful non-stub atomic \
494 prepare for it is tracked as ADR-099 follow-up work, not implemented \
495 in slice B3. No write was attempted."
496 )))
497}
498
499pub async fn prepare_add_entity(
509 runtime: &KhiveRuntime,
510 token: &NamespaceToken,
511 args: &Value,
512) -> RuntimeResult<AtomicOpPlan> {
513 let kind = require_str(args, "kind")?;
514 let name = require_str(args, "name")?;
515 runtime.validate_entity_kind(kind)?;
516 if name.trim().is_empty() {
517 return Err(RuntimeError::InvalidInput(
518 "name must not be empty".to_string(),
519 ));
520 }
521
522 let description = optional_create_string(args, "description")?;
523 let properties = optional_properties(args, "properties")?;
524 let tags = optional_tags(args)?.unwrap_or_default();
525
526 crate::secret_gate::check_at(name, "entity", "name")?;
527 if let Some(ref d) = description {
528 crate::secret_gate::check_at(d, "entity", "description")?;
529 }
530 if let Some(ref p) = properties {
531 crate::secret_gate::check_json_at(p, "entity", "properties")?;
532 }
533 crate::secret_gate::check_tags_at(&tags, "entity", "tags")?;
534 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
535
536 let ns = token.namespace().as_str();
537 let mut entity = khive_storage::Entity::new(ns, kind, name);
538 if let Some(d) = description {
539 entity = entity.with_description(d);
540 }
541 if let Some(p) = properties {
542 entity = entity.with_properties(p);
543 }
544 if !tags.is_empty() {
545 entity = entity.with_tags(tags);
546 }
547
548 let mut statements = vec![PlanStatement {
549 statement: entity_upsert_statement(&entity),
550 guard: Some(AffectedRowGuard::exactly(1)),
551 }];
552 for statement in insert_document_statements("fts_entities", &entity_fts_document(&entity)) {
556 statements.push(PlanStatement {
557 statement,
558 guard: None,
559 });
560 }
561
562 Ok(AtomicOpPlan::AddEntity(AddEntityPlan {
563 entity_id: entity.id,
564 statements,
565 post_commit: PostCommitEffect::ReindexEntity {
566 entity_id: entity.id,
567 },
568 }))
569}
570
571pub async fn prepare_add_note(
577 runtime: &KhiveRuntime,
578 token: &NamespaceToken,
579 args: &Value,
580) -> RuntimeResult<AtomicOpPlan> {
581 let kind = require_str(args, "kind")?;
582 let content = require_str(args, "content")?;
583 runtime.validate_note_kind(kind)?;
584
585 let name = optional_create_string(args, "name")?;
586 let properties = optional_properties(args, "properties")?;
587 let properties = runtime.derive_note_write_properties(kind, token, properties)?;
593
594 crate::secret_gate::check_at(content, "note", "content")?;
595 if let Some(ref n) = name {
596 crate::secret_gate::check_at(n, "note", "name")?;
597 }
598 if let Some(ref p) = properties {
599 crate::secret_gate::check_json_at(p, "note", "properties")?;
600 }
601 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
602
603 let ns = token.namespace().as_str();
604 let mut note = khive_storage::note::Note::new(ns, kind, content);
605 if let Some(n) = name {
606 note = note.with_name(n);
607 }
608 if let Some(p) = properties {
609 note = note.with_properties(p);
610 }
611
612 let mut statements = vec![PlanStatement {
613 statement: note_upsert_statement(¬e),
614 guard: Some(AffectedRowGuard::exactly(1)),
615 }];
616 for statement in insert_document_statements("fts_notes", ¬e_fts_document(¬e)) {
619 statements.push(PlanStatement {
620 statement,
621 guard: None,
622 });
623 }
624
625 Ok(AtomicOpPlan::AddNote(Box::new(AddNotePlan {
626 note_guard: None,
627 note_id: note.id,
628 statements,
629 post_commit: PostCommitEffect::ReindexNote {
630 note_id: note.id,
631 version: note.version,
632 },
633 })))
634}
635
636fn reject_inapplicable_update_fields(args: &Value, substrate: &str) -> RuntimeResult<()> {
651 let o = obj(args)?;
652 if substrate == "edge" && o.get("expected_version").is_some_and(|v| !v.is_null()) {
653 return Err(RuntimeError::InvalidInput(
654 "expected_version applies only to entities and notes".into(),
655 ));
656 }
657 if substrate != "note"
658 && ["embed", "fence"]
659 .iter()
660 .any(|field| o.contains_key(*field))
661 {
662 return Err(RuntimeError::InvalidInput(
663 "embed and fence apply only to notes".into(),
664 ));
665 }
666 let present = |k: &str| o.get(k).is_some_and(|v| !v.is_null());
667 let (bad_field, valid): (Option<&str>, &str) = match substrate {
668 "entity" => {
669 let bad = if present("content") {
670 Some("content")
671 } else if present("salience") {
672 Some("salience")
673 } else if present("decay_factor") {
674 Some("decay_factor")
675 } else if present("relation") {
676 Some("relation")
677 } else if present("weight") {
678 Some("weight")
679 } else {
680 None
681 };
682 (bad, "name, description, tags, properties, entity_type")
683 }
684 "note" => {
685 let bad = if present("description") {
686 Some("description")
687 } else if present("relation") {
688 Some("relation")
689 } else if present("weight") {
690 Some("weight")
691 } else if o.contains_key("entity_type") {
692 Some("entity_type")
695 } else {
696 None
697 };
698 (
699 bad,
700 "name, content, salience, decay_factor, properties, tags",
701 )
702 }
703 "edge" => {
709 let bad = if present("name") {
710 Some("name")
711 } else if present("description") {
712 Some("description")
713 } else if present("content") {
714 Some("content")
715 } else if present("tags") {
716 Some("tags")
717 } else if present("salience") {
718 Some("salience")
719 } else if present("decay_factor") {
720 Some("decay_factor")
721 } else if o.contains_key("entity_type") {
722 Some("entity_type")
725 } else {
726 None
727 };
728 (bad, "relation, weight, properties")
729 }
730 _ => (None, ""),
731 };
732 if let Some(field) = bad_field {
733 let substrate_label = match substrate {
734 "entity" => "an entity",
735 "note" => "a note",
736 "edge" => "an edge",
737 other => other,
738 };
739 return Err(RuntimeError::InvalidInput(format!(
740 "field '{field}' is not valid for {substrate_label}; valid fields: {valid}"
741 )));
742 }
743 Ok(())
744}
745
746pub enum AtomicUpdateKind {
756 Entity { specific: Option<String> },
757 Note { specific: Option<String> },
758 Edge,
759}
760
761pub fn validate_note_update_expected_kind(
766 note: &khive_storage::note::Note,
767 expected_kind: &Option<AtomicUpdateKind>,
768) -> RuntimeResult<()> {
769 let id = note.id;
770 match expected_kind {
771 None => Ok(()),
772 Some(AtomicUpdateKind::Note {
773 specific: Some(expected),
774 }) if ¬e.kind != expected => Err(RuntimeError::NotFound(format!("note {id}"))),
775 Some(AtomicUpdateKind::Note { .. }) => Ok(()),
776 Some(AtomicUpdateKind::Entity { .. }) => {
777 Err(RuntimeError::NotFound(format!("entity {id}")))
778 }
779 Some(AtomicUpdateKind::Edge) => Err(RuntimeError::NotFound(format!("edge {id}"))),
780 }
781}
782
783async fn prepare_note_update_plan_from_snapshot(
784 runtime: &KhiveRuntime,
785 token: &NamespaceToken,
786 args: &Value,
787 expected_kind: &Option<AtomicUpdateKind>,
788 note: khive_storage::note::Note,
789 policy: crate::NoteUpdatePolicy,
790 registry: Option<&crate::VerbRegistry>,
791) -> RuntimeResult<(khive_storage::Note, UpdatePlan)> {
792 let id = require_uuid(args, "id")?;
793 if note.id != id {
794 return Err(RuntimeError::NotFound(format!("note {id}")));
795 }
796 validate_note_update_expected_kind(¬e, expected_kind)?;
797
798 reject_inapplicable_update_fields(args, "note")?;
799 let mut normalized_args = args.clone();
800 crate::curation::normalize_note_update_tags(&mut normalized_args)?;
801 let args = &normalized_args;
802 let name = optional_string_patch(args, "name")?;
803 let content = optional_str(args, "content").map(str::to_string);
804 let properties = optional_properties(args, "properties")?;
805 let salience = optional_f64_patch(args, "salience")?;
806 let decay_factor = optional_f64_patch(args, "decay_factor")?;
807 let options = crate::note_write::NoteWriteOptions {
808 expected_version: obj(args)?
809 .get("expected_version")
810 .filter(|v| !v.is_null())
811 .map(|v| {
812 v.as_i64().ok_or_else(|| {
813 RuntimeError::InvalidInput("expected_version must be an integer".into())
814 })
815 })
816 .transpose()?,
817 fence: obj(args)?
818 .get("fence")
819 .map(|v| {
820 serde_json::from_value(v.clone())
821 .map_err(|error| RuntimeError::InvalidInput(format!("invalid fence: {error}")))
822 })
823 .transpose()?,
824 embed: obj(args)?
825 .get("embed")
826 .filter(|v| !v.is_null())
827 .map(|v| {
828 v.as_bool()
829 .ok_or_else(|| RuntimeError::InvalidInput("embed must be boolean".into()))
830 })
831 .transpose()?,
832 key: None,
833 };
834 let patch = crate::curation::NotePatch::new(name, content, salience, decay_factor, properties)
835 .with_update_policy(policy)
836 .with_write_options(options);
837 let (updated, mut plan) = runtime
838 .prepare_versioned_note_update(token, note.clone(), patch.clone())
839 .await?;
840 if let Some(registry) = registry {
841 attach_note_update_effects(runtime, token, registry, ¬e, &patch, &mut plan).await?;
842 }
843 Ok((updated, plan))
844}
845
846pub async fn prepare_update_from_note_snapshot(
856 runtime: &KhiveRuntime,
857 token: &NamespaceToken,
858 args: &Value,
859 expected_kind: Option<AtomicUpdateKind>,
860 note: khive_storage::note::Note,
861 policy: crate::NoteUpdatePolicy,
862 registry: &crate::VerbRegistry,
863) -> RuntimeResult<(khive_storage::Note, AtomicOpPlan)> {
864 if obj(args)?.get("entity_kind").is_some_and(|v| !v.is_null()) {
865 return Err(RuntimeError::InvalidInput(
866 "entity_kind is immutable; to change kind, delete then re-create the entity, \
867 or use merge() if this is a deduplication correction"
868 .into(),
869 ));
870 }
871 let (note, plan) = prepare_note_update_plan_from_snapshot(
872 runtime,
873 token,
874 args,
875 &expected_kind,
876 note,
877 policy,
878 Some(registry),
879 )
880 .await?;
881 Ok((note, AtomicOpPlan::Update(Box::new(plan))))
882}
883
884async fn attach_note_update_effects(
888 runtime: &KhiveRuntime,
889 token: &NamespaceToken,
890 registry: &crate::VerbRegistry,
891 snapshot: &khive_storage::Note,
892 patch: &crate::curation::NotePatch,
893 plan: &mut UpdatePlan,
894) -> RuntimeResult<()> {
895 use crate::atomic_plan::NoteUpdateStatement;
896 use crate::NoteUpdateEffect;
897
898 let Some(hook) = registry.find_kind_hook(&snapshot.kind) else {
899 return Ok(());
900 };
901 if let Some(expected) = patch.write_options.expected_version {
904 if expected != snapshot.version {
905 return Err(crate::note_write::NoteWriteConflict::Version {
906 expected,
907 current: snapshot.version,
908 }
909 .into_error()
910 .into());
911 }
912 }
913 let current = runtime.notes(token)?.get_note(snapshot.id).await?;
914 if !current.is_some_and(|current| {
915 current.updated_at == snapshot.updated_at
916 && current.deleted_at == snapshot.deleted_at
917 && current.version == snapshot.version
918 }) {
919 return Err(crate::curation::stale_note_snapshot_error(snapshot.id));
920 }
921 let effects = hook
922 .note_update_effects(runtime, token, snapshot, patch)
923 .await?;
924 if plan.idempotent_noop && !effects.is_empty() {
925 return Err(RuntimeError::InvalidInput(
926 "an unchanged note update cannot carry graph effects".into(),
927 ));
928 }
929 let edge_token = token.with_namespace(
930 crate::Namespace::parse(&snapshot.namespace)
931 .map_err(|error| RuntimeError::Internal(format!("invalid note namespace: {error}")))?,
932 );
933 for effect in effects {
934 match effect {
935 NoteUpdateEffect::Link(spec) => {
936 if spec.source_id != snapshot.id
937 || spec
938 .namespace
939 .as_deref()
940 .is_some_and(|ns| ns != snapshot.namespace)
941 {
942 return Err(RuntimeError::InvalidInput(
943 "note update links must originate from the note in its namespace".into(),
944 ));
945 }
946 let mut args = serde_json::json!({
947 "source_id": spec.source_id, "target_id": spec.target_id,
948 "relation": spec.relation, "weight": spec.weight,
949 "resurrect": spec.resurrect,
950 });
951 if let Some(metadata) = spec.metadata {
952 args["metadata"] = metadata;
953 }
954 let AtomicOpPlan::Link(link) = prepare_link(runtime, &edge_token, &args).await?
955 else {
956 return Err(RuntimeError::Internal("expected a link plan".into()));
957 };
958 if link.disposition == EdgeUpsertDisposition::Updated {
962 return Err(khive_types::KhiveError::conflict(
963 "a live edge appeared while preparing the note update; retry with fresh state",
964 ).into());
965 }
966 plan.graph_effects
967 .extend(link.statements.into_iter().map(NoteUpdateStatement::Write));
968 }
969 NoteUpdateEffect::DeleteEdge(edge) => {
970 let id = Uuid::from(edge.id);
971 if edge.source_id != snapshot.id
972 || edge.namespace != snapshot.namespace
973 || edge.deleted_at.is_some()
974 {
975 return Err(RuntimeError::InvalidInput(
976 "note update deletes must select an outgoing edge in the note namespace"
977 .into(),
978 ));
979 }
980 plan.graph_effects
981 .push(NoteUpdateStatement::Assert(PlanStatement {
982 statement: khive_db::stores::graph::edge_snapshot_assertion_statement(
983 &edge, false,
984 ),
985 guard: Some(AffectedRowGuard::exactly(1)),
986 }));
987 let actor = format!("{}:{}", token.actor().kind, token.actor().id);
988 let AtomicOpPlan::Delete(delete) =
989 prepare_delete_edge(&edge_token, id, edge, false, &actor).await?
990 else {
991 return Err(RuntimeError::Internal(
992 "expected an edge delete plan".into(),
993 ));
994 };
995 if delete.post_commit != PostCommitEffect::None {
996 return Err(RuntimeError::Internal(
997 "edge delete has a deferred effect".into(),
998 ));
999 }
1000 plan.graph_effects.extend(
1001 delete
1002 .statements
1003 .into_iter()
1004 .map(NoteUpdateStatement::Write),
1005 );
1006 }
1007 NoteUpdateEffect::AssertLink(edge) => {
1008 if edge.source_id != snapshot.id
1009 || edge.namespace != snapshot.namespace
1010 || edge.deleted_at.is_some()
1011 {
1012 return Err(RuntimeError::InvalidInput(
1013 "note update assertions must select a live outgoing edge in the note namespace".into(),
1014 ));
1015 }
1016 plan.graph_effects
1017 .push(NoteUpdateStatement::Assert(PlanStatement {
1018 statement: khive_db::stores::graph::edge_snapshot_assertion_statement(
1019 &edge, true,
1020 ),
1021 guard: Some(AffectedRowGuard::exactly(1)),
1022 }));
1023 }
1024 }
1025 }
1026 Ok(())
1027}
1028
1029pub async fn prepare_update(
1037 runtime: &KhiveRuntime,
1038 token: &NamespaceToken,
1039 args: &Value,
1040 expected_kind: Option<AtomicUpdateKind>,
1041) -> RuntimeResult<AtomicOpPlan> {
1042 let id = require_uuid(args, "id")?;
1043
1044 if obj(args)?.get("entity_kind").is_some_and(|v| !v.is_null()) {
1048 return Err(RuntimeError::InvalidInput(
1049 "entity_kind is immutable; to change kind, delete then re-create the entity, \
1050 or use merge() if this is a deduplication correction"
1051 .into(),
1052 ));
1053 }
1054
1055 match runtime.resolve_by_id(token, id).await? {
1056 Some(Resolved::Entity(entity)) => {
1057 match &expected_kind {
1058 None => {}
1059 Some(AtomicUpdateKind::Entity {
1060 specific: Some(expected),
1061 }) if &entity.kind != expected => {
1062 return Err(RuntimeError::NotFound(format!("entity {id}")));
1063 }
1064 Some(AtomicUpdateKind::Entity { .. }) => {}
1065 Some(AtomicUpdateKind::Note { .. }) => {
1066 return Err(RuntimeError::NotFound(format!("note {id}")));
1067 }
1068 Some(AtomicUpdateKind::Edge) => {
1069 return Err(RuntimeError::NotFound(format!("edge {id}")));
1070 }
1071 }
1072 reject_inapplicable_update_fields(args, "entity")?;
1080 let name = entity_name_patch(args)?;
1081 let description = optional_string_patch(args, "description")?;
1082 let properties = optional_properties(args, "properties")?;
1083 let tags = optional_tags(args)?;
1084 let entity_type = optional_entity_type_patch(args, "entity_type")?;
1085
1086 let expected_version = obj(args)?
1087 .get("expected_version")
1088 .filter(|v| !v.is_null())
1089 .map(|value| {
1090 value.as_i64().ok_or_else(|| {
1091 RuntimeError::InvalidInput("expected_version must be an integer".into())
1092 })
1093 })
1094 .transpose()?;
1095 prepare_update_entity_plan_with_version(
1096 runtime,
1097 token,
1098 id,
1099 crate::curation::EntityPatch {
1100 name,
1101 description,
1102 properties,
1103 tags,
1104 entity_type,
1105 },
1106 expected_version,
1107 )
1108 .await
1109 }
1110 Some(Resolved::Note(note)) => {
1111 let (_, plan) = prepare_note_update_plan_from_snapshot(
1116 runtime,
1117 token,
1118 args,
1119 &expected_kind,
1120 note,
1121 crate::NoteUpdatePolicy::default(),
1122 None,
1123 )
1124 .await?;
1125 Ok(AtomicOpPlan::Update(Box::new(plan)))
1126 }
1127 Some(_) => Err(RuntimeError::InvalidInput(format!(
1128 "update target {id} must be an entity, note, or edge"
1129 ))),
1130 None => match &expected_kind {
1137 Some(AtomicUpdateKind::Entity { .. }) => {
1138 Err(RuntimeError::NotFound(format!("entity/note {id}")))
1139 }
1140 Some(AtomicUpdateKind::Note { .. }) => {
1141 Err(RuntimeError::NotFound(format!("entity/note {id}")))
1142 }
1143 Some(AtomicUpdateKind::Edge) | None => match runtime.get_edge(token, id).await? {
1144 Some(edge) => prepare_update_edge(runtime, token, id, edge, args).await,
1145 None => Err(RuntimeError::NotFound(format!("entity/note/edge {id}"))),
1146 },
1147 },
1148 }
1149}
1150
1151pub async fn prepare_update_entity_plan(
1155 runtime: &KhiveRuntime,
1156 token: &NamespaceToken,
1157 id: Uuid,
1158 patch: crate::curation::EntityPatch,
1159) -> RuntimeResult<AtomicOpPlan> {
1160 prepare_update_entity_plan_with_version(runtime, token, id, patch, None).await
1161}
1162
1163pub(crate) async fn prepare_update_entity_plan_with_version(
1164 runtime: &KhiveRuntime,
1165 token: &NamespaceToken,
1166 id: Uuid,
1167 patch: crate::curation::EntityPatch,
1168 expected_version: Option<i64>,
1169) -> RuntimeResult<AtomicOpPlan> {
1170 crate::entity_write::validate_expected_version(expected_version)?;
1171 let (entity, reindex_required, changed_fields, expected_updated_at, expected_deleted_at) =
1172 runtime.prepare_update_entity(token, id, patch).await?;
1173 let mut statements = vec![PlanStatement {
1174 statement: entity_replace_if_unchanged_statement(
1175 &entity,
1176 expected_updated_at,
1177 expected_deleted_at,
1178 ),
1179 guard: Some(AffectedRowGuard::exactly(1)),
1180 }];
1181 statements.extend(event_append_statements(
1182 token,
1183 &entity.namespace,
1184 "update",
1185 EventKind::EntityUpdated,
1186 SubstrateKind::Entity,
1187 id,
1188 serde_json::json!({
1189 "id": id,
1190 "namespace": entity.namespace,
1191 "changed_fields": changed_fields,
1192 }),
1193 )?);
1194 let post_commit = if reindex_required {
1195 PostCommitEffect::ReindexEntity { entity_id: id }
1196 } else {
1197 PostCommitEffect::None
1198 };
1199 Ok(AtomicOpPlan::Update(Box::new(UpdatePlan {
1200 graph_effects: Vec::new(),
1201 note_vector_purge: None,
1202 note_embedding_inheritance: None,
1203 entity_guard: expected_version.map(|expected_version| {
1204 crate::entity_write::EntityWriteGuard {
1205 id,
1206 expected_version,
1207 }
1208 }),
1209 note_guard: None,
1210 target_id: id,
1211 statements,
1212 post_commit,
1213 edge_natural_key: None,
1214 idempotent_noop: false,
1215 })))
1216}
1217
1218async fn prepare_update_edge(
1238 runtime: &KhiveRuntime,
1239 token: &NamespaceToken,
1240 id: Uuid,
1241 mut edge: khive_storage::types::Edge,
1242 args: &Value,
1243) -> RuntimeResult<AtomicOpPlan> {
1244 reject_inapplicable_update_fields(args, "edge")?;
1245
1246 let expected_updated_at = edge.updated_at;
1247 let expected_deleted_at = edge.deleted_at;
1248
1249 let relation_raw = optional_str(args, "relation");
1250 let weight = optional_f64(args, "weight")?;
1251 let properties = optional_properties(args, "properties")?;
1252
1253 if let Some(ref p) = properties {
1254 crate::secret_gate::check_json_at(p, "edge", "properties")?;
1255 }
1256 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
1257
1258 let namespace = edge.namespace.clone();
1259 let record_tok = token.with_namespace(
1260 khive_types::Namespace::parse(&namespace)
1261 .map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
1262 );
1263
1264 let mut changed_fields: Vec<&'static str> = Vec::new();
1265 if let Some(raw) = relation_raw {
1266 let relation = parse_edge_relation(raw)?;
1267 runtime
1268 .validate_edge_relation_endpoints(&record_tok, edge.source_id, edge.target_id, relation)
1269 .await?;
1270 edge.relation = relation;
1271 changed_fields.push("relation");
1272 }
1273 if let Some(w) = weight {
1274 if !w.is_finite() || !(0.0..=1.0).contains(&w) {
1275 return Err(RuntimeError::InvalidInput(format!(
1276 "edge weight must be a finite value in [0.0, 1.0]; got {w}"
1277 )));
1278 }
1279 edge.weight = w;
1280 changed_fields.push("weight");
1281 }
1282 if let Some(p) = properties {
1283 edge.metadata = Some(p);
1284 changed_fields.push("properties");
1285 }
1286
1287 let (canon_src, canon_tgt) =
1288 canonical_edge_endpoints(edge.relation, edge.source_id, edge.target_id);
1289 let now = chrono::Utc::now();
1290
1291 let mut statements: Vec<PlanStatement> = Vec::new();
1292 let mut edge_natural_key: Option<EdgeNaturalKey> = None;
1293
1294 if edge.relation.is_symmetric() {
1295 let metadata_str = edge
1305 .metadata
1306 .as_ref()
1307 .map(|v| serde_json::to_string(v).unwrap_or_default());
1308
1309 let minimum_updated_at_micros = expected_updated_at
1314 .timestamp_micros()
1315 .checked_add(1)
1316 .ok_or_else(|| {
1317 RuntimeError::Internal(format!(
1318 "edge {id} updated_at is already at i64::MAX and cannot advance"
1319 ))
1320 })?;
1321 let symmetric_updated_at_micros = now.timestamp_micros().max(minimum_updated_at_micros);
1322 let expected_deleted_at_micros = expected_deleted_at.map(|v| v.timestamp_micros());
1323
1324 statements.push(PlanStatement {
1325 statement: edge_symmetric_delete_if_conflict_statement(
1326 &namespace,
1327 id,
1328 canon_src,
1329 canon_tgt,
1330 edge.relation,
1331 expected_updated_at.timestamp_micros(),
1332 expected_deleted_at_micros,
1333 ),
1334 guard: Some(AffectedRowGuard {
1335 expected_min: 0,
1336 expected_max: Some(1),
1337 }),
1338 });
1339 statements.push(PlanStatement {
1340 statement: edge_symmetric_absorb_or_update_inplace_statement(
1341 &namespace,
1342 id,
1343 canon_src,
1344 canon_tgt,
1345 edge.relation,
1346 edge.weight,
1347 symmetric_updated_at_micros,
1348 metadata_str.as_deref(),
1349 edge.target_backend.as_deref(),
1350 expected_updated_at.timestamp_micros(),
1351 expected_deleted_at_micros,
1352 ),
1353 guard: Some(AffectedRowGuard::exactly(1)),
1354 });
1355
1356 edge_natural_key = Some(EdgeNaturalKey {
1361 namespace: namespace.clone(),
1362 canon_source_id: canon_src,
1363 canon_target_id: canon_tgt,
1364 relation: edge.relation,
1365 });
1366 } else {
1367 let minimum_updated_at_micros = expected_updated_at
1378 .timestamp_micros()
1379 .checked_add(1)
1380 .ok_or_else(|| {
1381 RuntimeError::Internal(format!(
1382 "edge {id} updated_at is already at i64::MAX and cannot advance"
1383 ))
1384 })?;
1385 let now_micros = now.timestamp_micros().max(minimum_updated_at_micros);
1386 edge.updated_at = chrono::DateTime::from_timestamp_micros(now_micros).ok_or_else(|| {
1387 RuntimeError::Internal(format!(
1388 "edge {id}: computed updated_at {now_micros} is not a valid timestamp"
1389 ))
1390 })?;
1391 statements.push(PlanStatement {
1392 statement: edge_replace_if_unchanged_statement(
1393 &edge,
1394 expected_updated_at,
1395 expected_deleted_at,
1396 ),
1397 guard: Some(AffectedRowGuard::exactly(1)),
1398 });
1399 }
1400
1401 statements.extend(event_append_statements(
1406 token,
1407 &namespace,
1408 "update",
1409 EventKind::EdgeUpdated,
1410 SubstrateKind::Entity,
1411 id,
1412 serde_json::json!({"id": id, "namespace": namespace, "changed_fields": changed_fields}),
1413 )?);
1414
1415 Ok(AtomicOpPlan::Update(Box::new(UpdatePlan {
1416 graph_effects: Vec::new(),
1417 note_vector_purge: None,
1418 note_embedding_inheritance: None,
1419 entity_guard: None,
1420 note_guard: None,
1421 target_id: id,
1422 statements,
1423 post_commit: PostCommitEffect::None,
1424 edge_natural_key,
1425 idempotent_noop: false,
1426 })))
1427}
1428
1429pub enum AtomicDeleteKind {
1446 Entity { specific: Option<String> },
1447 Note { specific: Option<String> },
1448 Edge,
1449}
1450
1451pub async fn prepare_delete(
1457 runtime: &KhiveRuntime,
1458 token: &NamespaceToken,
1459 args: &Value,
1460 expected_kind: Option<AtomicDeleteKind>,
1461) -> RuntimeResult<AtomicOpPlan> {
1462 let id = require_uuid(args, "id")?;
1463 let actor = format!("{}:{}", token.actor().kind, token.actor().id);
1464 let hard = obj(args)?
1465 .get("hard")
1466 .and_then(|v| v.as_bool())
1467 .unwrap_or(false);
1468
1469 let resolved = if hard {
1475 runtime.resolve_by_id_including_deleted(token, id).await?
1476 } else {
1477 runtime.resolve_by_id(token, id).await?
1478 };
1479
1480 match resolved {
1481 Some(Resolved::Entity(entity)) => {
1482 match &expected_kind {
1483 None => {}
1484 Some(AtomicDeleteKind::Entity {
1485 specific: Some(expected),
1486 }) if &entity.kind != expected => {
1487 return Err(RuntimeError::NotFound(format!("{expected} {id}")));
1488 }
1489 Some(AtomicDeleteKind::Entity { .. }) => {}
1490 Some(AtomicDeleteKind::Note { .. }) => {
1491 return Err(RuntimeError::NotFound(format!("note {id}")));
1492 }
1493 Some(AtomicDeleteKind::Edge) => {
1494 return Err(RuntimeError::NotFound(format!("edge {id}")));
1495 }
1496 }
1497 let namespace = entity.namespace.clone();
1498 let mut statements = if hard {
1503 vec![
1504 PlanStatement {
1505 statement: delete_record_attachments_statement(
1506 id,
1507 AttachmentSubstrate::Entity,
1508 ),
1509 guard: None,
1510 },
1511 PlanStatement {
1512 statement: entity_hard_delete_statement(id),
1513 guard: Some(AffectedRowGuard::exactly(1)),
1514 },
1515 ]
1516 } else {
1517 let deleted_at = chrono::Utc::now().timestamp_micros();
1518 vec![PlanStatement {
1519 statement: entity_soft_delete_statement(id, deleted_at),
1520 guard: Some(AffectedRowGuard::exactly(1)),
1521 }]
1522 };
1523 if hard {
1524 statements.extend(
1525 hard_delete_lineage_warning_statements(
1526 &namespace,
1527 &actor,
1528 id,
1529 SubstrateKind::Entity,
1530 )
1531 .into_iter()
1532 .map(|statement| PlanStatement {
1533 statement,
1534 guard: None,
1535 }),
1536 );
1537 statements.push(PlanStatement {
1540 statement: purge_incident_edges_statement(id),
1541 guard: None,
1542 });
1543 }
1544 push_index_purge_statements(
1549 runtime,
1550 &mut statements,
1551 "fts_entities",
1552 &namespace,
1553 id,
1554 "atomic-delete-entity",
1555 )
1556 .await?;
1557 statements.extend(event_append_statements(
1563 token,
1564 &namespace,
1565 "delete",
1566 EventKind::EntityDeleted,
1567 SubstrateKind::Entity,
1568 id,
1569 serde_json::json!({"id": id, "namespace": namespace, "hard": hard}),
1570 )?);
1571 Ok(AtomicOpPlan::Delete(DeletePlan {
1572 target_id: id,
1573 statements,
1574 post_commit: PostCommitEffect::None,
1575 }))
1576 }
1577 Some(Resolved::Note(note)) => {
1578 match &expected_kind {
1579 None => {}
1580 Some(AtomicDeleteKind::Note {
1581 specific: Some(expected),
1582 }) if ¬e.kind != expected => {
1583 return Err(RuntimeError::NotFound(format!("{expected} {id}")));
1584 }
1585 Some(AtomicDeleteKind::Note { .. }) => {}
1586 Some(AtomicDeleteKind::Entity { .. }) => {
1587 return Err(RuntimeError::NotFound(format!("entity {id}")));
1588 }
1589 Some(AtomicDeleteKind::Edge) => {
1590 return Err(RuntimeError::NotFound(format!("edge {id}")));
1591 }
1592 }
1593 if let Some(error) = runtime.stream_member_error(¬e).await? {
1594 return Err(error);
1595 }
1596 let namespace = note.namespace.clone();
1597 let mut statements = if hard {
1601 vec![
1602 PlanStatement {
1603 statement: delete_record_attachments_statement(
1604 id,
1605 AttachmentSubstrate::Note,
1606 ),
1607 guard: None,
1608 },
1609 PlanStatement {
1610 statement: note_hard_delete_statement(id),
1611 guard: Some(AffectedRowGuard::exactly(1)),
1612 },
1613 ]
1614 } else {
1615 let deleted_at = chrono::Utc::now().timestamp_micros();
1616 vec![PlanStatement {
1617 statement: note_soft_delete_statement(id, deleted_at),
1618 guard: Some(AffectedRowGuard::exactly(1)),
1619 }]
1620 };
1621 if hard {
1622 statements.extend(
1623 hard_delete_lineage_warning_statements(
1624 &namespace,
1625 &actor,
1626 id,
1627 SubstrateKind::Note,
1628 )
1629 .into_iter()
1630 .map(|statement| PlanStatement {
1631 statement,
1632 guard: None,
1633 }),
1634 );
1635 statements.push(PlanStatement {
1636 statement: purge_incident_edges_statement(id),
1637 guard: None,
1638 });
1639 }
1640 push_index_purge_statements(
1645 runtime,
1646 &mut statements,
1647 "fts_notes",
1648 &namespace,
1649 id,
1650 "atomic-delete-note",
1651 )
1652 .await?;
1653 statements.extend(event_append_statements(
1657 token,
1658 &namespace,
1659 "delete",
1660 EventKind::NoteDeleted,
1661 SubstrateKind::Note,
1662 id,
1663 serde_json::json!({"id": id, "namespace": namespace, "hard": hard}),
1664 )?);
1665 Ok(AtomicOpPlan::Delete(DeletePlan {
1666 target_id: id,
1667 statements,
1668 post_commit: PostCommitEffect::NoteDeleted {
1673 note_id: id,
1674 kind: note.kind.clone(),
1675 },
1676 }))
1677 }
1678 Some(_) => Err(RuntimeError::InvalidInput(format!(
1679 "delete target {id} must be an entity, note, or edge"
1680 ))),
1681 None => match &expected_kind {
1685 Some(AtomicDeleteKind::Entity { .. }) => {
1686 Err(RuntimeError::NotFound(format!("entity/note {id}")))
1687 }
1688 Some(AtomicDeleteKind::Note { .. }) => {
1689 Err(RuntimeError::NotFound(format!("entity/note {id}")))
1690 }
1691 Some(AtomicDeleteKind::Edge) | None => {
1692 let edge = if hard {
1693 runtime.get_edge_including_deleted(token, id).await?
1694 } else {
1695 runtime.get_edge(token, id).await?
1696 };
1697 match edge {
1698 Some(edge) => prepare_delete_edge(token, id, edge, hard, &actor).await,
1699 None => Err(RuntimeError::NotFound(format!("entity/note/edge {id}"))),
1700 }
1701 }
1702 },
1703 }
1704}
1705
1706async fn prepare_delete_edge(
1715 token: &NamespaceToken,
1716 id: Uuid,
1717 edge: khive_storage::types::Edge,
1718 hard: bool,
1719 actor: &str,
1720) -> RuntimeResult<AtomicOpPlan> {
1721 let namespace = edge.namespace.clone();
1722 let mut statements: Vec<PlanStatement> = Vec::new();
1723
1724 if hard {
1725 statements.extend(
1726 hard_delete_lineage_warning_statements(&namespace, actor, id, SubstrateKind::Entity)
1727 .into_iter()
1728 .map(|statement| PlanStatement {
1729 statement,
1730 guard: None,
1731 }),
1732 );
1733 statements.push(PlanStatement {
1738 statement: purge_incident_edges_statement(id),
1739 guard: None,
1740 });
1741 statements.push(PlanStatement {
1742 statement: edge_hard_delete_statement(id),
1743 guard: Some(AffectedRowGuard::exactly(1)),
1744 });
1745 } else {
1746 let now = chrono::Utc::now().timestamp_micros();
1747 statements.push(PlanStatement {
1748 statement: edge_soft_delete_statement(id, now),
1749 guard: Some(AffectedRowGuard::exactly(1)),
1750 });
1751 }
1752
1753 statements.extend(event_append_statements(
1754 token,
1755 &namespace,
1756 "delete",
1757 EventKind::EdgeDeleted,
1758 SubstrateKind::Entity,
1759 id,
1760 serde_json::json!({"id": id, "namespace": namespace, "hard": hard}),
1761 )?);
1762
1763 Ok(AtomicOpPlan::Delete(DeletePlan {
1764 target_id: id,
1765 statements,
1766 post_commit: PostCommitEffect::None,
1767 }))
1768}
1769
1770fn parse_edge_relation(raw: &str) -> RuntimeResult<EdgeRelation> {
1775 raw.parse::<EdgeRelation>()
1776 .map_err(|e| RuntimeError::InvalidInput(format!("unknown edge relation {raw:?}: {e}")))
1777}
1778
1779async fn prepare_link(
1780 runtime: &KhiveRuntime,
1781 token: &NamespaceToken,
1782 args: &Value,
1783) -> RuntimeResult<AtomicOpPlan> {
1784 let source_id = require_uuid(args, "source_id")?;
1785 let target_id = require_uuid(args, "target_id")?;
1786 let relation = parse_edge_relation(require_str(args, "relation")?)?;
1787 let weight = optional_f64(args, "weight")?.unwrap_or(1.0);
1788 let metadata = obj(args)?.get("metadata").cloned();
1789 let resurrect = match obj(args)?.get("resurrect") {
1790 None => false,
1791 Some(Value::Bool(value)) => *value,
1792 Some(other) => {
1793 return Err(RuntimeError::InvalidInput(format!(
1794 "resurrect must be a boolean, got: {other}"
1795 )))
1796 }
1797 };
1798
1799 let mut metadata = crate::merge_entry_metadata(
1805 metadata,
1806 optional_str(args, "dependency_kind").map(String::from),
1807 )?;
1808
1809 validate_edge_weight(weight)?;
1810 runtime
1811 .validate_edge_relation_endpoints(token, source_id, target_id, relation)
1812 .await?;
1813
1814 let (canon_source, canon_target) = canonical_edge_endpoints(relation, source_id, target_id);
1815
1816 if relation == EdgeRelation::DependsOn {
1823 metadata = match (
1824 runtime.resolve_edge_endpoint(token, canon_source).await?,
1825 runtime.resolve_edge_endpoint(token, canon_target).await?,
1826 ) {
1827 (Some(Resolved::Entity(src_e)), Some(Resolved::Entity(tgt_e))) => {
1828 merge_dependency_kind(&src_e.kind, &tgt_e.kind, metadata)
1829 }
1830 _ => metadata,
1831 };
1832 }
1833
1834 validate_edge_metadata(relation, metadata.as_ref())?;
1835 let namespace = token.namespace().as_str().to_string();
1836 let previous = runtime
1837 .get_edge_by_natural_key_including_deleted(
1838 token,
1839 &namespace,
1840 canon_source,
1841 canon_target,
1842 relation,
1843 )
1844 .await?;
1845 if let Some(edge) = previous.as_ref() {
1846 if edge.deleted_at.is_some() && !resurrect {
1847 return Err(RuntimeError::InvalidInput(format!(
1848 "edge natural key is soft-deleted; pass resurrect=true to link explicitly: {}",
1849 Uuid::from(edge.id)
1850 )));
1851 }
1852 }
1853
1854 let disposition = match previous.as_ref() {
1855 None => EdgeUpsertDisposition::Created,
1856 Some(edge) if edge.deleted_at.is_some() => EdgeUpsertDisposition::Resurrected,
1857 Some(_) => EdgeUpsertDisposition::Updated,
1858 };
1859 let edge_id = previous
1860 .as_ref()
1861 .map(|edge| Uuid::from(edge.id))
1862 .unwrap_or_else(Uuid::new_v4);
1863 let now = previous.as_ref().map_or_else(
1864 || chrono::Utc::now().timestamp_micros(),
1865 |edge| {
1866 chrono::Utc::now()
1867 .timestamp_micros()
1868 .max(edge.updated_at.timestamp_micros().saturating_add(1))
1869 },
1870 );
1871 let metadata_str = metadata
1872 .as_ref()
1873 .map(|value| serde_json::to_string(value).unwrap_or_default());
1874
1875 let statement = match previous.as_ref() {
1879 None => edge_insert_new_guarded_by_endpoints_statement(
1880 &namespace,
1881 edge_id,
1882 canon_source,
1883 canon_target,
1884 relation,
1885 weight,
1886 now,
1887 metadata_str.as_deref(),
1888 ),
1889 Some(edge) => edge_link_replace_if_unchanged_and_endpoints_exist_statement(
1890 edge,
1891 weight,
1892 now,
1893 metadata_str.as_deref(),
1894 ),
1895 };
1896 let mut statements = vec![PlanStatement {
1897 statement,
1898 guard: Some(AffectedRowGuard::exactly(1)),
1899 }];
1900 let kind = match disposition {
1901 EdgeUpsertDisposition::Created => EventKind::LinkCreated,
1902 EdgeUpsertDisposition::Updated | EdgeUpsertDisposition::Resurrected => {
1903 EventKind::EdgeUpdated
1904 }
1905 };
1906 statements.extend(event_append_statements(
1907 token,
1908 &namespace,
1909 "link",
1910 kind,
1911 SubstrateKind::Entity,
1912 edge_id,
1913 serde_json::json!({
1914 "id": edge_id,
1915 "namespace": namespace,
1916 "mutation": disposition.name(),
1917 "source_id": canon_source,
1918 "target_id": canon_target,
1919 "relation": relation,
1920 "weight": weight,
1921 "metadata": metadata,
1922 "previous": previous,
1923 }),
1924 )?);
1925
1926 Ok(AtomicOpPlan::Link(LinkPlan {
1927 source_id: canon_source,
1928 target_id: canon_target,
1929 statements,
1930 disposition,
1931 }))
1932}
1933
1934async fn prepare_merge(
1948 runtime: &KhiveRuntime,
1949 token: &NamespaceToken,
1950 args: &Value,
1951) -> RuntimeResult<AtomicOpPlan> {
1952 let into_id = require_uuid(args, "into_id")?;
1953 let from_id = require_uuid(args, "from_id")?;
1954 if into_id == from_id {
1955 return Err(RuntimeError::InvalidInput(
1956 "cannot merge an entity into itself".into(),
1957 ));
1958 }
1959
1960 let entities = runtime.entities(token)?;
1961 entities
1962 .get_entity(into_id)
1963 .await?
1964 .ok_or_else(|| RuntimeError::NotFound(format!("entity {into_id}")))?;
1965 entities
1966 .get_entity(from_id)
1967 .await?
1968 .ok_or_else(|| RuntimeError::NotFound(format!("entity {from_id}")))?;
1969
1970 let now = chrono::Utc::now().timestamp_micros();
1971 let rewires = vec![
1972 crate::atomic_plan::PlanPredicate {
1973 description: "source_id = :from".to_string(),
1974 statement: SqlStatement {
1975 sql: "UPDATE graph_edges SET source_id = ?1, updated_at = ?2 WHERE source_id = ?3"
1976 .to_string(),
1977 params: vec![
1978 SqlValue::Text(into_id.to_string()),
1979 SqlValue::Integer(now),
1980 SqlValue::Text(from_id.to_string()),
1981 ],
1982 label: Some("atomic-merge-rewire-source".to_string()),
1983 },
1984 },
1985 crate::atomic_plan::PlanPredicate {
1986 description: "target_id = :from".to_string(),
1987 statement: SqlStatement {
1988 sql: "UPDATE graph_edges SET target_id = ?1, updated_at = ?2 WHERE target_id = ?3"
1989 .to_string(),
1990 params: vec![
1991 SqlValue::Text(into_id.to_string()),
1992 SqlValue::Integer(now),
1993 SqlValue::Text(from_id.to_string()),
1994 ],
1995 label: Some("atomic-merge-rewire-target".to_string()),
1996 },
1997 },
1998 ];
1999 let lifecycle = vec![PlanStatement {
2000 statement: SqlStatement {
2001 sql: "UPDATE entities SET deleted_at = ?1, merged_into = ?2, version = version + 1 \
2002 WHERE id = ?3 AND deleted_at IS NULL"
2003 .to_string(),
2004 params: vec![
2005 SqlValue::Integer(now),
2006 SqlValue::Text(into_id.to_string()),
2007 SqlValue::Text(from_id.to_string()),
2008 ],
2009 label: Some("atomic-merge-tombstone-from-entity".to_string()),
2010 },
2011 guard: Some(AffectedRowGuard::exactly(1)),
2012 }];
2013
2014 Ok(AtomicOpPlan::Merge(MergePlan {
2015 into_id,
2016 from_id,
2017 rewires,
2018 lifecycle,
2019 }))
2020}
2021
2022#[derive(Clone, Debug, PartialEq, Eq)]
2032pub struct PostCommitEmbeddingOutcome {
2033 pub effect: PostCommitEffect,
2035 pub truncation: crate::retrieval::EmbeddingTruncationReport,
2037}
2038
2039pub async fn apply_post_commit_effects(
2041 runtime: &KhiveRuntime,
2042 token: &NamespaceToken,
2043 effects: CommittedPostCommitEffects,
2044) -> RuntimeResult<()> {
2045 apply_post_commit_effects_with_report(runtime, token, effects)
2046 .await
2047 .map(|_| ())
2048}
2049
2050pub async fn apply_post_commit_effects_with_report(
2057 runtime: &KhiveRuntime,
2058 token: &NamespaceToken,
2059 effects: CommittedPostCommitEffects,
2060) -> RuntimeResult<Vec<PostCommitEmbeddingOutcome>> {
2061 let mut embedding_outcomes = Vec::new();
2062 let mut failures = Vec::new();
2063 for (index, effect) in effects.into_effects().into_iter().enumerate() {
2064 let identity = format!("{effect:?}");
2065 match apply_one_post_commit_effect(runtime, token, effect).await {
2066 Ok(Some(outcome)) => embedding_outcomes.push(outcome),
2067 Ok(None) => {}
2068 Err(error) => failures.push(format!("effect[{index}] {identity}: {error}")),
2069 }
2070 }
2071 if failures.is_empty() {
2072 Ok(embedding_outcomes)
2073 } else {
2074 Err(RuntimeError::Internal(format!(
2075 "post-commit effects failed after commit: {}",
2076 failures.join("; ")
2077 )))
2078 }
2079}
2080
2081async fn apply_one_post_commit_effect(
2082 runtime: &KhiveRuntime,
2083 token: &NamespaceToken,
2084 effect: PostCommitEffect,
2085) -> RuntimeResult<Option<PostCommitEmbeddingOutcome>> {
2086 match effect {
2087 PostCommitEffect::None => Ok(None),
2088 PostCommitEffect::NoteChanged { note_id, kind } => {
2089 runtime.fire_note_mutation_hook(&kind, note_id).await;
2090 Ok(None)
2091 }
2092 PostCommitEffect::ReindexEntity { entity_id } => {
2093 let Some(entity) = runtime.entities(token)?.get_entity(entity_id).await? else {
2094 return Ok(None);
2095 };
2096 let truncation = runtime.reindex_entity(token, &entity).await?;
2097 Ok(Some(PostCommitEmbeddingOutcome {
2098 effect: PostCommitEffect::ReindexEntity { entity_id },
2099 truncation,
2100 }))
2101 }
2102 PostCommitEffect::ReindexNote { note_id, version } => {
2103 let Some(note) = runtime.notes(token)?.get_note(note_id).await? else {
2104 return Ok(None);
2105 };
2106 if note.version != version {
2107 return Ok(None);
2108 }
2109 let truncation = runtime.reindex_note(token, ¬e).await?;
2110 if runtime
2111 .notes(token)?
2112 .get_note(note_id)
2113 .await?
2114 .is_none_or(|current| current.version != version)
2115 {
2116 return Ok(None);
2117 }
2118 runtime.fire_note_mutation_hook(¬e.kind, note.id).await;
2121 Ok(Some(PostCommitEmbeddingOutcome {
2122 effect: PostCommitEffect::ReindexNote { note_id, version },
2123 truncation,
2124 }))
2125 }
2126 PostCommitEffect::NoteDeleted { note_id, kind } => {
2127 runtime.fire_note_mutation_hook(&kind, note_id).await;
2129 Ok(None)
2130 }
2131 PostCommitEffect::GtdAudit { .. } => {
2132 Ok(None)
2134 }
2135 }
2136}
2137
2138#[cfg(test)]
2139mod tests {
2140 use super::*;
2141
2142 use async_trait::async_trait;
2143 use lattice_embed::{EmbedError, EmbeddingModel, EmbeddingService, MAX_TEXT_BYTES};
2144 use serde_json::json;
2145
2146 use khive_types::Namespace;
2147
2148 use crate::embedder_registry::EmbedderProvider;
2149 use crate::runtime::RuntimeConfig;
2150
2151 struct TestRuntime {
2153 runtime: KhiveRuntime,
2154 _temp_dir: tempfile::TempDir,
2155 }
2156
2157 impl std::ops::Deref for TestRuntime {
2158 type Target = KhiveRuntime;
2159
2160 fn deref(&self) -> &Self::Target {
2161 &self.runtime
2162 }
2163 }
2164
2165 const STUB_MODEL: &str = "stub-adr099-b3";
2166 const STUB_DIMS: usize = 4;
2167
2168 struct StubService;
2169
2170 #[async_trait]
2171 impl EmbeddingService for StubService {
2172 async fn embed(
2173 &self,
2174 texts: &[String],
2175 _model: EmbeddingModel,
2176 ) -> Result<Vec<Vec<f32>>, EmbedError> {
2177 Ok(texts.iter().map(|_| vec![0.5_f32; STUB_DIMS]).collect())
2178 }
2179
2180 fn supports_model(&self, _model: EmbeddingModel) -> bool {
2181 true
2182 }
2183
2184 fn name(&self) -> &'static str {
2185 STUB_MODEL
2186 }
2187 }
2188
2189 struct StubProvider;
2190
2191 #[async_trait]
2192 impl EmbedderProvider for StubProvider {
2193 fn name(&self) -> &str {
2194 STUB_MODEL
2195 }
2196
2197 fn dimensions(&self) -> usize {
2198 STUB_DIMS
2199 }
2200
2201 async fn build(&self) -> RuntimeResult<std::sync::Arc<dyn EmbeddingService>> {
2202 Ok(std::sync::Arc::new(StubService))
2203 }
2204 }
2205
2206 fn scratch_runtime() -> TestRuntime {
2207 let dir = tempfile::tempdir().expect("tempdir");
2208 let path = dir.path().join("atomic_prepare_reindex.db");
2209 let runtime = KhiveRuntime::new_for_test(RuntimeConfig {
2210 db_path: Some(path),
2211 embedding_model: None,
2212 additional_embedding_models: vec![],
2213 ..RuntimeConfig::default()
2214 })
2215 .expect("runtime");
2216 TestRuntime {
2217 runtime,
2218 _temp_dir: dir,
2219 }
2220 }
2221
2222 #[tokio::test]
2230 async fn atomic_update_entity_rejects_note_only_field_salience() {
2231 let runtime = scratch_runtime();
2232 let token = runtime
2233 .authorize(Namespace::parse("local").expect("ns"))
2234 .expect("authorize");
2235 let entity = khive_storage::Entity::new("local", "concept", "GapFourEntity");
2236 let entity_id = entity.id;
2237 runtime
2238 .entities(&token)
2239 .expect("entities store")
2240 .upsert_entity(entity)
2241 .await
2242 .expect("seed entity");
2243
2244 let err = prepare_update(
2245 &runtime,
2246 &token,
2247 &json!({"id": entity_id.to_string(), "salience": 0.9}),
2248 None,
2249 )
2250 .await
2251 .expect_err("salience on an entity must be rejected, not silently accepted");
2252 assert!(
2253 matches!(err, RuntimeError::InvalidInput(ref msg) if msg.contains("salience") && msg.contains("not valid for an entity")),
2254 "expected an InvalidInput naming the offending field, got: {err:?}"
2255 );
2256
2257 let plan = prepare_update(
2259 &runtime,
2260 &token,
2261 &json!({
2262 "id": entity_id.to_string(),
2263 "name": "GapFourEntity-renamed",
2264 "description": "updated description",
2265 "tags": ["a", "b"],
2266 }),
2267 None,
2268 )
2269 .await
2270 .expect("a valid entity field set must still be accepted");
2271 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
2272 .await
2273 .expect("seam call ok");
2274 assert!(matches!(
2275 outcome,
2276 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
2277 ));
2278 let entity = runtime
2279 .get_entity(&token, entity_id)
2280 .await
2281 .expect("get_entity");
2282 assert_eq!(entity.name, "GapFourEntity-renamed");
2283 }
2284
2285 #[tokio::test]
2292 async fn atomic_update_short_prefix_id_rejected_with_namespace_explanation() {
2293 let runtime = scratch_runtime();
2294 let token = runtime
2295 .authorize(Namespace::parse("local").expect("ns"))
2296 .expect("authorize");
2297
2298 let err = prepare_update(
2299 &runtime,
2300 &token,
2301 &json!({"id": "deadbeef", "name": "whatever"}),
2302 None,
2303 )
2304 .await
2305 .expect_err("a short hex prefix must be rejected by the atomic-plan seam");
2306 let msg = match err {
2307 RuntimeError::InvalidInput(ref msg) => msg.clone(),
2308 other => panic!("expected InvalidInput, got: {other:?}"),
2309 };
2310 assert!(
2311 msg.contains("full UUID"),
2312 "message must still state the rule; got: {msg}"
2313 );
2314 assert!(
2315 msg.to_ascii_lowercase().contains("namespace"),
2316 "message must explain the namespace-scoping consequence, not just restate the \
2317 rule; got: {msg}"
2318 );
2319 }
2320
2321 #[tokio::test]
2322 async fn atomic_update_entity_type_persists_patch_and_schedules_reindex() {
2323 let runtime = scratch_runtime();
2324 runtime.install_entity_type_validator(std::sync::Arc::new(|kind, entity_type| {
2325 let Some(raw) = entity_type else {
2326 return Ok(None);
2327 };
2328 let normalized = raw.trim().to_ascii_lowercase();
2329 if kind == "concept" && normalized == "algorithm" {
2330 Ok(Some(normalized))
2331 } else {
2332 Err(RuntimeError::InvalidInput(format!(
2333 "unknown entity_type {raw:?} for {kind:?}; valid: algorithm"
2334 )))
2335 }
2336 }));
2337 let token = runtime
2338 .authorize(Namespace::parse("local").expect("ns"))
2339 .expect("authorize");
2340 let mut entity = khive_storage::Entity::new("local", "concept", "AtomicHistorical");
2341 entity.description = Some("keep description".to_string());
2342 entity.properties = Some(json!({"type": "algorithm", "keep": true}));
2343 entity.tags = vec!["keep-tag".to_string()];
2344 let entity_id = entity.id;
2345 runtime
2346 .entities(&token)
2347 .expect("entities store")
2348 .upsert_entity(entity)
2349 .await
2350 .expect("seed entity");
2351
2352 let plan = prepare_update(
2353 &runtime,
2354 &token,
2355 &json!({"id": entity_id.to_string(), "entity_type": " Algorithm "}),
2356 None,
2357 )
2358 .await
2359 .expect("atomic prepare must accept a registered entity_type");
2360 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
2361 .await
2362 .expect("atomic update must run");
2363 let post_commit = match outcome {
2364 crate::atomic_runner::AtomicRunOutcome::Committed { post_commit } => post_commit,
2365 other => panic!("expected Committed, got {other:?}"),
2366 };
2367 assert_eq!(
2368 post_commit.as_slice(),
2369 &[PostCommitEffect::ReindexEntity { entity_id }],
2370 "entity_type patch must schedule the normal entity reindex path"
2371 );
2372
2373 let updated = runtime
2374 .get_entity(&token, entity_id)
2375 .await
2376 .expect("read updated entity");
2377 assert_eq!(updated.entity_type.as_deref(), Some("algorithm"));
2378 assert_eq!(updated.name, "AtomicHistorical");
2379 assert_eq!(updated.description.as_deref(), Some("keep description"));
2380 assert_eq!(
2381 updated.properties,
2382 Some(json!({"type": "algorithm", "keep": true}))
2383 );
2384 assert_eq!(updated.tags, vec!["keep-tag"]);
2385 }
2386
2387 #[tokio::test]
2390 async fn atomic_update_note_rejects_entity_only_field_description() {
2391 let runtime = scratch_runtime();
2392 let token = runtime
2393 .authorize(Namespace::parse("local").expect("ns"))
2394 .expect("authorize");
2395 let mut note = khive_storage::note::Note::new("local", "observation", "gap-4 note content");
2396 note.name = Some("gap-four-note".to_string());
2397 let note_id = note.id;
2398 runtime
2399 .notes(&token)
2400 .expect("notes store")
2401 .upsert_note(note)
2402 .await
2403 .expect("seed note");
2404
2405 let err = prepare_update(
2406 &runtime,
2407 &token,
2408 &json!({"id": note_id.to_string(), "description": "entities have descriptions, notes don't"}),
2409 None,
2410 )
2411 .await
2412 .expect_err("description on a note must be rejected, not silently accepted");
2413 assert!(
2414 matches!(err, RuntimeError::InvalidInput(ref msg) if msg.contains("description") && msg.contains("not valid for a note")),
2415 "expected an InvalidInput naming the offending field, got: {err:?}"
2416 );
2417
2418 let plan = prepare_update(
2420 &runtime,
2421 &token,
2422 &json!({"id": note_id.to_string(), "content": "gap-4 note content, revised"}),
2423 None,
2424 )
2425 .await
2426 .expect("a valid note field must still be accepted");
2427 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
2428 .await
2429 .expect("seam call ok");
2430 assert!(matches!(
2431 outcome,
2432 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
2433 ));
2434 }
2435
2436 #[tokio::test]
2437 async fn atomic_update_note_tags_replace_preserve_clear_and_override_nested_tags() {
2438 let runtime = scratch_runtime();
2439 let token = runtime
2440 .authorize(Namespace::parse("local").expect("ns"))
2441 .expect("authorize");
2442 let mut note = khive_storage::note::Note::new("local", "observation", "tagged note");
2443 note.properties = Some(json!({"tags": ["old"], "keep": {"value": 1}}));
2444 let note_id = note.id;
2445 runtime
2446 .notes(&token)
2447 .expect("notes store")
2448 .upsert_note(note)
2449 .await
2450 .expect("seed note");
2451
2452 for (mut args, expected_tags) in [
2453 (
2454 json!({"tags": ["new", "shared"], "properties": {"tags": ["nested"], "added": true}}),
2455 json!(["new", "shared"]),
2456 ),
2457 (
2458 json!({"name": "renamed note", "properties": {"omitted": true}}),
2459 json!(["new", "shared"]),
2460 ),
2461 (
2462 json!({"tags": null, "properties": null}),
2463 json!(["new", "shared"]),
2464 ),
2465 (
2466 json!({"tags": [], "properties": {"tags": ["nested-after-clear"]}}),
2467 json!([]),
2468 ),
2469 (
2470 json!({"tags": ["after-null-properties"], "properties": null}),
2471 json!(["after-null-properties"]),
2472 ),
2473 ] {
2474 args["id"] = json!(note_id.to_string());
2475 let original_args = args.clone();
2476 let plan = prepare_update(&runtime, &token, &args, None)
2477 .await
2478 .expect("valid atomic note tags patch");
2479 assert_eq!(
2480 args, original_args,
2481 "preparation must not mutate caller args"
2482 );
2483 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
2484 .await
2485 .expect("atomic note update");
2486 assert!(matches!(
2487 outcome,
2488 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
2489 ));
2490 let updated = runtime
2491 .notes(&token)
2492 .expect("notes store")
2493 .get_note(note_id)
2494 .await
2495 .expect("read note")
2496 .expect("note exists");
2497 let properties = updated.properties.expect("note properties");
2498 assert_eq!(properties["tags"], expected_tags);
2499 assert_eq!(properties["keep"], json!({"value": 1}));
2500 assert_eq!(properties["added"], json!(true));
2501 assert_eq!(updated.content, "tagged note");
2502 if original_args.get("name").is_some() {
2503 assert_eq!(updated.name.as_deref(), Some("renamed note"));
2504 assert_eq!(properties["omitted"], json!(true));
2505 }
2506 }
2507 }
2508
2509 #[tokio::test]
2510 async fn atomic_update_entity_tags_keep_replace_preserve_and_clear_semantics() {
2511 let runtime = scratch_runtime();
2512 let token = runtime
2513 .authorize(Namespace::parse("local").expect("ns"))
2514 .expect("authorize");
2515 let mut entity = khive_storage::Entity::new("local", "concept", "tagged entity");
2516 entity.tags = vec!["old".to_string()];
2517 entity.properties = Some(json!({"keep": true}));
2518 let entity_id = entity.id;
2519 runtime
2520 .entities(&token)
2521 .expect("entities store")
2522 .upsert_entity(entity)
2523 .await
2524 .expect("seed entity");
2525 for (mut args, expected_tags) in [
2526 (json!({"tags": ["new", "shared"]}), json!(["new", "shared"])),
2527 (json!({"name": "renamed entity"}), json!(["new", "shared"])),
2528 (json!({"tags": null}), json!(["new", "shared"])),
2529 (json!({"tags": []}), json!([])),
2530 ] {
2531 args["id"] = json!(entity_id.to_string());
2532 let plan = prepare_update(&runtime, &token, &args, None)
2533 .await
2534 .expect("valid atomic entity tags patch");
2535 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
2536 .await
2537 .expect("atomic entity update");
2538 assert!(matches!(
2539 outcome,
2540 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
2541 ));
2542 let updated = runtime
2543 .get_entity(&token, entity_id)
2544 .await
2545 .expect("read entity");
2546 assert_eq!(json!(updated.tags), expected_tags);
2547 assert_eq!(updated.properties, Some(json!({"keep": true})));
2548 }
2549 }
2550
2551 #[tokio::test]
2552 async fn atomic_update_note_invalid_tags_leave_snapshot_unchanged() {
2553 let runtime = scratch_runtime();
2554 let token = runtime
2555 .authorize(Namespace::parse("local").expect("ns"))
2556 .expect("authorize");
2557 let mut note = khive_storage::note::Note::new("local", "observation", "unchanged note");
2558 note.properties = Some(json!({"tags": ["keep"], "other": true}));
2559 let note_id = note.id;
2560 runtime
2561 .notes(&token)
2562 .expect("notes store")
2563 .upsert_note(note)
2564 .await
2565 .expect("seed note");
2566 let before = runtime
2567 .notes(&token)
2568 .expect("notes store")
2569 .get_note(note_id)
2570 .await
2571 .expect("read note")
2572 .expect("note exists");
2573
2574 for (mut args, expected_error) in [
2575 (
2576 json!({"tags": "invalid"}),
2577 "tags must be an array of strings",
2578 ),
2579 (
2580 json!({"tags": ["valid", 1]}),
2581 "tags must be an array of strings",
2582 ),
2583 (
2584 json!({"tags": {"nested": true}}),
2585 "tags must be an array of strings",
2586 ),
2587 (
2588 json!({"tags": ["valid"], "properties": []}),
2589 "properties must be an object",
2590 ),
2591 (
2592 json!({"tags": [], "properties": "invalid"}),
2593 "properties must be an object",
2594 ),
2595 (
2596 json!({"tags": "invalid", "description": "entity field"}),
2597 "field 'description' is not valid for a note",
2598 ),
2599 ] {
2600 args["id"] = json!(note_id.to_string());
2601 args["content"] = json!("must not persist");
2602 let error = prepare_update(&runtime, &token, &args, None)
2603 .await
2604 .expect_err("invalid tags patch must not produce a plan");
2605 assert!(
2606 matches!(error, RuntimeError::InvalidInput(ref message) if message.contains(expected_error)),
2607 "unexpected error: {error:?}"
2608 );
2609 let after = runtime
2610 .notes(&token)
2611 .expect("notes store")
2612 .get_note(note_id)
2613 .await
2614 .expect("read note")
2615 .expect("note exists");
2616 assert_eq!(after, before);
2617 }
2618 }
2619
2620 #[tokio::test]
2625 async fn atomic_update_note_content_is_fts_and_vector_reindexed_post_commit() {
2626 let runtime = scratch_runtime();
2627 runtime.register_embedder(StubProvider);
2628 let token = runtime
2629 .authorize(Namespace::parse("local").expect("ns"))
2630 .expect("authorize");
2631
2632 let mut note = khive_storage::note::Note::new("local", "observation", "original content");
2633 note.name = Some("reindex-target".to_string());
2634 let note_id = note.id;
2635 runtime
2636 .notes(&token)
2637 .expect("notes store")
2638 .upsert_note(note)
2639 .await
2640 .expect("seed note");
2641
2642 let vec_store = runtime
2644 .vectors_for_model(&token, STUB_MODEL)
2645 .expect("vec store");
2646 assert_eq!(vec_store.count().await.expect("count before"), 0);
2647
2648 let updated_content = format!("freshly-updated-content-xyz{}", "x".repeat(MAX_TEXT_BYTES));
2649 let plan = prepare_update(
2650 &runtime,
2651 &token,
2652 &json!({"id": note_id.to_string(), "content": updated_content, "embed": true}),
2653 None,
2654 )
2655 .await
2656 .expect("prepare update");
2657
2658 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
2659 .await
2660 .expect("seam call ok");
2661 let post_commit = match outcome {
2662 crate::atomic_runner::AtomicRunOutcome::Committed { post_commit } => post_commit,
2663 other => panic!("expected Committed, got {other:?}"),
2664 };
2665 assert_eq!(
2666 post_commit.as_slice(),
2667 &[PostCommitEffect::ReindexNote {
2668 note_id,
2669 version: 2
2670 }],
2671 "content change must schedule exactly one ReindexNote post-commit effect"
2672 );
2673
2674 let embedding_outcomes =
2675 apply_post_commit_effects_with_report(&runtime, &token, post_commit)
2676 .await
2677 .expect("apply post-commit effects");
2678 assert_eq!(embedding_outcomes.len(), 1);
2679 assert_eq!(
2680 embedding_outcomes[0].effect,
2681 PostCommitEffect::ReindexNote {
2682 note_id,
2683 version: 2
2684 }
2685 );
2686 assert_eq!(embedding_outcomes[0].truncation.truncated, 1);
2687 assert!(embedding_outcomes[0].truncation.discarded_bytes > 0);
2688
2689 let doc = runtime
2691 .text_for_notes(&token)
2692 .expect("text store")
2693 .get_document("local", note_id)
2694 .await
2695 .expect("get_document")
2696 .expect("document must be indexed after post-commit reindex");
2697 assert!(
2698 doc.body.contains("freshly-updated-content-xyz"),
2699 "FTS body must reflect the committed content: {:?}",
2700 doc.body
2701 );
2702
2703 assert_eq!(
2705 vec_store.count().await.expect("count after"),
2706 1,
2707 "post-commit reindex must have inserted a vector row for the stub model"
2708 );
2709 }
2710
2711 #[tokio::test]
2723 async fn atomic_note_update_and_delete_post_commit_effects_execute_exactly_once() {
2724 let runtime = scratch_runtime();
2725 let token = runtime
2726 .authorize(Namespace::parse("local").expect("ns"))
2727 .expect("authorize");
2728
2729 let fired: std::sync::Arc<std::sync::Mutex<Vec<(String, uuid::Uuid)>>> =
2730 std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2731 let fired_for_hook = fired.clone();
2732 runtime.install_note_mutation_hook(std::sync::Arc::new(
2733 move |kind: String, id: uuid::Uuid| {
2734 let fired = fired_for_hook.clone();
2735 Box::pin(async move {
2736 fired.lock().expect("lock").push((kind, id));
2737 })
2738 },
2739 ));
2740
2741 let mut note = khive_storage::note::Note::new("local", "observation", "hook-update-target");
2743 note.name = Some("hook-update-target".to_string());
2744 let update_note_id = note.id;
2745 runtime
2746 .notes(&token)
2747 .expect("notes store")
2748 .upsert_note(note)
2749 .await
2750 .expect("seed update-target note");
2751
2752 let plan = prepare_update(
2753 &runtime,
2754 &token,
2755 &json!({"id": update_note_id.to_string(), "content": "hook-update-target, revised"}),
2756 None,
2757 )
2758 .await
2759 .expect("prepare update");
2760 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
2761 .await
2762 .expect("seam call ok");
2763 let post_commit = match outcome {
2764 crate::atomic_runner::AtomicRunOutcome::Committed { post_commit } => post_commit,
2765 other => panic!("expected Committed, got {other:?}"),
2766 };
2767 apply_post_commit_effects(&runtime, &token, post_commit)
2768 .await
2769 .expect("apply post-commit effects (update)");
2770
2771 let mut del_note =
2774 khive_storage::note::Note::new("local", "observation", "hook-delete-target");
2775 del_note.name = Some("hook-delete-target".to_string());
2776 let delete_note_id = del_note.id;
2777 runtime
2778 .notes(&token)
2779 .expect("notes store")
2780 .upsert_note(del_note)
2781 .await
2782 .expect("seed delete-target note");
2783
2784 let plan = prepare_delete(
2785 &runtime,
2786 &token,
2787 &json!({"id": delete_note_id.to_string(), "hard": false}),
2788 None,
2789 )
2790 .await
2791 .expect("prepare delete");
2792 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
2793 .await
2794 .expect("seam call ok");
2795 let post_commit = match outcome {
2796 crate::atomic_runner::AtomicRunOutcome::Committed { post_commit } => post_commit,
2797 other => panic!("expected Committed, got {other:?}"),
2798 };
2799 assert_eq!(
2800 post_commit.as_slice(),
2801 &[PostCommitEffect::NoteDeleted {
2802 note_id: delete_note_id,
2803 kind: "observation".to_string(),
2804 }],
2805 "a committed note delete must schedule exactly one NoteDeleted post-commit effect"
2806 );
2807 apply_post_commit_effects(&runtime, &token, post_commit)
2808 .await
2809 .expect("apply post-commit effects (delete)");
2810
2811 assert_eq!(
2812 *fired.lock().expect("lock"),
2813 vec![
2814 ("observation".to_string(), update_note_id),
2815 ("observation".to_string(), delete_note_id),
2816 ],
2817 "each committed token must execute its note-mutation effect exactly once"
2818 );
2819 }
2820
2821 #[tokio::test]
2822 async fn failed_post_commit_reindexes_do_not_skip_later_note_mutation_hook() {
2823 let runtime = scratch_runtime();
2824 let token = runtime
2825 .authorize(Namespace::parse("local").expect("ns"))
2826 .expect("authorize");
2827 let fired: std::sync::Arc<std::sync::Mutex<Vec<uuid::Uuid>>> =
2828 std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2829 let fired_for_hook = fired.clone();
2830 runtime.install_note_mutation_hook(std::sync::Arc::new(
2831 move |_kind: String, id: uuid::Uuid| {
2832 let fired = fired_for_hook.clone();
2833 Box::pin(async move {
2834 fired.lock().expect("lock").push(id);
2835 })
2836 },
2837 ));
2838
2839 let mut plans = Vec::new();
2840 let mut failed_ids = Vec::new();
2841 for name in ["first", "second"] {
2842 let note = khive_storage::note::Note::new("local", "observation", name);
2843 let id = note.id;
2844 runtime
2845 .notes(&token)
2846 .expect("notes store")
2847 .upsert_note(note)
2848 .await
2849 .expect("seed reindex target");
2850 plans.push(
2851 prepare_update(
2852 &runtime,
2853 &token,
2854 &json!({"id": id.to_string(), "content": format!("{name} revised"), "embed": true}),
2859 None,
2860 )
2861 .await
2862 .expect("prepare note update"),
2863 );
2864 failed_ids.push(id);
2865 }
2866 let deleted = khive_storage::note::Note::new("local", "observation", "deleted");
2867 let deleted_id = deleted.id;
2868 runtime
2869 .notes(&token)
2870 .expect("notes store")
2871 .upsert_note(deleted)
2872 .await
2873 .expect("seed delete target");
2874 plans.push(
2875 prepare_delete(
2876 &runtime,
2877 &token,
2878 &json!({"id": deleted_id.to_string(), "hard": false}),
2879 None,
2880 )
2881 .await
2882 .expect("prepare note delete"),
2883 );
2884
2885 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), plans)
2886 .await
2887 .expect("commit atomic unit");
2888 let post_commit = match outcome {
2889 crate::atomic_runner::AtomicRunOutcome::Committed { post_commit } => post_commit,
2890 other => panic!("expected Committed, got {other:?}"),
2891 };
2892 assert_eq!(
2893 post_commit.as_slice(),
2894 &[
2895 PostCommitEffect::ReindexNote {
2896 note_id: failed_ids[0],
2897 version: 2,
2898 },
2899 PostCommitEffect::ReindexNote {
2900 note_id: failed_ids[1],
2901 version: 2,
2902 },
2903 PostCommitEffect::NoteDeleted {
2904 note_id: deleted_id,
2905 kind: "observation".into(),
2906 },
2907 ],
2908 "the fixture must commit two reindexes before the deletion hook"
2909 );
2910
2911 let mut writer = runtime.sql().writer().await.expect("writer");
2914 writer
2915 .execute(SqlStatement {
2916 sql: "DROP TABLE fts_notes".into(),
2917 params: Vec::new(),
2918 label: Some("test_remove_fts_notes_after_commit".into()),
2919 })
2920 .await
2921 .expect("remove FTS table");
2922 drop(writer);
2923
2924 let error = apply_post_commit_effects_with_report(&runtime, &token, post_commit)
2925 .await
2926 .expect_err("both reindexes must fail")
2927 .to_string();
2928 assert!(error.contains("effect[0]"), "{error}");
2929 assert!(error.contains("effect[1]"), "{error}");
2930 for id in failed_ids {
2931 assert!(error.contains(&id.to_string()), "{error}");
2932 }
2933 assert_eq!(*fired.lock().expect("lock"), vec![deleted_id]);
2934 }
2935
2936 #[tokio::test]
2940 async fn atomic_delete_note_purges_fts_and_vector_indexes_soft_and_hard() {
2941 let runtime = scratch_runtime();
2942 runtime.register_embedder(StubProvider);
2943 let token = runtime
2944 .authorize(Namespace::parse("local").expect("ns"))
2945 .expect("authorize");
2946
2947 for hard in [false, true] {
2948 let mut note =
2949 khive_storage::note::Note::new("local", "observation", "purge-target content");
2950 note.name = Some(format!("purge-target-hard-{hard}"));
2951 let note_id = note.id;
2952 runtime
2953 .notes(&token)
2954 .expect("notes store")
2955 .upsert_note(note.clone())
2956 .await
2957 .expect("seed note");
2958 runtime
2959 .reindex_note(&token, ¬e)
2960 .await
2961 .expect("seed index rows");
2962
2963 let vec_store = runtime
2964 .vectors_for_model(&token, STUB_MODEL)
2965 .expect("vec store");
2966 assert_eq!(
2967 vec_store.count().await.expect("count before"),
2968 1,
2969 "seeded note must have a vector row before delete (hard={hard})"
2970 );
2971 assert!(
2972 runtime
2973 .text_for_notes(&token)
2974 .expect("text store")
2975 .get_document("local", note_id)
2976 .await
2977 .expect("get_document")
2978 .is_some(),
2979 "seeded note must have an FTS row before delete (hard={hard})"
2980 );
2981
2982 let plan = prepare_delete(
2983 &runtime,
2984 &token,
2985 &json!({"id": note_id.to_string(), "hard": hard}),
2986 None,
2987 )
2988 .await
2989 .expect("prepare delete");
2990 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
2991 .await
2992 .expect("seam call ok");
2993 assert!(
2994 matches!(
2995 outcome,
2996 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
2997 ),
2998 "expected commit (hard={hard}): {outcome:?}"
2999 );
3000
3001 assert!(
3002 runtime
3003 .text_for_notes(&token)
3004 .expect("text store")
3005 .get_document("local", note_id)
3006 .await
3007 .expect("get_document")
3008 .is_none(),
3009 "FTS row must be purged after atomic delete (hard={hard})"
3010 );
3011 assert_eq!(
3012 vec_store.count().await.expect("count after"),
3013 0,
3014 "vector row must be purged after atomic delete (hard={hard})"
3015 );
3016 }
3017 }
3018
3019 #[tokio::test]
3023 async fn atomic_delete_entity_purges_fts_and_vector_indexes_soft_and_hard() {
3024 let runtime = scratch_runtime();
3025 runtime.register_embedder(StubProvider);
3026 let token = runtime
3027 .authorize(Namespace::parse("local").expect("ns"))
3028 .expect("authorize");
3029
3030 for hard in [false, true] {
3031 let entity =
3032 khive_storage::Entity::new("local", "concept", format!("purge-target-hard-{hard}"));
3033 let entity_id = entity.id;
3034 runtime
3035 .entities(&token)
3036 .expect("entities store")
3037 .upsert_entity(entity.clone())
3038 .await
3039 .expect("seed entity");
3040 runtime
3041 .reindex_entity(&token, &entity)
3042 .await
3043 .expect("seed index rows");
3044
3045 let vec_store = runtime
3046 .vectors_for_model(&token, STUB_MODEL)
3047 .expect("vec store");
3048 assert_eq!(
3049 vec_store.count().await.expect("count before"),
3050 1,
3051 "seeded entity must have a vector row before delete (hard={hard})"
3052 );
3053 assert!(
3054 runtime
3055 .text(&token)
3056 .expect("text store")
3057 .get_document("local", entity_id)
3058 .await
3059 .expect("get_document")
3060 .is_some(),
3061 "seeded entity must have an FTS row before delete (hard={hard})"
3062 );
3063
3064 let plan = prepare_delete(
3065 &runtime,
3066 &token,
3067 &json!({"id": entity_id.to_string(), "hard": hard}),
3068 None,
3069 )
3070 .await
3071 .expect("prepare delete");
3072 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
3073 .await
3074 .expect("seam call ok");
3075 assert!(
3076 matches!(
3077 outcome,
3078 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3079 ),
3080 "expected commit (hard={hard}): {outcome:?}"
3081 );
3082
3083 assert!(
3084 runtime
3085 .text(&token)
3086 .expect("text store")
3087 .get_document("local", entity_id)
3088 .await
3089 .expect("get_document")
3090 .is_none(),
3091 "FTS row must be purged after atomic delete (hard={hard})"
3092 );
3093 assert_eq!(
3094 vec_store.count().await.expect("count after"),
3095 0,
3096 "vector row must be purged after atomic delete (hard={hard})"
3097 );
3098 }
3099 }
3100
3101 #[tokio::test]
3102 async fn atomic_delete_entity_and_note_logs_vector_delete_rows() {
3103 let runtime = scratch_runtime();
3104 runtime.register_embedder(StubProvider);
3105 let token = runtime
3106 .authorize(Namespace::parse("local").expect("ns"))
3107 .expect("authorize");
3108
3109 let entity = khive_storage::Entity::new("local", "concept", "ann-delete-entity");
3110 let entity_id = entity.id;
3111 runtime
3112 .entities(&token)
3113 .expect("entities store")
3114 .upsert_entity(entity.clone())
3115 .await
3116 .expect("seed entity");
3117 runtime
3118 .reindex_entity(&token, &entity)
3119 .await
3120 .expect("seed entity vector row");
3121 let note = khive_storage::note::Note::new("local", "observation", "ann-delete-note");
3122 let note_id = note.id;
3123 runtime
3124 .notes(&token)
3125 .expect("notes store")
3126 .upsert_note(note.clone())
3127 .await
3128 .expect("seed note");
3129 runtime
3130 .reindex_note(&token, ¬e)
3131 .await
3132 .expect("seed note vector row");
3133
3134 for id in [entity_id, note_id] {
3135 let plan = prepare_delete(&runtime, &token, &json!({"id": id.to_string()}), None)
3136 .await
3137 .expect("prepare delete");
3138 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
3139 .await
3140 .expect("seam call ok");
3141 assert!(matches!(
3142 outcome,
3143 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3144 ));
3145 }
3146
3147 let mut reader = runtime.sql().reader().await.expect("sql reader");
3148 let count = reader
3149 .query_scalar(SqlStatement {
3150 sql: "SELECT COUNT(*) FROM ann_write_log \
3151 WHERE namespace = 'local' AND embedding_model = ?1 AND op = 'delete' \
3152 AND ((kind = 'entity' AND field = 'entity.body' AND subject_id = ?2) \
3153 OR (kind = 'note' AND field = 'note.content' AND subject_id = ?3))"
3154 .to_string(),
3155 params: vec![
3156 SqlValue::Text(STUB_MODEL.to_string()),
3157 SqlValue::Text(entity_id.to_string()),
3158 SqlValue::Text(note_id.to_string()),
3159 ],
3160 label: Some("test-atomic-delete-ann-write-log".to_string()),
3161 })
3162 .await
3163 .expect("query ANN write log");
3164 assert!(matches!(count, Some(SqlValue::Integer(2))));
3165 }
3166
3167 #[tokio::test]
3171 async fn atomic_link_symmetric_rejection_preserves_requested_pair() {
3172 let runtime = scratch_runtime();
3173 let token = runtime
3174 .authorize(Namespace::parse("local").expect("ns"))
3175 .expect("authorize");
3176 let entities = runtime.entities(&token).expect("entities store");
3177 let concept_id =
3178 Uuid::parse_str("ffffffff-ffff-ffff-ffff-ffffffffffff").expect("high UUID");
3179 let project_id = Uuid::nil();
3180 assert!(project_id < concept_id, "test must exercise UUID reversal");
3181
3182 let mut concept = khive_storage::Entity::new("local", "concept", "Concept source");
3183 concept.id = concept_id;
3184 let mut project = khive_storage::Entity::new("local", "project", "Project target");
3185 project.id = project_id;
3186 entities.upsert_entity(concept).await.expect("seed concept");
3187 entities.upsert_entity(project).await.expect("seed project");
3188
3189 let error = prepare_link(
3190 &runtime,
3191 &token,
3192 &json!({
3193 "source_id": concept_id.to_string(),
3194 "target_id": project_id.to_string(),
3195 "relation": "competes_with",
3196 }),
3197 )
3198 .await
3199 .expect_err("atomic link must reject concept competes_with project");
3200 let message = error.to_string();
3201 assert!(
3202 message.contains(
3203 "currently legal relations for concept -> project under the loaded endpoint rules: none"
3204 ),
3205 "atomic validation must diagnose caller order before persistence canonicalization; got: {message}"
3206 );
3207 assert!(
3208 !message.contains("currently legal relations for project -> concept"),
3209 "atomic validation must not diagnose the UUID-canonical reverse pair; got: {message}"
3210 );
3211 }
3212
3213 #[tokio::test]
3219 async fn atomic_link_persists_explicit_dependency_kind_and_infers_when_absent() {
3220 let runtime = scratch_runtime();
3221 let token = runtime
3222 .authorize(Namespace::parse("local").expect("ns"))
3223 .expect("authorize");
3224 let entities = runtime.entities(&token).expect("entities store");
3225
3226 fn metadata_json(plan: &AtomicOpPlan) -> String {
3227 let link_plan = match plan {
3228 AtomicOpPlan::Link(p) => p,
3229 other => panic!("expected an AtomicOpPlan::Link, got {other:?}"),
3230 };
3231 match link_plan.statements[0].statement.params.last() {
3232 Some(SqlValue::Text(s)) => s.clone(),
3233 other => panic!("expected the metadata param to be SqlValue::Text, got {other:?}"),
3234 }
3235 }
3236
3237 {
3239 let svc = khive_storage::Entity::new("local", "service", "SvcA");
3240 let proj = khive_storage::Entity::new("local", "project", "ProjB");
3241 let (svc_id, proj_id) = (svc.id, proj.id);
3242 entities.upsert_entity(svc).await.expect("seed svc");
3243 entities.upsert_entity(proj).await.expect("seed proj");
3244
3245 let plan = prepare_link(
3246 &runtime,
3247 &token,
3248 &json!({
3249 "source_id": svc_id.to_string(),
3250 "target_id": proj_id.to_string(),
3251 "relation": "depends_on",
3252 "dependency_kind": "artifact",
3253 }),
3254 )
3255 .await
3256 .expect("prepare link");
3257 let json_str = metadata_json(&plan);
3258 assert!(
3259 json_str.contains(r#""dependency_kind":"artifact""#),
3260 "explicit dependency_kind param must persist: {json_str}"
3261 );
3262 }
3263
3264 {
3267 let svc_a = khive_storage::Entity::new("local", "service", "SvcC");
3268 let svc_b = khive_storage::Entity::new("local", "service", "SvcD");
3269 let (a_id, b_id) = (svc_a.id, svc_b.id);
3270 entities.upsert_entity(svc_a).await.expect("seed svc a");
3271 entities.upsert_entity(svc_b).await.expect("seed svc b");
3272
3273 let plan = prepare_link(
3274 &runtime,
3275 &token,
3276 &json!({
3277 "source_id": a_id.to_string(),
3278 "target_id": b_id.to_string(),
3279 "relation": "depends_on",
3280 }),
3281 )
3282 .await
3283 .expect("prepare link");
3284 let json_str = metadata_json(&plan);
3285 assert!(
3286 json_str.contains(r#""dependency_kind":"runtime""#),
3287 "inferred dependency_kind for (service, service) must persist: {json_str}"
3288 );
3289 }
3290 }
3291
3292 async fn probe_edge_natural_key(
3297 runtime: &KhiveRuntime,
3298 namespace: &str,
3299 source_id: Uuid,
3300 target_id: Uuid,
3301 relation: &str,
3302 ) -> (usize, Option<f64>, Option<String>, Option<i64>) {
3303 let mut reader = runtime.sql().reader().await.expect("reader");
3304 let rows = reader
3305 .query_all(SqlStatement {
3306 sql: "SELECT weight, metadata, deleted_at FROM graph_edges \
3307 WHERE namespace = ?1 AND source_id = ?2 AND target_id = ?3 AND relation = ?4"
3308 .to_string(),
3309 params: vec![
3310 SqlValue::Text(namespace.to_string()),
3311 SqlValue::Text(source_id.to_string()),
3312 SqlValue::Text(target_id.to_string()),
3313 SqlValue::Text(relation.to_string()),
3314 ],
3315 label: Some("test-probe-edge-natural-key".to_string()),
3316 })
3317 .await
3318 .expect("probe edge natural key");
3319 let count = rows.len();
3320 let Some(row) = rows.into_iter().next() else {
3321 return (count, None, None, None);
3322 };
3323 let weight = match row.get("weight") {
3324 Some(SqlValue::Float(f)) => Some(*f),
3325 Some(SqlValue::Integer(i)) => Some(*i as f64),
3326 _ => None,
3327 };
3328 let metadata = match row.get("metadata") {
3329 Some(SqlValue::Text(s)) => Some(s.clone()),
3330 _ => None,
3331 };
3332 let deleted_at = match row.get("deleted_at") {
3333 Some(SqlValue::Integer(i)) => Some(*i),
3334 _ => None,
3335 };
3336 (count, weight, metadata, deleted_at)
3337 }
3338
3339 #[tokio::test]
3345 async fn atomic_link_of_already_linked_triple_upserts_weight_and_metadata() {
3346 let runtime = scratch_runtime();
3347 let token = runtime
3348 .authorize(Namespace::parse("local").expect("ns"))
3349 .expect("authorize");
3350 let entities = runtime.entities(&token).expect("entities store");
3351
3352 let a = khive_storage::Entity::new("local", "concept", "GapTwoA");
3353 let b = khive_storage::Entity::new("local", "concept", "GapTwoB");
3354 let (a_id, b_id) = (a.id, b.id);
3355 entities.upsert_entity(a).await.expect("seed a");
3356 entities.upsert_entity(b).await.expect("seed b");
3357
3358 let plan1 = prepare_link(
3360 &runtime,
3361 &token,
3362 &json!({
3363 "source_id": a_id.to_string(),
3364 "target_id": b_id.to_string(),
3365 "relation": "extends",
3366 "weight": 0.5,
3367 }),
3368 )
3369 .await
3370 .expect("prepare first link");
3371 let outcome1 = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan1])
3372 .await
3373 .expect("seam call ok");
3374 assert!(
3375 matches!(
3376 outcome1,
3377 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3378 ),
3379 "fresh link must commit: {outcome1:?}"
3380 );
3381 let (count, weight, _metadata, deleted_at) =
3382 probe_edge_natural_key(&runtime, "local", a_id, b_id, "extends").await;
3383 assert_eq!(count, 1, "exactly one edge row after the fresh link");
3384 assert_eq!(weight, Some(0.5));
3385 assert!(deleted_at.is_none());
3386
3387 let plan2 = prepare_link(
3391 &runtime,
3392 &token,
3393 &json!({
3394 "source_id": a_id.to_string(),
3395 "target_id": b_id.to_string(),
3396 "relation": "extends",
3397 "weight": 0.9,
3398 "metadata": {"note": "relinked"},
3399 }),
3400 )
3401 .await
3402 .expect("prepare second link");
3403 let outcome2 = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan2])
3404 .await
3405 .expect("seam call ok");
3406 assert!(
3407 matches!(
3408 outcome2,
3409 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3410 ),
3411 "re-link of an already-linked triple must upsert, not roll back: {outcome2:?}"
3412 );
3413 let (count, weight, metadata, deleted_at) =
3414 probe_edge_natural_key(&runtime, "local", a_id, b_id, "extends").await;
3415 assert_eq!(
3416 count, 1,
3417 "the natural-key UNIQUE constraint must still hold exactly one row (upsert, not a second insert)"
3418 );
3419 assert_eq!(weight, Some(0.9), "weight must be updated to the new value");
3420 assert!(
3421 metadata
3422 .as_deref()
3423 .is_some_and(|m| m.contains(r#""note":"relinked""#)),
3424 "metadata must be updated to the new value: {metadata:?}"
3425 );
3426 assert!(deleted_at.is_none());
3427 }
3428
3429 #[tokio::test]
3432 async fn atomic_link_of_soft_deleted_triple_resurrects_it() {
3433 let runtime = scratch_runtime();
3434 let token = runtime
3435 .authorize(Namespace::parse("local").expect("ns"))
3436 .expect("authorize");
3437 let entities = runtime.entities(&token).expect("entities store");
3438
3439 let a = khive_storage::Entity::new("local", "concept", "GapTwoResurrectA");
3440 let b = khive_storage::Entity::new("local", "concept", "GapTwoResurrectB");
3441 let (a_id, b_id) = (a.id, b.id);
3442 entities.upsert_entity(a).await.expect("seed a");
3443 entities.upsert_entity(b).await.expect("seed b");
3444
3445 let plan = prepare_link(
3446 &runtime,
3447 &token,
3448 &json!({
3449 "source_id": a_id.to_string(),
3450 "target_id": b_id.to_string(),
3451 "relation": "extends",
3452 }),
3453 )
3454 .await
3455 .expect("prepare link");
3456 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
3457 .await
3458 .expect("seam call ok");
3459 assert!(matches!(
3460 outcome,
3461 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3462 ));
3463
3464 {
3467 let mut writer = runtime.sql().writer().await.expect("writer");
3468 let affected = writer
3469 .execute(SqlStatement {
3470 sql: "UPDATE graph_edges SET deleted_at = ?1 \
3471 WHERE namespace = ?2 AND source_id = ?3 AND target_id = ?4 AND relation = ?5"
3472 .to_string(),
3473 params: vec![
3474 SqlValue::Integer(chrono::Utc::now().timestamp_micros()),
3475 SqlValue::Text("local".to_string()),
3476 SqlValue::Text(a_id.to_string()),
3477 SqlValue::Text(b_id.to_string()),
3478 SqlValue::Text("extends".to_string()),
3479 ],
3480 label: Some("test-soft-delete-edge".to_string()),
3481 })
3482 .await
3483 .expect("soft delete edge");
3484 assert_eq!(affected, 1, "soft-delete must touch exactly the seeded row");
3485 }
3486 let (_, _, _, deleted_at) =
3487 probe_edge_natural_key(&runtime, "local", a_id, b_id, "extends").await;
3488 assert!(
3489 deleted_at.is_some(),
3490 "row must be soft-deleted before the resurrect attempt"
3491 );
3492
3493 let refusal = prepare_link(
3494 &runtime,
3495 &token,
3496 &json!({
3497 "source_id": a_id.to_string(),
3498 "target_id": b_id.to_string(),
3499 "relation": "extends",
3500 "weight": 0.75,
3501 }),
3502 )
3503 .await
3504 .expect_err("implicit resurrection must be refused at prepare time");
3505 assert!(matches!(
3506 refusal,
3507 RuntimeError::InvalidInput(message) if message.contains("resurrect=true")
3508 ));
3509 let (_, weight, _, deleted_at) =
3510 probe_edge_natural_key(&runtime, "local", a_id, b_id, "extends").await;
3511 assert_eq!(weight, Some(1.0), "refusal must preserve the tombstone");
3512 assert!(deleted_at.is_some());
3513
3514 let plan_relink = prepare_link(
3515 &runtime,
3516 &token,
3517 &json!({
3518 "source_id": a_id.to_string(),
3519 "target_id": b_id.to_string(),
3520 "relation": "extends",
3521 "weight": 0.75,
3522 "resurrect": true,
3523 }),
3524 )
3525 .await
3526 .expect("prepare explicitly resurrecting link");
3527 let outcome_relink =
3528 crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan_relink])
3529 .await
3530 .expect("seam call ok");
3531 assert!(
3532 matches!(
3533 outcome_relink,
3534 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3535 ),
3536 "re-linking a soft-deleted triple must resurrect it, not roll back: {outcome_relink:?}"
3537 );
3538 let (count, weight, _, deleted_at) =
3539 probe_edge_natural_key(&runtime, "local", a_id, b_id, "extends").await;
3540 assert_eq!(count, 1);
3541 assert_eq!(weight, Some(0.75));
3542 assert!(
3543 deleted_at.is_none(),
3544 "re-link must resurrect the soft-deleted row (deleted_at -> NULL)"
3545 );
3546 }
3547
3548 #[tokio::test]
3556 async fn atomic_delete_succeeds_when_vec_table_never_created() {
3557 let runtime = scratch_runtime();
3558 runtime.register_embedder(StubProvider);
3559 let token = runtime
3560 .authorize(Namespace::parse("local").expect("ns"))
3561 .expect("authorize");
3562
3563 let entity = khive_storage::Entity::new("local", "concept", "no-vec-table-entity");
3567 let entity_id = entity.id;
3568 runtime
3569 .entities(&token)
3570 .expect("entities store")
3571 .upsert_entity(entity)
3572 .await
3573 .expect("seed entity");
3574
3575 let mut note = khive_storage::note::Note::new("local", "observation", "no-vec-table-note");
3576 note.name = Some("no-vec-table-note".to_string());
3577 let note_id = note.id;
3578 runtime
3579 .notes(&token)
3580 .expect("notes store")
3581 .upsert_note(note)
3582 .await
3583 .expect("seed note");
3584
3585 for (id, kind) in [(entity_id, "entity"), (note_id, "note")] {
3586 let plan = prepare_delete(&runtime, &token, &json!({"id": id.to_string()}), None)
3587 .await
3588 .unwrap_or_else(|e| panic!("prepare delete ({kind}) must not fail: {e}"));
3589 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
3590 .await
3591 .unwrap_or_else(|e| {
3592 panic!("atomic delete ({kind}) must not hit `no such table`: {e}")
3593 });
3594 assert!(
3595 matches!(
3596 outcome,
3597 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3598 ),
3599 "expected a clean commit ({kind}): {outcome:?}"
3600 );
3601 }
3602
3603 assert!(
3604 runtime
3605 .get_entity_including_deleted(&token, entity_id)
3606 .await
3607 .expect("get entity")
3608 .expect("entity row still present (soft delete)")
3609 .deleted_at
3610 .is_some(),
3611 "entity must be soft-deleted"
3612 );
3613 assert!(
3614 runtime
3615 .get_note_including_deleted(&token, note_id)
3616 .await
3617 .expect("get note")
3618 .expect("note row still present (soft delete)")
3619 .deleted_at
3620 .is_some(),
3621 "note must be soft-deleted"
3622 );
3623 }
3624
3625 #[tokio::test]
3631 async fn atomic_hard_delete_purges_already_soft_deleted_entity_and_note() {
3632 let runtime = scratch_runtime();
3633 runtime.register_embedder(StubProvider);
3634 let token = runtime
3635 .authorize(Namespace::parse("local").expect("ns"))
3636 .expect("authorize");
3637
3638 let entity =
3639 khive_storage::Entity::new("local", "concept", "tombstoned-entity-hard-delete");
3640 let entity_id = entity.id;
3641 runtime
3642 .entities(&token)
3643 .expect("entities store")
3644 .upsert_entity(entity.clone())
3645 .await
3646 .expect("seed entity");
3647 runtime
3648 .reindex_entity(&token, &entity)
3649 .await
3650 .expect("seed index rows");
3651
3652 let mut note =
3653 khive_storage::note::Note::new("local", "observation", "tombstoned-note-hard-delete");
3654 note.name = Some("tombstoned-note-hard-delete".to_string());
3655 let note_id = note.id;
3656 runtime
3657 .notes(&token)
3658 .expect("notes store")
3659 .upsert_note(note.clone())
3660 .await
3661 .expect("seed note");
3662 runtime
3663 .reindex_note(&token, ¬e)
3664 .await
3665 .expect("seed index rows");
3666
3667 for id in [entity_id, note_id] {
3670 let plan = prepare_delete(&runtime, &token, &json!({"id": id.to_string()}), None)
3671 .await
3672 .expect("prepare soft delete");
3673 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
3674 .await
3675 .expect("soft delete commit");
3676 assert!(matches!(
3677 outcome,
3678 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3679 ));
3680 }
3681 assert!(
3682 runtime
3683 .get_entity_including_deleted(&token, entity_id)
3684 .await
3685 .expect("get entity")
3686 .expect("entity present")
3687 .deleted_at
3688 .is_some(),
3689 "entity must be soft-deleted before the hard-delete attempt"
3690 );
3691 assert!(
3692 runtime
3693 .get_note_including_deleted(&token, note_id)
3694 .await
3695 .expect("get note")
3696 .expect("note present")
3697 .deleted_at
3698 .is_some(),
3699 "note must be soft-deleted before the hard-delete attempt"
3700 );
3701
3702 for (id, kind) in [(entity_id, "entity"), (note_id, "note")] {
3704 let plan = prepare_delete(
3705 &runtime,
3706 &token,
3707 &json!({"id": id.to_string(), "hard": true}),
3708 None,
3709 )
3710 .await
3711 .unwrap_or_else(|e| {
3712 panic!(
3713 "prepare hard delete ({kind}) of an already-soft-deleted record \
3714 must resolve it: {e}"
3715 )
3716 });
3717 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
3718 .await
3719 .unwrap_or_else(|e| panic!("hard delete ({kind}) commit failed: {e}"));
3720 assert!(
3721 matches!(
3722 outcome,
3723 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3724 ),
3725 "expected a clean hard-delete commit ({kind}): {outcome:?}"
3726 );
3727 }
3728
3729 assert!(
3730 runtime
3731 .get_entity_including_deleted(&token, entity_id)
3732 .await
3733 .expect("get entity")
3734 .is_none(),
3735 "entity row must be fully purged after hard delete"
3736 );
3737 assert!(
3738 runtime
3739 .get_note_including_deleted(&token, note_id)
3740 .await
3741 .expect("get note")
3742 .is_none(),
3743 "note row must be fully purged after hard delete"
3744 );
3745 assert!(
3746 runtime
3747 .text(&token)
3748 .expect("text store")
3749 .get_document("local", entity_id)
3750 .await
3751 .expect("get_document")
3752 .is_none(),
3753 "entity FTS row must be purged after hard delete"
3754 );
3755 assert!(
3756 runtime
3757 .text_for_notes(&token)
3758 .expect("text store")
3759 .get_document("local", note_id)
3760 .await
3761 .expect("get_document")
3762 .is_none(),
3763 "note FTS row must be purged after hard delete"
3764 );
3765 let vec_store = runtime
3766 .vectors_for_model(&token, STUB_MODEL)
3767 .expect("vec store");
3768 assert_eq!(
3769 vec_store.count().await.expect("count after"),
3770 0,
3771 "vector rows for both records must be purged after hard delete"
3772 );
3773 }
3774
3775 async fn events_for_target(
3783 runtime: &KhiveRuntime,
3784 token: &NamespaceToken,
3785 target_id: Uuid,
3786 kind: EventKind,
3787 ) -> Vec<khive_storage::Event> {
3788 let event_store = runtime.events(token).expect("event store");
3789 let filter = khive_storage::EventFilter {
3790 kinds: vec![kind],
3791 ..Default::default()
3792 };
3793 let page = event_store
3794 .query_events(filter, khive_storage::types::PageRequest::default())
3795 .await
3796 .expect("query_events");
3797 page.items
3798 .into_iter()
3799 .filter(|e| e.target_id == Some(target_id))
3800 .collect()
3801 }
3802
3803 #[tokio::test]
3808 async fn atomic_update_entity_appends_entity_updated_event() {
3809 let runtime = scratch_runtime();
3810 let token = runtime
3811 .authorize(Namespace::parse("local").expect("ns"))
3812 .expect("authorize");
3813 let entity = khive_storage::Entity::new("local", "concept", "gap1-entity");
3814 let entity_id = entity.id;
3815 runtime
3816 .entities(&token)
3817 .expect("entities store")
3818 .upsert_entity(entity)
3819 .await
3820 .expect("seed entity");
3821
3822 let plan = prepare_update(
3823 &runtime,
3824 &token,
3825 &json!({"id": entity_id.to_string(), "name": "gap1-entity-renamed"}),
3826 None,
3827 )
3828 .await
3829 .expect("prepare update");
3830 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
3831 .await
3832 .expect("seam call ok");
3833 assert!(matches!(
3834 outcome,
3835 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3836 ));
3837
3838 let events = events_for_target(&runtime, &token, entity_id, EventKind::EntityUpdated).await;
3839 assert_eq!(
3840 events.len(),
3841 1,
3842 "expected exactly one EntityUpdated event for {entity_id}"
3843 );
3844 assert_eq!(events[0].namespace, "local");
3845 assert_eq!(
3846 events[0].actor, "anonymous:local",
3847 "atomic event attribution must come from the authorized token"
3848 );
3849 assert_eq!(events[0].payload["id"], json!(entity_id.to_string()));
3850 assert_eq!(
3851 events[0].payload["changed_fields"],
3852 json!(["name"]),
3853 "changed_fields must name exactly the patched fields"
3854 );
3855 }
3856
3857 #[tokio::test]
3861 async fn atomic_delete_entity_appends_entity_deleted_event_soft_and_hard() {
3862 let runtime = scratch_runtime();
3863 let token = runtime
3864 .authorize(Namespace::parse("local").expect("ns"))
3865 .expect("authorize");
3866
3867 for hard in [false, true] {
3868 let entity =
3869 khive_storage::Entity::new("local", "concept", format!("gap1-entity-hard-{hard}"));
3870 let entity_id = entity.id;
3871 runtime
3872 .entities(&token)
3873 .expect("entities store")
3874 .upsert_entity(entity)
3875 .await
3876 .expect("seed entity");
3877
3878 let args = if hard {
3879 json!({"id": entity_id.to_string(), "hard": true})
3880 } else {
3881 json!({"id": entity_id.to_string()})
3882 };
3883 let plan = prepare_delete(&runtime, &token, &args, None)
3884 .await
3885 .unwrap_or_else(|e| panic!("prepare delete (hard={hard}): {e}"));
3886 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
3887 .await
3888 .unwrap_or_else(|e| panic!("delete commit (hard={hard}): {e}"));
3889 assert!(
3890 matches!(
3891 outcome,
3892 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3893 ),
3894 "expected a clean delete commit (hard={hard}): {outcome:?}"
3895 );
3896
3897 let events =
3898 events_for_target(&runtime, &token, entity_id, EventKind::EntityDeleted).await;
3899 assert_eq!(
3900 events.len(),
3901 1,
3902 "expected exactly one EntityDeleted event for {entity_id} (hard={hard})"
3903 );
3904 assert_eq!(events[0].payload["hard"], json!(hard));
3905 }
3906 }
3907
3908 #[tokio::test]
3909 async fn atomic_hard_delete_emits_lineage_warning_in_commit_unit() {
3910 let runtime = scratch_runtime();
3911 let token = NamespaceToken::mint_authorized(
3912 Namespace::local(),
3913 crate::ActorRef::new("agent", "atomic-lineage-deleter"),
3914 );
3915 let doomed = khive_storage::Entity::new("local", "document", "atomic-doomed");
3916 let source = khive_storage::Entity::new("local", "artifact", "atomic-source");
3917 let doomed_id = doomed.id;
3918 runtime
3919 .entities(&token)
3920 .expect("entities store")
3921 .upsert_entity(doomed)
3922 .await
3923 .expect("seed doomed entity");
3924 runtime
3925 .entities(&token)
3926 .expect("entities store")
3927 .upsert_entity(source.clone())
3928 .await
3929 .expect("seed source entity");
3930 runtime
3931 .link(
3932 &token,
3933 source.id,
3934 doomed_id,
3935 EdgeRelation::DerivedFrom,
3936 1.0,
3937 None,
3938 )
3939 .await
3940 .expect("seed protected edge");
3941
3942 let plan = prepare_delete(
3943 &runtime,
3944 &token,
3945 &json!({"id": doomed_id.to_string(), "hard": true}),
3946 None,
3947 )
3948 .await
3949 .expect("prepare hard delete");
3950 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
3951 .await
3952 .expect("delete commit");
3953 assert!(matches!(
3954 outcome,
3955 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
3956 ));
3957
3958 let warnings = events_for_target(&runtime, &token, doomed_id, EventKind::Audit).await;
3959 assert_eq!(warnings.len(), 1);
3960 assert_eq!(warnings[0].actor, "agent:atomic-lineage-deleter");
3961 assert_eq!(warnings[0].payload["relation"], "derived_from");
3962 assert_eq!(warnings[0].payload["warning"], "provenance_loss");
3963 }
3964
3965 #[tokio::test]
3969 async fn atomic_delete_note_appends_note_deleted_event_soft_and_hard() {
3970 let runtime = scratch_runtime();
3971 let token = runtime
3972 .authorize(Namespace::parse("local").expect("ns"))
3973 .expect("authorize");
3974
3975 for hard in [false, true] {
3976 let mut note = khive_storage::note::Note::new(
3977 "local",
3978 "observation",
3979 format!("gap1-note-content-hard-{hard}"),
3980 );
3981 note.name = Some(format!("gap1-note-hard-{hard}"));
3982 let note_id = note.id;
3983 runtime
3984 .notes(&token)
3985 .expect("notes store")
3986 .upsert_note(note)
3987 .await
3988 .expect("seed note");
3989
3990 let args = if hard {
3991 json!({"id": note_id.to_string(), "hard": true})
3992 } else {
3993 json!({"id": note_id.to_string()})
3994 };
3995 let plan = prepare_delete(&runtime, &token, &args, None)
3996 .await
3997 .unwrap_or_else(|e| panic!("prepare delete (hard={hard}): {e}"));
3998 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
3999 .await
4000 .unwrap_or_else(|e| panic!("delete commit (hard={hard}): {e}"));
4001 assert!(
4002 matches!(
4003 outcome,
4004 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
4005 ),
4006 "expected a clean delete commit (hard={hard}): {outcome:?}"
4007 );
4008
4009 let events = events_for_target(&runtime, &token, note_id, EventKind::NoteDeleted).await;
4010 assert_eq!(
4011 events.len(),
4012 1,
4013 "expected exactly one NoteDeleted event for {note_id} (hard={hard})"
4014 );
4015 assert_eq!(events[0].payload["hard"], json!(hard));
4016 }
4017 }
4018
4019 #[tokio::test]
4026 async fn atomic_update_edge_patches_weight_and_appends_edge_updated_event() {
4027 let runtime = scratch_runtime();
4028 let token = runtime
4029 .authorize(Namespace::parse("local").expect("ns"))
4030 .expect("authorize");
4031 let entities = runtime.entities(&token).expect("entities store");
4032 let a = khive_storage::Entity::new("local", "concept", "GapEdgeA");
4033 let b = khive_storage::Entity::new("local", "concept", "GapEdgeB");
4034 let (a_id, b_id) = (a.id, b.id);
4035 entities.upsert_entity(a).await.expect("seed a");
4036 entities.upsert_entity(b).await.expect("seed b");
4037
4038 let edge = runtime
4039 .link(&token, a_id, b_id, EdgeRelation::Extends, 0.4, None)
4040 .await
4041 .expect("seed edge");
4042 let edge_id = Uuid::from(edge.id);
4043
4044 let plan = prepare_update(
4045 &runtime,
4046 &token,
4047 &json!({"id": edge_id.to_string(), "weight": 0.75}),
4048 None,
4049 )
4050 .await
4051 .expect("prepare update edge");
4052 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
4053 .await
4054 .expect("seam call ok");
4055 assert!(
4056 matches!(
4057 outcome,
4058 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
4059 ),
4060 "expected a clean edge update commit: {outcome:?}"
4061 );
4062
4063 let updated = runtime
4064 .get_edge(&token, edge_id)
4065 .await
4066 .expect("get_edge")
4067 .expect("edge still present");
4068 assert_eq!(updated.weight, 0.75, "weight patch must persist");
4069 assert_eq!(updated.relation, EdgeRelation::Extends);
4070
4071 let events = events_for_target(&runtime, &token, edge_id, EventKind::EdgeUpdated).await;
4072 assert_eq!(
4073 events.len(),
4074 1,
4075 "expected exactly one EdgeUpdated event for {edge_id}"
4076 );
4077 assert_eq!(
4078 events[0].payload["changed_fields"],
4079 json!(["weight"]),
4080 "changed_fields must name exactly the patched field"
4081 );
4082 }
4083
4084 #[tokio::test]
4088 async fn atomic_update_edge_rejects_reserved_secret_gate_property() {
4089 let runtime = scratch_runtime();
4090 let token = runtime
4091 .authorize(Namespace::parse("local").expect("ns"))
4092 .expect("authorize");
4093 let entities = runtime.entities(&token).expect("entities store");
4094 let a = khive_storage::Entity::new("local", "concept", "ReservedEdgeA");
4095 let b = khive_storage::Entity::new("local", "concept", "ReservedEdgeB");
4096 let (a_id, b_id) = (a.id, b.id);
4097 entities.upsert_entity(a).await.expect("seed a");
4098 entities.upsert_entity(b).await.expect("seed b");
4099
4100 let edge = runtime
4101 .link(&token, a_id, b_id, EdgeRelation::Extends, 0.4, None)
4102 .await
4103 .expect("seed edge");
4104 let edge_id = Uuid::from(edge.id);
4105
4106 let err = prepare_update(
4107 &runtime,
4108 &token,
4109 &json!({
4110 "id": edge_id.to_string(),
4111 "properties": {"khive:secret_gate": "exempted:content-sha256-manifest-v1"},
4112 }),
4113 None,
4114 )
4115 .await
4116 .expect_err("a caller-supplied reserved key on edge metadata must be rejected");
4117 assert!(
4118 matches!(err, RuntimeError::InvalidInput(ref msg) if msg.contains("khive:secret_gate") && msg.contains("runtime-owned")),
4119 "expected a reservation rejection, got: {err:?}"
4120 );
4121
4122 let unchanged = runtime
4123 .get_edge(&token, edge_id)
4124 .await
4125 .expect("get_edge")
4126 .expect("edge still present");
4127 assert!(
4128 unchanged.metadata.is_none(),
4129 "rejected edge update must leave metadata untouched"
4130 );
4131 }
4132
4133 #[tokio::test]
4143 async fn atomic_update_edge_symmetric_conflict_absorbs_into_surviving_row() {
4144 let runtime = scratch_runtime();
4145 let token = runtime
4146 .authorize(Namespace::parse("local").expect("ns"))
4147 .expect("authorize");
4148 let entities = runtime.entities(&token).expect("entities store");
4149 let a = khive_storage::Entity::new("local", "concept", "GapEdgeSymA");
4150 let b = khive_storage::Entity::new("local", "concept", "GapEdgeSymB");
4151 let (a_id, b_id) = (a.id, b.id);
4152 entities.upsert_entity(a).await.expect("seed a");
4153 entities.upsert_entity(b).await.expect("seed b");
4154
4155 let requested_edge = runtime
4157 .link(&token, a_id, b_id, EdgeRelation::Extends, 0.2, None)
4158 .await
4159 .expect("seed requested edge");
4160 let requested_id = Uuid::from(requested_edge.id);
4161
4162 let canonical_edge = runtime
4165 .link(&token, a_id, b_id, EdgeRelation::CompetesWith, 0.6, None)
4166 .await
4167 .expect("seed canonical edge");
4168 let canonical_id = Uuid::from(canonical_edge.id);
4169 assert_ne!(requested_id, canonical_id);
4170
4171 let plan = prepare_update(
4172 &runtime,
4173 &token,
4174 &json!({"id": requested_id.to_string(), "relation": "competes_with", "weight": 0.9}),
4175 None,
4176 )
4177 .await
4178 .expect("prepare update edge (symmetric conflict)");
4179 let (canon_src, canon_tgt) =
4186 canonical_edge_endpoints(EdgeRelation::CompetesWith, a_id, b_id);
4187 match &plan {
4188 AtomicOpPlan::Update(p) => {
4189 assert_eq!(p.target_id, requested_id);
4190 let key = p
4191 .edge_natural_key
4192 .as_ref()
4193 .expect("symmetric edge update must carry edge_natural_key");
4194 assert_eq!(key.canon_source_id, canon_src);
4195 assert_eq!(key.canon_target_id, canon_tgt);
4196 assert_eq!(key.relation, EdgeRelation::CompetesWith);
4197 }
4198 other => panic!("expected an Update plan, got {other:?}"),
4199 }
4200
4201 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
4202 .await
4203 .expect("seam call ok");
4204 assert!(
4205 matches!(
4206 outcome,
4207 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
4208 ),
4209 "expected a clean symmetric-conflict-absorption commit: {outcome:?}"
4210 );
4211
4212 let requested_after = runtime
4214 .get_edge_including_deleted(&token, requested_id)
4215 .await
4216 .expect("get_edge_including_deleted");
4217 assert!(
4218 requested_after.is_none(),
4219 "the non-canonical requested row must be deleted, not just tombstoned"
4220 );
4221
4222 let surviving = runtime
4226 .get_edge(&token, canonical_id)
4227 .await
4228 .expect("get_edge")
4229 .expect("surviving canonical row must remain");
4230 assert_eq!(
4231 surviving.weight, 0.6,
4232 "survivor weight must not be overwritten by the discarded edge's patch"
4233 );
4234 assert_eq!(surviving.relation, EdgeRelation::CompetesWith);
4235
4236 let events =
4240 events_for_target(&runtime, &token, requested_id, EventKind::EdgeUpdated).await;
4241 assert_eq!(events.len(), 1);
4242 }
4243
4244 #[tokio::test]
4257 async fn atomic_entity_update_plan_stale_revision_rolls_back_unit() {
4258 let runtime = scratch_runtime();
4259 let token = runtime
4260 .authorize(Namespace::parse("local").expect("ns"))
4261 .expect("authorize");
4262 let entity = runtime
4263 .create_entity(
4264 &token,
4265 "concept",
4266 None,
4267 "StaleEntityPlanTarget",
4268 None,
4269 Some(json!({"a": 0})),
4270 vec![],
4271 )
4272 .await
4273 .expect("seed entity");
4274 let id = entity.id;
4275
4276 let mut plan = prepare_update(
4279 &runtime,
4280 &token,
4281 &json!({"id": id.to_string(), "description": "from the stale plan"}),
4282 None,
4283 )
4284 .await
4285 .expect("prepare update entity plan");
4286
4287 let (planned_replacement, plan_target_id, plan_expected, plan_expected_deleted) = {
4294 let statements = match &plan {
4295 AtomicOpPlan::Update(p) => p.statements.clone(),
4296 other => panic!("expected an Update plan, got {other:?}"),
4297 };
4298 let cas = statements
4303 .iter()
4304 .find(|s| s.statement.label.as_deref() == Some("entity-replace-if-unchanged"))
4305 .expect("the plan must carry the guarded entity replacement");
4306 let read = |i: usize| match &cas.statement.params[i] {
4307 SqlValue::Integer(v) => *v,
4308 other => panic!("param {i} must be an integer revision, got {other:?}"),
4309 };
4310 let read_marker = |i: usize| match &cas.statement.params[i] {
4311 SqlValue::Null => None,
4312 SqlValue::Integer(v) => Some(*v),
4313 other => panic!("param {i} must be a deletion marker, got {other:?}"),
4314 };
4315 let read_text = |i: usize| match &cas.statement.params[i] {
4316 SqlValue::Text(v) => v.clone(),
4317 other => panic!("param {i} must be a text id, got {other:?}"),
4318 };
4319 (read(7), read_text(11), read(12), read_marker(13))
4320 };
4321 assert_eq!(
4327 plan_target_id,
4328 id.to_string(),
4329 "fixture premise: the plan's `?12` must be the row under test, otherwise \
4330 `id = ?12` refuses on identity and the rollback stops being attributable to \
4331 the expected-revision guard"
4332 );
4333
4334 runtime
4337 .update_entity(
4338 &token,
4339 id,
4340 crate::curation::EntityPatch {
4341 name: Some("ConcurrentWriterWon".to_string()),
4342 ..Default::default()
4343 },
4344 )
4345 .await
4346 .expect("concurrent writer update");
4347
4348 let stored_pinned = planned_replacement - 1;
4359 assert_ne!(
4360 stored_pinned, plan_expected,
4361 "fixture premise: the pinned revision must differ from the plan's expected \
4362 revision, otherwise `updated_at = ?13` would MATCH and nothing would refuse \
4363 the plan"
4364 );
4365 {
4366 let mut writer = runtime.sql().writer().await.expect("writer");
4367 let affected = writer
4368 .execute(SqlStatement {
4369 sql: "UPDATE entities SET version = version + 1, updated_at = ?1 WHERE id = ?2"
4370 .to_string(),
4371 params: vec![
4372 SqlValue::Integer(stored_pinned),
4373 SqlValue::Text(id.to_string()),
4374 ],
4375 label: Some("test-pin-stored-revision".to_string()),
4376 })
4377 .await
4378 .expect("pin the stored revision");
4379 assert_eq!(affected, 1, "the pin must touch exactly the seeded row");
4380 }
4381 assert!(
4382 planned_replacement > stored_pinned,
4383 "fixture premise: the plan's replacement must still strictly advance past the \
4384 stored revision, otherwise `?8 > updated_at` would refuse and this stops being \
4385 a test of the expected-revision guard"
4386 );
4387 {
4393 let stored = runtime
4394 .get_entity_including_deleted(&token, id)
4395 .await
4396 .expect("read the stored row")
4397 .expect("the seeded row is present before the plan runs");
4398 assert_eq!(
4399 stored.updated_at, stored_pinned,
4400 "fixture premise: the pin must be what the guard reads, so the stored \
4401 revision is the pinned value and nothing re-advanced it"
4402 );
4403 assert_eq!(
4404 stored.deleted_at, plan_expected_deleted,
4405 "fixture premise: the stored deletion marker must MATCH the plan's `?14`, \
4406 otherwise `deleted_at IS ?14` refuses too and this stops being a test of \
4407 the expected-revision guard alone"
4408 );
4409 let AtomicOpPlan::Update(update) = &mut plan else {
4414 unreachable!()
4415 };
4416 let cas = update
4417 .statements
4418 .iter_mut()
4419 .find(|s| s.statement.label.as_deref() == Some("entity-replace-if-unchanged"))
4420 .unwrap();
4421 cas.statement.params[14] = SqlValue::Integer(stored.version);
4422 }
4423
4424 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
4425 .await
4426 .expect("the seam call itself must not error; the unit rolls back cleanly");
4427 match outcome {
4428 crate::atomic_runner::AtomicRunOutcome::RolledBack {
4429 failed_op_index,
4430 failure,
4431 } => {
4432 assert_eq!(
4433 failed_op_index, 0,
4434 "the sole op's guard must be the one that fails"
4435 );
4436 assert_eq!(
4443 failure,
4444 crate::atomic_runner::AtomicOpFailure::GuardFailed {
4445 statement_label: Some("entity-replace-if-unchanged".to_string()),
4446 expected: crate::atomic_plan::AffectedRowGuard::exactly(1),
4447 observed: 0,
4448 },
4449 "the guarded entity replacement must be the statement whose guard refused"
4450 );
4451 }
4452 other => panic!(
4453 "a stale entity plan must roll back, not silently overwrite the concurrent \
4454 writer's change: {other:?}"
4455 ),
4456 }
4457
4458 let after = runtime.get_entity(&token, id).await.expect("get_entity");
4459 assert_eq!(
4460 after.name, "ConcurrentWriterWon",
4461 "the concurrent writer's committed name must survive the rolled-back stale plan"
4462 );
4463 assert_eq!(
4464 after.description, None,
4465 "the stale plan's description patch must NOT have landed"
4466 );
4467 }
4468
4469 #[tokio::test]
4474 async fn atomic_edge_update_plan_stale_revision_rolls_back_unit() {
4475 let runtime = scratch_runtime();
4476 let token = runtime
4477 .authorize(Namespace::parse("local").expect("ns"))
4478 .expect("authorize");
4479 let entities = runtime.entities(&token).expect("entities store");
4480 let a = khive_storage::Entity::new("local", "concept", "StaleEdgePlanA");
4481 let b = khive_storage::Entity::new("local", "concept", "StaleEdgePlanB");
4482 let (a_id, b_id) = (a.id, b.id);
4483 entities.upsert_entity(a).await.expect("seed a");
4484 entities.upsert_entity(b).await.expect("seed b");
4485
4486 let edge = runtime
4487 .link(&token, a_id, b_id, EdgeRelation::Extends, 0.2, None)
4488 .await
4489 .expect("seed edge");
4490 let edge_id = Uuid::from(edge.id);
4491
4492 let plan = prepare_update(
4495 &runtime,
4496 &token,
4497 &json!({"id": edge_id.to_string(), "properties": {"note": "from the stale plan"}}),
4498 None,
4499 )
4500 .await
4501 .expect("prepare update edge plan");
4502
4503 let (planned_replacement, plan_target_id, plan_expected, plan_expected_deleted) = {
4510 let statements = match &plan {
4511 AtomicOpPlan::Update(p) => p.statements.clone(),
4512 other => panic!("expected an Update plan, got {other:?}"),
4513 };
4514 let cas = statements
4518 .iter()
4519 .find(|s| s.statement.label.as_deref() == Some("edge-replace-if-unchanged"))
4520 .expect("the plan must carry the guarded edge replacement");
4521 let read = |i: usize| match &cas.statement.params[i] {
4522 SqlValue::Integer(v) => *v,
4523 other => panic!("param {i} must be an integer revision, got {other:?}"),
4524 };
4525 let read_marker = |i: usize| match &cas.statement.params[i] {
4526 SqlValue::Null => None,
4527 SqlValue::Integer(v) => Some(*v),
4528 other => panic!("param {i} must be a deletion marker, got {other:?}"),
4529 };
4530 let read_text = |i: usize| match &cas.statement.params[i] {
4531 SqlValue::Text(v) => v.clone(),
4532 other => panic!("param {i} must be a text id, got {other:?}"),
4533 };
4534 (read(5), read_text(9), read(10), read_marker(11))
4535 };
4536 assert_eq!(
4540 plan_target_id,
4541 edge_id.to_string(),
4542 "fixture premise: the plan's `?10` must be the edge under test, otherwise \
4543 `id = ?10` refuses on identity and the rollback stops being attributable to \
4544 the expected-revision guard"
4545 );
4546
4547 runtime
4550 .update_edge(
4551 &token,
4552 edge_id,
4553 crate::curation::EdgePatch {
4554 weight: Some(0.75),
4555 ..Default::default()
4556 },
4557 )
4558 .await
4559 .expect("concurrent writer update");
4560
4561 let stored_pinned = planned_replacement - 1;
4569 assert_ne!(
4570 stored_pinned, plan_expected,
4571 "fixture premise: the pinned revision must differ from the plan's expected \
4572 revision, otherwise `updated_at = ?11` would MATCH and nothing would refuse \
4573 the plan"
4574 );
4575 {
4576 let mut writer = runtime.sql().writer().await.expect("writer");
4577 let affected = writer
4578 .execute(SqlStatement {
4579 sql: "UPDATE graph_edges SET updated_at = ?1 WHERE id = ?2".to_string(),
4580 params: vec![
4581 SqlValue::Integer(stored_pinned),
4582 SqlValue::Text(edge_id.to_string()),
4583 ],
4584 label: Some("test-pin-stored-revision".to_string()),
4585 })
4586 .await
4587 .expect("pin the stored revision");
4588 assert_eq!(affected, 1, "the pin must touch exactly the seeded edge");
4589 }
4590 assert!(
4591 planned_replacement > stored_pinned,
4592 "fixture premise: the plan's replacement must still strictly advance past the \
4593 stored revision, otherwise `?6 > updated_at` would refuse and this stops being \
4594 a test of the expected-revision guard"
4595 );
4596 {
4600 let stored = runtime
4601 .get_edge_including_deleted(&token, edge_id)
4602 .await
4603 .expect("read the stored edge")
4604 .expect("the seeded edge is present before the plan runs");
4605 assert_eq!(
4606 stored.updated_at.timestamp_micros(),
4607 stored_pinned,
4608 "fixture premise: the pin must be what the guard reads, so the stored \
4609 revision is the pinned value and nothing re-advanced it"
4610 );
4611 assert_eq!(
4612 stored.deleted_at.map(|d| d.timestamp_micros()),
4613 plan_expected_deleted,
4614 "fixture premise: the stored deletion marker must MATCH the plan's `?12`, \
4615 otherwise `deleted_at IS ?12` refuses too and this stops being a test of \
4616 the expected-revision guard alone"
4617 );
4618 }
4619
4620 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
4621 .await
4622 .expect("the seam call itself must not error; the unit rolls back cleanly");
4623 match outcome {
4624 crate::atomic_runner::AtomicRunOutcome::RolledBack {
4625 failed_op_index,
4626 failure,
4627 } => {
4628 assert_eq!(
4629 failed_op_index, 0,
4630 "the sole op's guard must be the one that fails"
4631 );
4632 assert_eq!(
4634 failure,
4635 crate::atomic_runner::AtomicOpFailure::GuardFailed {
4636 statement_label: Some("edge-replace-if-unchanged".to_string()),
4637 expected: crate::atomic_plan::AffectedRowGuard::exactly(1),
4638 observed: 0,
4639 },
4640 "the guarded edge replacement must be the statement whose guard refused"
4641 );
4642 }
4643 other => panic!(
4644 "a stale edge plan must roll back, not silently overwrite the concurrent \
4645 writer's change: {other:?}"
4646 ),
4647 }
4648
4649 let after = runtime
4650 .get_edge(&token, edge_id)
4651 .await
4652 .expect("get_edge")
4653 .expect("edge still exists");
4654 assert!(
4655 (after.weight - 0.75).abs() < 0.001,
4656 "the concurrent writer's committed weight must survive the rolled-back stale plan: {after:?}"
4657 );
4658 assert!(
4659 after.metadata.is_none(),
4660 "the stale plan's properties patch must NOT have landed: {after:?}"
4661 );
4662 }
4663
4664 #[tokio::test]
4669 async fn atomic_update_edge_symmetric_conflict_does_not_resurrect_tombstoned_survivor() {
4670 let runtime = scratch_runtime();
4671 let token = runtime
4672 .authorize(Namespace::parse("local").expect("ns"))
4673 .expect("authorize");
4674 let entities = runtime.entities(&token).expect("entities store");
4675 let a = khive_storage::Entity::new("local", "concept", "GapEdgeSymTombA");
4676 let b = khive_storage::Entity::new("local", "concept", "GapEdgeSymTombB");
4677 let (a_id, b_id) = (a.id, b.id);
4678 entities.upsert_entity(a).await.expect("seed a");
4679 entities.upsert_entity(b).await.expect("seed b");
4680
4681 let requested_edge = runtime
4682 .link(&token, a_id, b_id, EdgeRelation::Extends, 0.2, None)
4683 .await
4684 .expect("seed requested edge");
4685 let requested_id = Uuid::from(requested_edge.id);
4686
4687 let canonical_edge = runtime
4688 .link(&token, a_id, b_id, EdgeRelation::CompetesWith, 0.6, None)
4689 .await
4690 .expect("seed canonical edge");
4691 let canonical_id = Uuid::from(canonical_edge.id);
4692 runtime
4693 .delete_edge(&token, canonical_id, false)
4694 .await
4695 .expect("soft-delete canonical edge");
4696
4697 let plan = prepare_update(
4698 &runtime,
4699 &token,
4700 &json!({"id": requested_id.to_string(), "relation": "competes_with", "weight": 0.9}),
4701 None,
4702 )
4703 .await
4704 .expect("prepare update edge (symmetric conflict over tombstone)");
4705
4706 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
4707 .await
4708 .expect("seam call ok");
4709 assert!(
4710 matches!(
4711 outcome,
4712 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
4713 ),
4714 "expected a clean symmetric-conflict-absorption commit: {outcome:?}"
4715 );
4716
4717 let requested_after = runtime
4718 .get_edge_including_deleted(&token, requested_id)
4719 .await
4720 .expect("get_edge_including_deleted");
4721 assert!(
4722 requested_after.is_none(),
4723 "the non-canonical requested row must be deleted, not just tombstoned"
4724 );
4725
4726 let canonical_after = runtime
4727 .get_edge(&token, canonical_id)
4728 .await
4729 .expect("get_edge");
4730 assert!(
4731 canonical_after.is_none(),
4732 "a tombstoned survivor must not be resurrected by a conflicting update"
4733 );
4734 }
4735
4736 #[tokio::test]
4748 async fn atomic_update_edge_symmetric_same_unit_delete_race_aborts_the_unit() {
4749 let runtime = scratch_runtime();
4750 let token = runtime
4751 .authorize(Namespace::parse("local").expect("ns"))
4752 .expect("authorize");
4753 let entities = runtime.entities(&token).expect("entities store");
4754 let a = khive_storage::Entity::new("local", "concept", "GapEdgeRaceA");
4755 let b = khive_storage::Entity::new("local", "concept", "GapEdgeRaceB");
4756 let (a_id, b_id) = (a.id, b.id);
4757 entities.upsert_entity(a).await.expect("seed a");
4758 entities.upsert_entity(b).await.expect("seed b");
4759
4760 let requested_edge = runtime
4762 .link(&token, a_id, b_id, EdgeRelation::Extends, 0.2, None)
4763 .await
4764 .expect("seed requested edge");
4765 let requested_id = Uuid::from(requested_edge.id);
4766
4767 let canonical_edge = runtime
4770 .link(&token, a_id, b_id, EdgeRelation::CompetesWith, 0.6, None)
4771 .await
4772 .expect("seed canonical edge");
4773 let canonical_id = Uuid::from(canonical_edge.id);
4774
4775 let delete_plan = prepare_delete(
4776 &runtime,
4777 &token,
4778 &json!({"id": requested_id.to_string(), "hard": true}),
4779 None,
4780 )
4781 .await
4782 .expect("prepare delete edge");
4783 let update_plan = prepare_update(
4784 &runtime,
4785 &token,
4786 &json!({"id": requested_id.to_string(), "relation": "competes_with", "weight": 0.9}),
4787 None,
4788 )
4789 .await
4790 .expect("prepare update edge (both prepares run before either commits)");
4791
4792 let outcome = crate::atomic_runner::run_atomic_unit(
4793 runtime.sql().as_ref(),
4794 vec![delete_plan, update_plan],
4795 )
4796 .await
4797 .expect("the seam call itself must not error — the unit rolls back cleanly");
4798 match outcome {
4799 crate::atomic_runner::AtomicRunOutcome::RolledBack {
4800 failed_op_index, ..
4801 } => {
4802 assert_eq!(
4803 failed_op_index, 1,
4804 "op 1 (the update) must be the one whose guard fails"
4805 );
4806 }
4807 other => panic!("expected the whole unit to roll back, got {other:?}"),
4808 }
4809
4810 let requested_after = runtime
4812 .get_edge(&token, requested_id)
4813 .await
4814 .expect("get_edge");
4815 assert!(
4816 requested_after.is_some(),
4817 "delete(X) must have rolled back along with the failed update"
4818 );
4819 let canonical_after = runtime
4821 .get_edge(&token, canonical_id)
4822 .await
4823 .expect("get_edge")
4824 .expect("canonical row must still be present");
4825 assert_eq!(
4826 canonical_after.weight, 0.6,
4827 "the pre-existing canonical row must never have been touched by the aborted update"
4828 );
4829 }
4830
4831 #[tokio::test]
4842 async fn atomic_update_edge_symmetric_absorption_plan_stale_revision_rolls_back_unit() {
4843 let runtime = scratch_runtime();
4844 let token = runtime
4845 .authorize(Namespace::parse("local").expect("ns"))
4846 .expect("authorize");
4847 let entities = runtime.entities(&token).expect("entities store");
4848 let a = khive_storage::Entity::new("local", "concept", "StaleAbsorbA");
4849 let b = khive_storage::Entity::new("local", "concept", "StaleAbsorbB");
4850 let (a_id, b_id) = (a.id, b.id);
4851 entities.upsert_entity(a).await.expect("seed a");
4852 entities.upsert_entity(b).await.expect("seed b");
4853
4854 let canonical_edge = runtime
4856 .link(&token, a_id, b_id, EdgeRelation::CompetesWith, 0.6, None)
4857 .await
4858 .expect("seed canonical edge");
4859 let canonical_id = Uuid::from(canonical_edge.id);
4860
4861 let requested_edge = runtime
4863 .link(&token, a_id, b_id, EdgeRelation::Extends, 0.2, None)
4864 .await
4865 .expect("seed requested edge");
4866 let requested_id = Uuid::from(requested_edge.id);
4867
4868 let plan = prepare_update(
4871 &runtime,
4872 &token,
4873 &json!({"id": requested_id.to_string(), "relation": "competes_with", "weight": 0.9}),
4874 None,
4875 )
4876 .await
4877 .expect("prepare update edge (symmetric absorption)");
4878
4879 let concurrent = runtime
4884 .update_edge(
4885 &token,
4886 requested_id,
4887 crate::curation::EdgePatch {
4888 weight: Some(0.77),
4889 ..Default::default()
4890 },
4891 )
4892 .await
4893 .expect("concurrent writer update");
4894 assert!((concurrent.weight - 0.77).abs() < 1e-9);
4895
4896 let (absorb_replacement, absorb_expected, absorb_expected_deleted, absorb_ns, absorb_id) = {
4906 let statements = match &plan {
4907 AtomicOpPlan::Update(p) => p.statements.clone(),
4908 other => panic!("expected an Update plan, got {other:?}"),
4909 };
4910 let s = statements
4911 .iter()
4912 .find(|st| {
4913 st.statement.label.as_deref() == Some("edge-symmetric-absorb-or-update-inplace")
4914 })
4915 .expect("the plan must carry the guarded symmetric in-place absorb");
4916 let int = |i: usize| match &s.statement.params[i] {
4917 SqlValue::Integer(v) => *v,
4918 other => panic!("param {i} must be an integer revision, got {other:?}"),
4919 };
4920 let marker = |i: usize| match &s.statement.params[i] {
4921 SqlValue::Null => None,
4922 SqlValue::Integer(v) => Some(*v),
4923 other => panic!("param {i} must be a deletion marker, got {other:?}"),
4924 };
4925 let txt = |i: usize| match &s.statement.params[i] {
4926 SqlValue::Text(v) => v.clone(),
4927 other => panic!("param {i} must be text, got {other:?}"),
4928 };
4929 (int(6), int(9), marker(10), txt(0), txt(1))
4930 };
4931 let stored_pinned = absorb_replacement - 1;
4932 assert_ne!(
4933 stored_pinned, absorb_expected,
4934 "fixture premise: the pinned revision must differ from the plan's `?10`, otherwise \
4935 `updated_at = ?10` would MATCH and nothing would refuse the in-place arm (and the \
4936 guarded delete's `updated_at = ?6`, bound to the same value, would MATCH too and \
4937 fire, leaving `changes() = 1` and selecting the absorbed arm instead)"
4938 );
4939 {
4940 let mut writer = runtime.sql().writer().await.expect("writer");
4941 let affected = writer
4942 .execute(SqlStatement {
4943 sql: "UPDATE graph_edges SET updated_at = ?1 WHERE id = ?2".to_string(),
4944 params: vec![
4945 SqlValue::Integer(stored_pinned),
4946 SqlValue::Text(requested_id.to_string()),
4947 ],
4948 label: Some("test-pin-stored-revision".to_string()),
4949 })
4950 .await
4951 .expect("pin the stored revision");
4952 assert_eq!(affected, 1, "the pin must touch exactly the requested edge");
4953 }
4954 assert!(
4955 absorb_replacement > stored_pinned,
4956 "fixture premise: the plan's `?7` must still strictly advance past the stored \
4957 revision, otherwise `?7 > updated_at` refuses too and the refusal stops being \
4958 attributable to `updated_at = ?10`"
4959 );
4960
4961 {
4985 let statements = match &plan {
4986 AtomicOpPlan::Update(p) => p.statements.clone(),
4987 other => panic!("expected an Update plan, got {other:?}"),
4988 };
4989 let del = statements
4991 .iter()
4992 .find(|s| s.statement.label.as_deref() == Some("edge-symmetric-delete-if-conflict"))
4993 .expect("the plan must carry the guarded symmetric delete");
4994 let text = |i: usize| match &del.statement.params[i] {
4995 SqlValue::Text(v) => v.clone(),
4996 other => panic!("param {i} must be text, got {other:?}"),
4997 };
4998 let plan_expected_updated = match &del.statement.params[5] {
4999 SqlValue::Integer(v) => *v,
5000 other => panic!("param 5 must be the expected revision, got {other:?}"),
5001 };
5002 let plan_expected_deleted = match &del.statement.params[6] {
5003 SqlValue::Null => None,
5004 SqlValue::Integer(v) => Some(*v),
5005 other => panic!("param 6 must be a deletion marker, got {other:?}"),
5006 };
5007
5008 assert_eq!(
5016 text(0),
5017 "local",
5018 "fixture premise: the delete's `?1` must be the namespace the row lives in, \
5019 otherwise `namespace = ?1` refuses on its own"
5020 );
5021 assert_eq!(
5022 text(1),
5023 requested_id.to_string(),
5024 "fixture premise: the delete's `?2` must be the edge under test, otherwise \
5025 `id = ?2` refuses on identity and the delete never attempted the requested row"
5026 );
5027
5028 let requested_now = runtime
5029 .get_edge_including_deleted(&token, requested_id)
5030 .await
5031 .expect("read the requested edge")
5032 .expect("the requested edge is present before the plan runs");
5033 assert_ne!(
5034 requested_now.updated_at.timestamp_micros(),
5035 plan_expected_updated,
5036 "fixture premise: the concurrent writer must actually have moved the \
5037 revision past the plan's `?6`, otherwise nothing refuses the delete"
5038 );
5039 assert_eq!(
5040 requested_now.deleted_at.map(|d| d.timestamp_micros()),
5041 plan_expected_deleted,
5042 "fixture premise: the stored deletion marker must MATCH the plan's `?7`, \
5043 otherwise `deleted_at IS ?7` refuses too and the refusal is not \
5044 attributable to the revision guard"
5045 );
5046
5047 assert_ne!(
5056 canonical_id, requested_id,
5057 "fixture premise: the survivor must be a DIFFERENT row, since the EXISTS \
5058 arm excludes the requested id"
5059 );
5060 let survivors = {
5061 let mut reader = runtime.sql().reader().await.expect("sql reader");
5062 reader
5063 .query_scalar(SqlStatement {
5064 sql: "SELECT count(*) FROM graph_edges \
5065 WHERE namespace = ?1 AND source_id = ?3 AND target_id = ?4 \
5066 AND relation = ?5 AND id != ?2"
5067 .to_string(),
5068 params: del.statement.params[0..5].to_vec(),
5069 label: Some("test-survivor-exists-premise".to_string()),
5070 })
5071 .await
5072 .expect("run the plan's own EXISTS predicate")
5073 };
5074 let survivors = match survivors {
5075 Some(SqlValue::Integer(n)) => n,
5076 other => panic!("count(*) must come back as an integer, got {other:?}"),
5077 };
5078 assert_eq!(
5079 survivors,
5080 1,
5081 "fixture premise: exactly one survivor must satisfy the plan's own EXISTS \
5082 predicate (namespace {}, source {}, target {}, relation {}, id != {}), \
5083 otherwise the delete refuses on the survivor arm and the refusal is not \
5084 attributable to the revision guard",
5085 text(0),
5086 text(2),
5087 text(3),
5088 text(4),
5089 text(1),
5090 );
5091 }
5092
5093 assert_eq!(
5097 absorb_ns, "local",
5098 "fixture premise: the plan's `?1` must be the namespace the row lives in, \
5099 otherwise the outer `namespace = ?1` refuses on its own"
5100 );
5101 assert_eq!(
5102 absorb_id,
5103 requested_id.to_string(),
5104 "fixture premise: the plan's `?2` must be the edge under test, otherwise the \
5105 in-place arm's `id = ?2` refuses on identity"
5106 );
5107 {
5108 let stored = runtime
5109 .get_edge_including_deleted(&token, requested_id)
5110 .await
5111 .expect("read the requested edge")
5112 .expect("the requested edge is present before the plan runs");
5113 assert_eq!(
5114 stored.updated_at.timestamp_micros(),
5115 stored_pinned,
5116 "fixture premise: the pin must be what the guard reads at DML time, so the \
5117 stored revision is the pinned value and nothing re-advanced it"
5118 );
5119 assert_eq!(
5120 stored.deleted_at.map(|d| d.timestamp_micros()),
5121 absorb_expected_deleted,
5122 "fixture premise: the stored deletion marker must MATCH the plan's `?11`, \
5123 otherwise `deleted_at IS ?11` refuses too and the refusal stops being \
5124 attributable to `updated_at = ?10`"
5125 );
5126 }
5127
5128 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
5129 .await
5130 .expect("the seam call itself must not error; the unit rolls back cleanly");
5131 match outcome {
5132 crate::atomic_runner::AtomicRunOutcome::RolledBack {
5133 failed_op_index,
5134 failure,
5135 } => {
5136 assert_eq!(
5137 failed_op_index, 0,
5138 "the sole op's guard must be the one that fails"
5139 );
5140 assert_eq!(
5148 failure,
5149 crate::atomic_runner::AtomicOpFailure::GuardFailed {
5150 statement_label: Some(
5151 "edge-symmetric-absorb-or-update-inplace".to_string()
5152 ),
5153 expected: crate::atomic_plan::AffectedRowGuard::exactly(1),
5154 observed: 0,
5155 },
5156 "the guarded in-place absorb must be the statement whose guard refused"
5157 );
5158 }
5159 other => panic!(
5160 "a stale absorption plan must roll back, not silently absorb E into the \
5161 survivor and discard the concurrent writer's change: {other:?}"
5162 ),
5163 }
5164
5165 let requested_after = runtime
5166 .get_edge(&token, requested_id)
5167 .await
5168 .expect("get_edge")
5169 .expect("E must still exist after the rolled-back absorption");
5170 assert_eq!(
5171 requested_after.relation,
5172 EdgeRelation::Extends,
5173 "E must be untouched by the rolled-back absorption: {requested_after:?}"
5174 );
5175 assert!(
5176 (requested_after.weight - 0.77).abs() < 1e-9,
5177 "the concurrent writer's committed weight must survive the rolled-back plan: \
5178 {requested_after:?}"
5179 );
5180
5181 let canonical_after = runtime
5182 .get_edge(&token, canonical_id)
5183 .await
5184 .expect("get_edge")
5185 .expect("S must still exist after the rolled-back absorption");
5186 assert_eq!(
5187 canonical_after.weight, 0.6,
5188 "the survivor must never have been touched by the aborted absorption"
5189 );
5190 }
5191
5192 #[tokio::test]
5197 async fn atomic_update_edge_rejects_entity_only_field_name() {
5198 let runtime = scratch_runtime();
5199 let token = runtime
5200 .authorize(Namespace::parse("local").expect("ns"))
5201 .expect("authorize");
5202 let entities = runtime.entities(&token).expect("entities store");
5203 let a = khive_storage::Entity::new("local", "concept", "GapEdgeRejectA");
5204 let b = khive_storage::Entity::new("local", "concept", "GapEdgeRejectB");
5205 let (a_id, b_id) = (a.id, b.id);
5206 entities.upsert_entity(a).await.expect("seed a");
5207 entities.upsert_entity(b).await.expect("seed b");
5208 let edge = runtime
5209 .link(&token, a_id, b_id, EdgeRelation::Extends, 0.5, None)
5210 .await
5211 .expect("seed edge");
5212 let edge_id = Uuid::from(edge.id);
5213
5214 let err = prepare_update(
5215 &runtime,
5216 &token,
5217 &json!({"id": edge_id.to_string(), "name": "not-a-valid-edge-field"}),
5218 None,
5219 )
5220 .await
5221 .expect_err("edge update with an entity-only field must be rejected");
5222 let message = err.to_string();
5223 assert!(
5224 message.contains("name") && message.contains("edge"),
5225 "error must name the offending field and the substrate: {message}"
5226 );
5227 }
5228
5229 #[tokio::test]
5234 async fn atomic_delete_edge_soft_and_hard_appends_edge_deleted_event() {
5235 let runtime = scratch_runtime();
5236 let token = runtime
5237 .authorize(Namespace::parse("local").expect("ns"))
5238 .expect("authorize");
5239
5240 for hard in [false, true] {
5241 let entities = runtime.entities(&token).expect("entities store");
5242 let a = khive_storage::Entity::new("local", "concept", format!("GapEdgeDelA{hard}"));
5243 let b = khive_storage::Entity::new("local", "concept", format!("GapEdgeDelB{hard}"));
5244 let (a_id, b_id) = (a.id, b.id);
5245 entities.upsert_entity(a).await.expect("seed a");
5246 entities.upsert_entity(b).await.expect("seed b");
5247 let edge = runtime
5248 .link(&token, a_id, b_id, EdgeRelation::Extends, 0.5, None)
5249 .await
5250 .expect("seed edge");
5251 let edge_id = Uuid::from(edge.id);
5252
5253 let args = if hard {
5254 json!({"id": edge_id.to_string(), "hard": true})
5255 } else {
5256 json!({"id": edge_id.to_string()})
5257 };
5258 let plan = prepare_delete(&runtime, &token, &args, None)
5259 .await
5260 .unwrap_or_else(|e| panic!("prepare delete edge (hard={hard}): {e}"));
5261 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
5262 .await
5263 .unwrap_or_else(|e| panic!("edge delete commit (hard={hard}): {e}"));
5264 assert!(
5265 matches!(
5266 outcome,
5267 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
5268 ),
5269 "expected a clean edge delete commit (hard={hard}): {outcome:?}"
5270 );
5271
5272 let after = runtime
5273 .get_edge_including_deleted(&token, edge_id)
5274 .await
5275 .expect("get_edge_including_deleted");
5276 if hard {
5277 assert!(after.is_none(), "hard delete must purge the row entirely");
5278 } else {
5279 assert!(
5280 after.as_ref().is_some_and(|e| e.deleted_at.is_some()),
5281 "soft delete must tombstone, not purge"
5282 );
5283 }
5284
5285 let events = events_for_target(&runtime, &token, edge_id, EventKind::EdgeDeleted).await;
5286 assert_eq!(
5287 events.len(),
5288 1,
5289 "expected exactly one EdgeDeleted event for {edge_id} (hard={hard})"
5290 );
5291 assert_eq!(events[0].payload["hard"], json!(hard));
5292 }
5293 }
5294
5295 #[tokio::test]
5304 async fn atomic_update_note_appends_its_domain_event() {
5305 let runtime = scratch_runtime();
5306 let token = runtime
5307 .authorize(Namespace::parse("local").expect("ns"))
5308 .expect("authorize");
5309 let mut note = khive_storage::note::Note::new("local", "observation", "gap1-note-noevent");
5310 note.name = Some("gap1-note-noevent".to_string());
5311 let note_id = note.id;
5312 runtime
5313 .notes(&token)
5314 .expect("notes store")
5315 .upsert_note(note)
5316 .await
5317 .expect("seed note");
5318
5319 let plan = prepare_update(
5320 &runtime,
5321 &token,
5322 &json!({"id": note_id.to_string(), "content": "gap1-note-noevent, revised"}),
5323 None,
5324 )
5325 .await
5326 .expect("prepare update");
5327 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
5328 .await
5329 .expect("seam call ok");
5330 assert!(matches!(
5331 outcome,
5332 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
5333 ));
5334
5335 let event_store = runtime.events(&token).expect("event store");
5336 let page = event_store
5337 .query_events(
5338 khive_storage::EventFilter::default(),
5339 khive_storage::types::PageRequest::default(),
5340 )
5341 .await
5342 .expect("query_events");
5343 let for_note: Vec<_> = page
5344 .items
5345 .iter()
5346 .filter(|e| e.target_id == Some(note_id))
5347 .collect();
5348 assert_eq!(
5349 for_note.len(),
5350 1,
5351 "an atomic note update must append exactly one event; found: {for_note:?}"
5352 );
5353 assert_eq!(for_note[0].kind, EventKind::NoteUpdated);
5354 assert_eq!(for_note[0].substrate, SubstrateKind::Note);
5355 assert_eq!(for_note[0].verb, "update");
5356 assert_eq!(for_note[0].payload["id"], json!(note_id));
5357 assert_eq!(for_note[0].payload["text_changed"], json!(true));
5358 }
5359
5360 #[tokio::test]
5363 async fn atomic_link_appends_created_event_with_edge_observation() {
5364 let runtime = scratch_runtime();
5365 let token = runtime
5366 .authorize(Namespace::parse("local").expect("ns"))
5367 .expect("authorize");
5368 let source = khive_storage::Entity::new("local", "concept", "gap1-link-source");
5369 let target = khive_storage::Entity::new("local", "concept", "gap1-link-target");
5370 let (source_id, target_id) = (source.id, target.id);
5371 runtime
5372 .entities(&token)
5373 .expect("entities store")
5374 .upsert_entity(source)
5375 .await
5376 .expect("seed source");
5377 runtime
5378 .entities(&token)
5379 .expect("entities store")
5380 .upsert_entity(target)
5381 .await
5382 .expect("seed target");
5383
5384 let plan = prepare_link(
5385 &runtime,
5386 &token,
5387 &json!({
5388 "source_id": source_id.to_string(),
5389 "target_id": target_id.to_string(),
5390 "relation": "extends",
5391 }),
5392 )
5393 .await
5394 .expect("prepare link");
5395 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
5396 .await
5397 .expect("seam call ok");
5398 assert!(matches!(
5399 outcome,
5400 crate::atomic_runner::AtomicRunOutcome::Committed { .. }
5401 ));
5402
5403 let event_store = runtime.events(&token).expect("event store");
5404 let page = event_store
5405 .query_events(
5406 khive_storage::EventFilter {
5407 kinds: vec![EventKind::LinkCreated],
5408 ..khive_storage::EventFilter::default()
5409 },
5410 khive_storage::types::PageRequest::default(),
5411 )
5412 .await
5413 .expect("query_events");
5414 assert_eq!(page.items.len(), 1, "link must append one created event");
5415 let event = &page.items[0];
5416 assert_eq!(event.payload["mutation"], "created");
5417 let edge_id = event.target_id.expect("link event targets its edge");
5418 let observed = event_store
5419 .query_events(
5420 khive_storage::EventFilter {
5421 observed: vec![edge_id],
5422 ..khive_storage::EventFilter::default()
5423 },
5424 khive_storage::types::PageRequest::default(),
5425 )
5426 .await
5427 .expect("query observed edge");
5428 assert_eq!(observed.items.len(), 1);
5429 assert_eq!(observed.items[0].id, event.id);
5430 }
5431
5432 #[tokio::test]
5437 async fn prepare_add_entity_rejects_whitespace_only_name() {
5438 let runtime = scratch_runtime();
5439 let token = runtime
5440 .authorize(Namespace::parse("local").expect("ns"))
5441 .expect("authorize");
5442
5443 let err = prepare_add_entity(&runtime, &token, &json!({"kind": "concept", "name": " "}))
5444 .await
5445 .expect_err("whitespace-only entity name must fail prepare");
5446
5447 assert!(matches!(
5448 err,
5449 RuntimeError::InvalidInput(message) if message.contains("name must not be empty")
5450 ));
5451 }
5452
5453 #[tokio::test]
5454 async fn prepare_add_entity_rejects_non_string_description() {
5455 let runtime = scratch_runtime();
5456 let token = runtime
5457 .authorize(Namespace::parse("local").expect("ns"))
5458 .expect("authorize");
5459
5460 let err = prepare_add_entity(
5461 &runtime,
5462 &token,
5463 &json!({"kind": "concept", "name": "Valid", "description": 42}),
5464 )
5465 .await
5466 .expect_err("non-string entity description must fail prepare");
5467
5468 assert!(matches!(
5469 err,
5470 RuntimeError::InvalidInput(message)
5471 if message.contains("description must be a string or null")
5472 ));
5473 }
5474
5475 #[tokio::test]
5476 async fn prepare_add_note_rejects_non_string_name() {
5477 let runtime = scratch_runtime();
5478 let token = runtime
5479 .authorize(Namespace::parse("local").expect("ns"))
5480 .expect("authorize");
5481
5482 let err = prepare_add_note(
5483 &runtime,
5484 &token,
5485 &json!({"kind": "observation", "content": "Valid", "name": 42}),
5486 )
5487 .await
5488 .expect_err("non-string note name must fail prepare");
5489
5490 assert!(matches!(
5491 err,
5492 RuntimeError::InvalidInput(message) if message.contains("name must be a string or null")
5493 ));
5494 }
5495
5496 #[tokio::test]
5501 async fn prepare_add_entity_rejects_reserved_secret_gate_property() {
5502 let runtime = scratch_runtime();
5503 let token = runtime
5504 .authorize(Namespace::parse("local").expect("ns"))
5505 .expect("authorize");
5506
5507 let err = prepare_add_entity(
5508 &runtime,
5509 &token,
5510 &json!({
5511 "kind": "concept",
5512 "name": "ReservedKeyEntity",
5513 "properties": {"khive:secret_gate": "exempted:content-sha256-manifest-v1"},
5514 }),
5515 )
5516 .await
5517 .expect_err("a caller-supplied reserved key on a new entity must be rejected");
5518 assert!(
5519 matches!(err, RuntimeError::InvalidInput(ref msg) if msg.contains("khive:secret_gate") && msg.contains("runtime-owned")),
5520 "expected a reservation rejection, got: {err:?}"
5521 );
5522 }
5523
5524 #[tokio::test]
5527 async fn prepare_add_note_rejects_reserved_secret_gate_property() {
5528 let runtime = scratch_runtime();
5529 let token = runtime
5530 .authorize(Namespace::parse("local").expect("ns"))
5531 .expect("authorize");
5532
5533 let err = prepare_add_note(
5534 &runtime,
5535 &token,
5536 &json!({
5537 "kind": "observation",
5538 "content": "a note carrying a reserved property key",
5539 "properties": {"khive:secret_gate": "exempted:content-sha256-manifest-v1"},
5540 }),
5541 )
5542 .await
5543 .expect_err("a caller-supplied reserved key on a new note must be rejected");
5544 assert!(
5545 matches!(err, RuntimeError::InvalidInput(ref msg) if msg.contains("khive:secret_gate") && msg.contains("runtime-owned")),
5546 "expected a reservation rejection, got: {err:?}"
5547 );
5548 }
5549
5550 #[tokio::test]
5551 async fn atomic_proposal_vectors_materialize_only_after_successful_commit() {
5552 let runtime = scratch_runtime();
5553 runtime.register_embedder(StubProvider);
5554 let token = runtime
5555 .authorize(Namespace::parse("local").expect("ns"))
5556 .expect("authorize");
5557 let vec_store = runtime
5558 .vectors_for_model(&token, STUB_MODEL)
5559 .expect("vec store");
5560 let entities = runtime.entities(&token).expect("entities store");
5561 let a = khive_storage::Entity::new("local", "concept", "ProposalPlanLinkA");
5562 let b = khive_storage::Entity::new("local", "concept", "ProposalPlanLinkB");
5563 let (a_id, b_id) = (a.id, b.id);
5564 entities.upsert_entity(a).await.expect("seed a");
5565 entities.upsert_entity(b).await.expect("seed b");
5566
5567 let add_entity_plan = prepare_add_entity(
5568 &runtime,
5569 &token,
5570 &json!({"kind": "concept", "name": "ProposalPlanNewEntity", "description": "created atomically"}),
5571 )
5572 .await
5573 .expect("prepare add_entity");
5574 let link_plan = prepare_link(
5575 &runtime,
5576 &token,
5577 &json!({"source_id": a_id.to_string(), "target_id": b_id.to_string(), "relation": "extends"}),
5578 )
5579 .await
5580 .expect("prepare link");
5581 let add_note_plan = prepare_add_note(
5582 &runtime,
5583 &token,
5584 &json!({"kind": "observation", "content": "created atomically alongside the entity"}),
5585 )
5586 .await
5587 .expect("prepare add_note");
5588
5589 let entity_id = match &add_entity_plan {
5590 AtomicOpPlan::AddEntity(p) => p.entity_id,
5591 other => panic!("expected an AddEntity plan, got {other:?}"),
5592 };
5593 let note_id = match &add_note_plan {
5594 AtomicOpPlan::AddNote(p) => p.note_id,
5595 other => panic!("expected an AddNote plan, got {other:?}"),
5596 };
5597
5598 let outcome = crate::atomic_runner::run_atomic_unit(
5599 runtime.sql().as_ref(),
5600 vec![add_entity_plan, link_plan, add_note_plan],
5601 )
5602 .await
5603 .expect("seam call ok");
5604 let post_commit = match outcome {
5605 crate::atomic_runner::AtomicRunOutcome::Committed { post_commit } => post_commit,
5606 other => panic!("expected the whole unit to commit: {other:?}"),
5607 };
5608 assert_eq!(
5609 post_commit.as_slice(),
5610 &[
5611 PostCommitEffect::ReindexEntity { entity_id },
5612 PostCommitEffect::ReindexNote {
5613 note_id,
5614 version: 1
5615 },
5616 ],
5617 "prepare-derived effects must reach the committed token unchanged"
5618 );
5619 assert_eq!(
5620 vec_store.count().await.expect("count before effects"),
5621 0,
5622 "commit returns deferred effects without materializing vectors"
5623 );
5624 let entity = runtime
5625 .entities(&token)
5626 .expect("entities store")
5627 .get_entity(entity_id)
5628 .await
5629 .expect("get_entity")
5630 .expect("entity must exist after commit");
5631 assert_eq!(entity.name, "ProposalPlanNewEntity");
5632 assert!(
5633 runtime
5634 .text(&token)
5635 .expect("text store")
5636 .get_document("local", entity_id)
5637 .await
5638 .expect("get_document")
5639 .is_some(),
5640 "entity's FTS document must exist after commit"
5641 );
5642
5643 let (edge_count, _, _, edge_deleted_at) =
5644 probe_edge_natural_key(&runtime, "local", a_id, b_id, "extends").await;
5645 assert_eq!(
5646 edge_count, 1,
5647 "the edge must be committed alongside the entity/note"
5648 );
5649 assert!(edge_deleted_at.is_none());
5650
5651 let note = runtime
5652 .notes(&token)
5653 .expect("notes store")
5654 .get_note(note_id)
5655 .await
5656 .expect("get_note")
5657 .expect("note must exist after commit");
5658 assert_eq!(note.content, "created atomically alongside the entity");
5659 assert!(
5660 runtime
5661 .text_for_notes(&token)
5662 .expect("text store")
5663 .get_document("local", note_id)
5664 .await
5665 .expect("get_document")
5666 .is_some(),
5667 "note's FTS document must exist after commit"
5668 );
5669
5670 apply_post_commit_effects(&runtime, &token, post_commit)
5671 .await
5672 .expect("apply post-commit effects");
5673
5674 assert_eq!(
5675 vec_store.count().await.expect("count after"),
5676 2,
5677 "post-commit reindex must have embedded both the new entity and the new note"
5678 );
5679 }
5680
5681 #[tokio::test]
5682 async fn atomic_proposal_abort_leaves_zero_vector_rows() {
5683 let runtime = scratch_runtime();
5684 runtime.register_embedder(StubProvider);
5685 let token = runtime
5686 .authorize(Namespace::parse("local").expect("ns"))
5687 .expect("authorize");
5688 let vec_store = runtime
5689 .vectors_for_model(&token, STUB_MODEL)
5690 .expect("vec store");
5691 let entities = runtime.entities(&token).expect("entities store");
5692 let a = khive_storage::Entity::new("local", "concept", "ProposalPlanRollbackA");
5693 let x = khive_storage::Entity::new("local", "concept", "ProposalPlanRollbackX");
5694 let (a_id, x_id) = (a.id, x.id);
5695 entities.upsert_entity(a).await.expect("seed a");
5696 entities.upsert_entity(x.clone()).await.expect("seed x");
5697
5698 let add_entity_plan = prepare_add_entity(
5699 &runtime,
5700 &token,
5701 &json!({"kind": "concept", "name": "ProposalPlanRollbackNewEntity"}),
5702 )
5703 .await
5704 .expect("prepare add_entity");
5705 let add_note_plan = prepare_add_note(
5706 &runtime,
5707 &token,
5708 &json!({"kind": "observation", "content": "must not survive the rollback"}),
5709 )
5710 .await
5711 .expect("prepare add_note");
5712 let delete_plan = prepare_delete(
5713 &runtime,
5714 &token,
5715 &json!({"id": x_id.to_string(), "hard": true}),
5716 None,
5717 )
5718 .await
5719 .expect("prepare delete x");
5720 let link_plan = prepare_link(
5723 &runtime,
5724 &token,
5725 &json!({"source_id": a_id.to_string(), "target_id": x_id.to_string(), "relation": "extends"}),
5726 )
5727 .await
5728 .expect("prepare link (endpoint still exists at prepare time)");
5729
5730 let entity_id = match &add_entity_plan {
5731 AtomicOpPlan::AddEntity(p) => p.entity_id,
5732 other => panic!("expected an AddEntity plan, got {other:?}"),
5733 };
5734 let note_id = match &add_note_plan {
5735 AtomicOpPlan::AddNote(p) => p.note_id,
5736 other => panic!("expected an AddNote plan, got {other:?}"),
5737 };
5738
5739 let outcome = crate::atomic_runner::run_atomic_unit(
5740 runtime.sql().as_ref(),
5741 vec![add_entity_plan, add_note_plan, delete_plan, link_plan],
5742 )
5743 .await
5744 .expect("the seam call itself must not error; the unit rolls back cleanly");
5745 match outcome {
5746 crate::atomic_runner::AtomicRunOutcome::RolledBack {
5747 failed_op_index, ..
5748 } => {
5749 assert_eq!(
5750 failed_op_index, 3,
5751 "the trailing link (index 3) must be the op whose guard fails"
5752 );
5753 }
5754 other => panic!("expected the whole unit to roll back, got {other:?}"),
5755 }
5756
5757 assert_eq!(
5758 vec_store
5759 .count()
5760 .await
5761 .expect("vector count after rollback"),
5762 0,
5763 "a rolled-back atomic apply must not materialize vectors"
5764 );
5765
5766 assert!(
5767 runtime
5768 .get_entity_including_deleted(&token, entity_id)
5769 .await
5770 .expect("get_entity_including_deleted")
5771 .is_none(),
5772 "the new entity must leave zero trace after rollback"
5773 );
5774 assert!(
5775 runtime
5776 .text(&token)
5777 .expect("text store")
5778 .get_document("local", entity_id)
5779 .await
5780 .expect("get_document")
5781 .is_none(),
5782 "the new entity's FTS document must leave zero trace after rollback"
5783 );
5784 assert!(
5785 runtime
5786 .get_note_including_deleted(&token, note_id)
5787 .await
5788 .expect("get_note_including_deleted")
5789 .is_none(),
5790 "the new note must leave zero trace after rollback"
5791 );
5792 assert!(
5793 runtime
5794 .text_for_notes(&token)
5795 .expect("text store")
5796 .get_document("local", note_id)
5797 .await
5798 .expect("get_document")
5799 .is_none(),
5800 "the new note's FTS document must leave zero trace after rollback"
5801 );
5802
5803 let x_after = runtime
5804 .get_entity_including_deleted(&token, x_id)
5805 .await
5806 .expect("get_entity_including_deleted")
5807 .expect("x must still be present because its delete rolled back too");
5808 assert!(
5809 x_after.deleted_at.is_none(),
5810 "x's delete must have rolled back along with the failed link"
5811 );
5812
5813 let (edge_count, _, _, _) =
5814 probe_edge_natural_key(&runtime, "local", a_id, x_id, "extends").await;
5815 assert_eq!(edge_count, 0, "no edge may have been committed");
5816 }
5817
5818 #[tokio::test]
5822 async fn atomic_update_entity_type_null_clears_stored_type() {
5823 let runtime = scratch_runtime();
5824 runtime.install_entity_type_validator(std::sync::Arc::new(|kind, entity_type| {
5825 let Some(raw) = entity_type else {
5826 return Ok(None);
5827 };
5828 let normalized = raw.trim().to_ascii_lowercase();
5829 if kind == "concept" && normalized == "algorithm" {
5830 Ok(Some(normalized))
5831 } else {
5832 Err(RuntimeError::InvalidInput(format!(
5833 "unknown entity_type {raw:?} for {kind:?}; valid: algorithm"
5834 )))
5835 }
5836 }));
5837 let token = runtime
5838 .authorize(Namespace::parse("local").expect("ns"))
5839 .expect("authorize");
5840 let mut entity = khive_storage::Entity::new("local", "concept", "AtomicNullClear");
5841 entity.entity_type = Some("algorithm".to_string());
5842 let entity_id = entity.id;
5843 runtime
5844 .entities(&token)
5845 .expect("entities store")
5846 .upsert_entity(entity)
5847 .await
5848 .expect("seed entity");
5849
5850 let plan = prepare_update(
5851 &runtime,
5852 &token,
5853 &json!({"id": entity_id.to_string(), "entity_type": null}),
5854 None,
5855 )
5856 .await
5857 .expect("atomic prepare must accept entity_type: null");
5858 let outcome = crate::atomic_runner::run_atomic_unit(runtime.sql().as_ref(), vec![plan])
5859 .await
5860 .expect("atomic update must run");
5861 let post_commit = match outcome {
5862 crate::atomic_runner::AtomicRunOutcome::Committed { post_commit } => post_commit,
5863 other => panic!("expected Committed, got {other:?}"),
5864 };
5865 assert_eq!(
5866 post_commit.as_slice(),
5867 &[PostCommitEffect::ReindexEntity { entity_id }],
5868 "a type clear that differs from the prior value must reindex"
5869 );
5870
5871 let updated = runtime
5872 .get_entity(&token, entity_id)
5873 .await
5874 .expect("read updated entity");
5875 assert_eq!(
5876 updated.entity_type, None,
5877 "entity_type: null must clear the stored type"
5878 );
5879 assert_eq!(updated.name, "AtomicNullClear");
5880 }
5881
5882 #[tokio::test]
5886 async fn atomic_update_null_entity_type_rejected_for_note_and_edge() {
5887 let runtime = scratch_runtime();
5888 let token = runtime
5889 .authorize(Namespace::parse("local").expect("ns"))
5890 .expect("authorize");
5891 let note = runtime
5892 .create_note(
5893 &token,
5894 "observation",
5895 None,
5896 "note body for entity_type guard",
5897 Some(0.5),
5898 None,
5899 vec![],
5900 )
5901 .await
5902 .expect("create note");
5903
5904 let note_err = prepare_update(
5905 &runtime,
5906 &token,
5907 &json!({"id": note.id.to_string(), "entity_type": null}),
5908 Some(crate::atomic_prepare::AtomicUpdateKind::Note { specific: None }),
5909 )
5910 .await
5911 .expect_err("entity_type: null on a note must be rejected");
5912 assert!(
5913 matches!(note_err, RuntimeError::InvalidInput(ref msg) if msg.contains("entity_type") && msg.contains("not valid for a note")),
5914 "expected an InvalidInput naming entity_type for a note, got: {note_err:?}"
5915 );
5916
5917 let source = khive_storage::Entity::new("local", "concept", "AtomicNullTypeEdgeSource");
5918 let target = khive_storage::Entity::new("local", "concept", "AtomicNullTypeEdgeTarget");
5919 let source_id = source.id;
5920 let target_id = target.id;
5921 runtime
5922 .entities(&token)
5923 .expect("entities store")
5924 .upsert_entity(source)
5925 .await
5926 .expect("seed source");
5927 runtime
5928 .entities(&token)
5929 .expect("entities store")
5930 .upsert_entity(target)
5931 .await
5932 .expect("seed target");
5933 let edge = runtime
5934 .link(
5935 &token,
5936 source_id,
5937 target_id,
5938 "supports".parse().expect("relation"),
5939 0.5,
5940 None,
5941 )
5942 .await
5943 .expect("create edge");
5944
5945 let edge_err = prepare_update(
5946 &runtime,
5947 &token,
5948 &json!({"id": edge.id.to_string(), "entity_type": null}),
5949 Some(crate::atomic_prepare::AtomicUpdateKind::Edge),
5950 )
5951 .await
5952 .expect_err("entity_type: null on an edge must be rejected");
5953 assert!(
5954 matches!(edge_err, RuntimeError::InvalidInput(ref msg) if msg.contains("entity_type") && msg.contains("not valid for an edge")),
5955 "expected an InvalidInput naming entity_type for an edge, got: {edge_err:?}"
5956 );
5957 }
5958}