khive_runtime/atomic_prepare/edge_delete.rs
1use super::{
2 canonical_edge_endpoints, delete_record_attachments_statement, edge_hard_delete_statement,
3 edge_replace_if_unchanged_statement, edge_soft_delete_statement,
4 edge_symmetric_absorb_or_update_inplace_statement, edge_symmetric_delete_if_conflict_statement,
5 entity_hard_delete_statement, entity_soft_delete_statement, event_append_statements,
6 hard_delete_lineage_warning_statements, note_hard_delete_statement, note_soft_delete_statement,
7 obj, optional_f64, optional_properties, optional_str, parse_edge_relation,
8 purge_incident_edges_statement, push_index_purge_statements, refuse_pack_registry_tags,
9 reject_inapplicable_update_fields, require_uuid, AffectedRowGuard, AtomicOpPlan,
10 AttachmentSubstrate, DeletePlan, EdgeNaturalKey, EventKind, KhiveRuntime, NamespaceToken,
11 PlanStatement, PostCommitEffect, Resolved, RuntimeError, RuntimeResult, SubstrateKind,
12 UpdatePlan, Uuid, Value,
13};
14
15/// Edge branch of `prepare_update`. Mirrors `KhiveRuntime::update_edge`'s
16/// patch semantics: `relation`/`weight`/`properties` are the only applicable
17/// fields, a changed `relation` is endpoint-validated first, `weight` is
18/// range-checked, and `properties` REPLACES `metadata` wholesale (no merge).
19/// See `docs/api/atomic_prepare.md#prepare_update_edge` for the DML-shape parity
20/// detail with `update_edge`.
21///
22/// Invariant (symmetric relations `competes_with`/`composed_with`): this
23/// function must never branch on a prepare-time conflict probe — a different
24/// op in the same atomic unit could change the conflict landscape between
25/// probe and commit, making any such branch stale by construction. It always
26/// emits BOTH statements from [`edge_symmetric_delete_if_conflict_statement`]
27/// and [`edge_symmetric_absorb_or_update_inplace_statement`], each carrying
28/// its own commit-time `WHERE`/`CASE WHEN` predicate that re-evaluates the
29/// conflict condition fresh inside the transaction. This function reads no
30/// state to guess a surviving id; the plan instead carries `edge_natural_key`
31/// so a post-commit caller derives the actual surviving id from the
32/// committed row, never from a value computed before the rest of this atomic
33/// unit has even run.
34pub(super) async fn prepare_update_edge(
35 runtime: &KhiveRuntime,
36 token: &NamespaceToken,
37 id: Uuid,
38 mut edge: khive_storage::types::Edge,
39 args: &Value,
40) -> RuntimeResult<AtomicOpPlan> {
41 reject_inapplicable_update_fields(args, "edge")?;
42
43 let expected_updated_at = edge.updated_at;
44 let expected_deleted_at = edge.deleted_at;
45
46 let relation_raw = optional_str(args, "relation");
47 let weight = optional_f64(args, "weight")?;
48 let properties = optional_properties(args, "properties")?;
49
50 if let Some(ref p) = properties {
51 crate::secret_gate::check_json_at(p, "edge", "properties")?;
52 }
53 crate::secret_gate::reject_reserved_secret_gate_property(properties.as_ref())?;
54
55 let namespace = edge.namespace.clone();
56 let record_tok = token.with_namespace(
57 khive_types::Namespace::parse(&namespace)
58 .map_err(|e| RuntimeError::Internal(format!("edge namespace invalid: {e}")))?,
59 );
60
61 let mut changed_fields: Vec<&'static str> = Vec::new();
62 if let Some(raw) = relation_raw {
63 let relation = parse_edge_relation(raw)?;
64 runtime
65 .validate_edge_relation_endpoints(&record_tok, edge.source_id, edge.target_id, relation)
66 .await?;
67 edge.relation = relation;
68 changed_fields.push("relation");
69 }
70 if let Some(w) = weight {
71 if !w.is_finite() || !(0.0..=1.0).contains(&w) {
72 return Err(RuntimeError::InvalidInput(format!(
73 "edge weight must be a finite value in [0.0, 1.0]; got {w}"
74 )));
75 }
76 edge.weight = w;
77 changed_fields.push("weight");
78 }
79 if let Some(p) = properties {
80 edge.metadata = Some(p);
81 changed_fields.push("properties");
82 }
83
84 let (canon_src, canon_tgt) =
85 canonical_edge_endpoints(edge.relation, edge.source_id, edge.target_id);
86 let now = chrono::Utc::now();
87
88 let mut statements: Vec<PlanStatement> = Vec::new();
89 let mut edge_natural_key: Option<EdgeNaturalKey> = None;
90
91 if edge.relation.is_symmetric() {
92 // The write for a symmetric relation never branches on a
93 // prepare-time probe result: it always carries both self-guarding,
94 // commit-time-predicate statements (see their doc comment in
95 // khive-db's graph.rs for the full rationale). This avoids the
96 // staleness window a prepare-time probe would expose: an earlier op
97 // in the same atomic unit could change the conflict landscape before
98 // commit. Canonical's own probe-then-branch
99 // `update_edge_symmetric_dml` has no such exposure (single
100 // transaction, no interleaving) and is unaffected.
101 let metadata_str = edge
102 .metadata
103 .as_ref()
104 .map(|v| serde_json::to_string(v).unwrap_or_default());
105
106 // `updated_at` must strictly advance past the snapshot even when two
107 // operations land inside one clock microsecond; saturating to
108 // i64::MAX would let the CAS accept a write without advancing its
109 // revision, so that is not a valid fallback (mirrors the note path).
110 let minimum_updated_at_micros = expected_updated_at
111 .timestamp_micros()
112 .checked_add(1)
113 .ok_or_else(|| {
114 RuntimeError::Internal(format!(
115 "edge {id} updated_at is already at i64::MAX and cannot advance"
116 ))
117 })?;
118 let symmetric_updated_at_micros = now.timestamp_micros().max(minimum_updated_at_micros);
119 let expected_deleted_at_micros = expected_deleted_at.map(|v| v.timestamp_micros());
120
121 statements.push(PlanStatement {
122 statement: edge_symmetric_delete_if_conflict_statement(
123 &namespace,
124 id,
125 canon_src,
126 canon_tgt,
127 edge.relation,
128 expected_updated_at.timestamp_micros(),
129 expected_deleted_at_micros,
130 ),
131 guard: Some(AffectedRowGuard {
132 expected_min: 0,
133 expected_max: Some(1),
134 }),
135 });
136 statements.push(PlanStatement {
137 statement: edge_symmetric_absorb_or_update_inplace_statement(
138 &namespace,
139 id,
140 canon_src,
141 canon_tgt,
142 edge.relation,
143 edge.weight,
144 symmetric_updated_at_micros,
145 metadata_str.as_deref(),
146 edge.target_backend.as_deref(),
147 expected_updated_at.timestamp_micros(),
148 expected_deleted_at_micros,
149 ),
150 guard: Some(AffectedRowGuard::exactly(1)),
151 });
152
153 // No prepare-time read needed: the two statements above are
154 // self-guarding at commit time (see their doc comment). Post-commit
155 // result rendering derives the actual surviving id from THIS
156 // natural key, never from a value computed here.
157 edge_natural_key = Some(EdgeNaturalKey {
158 namespace: namespace.clone(),
159 canon_source_id: canon_src,
160 canon_target_id: canon_tgt,
161 relation: edge.relation,
162 });
163 } else {
164 // Non-symmetric: guarded replace of the read snapshot rather than
165 // `graph.upsert_edge`'s unconditional natural-key upsert — a zero
166 // affected-row result now means a concurrent writer moved this edge
167 // between PREPARE and commit, and the atomic unit must roll back
168 // instead of silently overwriting it.
169 //
170 // `updated_at` must strictly advance past the snapshot even when two
171 // operations land inside one clock microsecond; saturating to
172 // i64::MAX would let the CAS accept a write without advancing its
173 // revision, so that is not a valid fallback (mirrors the note path).
174 let minimum_updated_at_micros = expected_updated_at
175 .timestamp_micros()
176 .checked_add(1)
177 .ok_or_else(|| {
178 RuntimeError::Internal(format!(
179 "edge {id} updated_at is already at i64::MAX and cannot advance"
180 ))
181 })?;
182 let now_micros = now.timestamp_micros().max(minimum_updated_at_micros);
183 edge.updated_at = chrono::DateTime::from_timestamp_micros(now_micros).ok_or_else(|| {
184 RuntimeError::Internal(format!(
185 "edge {id}: computed updated_at {now_micros} is not a valid timestamp"
186 ))
187 })?;
188 statements.push(PlanStatement {
189 statement: edge_replace_if_unchanged_statement(
190 &edge,
191 expected_updated_at,
192 expected_deleted_at,
193 ),
194 guard: Some(AffectedRowGuard::exactly(1)),
195 });
196 }
197
198 // Mirrors `update_edge`'s unconditional post-mutation `EdgeUpdated`
199 // event append, keyed on the original `edge_id` the caller supplied:
200 // canonical does the same (the event target is `edge_id`, not the
201 // post-absorption surviving id).
202 statements.extend(event_append_statements(
203 token,
204 &namespace,
205 "update",
206 EventKind::EdgeUpdated,
207 SubstrateKind::Entity,
208 id,
209 serde_json::json!({"id": id, "namespace": namespace, "changed_fields": changed_fields}),
210 )?);
211
212 Ok(AtomicOpPlan::Update(Box::new(UpdatePlan {
213 graph_effects: Vec::new(),
214 note_vector_purge: None,
215 note_embedding_inheritance: None,
216 entity_guard: None,
217 note_guard: None,
218 target_id: id,
219 statements,
220 post_commit: PostCommitEffect::None,
221 edge_natural_key,
222 idempotent_noop: false,
223 })))
224}
225
226// ---------------------------------------------------------------------------
227// delete
228// ---------------------------------------------------------------------------
229
230/// Caller-supplied delete-kind expectation, resolved via the canonical
231/// `resolve_kind_spec` at the kkernel `--atomic` seam. `khive-runtime` must
232/// not depend on `khive-pack-kg` (packs depend on the runtime, not the other
233/// way around), so this is a plain substrate-level shape rather than
234/// `khive_pack_kg::handlers::KindSpec` itself: the kkernel seam does the
235/// pack-aware `resolve_kind_spec` resolution (which needs a `VerbRegistry`,
236/// unreachable from this crate) and passes down only what `prepare_delete`
237/// needs to enforce the mismatch check.
238///
239/// `delete` admits `kind="edge"` per `ATOMIC_ADMISSIBLE_VERBS`, hence the
240/// `Edge` variant. `Event`/`Proposal` remain rejected at the kkernel seam
241/// (not v1-admissible for atomic delete at all).
242pub enum AtomicDeleteKind {
243 Entity {
244 specific: Option<String>,
245 entity_type: Option<String>,
246 },
247 Note {
248 specific: Option<String>,
249 },
250 Edge,
251}
252
253/// `expected_kind`: `None` when the caller omitted `kind` (no check, parity
254/// with canonical's own optional discriminator); `Some(_)` enforces an
255/// exact-parity mismatch check against the resolved record's actual
256/// substrate/specific kind, mirroring `handle_delete`'s
257/// `entity.kind != *expected` / `note.kind != *expected` checks.
258pub async fn prepare_delete(
259 runtime: &KhiveRuntime,
260 token: &NamespaceToken,
261 args: &Value,
262 expected_kind: Option<AtomicDeleteKind>,
263) -> RuntimeResult<AtomicOpPlan> {
264 let id = require_uuid(args, "id")?;
265 let actor = format!("{}:{}", token.actor().kind, token.actor().id);
266 let hard = obj(args)?
267 .get("hard")
268 .and_then(|v| v.as_bool())
269 .unwrap_or(false);
270
271 // `delete(id, hard=true)` is the public purge route after a prior soft
272 // delete, so it must resolve including already-tombstoned rows (a
273 // live-only resolve would never find one). Soft delete keeps the
274 // live-only resolve: a soft delete of an already-tombstoned row is a
275 // no-op, matching non-atomic behavior.
276 let resolved = if hard {
277 runtime.resolve_by_id_including_deleted(token, id).await?
278 } else {
279 runtime.resolve_by_id(token, id).await?
280 };
281
282 match resolved {
283 Some(Resolved::Entity(entity)) => {
284 match &expected_kind {
285 None => {}
286 Some(AtomicDeleteKind::Entity {
287 specific: Some(expected),
288 ..
289 }) if &entity.kind != expected => {
290 return Err(RuntimeError::NotFound(format!("{expected} {id}")));
291 }
292 Some(AtomicDeleteKind::Entity { .. }) => {}
293 Some(AtomicDeleteKind::Note { .. }) => {
294 return Err(RuntimeError::NotFound(format!("note {id}")));
295 }
296 Some(AtomicDeleteKind::Edge) => {
297 return Err(RuntimeError::NotFound(format!("edge {id}")));
298 }
299 }
300 if let Some(AtomicDeleteKind::Entity {
301 entity_type: Some(expected),
302 ..
303 }) = &expected_kind
304 {
305 if entity
306 .entity_type
307 .as_deref()
308 .is_some_and(|actual| actual != expected.as_str())
309 {
310 return Err(RuntimeError::NotFound(format!("entity {id}")));
311 }
312 }
313 refuse_pack_registry_tags(&entity.tags, "delete")?;
314 let namespace = entity.namespace.clone();
315 // Storage parity: `entity_soft_delete_statement`/
316 // `entity_hard_delete_statement` are the SAME khive-db builders
317 // khive-db's own `SqlEntityStore::delete_entity` calls — no DML
318 // text is hand-duplicated here.
319 let mut statements = if hard {
320 vec![
321 PlanStatement {
322 statement: delete_record_attachments_statement(
323 id,
324 AttachmentSubstrate::Entity,
325 ),
326 guard: None,
327 },
328 PlanStatement {
329 statement: entity_hard_delete_statement(id),
330 guard: Some(AffectedRowGuard::exactly(1)),
331 },
332 ]
333 } else {
334 let deleted_at = chrono::Utc::now().timestamp_micros();
335 vec![PlanStatement {
336 statement: entity_soft_delete_statement(id, deleted_at),
337 guard: Some(AffectedRowGuard::exactly(1)),
338 }]
339 };
340 if hard {
341 statements.extend(
342 hard_delete_lineage_warning_statements(
343 &namespace,
344 &actor,
345 id,
346 SubstrateKind::Entity,
347 )
348 .into_iter()
349 .map(|statement| PlanStatement {
350 statement,
351 guard: None,
352 }),
353 );
354 // Same builder canonical `delete_entity`'s hard-delete
355 // cascade calls (`graph.purge_incident_edges`).
356 statements.push(PlanStatement {
357 statement: purge_incident_edges_statement(id),
358 guard: None,
359 });
360 }
361 // FTS + vector index purge, matching operations.rs
362 // `delete_entity`: both soft and hard delete clean indexes (a
363 // hard delete of an already-tombstoned record must still purge
364 // them); only hard additionally cascades edges above.
365 push_index_purge_statements(
366 runtime,
367 &mut statements,
368 "fts_entities",
369 &namespace,
370 id,
371 "atomic-delete-entity",
372 )
373 .await?;
374 // operations.rs's `delete_entity` appends an `EntityDeleted`
375 // event after a successful row delete, on both soft and hard
376 // delete. `apply_plan` never reaches this statement unless the
377 // guarded row statement above affected a row, so no extra `if`
378 // is needed here.
379 statements.extend(event_append_statements(
380 token,
381 &namespace,
382 "delete",
383 EventKind::EntityDeleted,
384 SubstrateKind::Entity,
385 id,
386 serde_json::json!({"id": id, "namespace": namespace, "hard": hard}),
387 )?);
388 Ok(AtomicOpPlan::Delete(DeletePlan {
389 target_id: id,
390 statements,
391 post_commit: PostCommitEffect::None,
392 }))
393 }
394 Some(Resolved::Note(note)) => {
395 match &expected_kind {
396 None => {}
397 Some(AtomicDeleteKind::Note {
398 specific: Some(expected),
399 }) if ¬e.kind != expected => {
400 return Err(RuntimeError::NotFound(format!("{expected} {id}")));
401 }
402 Some(AtomicDeleteKind::Note { .. }) => {}
403 Some(AtomicDeleteKind::Entity { .. }) => {
404 return Err(RuntimeError::NotFound(format!("entity {id}")));
405 }
406 Some(AtomicDeleteKind::Edge) => {
407 return Err(RuntimeError::NotFound(format!("edge {id}")));
408 }
409 }
410 if let Some(error) = runtime.stream_member_error(¬e).await? {
411 return Err(error);
412 }
413 let namespace = note.namespace.clone();
414 // Storage parity: `note_soft_delete_statement`/
415 // `note_hard_delete_statement` are the SAME khive-db builders
416 // khive-db's own `SqlNoteStore::delete_note` calls.
417 let mut statements = if hard {
418 vec![
419 PlanStatement {
420 statement: delete_record_attachments_statement(
421 id,
422 AttachmentSubstrate::Note,
423 ),
424 guard: None,
425 },
426 PlanStatement {
427 statement: note_hard_delete_statement(id),
428 guard: Some(AffectedRowGuard::exactly(1)),
429 },
430 ]
431 } else {
432 let deleted_at = chrono::Utc::now().timestamp_micros();
433 vec![PlanStatement {
434 statement: note_soft_delete_statement(id, deleted_at),
435 guard: Some(AffectedRowGuard::exactly(1)),
436 }]
437 };
438 if hard {
439 statements.extend(
440 hard_delete_lineage_warning_statements(
441 &namespace,
442 &actor,
443 id,
444 SubstrateKind::Note,
445 )
446 .into_iter()
447 .map(|statement| PlanStatement {
448 statement,
449 guard: None,
450 }),
451 );
452 statements.push(PlanStatement {
453 statement: purge_incident_edges_statement(id),
454 guard: None,
455 });
456 }
457 // FTS + vector index purge, matching operations.rs
458 // `delete_note`: both soft and hard delete clean indexes (a hard
459 // delete of an already-tombstoned record must still purge
460 // them); only hard additionally cascades edges above.
461 push_index_purge_statements(
462 runtime,
463 &mut statements,
464 "fts_notes",
465 &namespace,
466 id,
467 "atomic-delete-note",
468 )
469 .await?;
470 // operations.rs's `delete_note` appends a `NoteDeleted` event
471 // after a successful row delete, on both soft and hard delete:
472 // same reasoning as the entity branch above.
473 statements.extend(event_append_statements(
474 token,
475 &namespace,
476 "delete",
477 EventKind::NoteDeleted,
478 SubstrateKind::Note,
479 id,
480 serde_json::json!({"id": id, "namespace": namespace, "hard": hard}),
481 )?);
482 Ok(AtomicOpPlan::Delete(DeletePlan {
483 target_id: id,
484 statements,
485 // A committed atomic note delete must fire the same
486 // pack-installed note-mutation hook `operations.rs::
487 // delete_note` fires, so a warm ANN cache sees the deletion
488 // even when the mutation went through the atomic-plan path.
489 post_commit: PostCommitEffect::NoteDeleted {
490 note_id: id,
491 kind: note.kind.clone(),
492 },
493 }))
494 }
495 Some(_) => Err(RuntimeError::InvalidInput(format!(
496 "delete target {id} must be an entity, note, or edge"
497 ))),
498 // `Resolved` has no `Edge` variant (same reasoning as
499 // `prepare_update`'s fallback above) — probe the graph store
500 // directly.
501 None => match &expected_kind {
502 Some(AtomicDeleteKind::Entity { .. }) => {
503 Err(RuntimeError::NotFound(format!("entity/note {id}")))
504 }
505 Some(AtomicDeleteKind::Note { .. }) => {
506 Err(RuntimeError::NotFound(format!("entity/note {id}")))
507 }
508 Some(AtomicDeleteKind::Edge) | None => {
509 let edge = if hard {
510 runtime.get_edge_including_deleted(token, id).await?
511 } else {
512 runtime.get_edge(token, id).await?
513 };
514 match edge {
515 Some(edge) => prepare_delete_edge(token, id, edge, hard, &actor).await,
516 None => Err(RuntimeError::NotFound(format!("entity/note/edge {id}"))),
517 }
518 }
519 },
520 }
521}
522
523/// Edge branch of `prepare_delete`. Mirrors
524/// `khive-runtime::operations::KhiveRuntime::delete_edge` exactly: hard
525/// delete cascades `purge_incident_edges` (any `annotates` edge — or any
526/// other edge — pointing AT this edge as a node) BEFORE deleting the edge
527/// row itself, then a soft or hard delete statement, then an unconditional
528/// `EdgeDeleted` event (edges are never FTS/vector-indexed, so unlike the
529/// entity/note branches there is no index purge here — `delete_edge` has
530/// none either).
531pub(super) async fn prepare_delete_edge(
532 token: &NamespaceToken,
533 id: Uuid,
534 edge: khive_storage::types::Edge,
535 hard: bool,
536 actor: &str,
537) -> RuntimeResult<AtomicOpPlan> {
538 let namespace = edge.namespace.clone();
539 let mut statements: Vec<PlanStatement> = Vec::new();
540
541 if hard {
542 statements.extend(
543 hard_delete_lineage_warning_statements(&namespace, actor, id, SubstrateKind::Entity)
544 .into_iter()
545 .map(|statement| PlanStatement {
546 statement,
547 guard: None,
548 }),
549 );
550 // Mirrors `delete_edge`'s `graph.purge_incident_edges(edge_id)` —
551 // unguarded: zero incident edges is a legitimate outcome, not a
552 // failure (same reasoning as the entity/note cascade-edges
553 // statements above).
554 statements.push(PlanStatement {
555 statement: purge_incident_edges_statement(id),
556 guard: None,
557 });
558 statements.push(PlanStatement {
559 statement: edge_hard_delete_statement(id),
560 guard: Some(AffectedRowGuard::exactly(1)),
561 });
562 } else {
563 let now = chrono::Utc::now().timestamp_micros();
564 statements.push(PlanStatement {
565 statement: edge_soft_delete_statement(id, now),
566 guard: Some(AffectedRowGuard::exactly(1)),
567 });
568 }
569
570 statements.extend(event_append_statements(
571 token,
572 &namespace,
573 "delete",
574 EventKind::EdgeDeleted,
575 SubstrateKind::Entity,
576 id,
577 serde_json::json!({"id": id, "namespace": namespace, "hard": hard}),
578 )?);
579
580 Ok(AtomicOpPlan::Delete(DeletePlan {
581 target_id: id,
582 statements,
583 post_commit: PostCommitEffect::None,
584 }))
585}