1#[cfg(doc)]
2use super::AtomicDeleteKind;
3use super::{
4 entity_fts_document, entity_name_patch, entity_replace_if_unchanged_statement,
5 entity_upsert_statement, event_append_statements, insert_document_statements,
6 note_fts_document, note_upsert_statement, obj, optional_create_string,
7 optional_entity_type_patch, optional_f64_patch, optional_properties, optional_str,
8 optional_string_patch, optional_tags, prepare_delete, prepare_delete_edge, prepare_link,
9 prepare_merge, prepare_update_edge, prepare_update_entity_plan_with_version,
10 refuse_pack_registry_tags, require_str, require_uuid, AddEntityPlan, AddNotePlan,
11 AffectedRowGuard, AtomicOpPlan, EdgeUpsertDisposition, EventKind, KhiveRuntime, NamespaceToken,
12 PlanStatement, PostCommitEffect, Resolved, RuntimeError, RuntimeResult, SubstrateKind,
13 UpdatePlan, Uuid, Value,
14};
15
16pub async fn prepare_op(
26 runtime: &KhiveRuntime,
27 token: &NamespaceToken,
28 tool: &str,
29 args: &Value,
30) -> RuntimeResult<AtomicOpPlan> {
31 match tool {
32 "update" => prepare_update(runtime, token, args, None).await,
41 "delete" => prepare_delete(runtime, token, args, None).await,
50 "link" => prepare_link(runtime, token, args).await,
51 "merge" => prepare_merge(runtime, token, args).await,
52 "propose" | "review" | "withdraw" => prepare_governance_unimplemented(tool),
53 other => Err(RuntimeError::InvalidInput(format!(
54 "{other:?} has no atomic_prepare::prepare_op implementation; the CLI \
55 admissibility check should have rejected this before prepare"
56 ))),
57 }
58}
59
60fn prepare_governance_unimplemented(tool: &str) -> RuntimeResult<AtomicOpPlan> {
61 Err(RuntimeError::InvalidInput(format!(
62 "{tool:?} is on the ADR-099 v1 admissible verb list but has no --atomic \
63 prepare/apply implementation yet: its lifecycle (ADR-046) is an \
64 event-sourced changeset-interpreter over a dedicated `proposals_open` \
65 table, not a small guarded-DML plan — a faithful non-stub atomic \
66 prepare for it is tracked as ADR-099 follow-up work, not implemented \
67 in slice B3. No write was attempted."
68 )))
69}
70
71pub async fn prepare_add_entity(
81 runtime: &KhiveRuntime,
82 token: &NamespaceToken,
83 args: &Value,
84) -> RuntimeResult<AtomicOpPlan> {
85 let kind = require_str(args, "kind")?;
86 let name = require_str(args, "name")?;
87 runtime.validate_entity_kind(kind)?;
88 if name.trim().is_empty() {
89 return Err(RuntimeError::InvalidInput(
90 "name must not be empty".to_string(),
91 ));
92 }
93
94 let description = optional_create_string(args, "description")?;
95 let properties = optional_properties(args, "properties")?;
96 let tags = optional_tags(args)?.unwrap_or_default();
97
98 crate::secret_gate::check_at(name, "entity", "name")?;
99 if let Some(ref d) = description {
100 crate::secret_gate::check_at(d, "entity", "description")?;
101 }
102 if let Some(ref p) = properties {
103 crate::secret_gate::check_json_at(p, "entity", "properties")?;
104 }
105 crate::secret_gate::check_tags_at(&tags, "entity", "tags")?;
106 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
107
108 let ns = token.namespace().as_str();
109 let mut entity = khive_storage::Entity::new(ns, kind, name);
110 if let Some(d) = description {
111 entity = entity.with_description(d);
112 }
113 if let Some(p) = properties {
114 entity = entity.with_properties(p);
115 }
116 if !tags.is_empty() {
117 entity = entity.with_tags(tags);
118 }
119
120 let mut statements = vec![PlanStatement {
121 statement: entity_upsert_statement(&entity),
122 guard: Some(AffectedRowGuard::exactly(1)),
123 }];
124 for statement in insert_document_statements("fts_entities", &entity_fts_document(&entity)) {
128 statements.push(PlanStatement {
129 statement,
130 guard: None,
131 });
132 }
133
134 Ok(AtomicOpPlan::AddEntity(AddEntityPlan {
135 entity_id: entity.id,
136 statements,
137 post_commit: PostCommitEffect::ReindexEntity {
138 entity_id: entity.id,
139 },
140 }))
141}
142
143pub async fn prepare_add_note(
149 runtime: &KhiveRuntime,
150 token: &NamespaceToken,
151 args: &Value,
152) -> RuntimeResult<AtomicOpPlan> {
153 let kind = require_str(args, "kind")?;
154 let content = require_str(args, "content")?;
155 runtime.validate_note_kind(kind)?;
156
157 let name = optional_create_string(args, "name")?;
158 let properties = optional_properties(args, "properties")?;
159 let properties = runtime.derive_note_write_properties(kind, token, properties)?;
165
166 crate::secret_gate::check_at(content, "note", "content")?;
167 if let Some(ref n) = name {
168 crate::secret_gate::check_at(n, "note", "name")?;
169 }
170 if let Some(ref p) = properties {
171 crate::secret_gate::check_json_at(p, "note", "properties")?;
172 }
173 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
174
175 let ns = token.namespace().as_str();
176 let mut note = khive_storage::note::Note::new(ns, kind, content);
177 if let Some(n) = name {
178 note = note.with_name(n);
179 }
180 if let Some(p) = properties {
181 note = note.with_properties(p);
182 }
183
184 let mut statements = vec![PlanStatement {
185 statement: note_upsert_statement(¬e),
186 guard: Some(AffectedRowGuard::exactly(1)),
187 }];
188 for statement in insert_document_statements("fts_notes", ¬e_fts_document(¬e)) {
191 statements.push(PlanStatement {
192 statement,
193 guard: None,
194 });
195 }
196
197 Ok(AtomicOpPlan::AddNote(Box::new(AddNotePlan {
198 note_guard: None,
199 note_id: note.id,
200 statements,
201 post_commit: PostCommitEffect::ReindexNote {
202 note_id: note.id,
203 version: note.version,
204 },
205 })))
206}
207
208pub(super) fn reject_inapplicable_update_fields(
223 args: &Value,
224 substrate: &str,
225) -> RuntimeResult<()> {
226 let o = obj(args)?;
227 if substrate == "edge" && o.get("expected_version").is_some_and(|v| !v.is_null()) {
228 return Err(RuntimeError::InvalidInput(
229 "expected_version applies only to entities and notes".into(),
230 ));
231 }
232 if substrate != "note"
233 && ["embed", "fence"]
234 .iter()
235 .any(|field| o.contains_key(*field))
236 {
237 return Err(RuntimeError::InvalidInput(
238 "embed and fence apply only to notes".into(),
239 ));
240 }
241 let present = |k: &str| o.get(k).is_some_and(|v| !v.is_null());
242 let (bad_field, valid): (Option<&str>, &str) = match substrate {
243 "entity" => {
244 let bad = if present("content") {
245 Some("content")
246 } else if present("salience") {
247 Some("salience")
248 } else if present("decay_factor") {
249 Some("decay_factor")
250 } else if present("relation") {
251 Some("relation")
252 } else if present("weight") {
253 Some("weight")
254 } else {
255 None
256 };
257 (bad, "name, description, tags, properties, entity_type")
258 }
259 "note" => {
260 let bad = if present("description") {
261 Some("description")
262 } else if present("relation") {
263 Some("relation")
264 } else if present("weight") {
265 Some("weight")
266 } else if o.contains_key("entity_type") {
267 Some("entity_type")
270 } else {
271 None
272 };
273 (
274 bad,
275 "name, content, salience, decay_factor, properties, tags",
276 )
277 }
278 "edge" => {
284 let bad = if present("name") {
285 Some("name")
286 } else if present("description") {
287 Some("description")
288 } else if present("content") {
289 Some("content")
290 } else if present("tags") {
291 Some("tags")
292 } else if present("salience") {
293 Some("salience")
294 } else if present("decay_factor") {
295 Some("decay_factor")
296 } else if o.contains_key("entity_type") {
297 Some("entity_type")
300 } else {
301 None
302 };
303 (bad, "relation, weight, properties")
304 }
305 _ => (None, ""),
306 };
307 if let Some(field) = bad_field {
308 let substrate_label = match substrate {
309 "entity" => "an entity",
310 "note" => "a note",
311 "edge" => "an edge",
312 other => other,
313 };
314 return Err(RuntimeError::InvalidInput(format!(
315 "field '{field}' is not valid for {substrate_label}; valid fields: {valid}"
316 )));
317 }
318 Ok(())
319}
320
321pub enum AtomicUpdateKind {
331 Entity {
332 specific: Option<String>,
333 entity_type: Option<String>,
334 },
335 Note {
336 specific: Option<String>,
337 },
338 Edge,
339}
340
341pub fn validate_note_update_expected_kind(
346 note: &khive_storage::note::Note,
347 expected_kind: &Option<AtomicUpdateKind>,
348) -> RuntimeResult<()> {
349 let id = note.id;
350 match expected_kind {
351 None => Ok(()),
352 Some(AtomicUpdateKind::Note {
353 specific: Some(expected),
354 }) if ¬e.kind != expected => Err(RuntimeError::NotFound(format!("note {id}"))),
355 Some(AtomicUpdateKind::Note { .. }) => Ok(()),
356 Some(AtomicUpdateKind::Entity { .. }) => {
357 Err(RuntimeError::NotFound(format!("entity {id}")))
358 }
359 Some(AtomicUpdateKind::Edge) => Err(RuntimeError::NotFound(format!("edge {id}"))),
360 }
361}
362
363async fn prepare_note_update_plan_from_snapshot(
364 runtime: &KhiveRuntime,
365 token: &NamespaceToken,
366 args: &Value,
367 expected_kind: &Option<AtomicUpdateKind>,
368 note: khive_storage::note::Note,
369 policy: crate::NoteUpdatePolicy,
370 registry: Option<&crate::VerbRegistry>,
371) -> RuntimeResult<(khive_storage::Note, UpdatePlan)> {
372 let id = require_uuid(args, "id")?;
373 if note.id != id {
374 return Err(RuntimeError::NotFound(format!("note {id}")));
375 }
376 validate_note_update_expected_kind(¬e, expected_kind)?;
377
378 reject_inapplicable_update_fields(args, "note")?;
379 let mut normalized_args = args.clone();
380 crate::curation::normalize_note_update_tags(&mut normalized_args)?;
381 let args = &normalized_args;
382 let name = optional_string_patch(args, "name")?;
383 let content = optional_str(args, "content").map(str::to_string);
384 let properties = optional_properties(args, "properties")?;
385 let salience = optional_f64_patch(args, "salience")?;
386 let decay_factor = optional_f64_patch(args, "decay_factor")?;
387 let options = crate::note_write::NoteWriteOptions {
388 expected_version: obj(args)?
389 .get("expected_version")
390 .filter(|v| !v.is_null())
391 .map(|v| {
392 v.as_i64().ok_or_else(|| {
393 RuntimeError::InvalidInput("expected_version must be an integer".into())
394 })
395 })
396 .transpose()?,
397 fence: obj(args)?
398 .get("fence")
399 .map(|v| {
400 serde_json::from_value(v.clone())
401 .map_err(|error| RuntimeError::InvalidInput(format!("invalid fence: {error}")))
402 })
403 .transpose()?,
404 embed: obj(args)?
405 .get("embed")
406 .filter(|v| !v.is_null())
407 .map(|v| {
408 v.as_bool()
409 .ok_or_else(|| RuntimeError::InvalidInput("embed must be boolean".into()))
410 })
411 .transpose()?,
412 key: None,
413 };
414 let patch = crate::curation::NotePatch::new(name, content, salience, decay_factor, properties)
415 .with_update_policy(policy)
416 .with_write_options(options);
417 let (updated, mut plan) = runtime
418 .prepare_versioned_note_update(token, note.clone(), patch.clone())
419 .await?;
420 if let Some(registry) = registry {
421 attach_note_update_effects(runtime, token, registry, ¬e, &patch, &mut plan).await?;
422 }
423 Ok((updated, plan))
424}
425
426pub async fn prepare_update_from_note_snapshot(
436 runtime: &KhiveRuntime,
437 token: &NamespaceToken,
438 args: &Value,
439 expected_kind: Option<AtomicUpdateKind>,
440 note: khive_storage::note::Note,
441 policy: crate::NoteUpdatePolicy,
442 registry: &crate::VerbRegistry,
443) -> RuntimeResult<(khive_storage::Note, AtomicOpPlan)> {
444 if obj(args)?.get("entity_kind").is_some_and(|v| !v.is_null()) {
445 return Err(RuntimeError::InvalidInput(
446 "entity_kind is immutable; to change kind, delete then re-create the entity, \
447 or use merge() if this is a deduplication correction"
448 .into(),
449 ));
450 }
451 let (note, plan) = prepare_note_update_plan_from_snapshot(
452 runtime,
453 token,
454 args,
455 &expected_kind,
456 note,
457 policy,
458 Some(registry),
459 )
460 .await?;
461 Ok((note, AtomicOpPlan::Update(Box::new(plan))))
462}
463
464async fn attach_note_update_effects(
468 runtime: &KhiveRuntime,
469 token: &NamespaceToken,
470 registry: &crate::VerbRegistry,
471 snapshot: &khive_storage::Note,
472 patch: &crate::curation::NotePatch,
473 plan: &mut UpdatePlan,
474) -> RuntimeResult<()> {
475 use crate::atomic_plan::NoteUpdateStatement;
476 use crate::NoteUpdateEffect;
477
478 let Some(hook) = registry.find_kind_hook(&snapshot.kind) else {
479 return Ok(());
480 };
481 if let Some(expected) = patch.write_options.expected_version {
484 if expected != snapshot.version {
485 return Err(crate::note_write::NoteWriteConflict::Version {
486 expected,
487 current: snapshot.version,
488 }
489 .into_error()
490 .into());
491 }
492 }
493 let current = runtime.notes(token)?.get_note(snapshot.id).await?;
494 if !current.is_some_and(|current| {
495 current.updated_at == snapshot.updated_at
496 && current.deleted_at == snapshot.deleted_at
497 && current.version == snapshot.version
498 }) {
499 return Err(crate::curation::stale_note_snapshot_error(snapshot.id));
500 }
501 let effects = hook
502 .note_update_effects(runtime, token, snapshot, patch)
503 .await?;
504 if plan.idempotent_noop && !effects.is_empty() {
505 return Err(RuntimeError::InvalidInput(
506 "an unchanged note update cannot carry graph effects".into(),
507 ));
508 }
509 let edge_token = token.with_namespace(
510 crate::Namespace::parse(&snapshot.namespace)
511 .map_err(|error| RuntimeError::Internal(format!("invalid note namespace: {error}")))?,
512 );
513 for effect in effects {
514 match effect {
515 NoteUpdateEffect::Link(spec) => {
516 if spec.source_id != snapshot.id
517 || spec
518 .namespace
519 .as_deref()
520 .is_some_and(|ns| ns != snapshot.namespace)
521 {
522 return Err(RuntimeError::InvalidInput(
523 "note update links must originate from the note in its namespace".into(),
524 ));
525 }
526 let mut args = serde_json::json!({
527 "source_id": spec.source_id, "target_id": spec.target_id,
528 "relation": spec.relation, "weight": spec.weight,
529 "resurrect": spec.resurrect,
530 });
531 if let Some(metadata) = spec.metadata {
532 args["metadata"] = metadata;
533 }
534 let AtomicOpPlan::Link(link) = prepare_link(runtime, &edge_token, &args).await?
535 else {
536 return Err(RuntimeError::Internal("expected a link plan".into()));
537 };
538 if link.disposition == EdgeUpsertDisposition::Updated {
542 return Err(khive_types::KhiveError::conflict(
543 "a live edge appeared while preparing the note update; retry with fresh state",
544 ).into());
545 }
546 plan.graph_effects
547 .extend(link.statements.into_iter().map(NoteUpdateStatement::Write));
548 }
549 NoteUpdateEffect::DeleteEdge(edge) => {
550 let id = Uuid::from(edge.id);
551 if edge.source_id != snapshot.id
552 || edge.namespace != snapshot.namespace
553 || edge.deleted_at.is_some()
554 {
555 return Err(RuntimeError::InvalidInput(
556 "note update deletes must select an outgoing edge in the note namespace"
557 .into(),
558 ));
559 }
560 plan.graph_effects
561 .push(NoteUpdateStatement::Assert(PlanStatement {
562 statement: khive_db::stores::graph::edge_snapshot_assertion_statement(
563 &edge, false,
564 ),
565 guard: Some(AffectedRowGuard::exactly(1)),
566 }));
567 let actor = format!("{}:{}", token.actor().kind, token.actor().id);
568 let AtomicOpPlan::Delete(delete) =
569 prepare_delete_edge(&edge_token, id, edge, false, &actor).await?
570 else {
571 return Err(RuntimeError::Internal(
572 "expected an edge delete plan".into(),
573 ));
574 };
575 if delete.post_commit != PostCommitEffect::None {
576 return Err(RuntimeError::Internal(
577 "edge delete has a deferred effect".into(),
578 ));
579 }
580 plan.graph_effects.extend(
581 delete
582 .statements
583 .into_iter()
584 .map(NoteUpdateStatement::Write),
585 );
586 }
587 NoteUpdateEffect::AssertLink(edge) => {
588 if edge.source_id != snapshot.id
589 || edge.namespace != snapshot.namespace
590 || edge.deleted_at.is_some()
591 {
592 return Err(RuntimeError::InvalidInput(
593 "note update assertions must select a live outgoing edge in the note namespace".into(),
594 ));
595 }
596 plan.graph_effects
597 .push(NoteUpdateStatement::Assert(PlanStatement {
598 statement: khive_db::stores::graph::edge_snapshot_assertion_statement(
599 &edge, true,
600 ),
601 guard: Some(AffectedRowGuard::exactly(1)),
602 }));
603 }
604 }
605 }
606 Ok(())
607}
608
609pub async fn prepare_update(
617 runtime: &KhiveRuntime,
618 token: &NamespaceToken,
619 args: &Value,
620 expected_kind: Option<AtomicUpdateKind>,
621) -> RuntimeResult<AtomicOpPlan> {
622 let id = require_uuid(args, "id")?;
623
624 if obj(args)?.get("entity_kind").is_some_and(|v| !v.is_null()) {
628 return Err(RuntimeError::InvalidInput(
629 "entity_kind is immutable; to change kind, delete then re-create the entity, \
630 or use merge() if this is a deduplication correction"
631 .into(),
632 ));
633 }
634
635 match runtime.resolve_by_id(token, id).await? {
636 Some(Resolved::Entity(entity)) => {
637 match &expected_kind {
638 None => {}
639 Some(AtomicUpdateKind::Entity {
640 specific: Some(expected),
641 ..
642 }) if &entity.kind != expected => {
643 return Err(RuntimeError::NotFound(format!("entity {id}")));
644 }
645 Some(AtomicUpdateKind::Entity { .. }) => {}
646 Some(AtomicUpdateKind::Note { .. }) => {
647 return Err(RuntimeError::NotFound(format!("note {id}")));
648 }
649 Some(AtomicUpdateKind::Edge) => {
650 return Err(RuntimeError::NotFound(format!("edge {id}")));
651 }
652 }
653 if let Some(AtomicUpdateKind::Entity {
654 entity_type: Some(expected),
655 ..
656 }) = &expected_kind
657 {
658 if entity
659 .entity_type
660 .as_deref()
661 .is_some_and(|actual| actual != expected.as_str())
662 {
663 return Err(RuntimeError::NotFound(format!("entity {id}")));
664 }
665 }
666 refuse_pack_registry_tags(&entity.tags, "update")?;
667 reject_inapplicable_update_fields(args, "entity")?;
675 let name = entity_name_patch(args)?;
676 let description = optional_string_patch(args, "description")?;
677 let properties = optional_properties(args, "properties")?;
678 let tags = optional_tags(args)?;
679 if let Some(ref tags) = tags {
680 refuse_pack_registry_tags(tags, "update")?;
681 }
682 let entity_type = optional_entity_type_patch(args, "entity_type")?;
683
684 let expected_version = obj(args)?
685 .get("expected_version")
686 .filter(|v| !v.is_null())
687 .map(|value| {
688 value.as_i64().ok_or_else(|| {
689 RuntimeError::InvalidInput("expected_version must be an integer".into())
690 })
691 })
692 .transpose()?;
693 let required_entity_type = match &expected_kind {
694 Some(AtomicUpdateKind::Entity { entity_type, .. }) => entity_type.as_deref(),
695 _ => None,
696 };
697 prepare_update_entity_plan_with_version_and_type(
698 runtime,
699 token,
700 id,
701 crate::curation::EntityPatch {
702 name,
703 description,
704 properties,
705 tags,
706 entity_type,
707 },
708 expected_version,
709 required_entity_type,
710 )
711 .await
712 }
713 Some(Resolved::Note(note)) => {
714 let (_, plan) = prepare_note_update_plan_from_snapshot(
719 runtime,
720 token,
721 args,
722 &expected_kind,
723 note,
724 crate::NoteUpdatePolicy::default(),
725 None,
726 )
727 .await?;
728 Ok(AtomicOpPlan::Update(Box::new(plan)))
729 }
730 Some(_) => Err(RuntimeError::InvalidInput(format!(
731 "update target {id} must be an entity, note, or edge"
732 ))),
733 None => match &expected_kind {
740 Some(AtomicUpdateKind::Entity { .. }) => {
741 Err(RuntimeError::NotFound(format!("entity/note {id}")))
742 }
743 Some(AtomicUpdateKind::Note { .. }) => {
744 Err(RuntimeError::NotFound(format!("entity/note {id}")))
745 }
746 Some(AtomicUpdateKind::Edge) | None => match runtime.get_edge(token, id).await? {
747 Some(edge) => prepare_update_edge(runtime, token, id, edge, args).await,
748 None => Err(RuntimeError::NotFound(format!("entity/note/edge {id}"))),
749 },
750 },
751 }
752}
753
754pub async fn prepare_update_entity_plan(
758 runtime: &KhiveRuntime,
759 token: &NamespaceToken,
760 id: Uuid,
761 patch: crate::curation::EntityPatch,
762) -> RuntimeResult<AtomicOpPlan> {
763 prepare_update_entity_plan_with_version(runtime, token, id, patch, None).await
764}
765
766pub(super) async fn prepare_update_entity_plan_with_version_and_type(
767 runtime: &KhiveRuntime,
768 token: &NamespaceToken,
769 id: Uuid,
770 patch: crate::curation::EntityPatch,
771 expected_version: Option<i64>,
772 required_entity_type: Option<&str>,
773) -> RuntimeResult<AtomicOpPlan> {
774 crate::entity_write::validate_expected_version(expected_version)?;
775 let explicit_entity_type_patch = patch.entity_type.is_some();
776 let (entity, reindex_required, changed_fields, expected_updated_at, expected_deleted_at) =
777 runtime.prepare_update_entity(token, id, patch).await?;
778 if required_entity_type.is_some_and(|expected| {
779 (explicit_entity_type_patch || entity.entity_type.is_some())
780 && entity.entity_type.as_deref() != Some(expected)
781 }) {
782 return Err(RuntimeError::InvalidInput(
783 "kind subtype contradicts the requested entity_type update".into(),
784 ));
785 }
786 let mut statements = vec![PlanStatement {
787 statement: entity_replace_if_unchanged_statement(
788 &entity,
789 expected_updated_at,
790 expected_deleted_at,
791 ),
792 guard: Some(AffectedRowGuard::exactly(1)),
793 }];
794 statements.extend(event_append_statements(
795 token,
796 &entity.namespace,
797 "update",
798 EventKind::EntityUpdated,
799 SubstrateKind::Entity,
800 id,
801 serde_json::json!({
802 "id": id,
803 "namespace": entity.namespace,
804 "changed_fields": changed_fields,
805 }),
806 )?);
807 let post_commit = if reindex_required {
808 PostCommitEffect::ReindexEntity { entity_id: id }
809 } else {
810 PostCommitEffect::None
811 };
812 Ok(AtomicOpPlan::Update(Box::new(UpdatePlan {
813 graph_effects: Vec::new(),
814 note_vector_purge: None,
815 note_embedding_inheritance: None,
816 entity_guard: expected_version.map(|expected_version| {
817 crate::entity_write::EntityWriteGuard {
818 id,
819 expected_version,
820 }
821 }),
822 note_guard: None,
823 target_id: id,
824 statements,
825 post_commit,
826 edge_natural_key: None,
827 idempotent_noop: false,
828 })))
829}