Skip to main content

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 &note.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(&note).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}