Skip to main content

khive_runtime/
atomic_prepare.rs

1//! ADR-099: the per-verb async prepare pass for the KG-substrate v1
2//! admissible verbs (`update`, `delete`, `link`, `merge`), plus
3//! [`prepare_add_entity`]/[`prepare_add_note`] for the ADR-046 proposal
4//! changeset `AddEntity`/`AddNote` arms. Each `prepare_*` function reads
5//! current state (async, outside any transaction) and returns a plain-data
6//! [`crate::atomic_runner::AtomicOpPlan`] ([`crate::atomic_plan`]) for the
7//! synchronous commit pass ([`crate::atomic_runner::run_atomic_unit`]) to
8//! apply.
9//!
10//! `gtd.transition`/`gtd.complete` prepare is deliberately not here (lives in
11//! `kkernel` instead), and `propose`/`review`/`withdraw`/`merge` are on the
12//! v1 admissible list but have no working prepare implementation in this
13//! module (`prepare_governance_unimplemented` fails loudly rather than
14//! silently no-opping; `prepare_merge` is unreachable through `--atomic` and
15//! kept only for its own tests and as defense in depth). See
16//! `docs/api/atomic_prepare.md#scope-what-is-excluded-and-why` for why each of these is excluded
17//! and what would be required to admit them.
18
19use 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
59// ---------------------------------------------------------------------------
60// arg extraction helpers
61// ---------------------------------------------------------------------------
62
63fn 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
75/// Parse `key` as a bare UUID — never a short hex prefix.
76///
77/// A short prefix is a *resolution*: an unfiltered search that applies no
78/// namespace predicate, so it can match nothing, match exactly one record
79/// across every namespace, or match ambiguously. That search already
80/// happened upstream, at the `kkernel` CLI boundary that has namespace
81/// context and calls `resolve_uuid_unfiltered` before handing args down to
82/// this module
83/// (`crates/kkernel/src/atomic_apply.rs::resolve_kg_ids_in_args`). By the
84/// time an id reaches this plan-preparation stage it must already name one
85/// specific, already-identified record — which is exactly what a full UUID
86/// demonstrates and a prefix does not.
87fn 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
113/// ADR-014 tri-state patch for entity `entity_type`, read from raw JSON to
114/// mirror `UpdateParams.entity_type`'s `tri_string` deserializer (the
115/// kkernel `--atomic` seam deserializes through that struct first, so the
116/// two surfaces cannot diverge): key absent -> `None` (unchanged), key
117/// present as `null` -> `Some(None)` (explicit clear), key present as a
118/// string -> `Some(Some(s))` (set); any other JSON type -> a hard error.
119fn 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
130/// Nullable-string patch semantics shared by note updates and mirroring the
131/// actually-reachable entity description behavior of
132/// `khive-pack-kg::handlers::common::description_patch`. Canonical's field
133/// type is `Option<Value>` (`UpdateParams.name`/`.description`); serde_json's
134/// derived `Deserialize` for `Option<T>` intercepts a literal JSON `null` at
135/// the outer `Option` boundary and maps it straight to Rust `None`
136/// regardless of the inner type, so canonical's own "clear" arm is
137/// unreachable through normal struct deserialization: `update(name=null)` /
138/// `update(description=null)` are no-ops, not clears. This module reads raw,
139/// un-deserialized JSON, so it must replicate that collapse explicitly: key
140/// absent OR JSON `null` -> `None` (leave unchanged, no-op); key present as a
141/// string -> `Some(Some(s))` (set); any other JSON type -> a hard error.
142fn 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
152/// Strict string-or-absent-or-null patch for entity `name`. Unlike
153/// `optional_str`'s `.as_str()`, this does not silently drop a non-string,
154/// non-null value like `name: 123` as absent: it rejects it instead of
155/// reporting success for an invalid update. Canonical validates entity
156/// `name` via `string_value` on `UpdateParams.name: Option<Value>`: null
157/// collapses to absent at the struct-deserialize boundary (see
158/// `optional_string_patch` doc above), so the reachable behavior is:
159/// absent/null -> unchanged; non-null string -> set; any other JSON type ->
160/// hard error. This mirrors that exactly, reading raw JSON instead of a
161/// deserialized struct.
162fn 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
172/// Nullable-JSON-value patch for `properties`: canonical
173/// `properties: Option<Value>` on `UpdateParams` collapses a literal JSON
174/// `null` to Rust `None` at the struct-deserialize boundary (same collapse
175/// as `optional_string_patch` above), so `properties=null` is canonically a
176/// no-op (leave existing properties unchanged): not a stored JSON `null`.
177/// This module reads raw JSON, so it must replicate that collapse: key
178/// absent OR JSON `null` -> `None` (no merge); any other JSON value ->
179/// `Some(value)` (merge), with no further shape validation at this layer.
180fn 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
187/// `tags` patch: canonical `tags: Option<Vec<String>>` on `UpdateParams`
188/// collapses a literal JSON `null` to Rust `None` at the struct-deserialize
189/// boundary (same collapse as above), so `tags=null` is canonically a no-op
190/// (leave existing tags unchanged). A non-array, non-null value is still a
191/// hard error (mirrors the type failure `UpdateParams` deserialization would
192/// itself produce for a malformed `tags`).
193fn 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
223/// Tri-state patch extraction for `Option<Option<f64>>`-shaped fields
224/// (`NotePatch::salience` / `NotePatch::decay_factor`): key absent -> `None`
225/// (untouched), key present as JSON `null` -> `Some(None)` (clear), key
226/// present as a number -> `Some(Some(v))` (set). Range validation lives in
227/// curation.rs's `prepare_update_note_from_snapshot`, not here.
228fn 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
239/// Every registered embedding model's vector table name, in the exact format
240/// `curation::merge_entity_sql` uses (`"vec_{sanitize_key(model_name)}"`) —
241/// reused here so atomic delete/merge purge the same tables the non-atomic
242/// paths do.
243fn 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
251/// A guarded (`guard: None` — best-effort mirror, matching the non-atomic
252/// index-cleanup calls which don't assert a row existed) `DELETE` statement
253/// against one vector table for a single subject, scoped by namespace.
254///
255/// Vector tables carry a real index on `(subject_id, namespace)`
256/// ([`khive_db::stores::vectors`]) — this row-scan predicate is not the FTS
257/// full-table-scan class this module's `purge_fts_document_statements`
258/// exists to avoid, so it is left as a direct `namespace = ? AND subject_id =
259/// ?` predicate.
260fn 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
279/// The FTS-document half of an index purge: `fts_table`'s row for `subject_id`
280/// (looked up via `khive_db::stores::text::rowid_map_table`, not a
281/// `namespace`/`subject_id` scan — those columns are `UNINDEXED` in every
282/// FTS5 DDL) plus that row's own entry in the sidecar map. Order-sensitive:
283/// index 0 must run before index 1 — see `delete_document_statements`'s
284/// adjacency contract.
285fn 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
330/// `true` iff a table named `table` currently exists in the backing SQLite
331/// database (`sqlite_master` probe, read-only — safe in async prepare, does
332/// NOT open/create the vector store, so it cannot lazily create the table
333/// itself).
334async 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
351/// Append the FTS + every registered model's vector-row purge for `subject_id`
352/// (scoped to the RECORD's own namespace, matching `delete_entity`/
353/// `delete_note`'s `record_tok`/`record_ns` convention: not the caller
354/// token's namespace, per by-ID namespace-agnosticism) onto `statements`.
355///
356/// FTS tables (`fts_entities`/`fts_notes`) always exist (created at schema
357/// migration time) so their purge is unconditional. `vec_*` tables are
358/// created lazily on first vector-store open, so a default runtime can
359/// register embedding models before any vector table necessarily exists:
360/// a raw unconditional `DELETE FROM vec_*` can hit `no such table` on a
361/// fresh DB. Only push the vec purge for tables that actually exist:
362/// absence means the record definitionally has no vector row for that
363/// model, so skipping is data-parity-correct (the non-atomic path would
364/// lazily create the table then delete zero rows: same data outcome,
365/// without this read-only prepare pass performing an init side effect).
366async 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
399/// Event-store append parity for the canonical handlers that emit a
400/// lifecycle event after their row mutation: `update_entity` ->
401/// `EntityUpdated`, `delete_entity` -> `EntityDeleted`, `delete_note` ->
402/// `NoteDeleted`, `update_edge` -> `EdgeUpdated`, `delete_edge` ->
403/// `EdgeDeleted`, `link` -> `LinkCreated`/`EdgeUpdated`, and
404/// `update_note` -> `NoteUpdated`. See
405/// `docs/api/atomic_prepare.md#event_append_statements` for why
406/// this is a `PlanStatement` rather than a `PostCommitEffect`.
407///
408/// Invariant: returned statements are unguarded — appended after the plan's
409/// own guarded row statement, so [`apply_plan`]'s stop-on-first-failure
410/// contract means they are only reached once that row mutation's guard has
411/// already held. Committing the event row atomically with the mutation it
412/// describes strengthens canonical's guarantee: the non-atomic handlers write
413/// the event in a separate transaction, ordered but not atomic with the row
414/// mutation.
415pub(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
444// ---------------------------------------------------------------------------
445// dispatch
446// ---------------------------------------------------------------------------
447
448/// Build the prepared [`AtomicOpPlan`] for one KG-substrate admissible op
449/// (`update`, `delete`, `link`, `merge`). Returns a loud [`RuntimeError`] for
450/// `propose`/`review`/`withdraw` (known scope gap, see module doc) and any
451/// other verb (the CLI boundary must reject those before calling this — a
452/// verb reaching here is either KG-substrate-admissible or a bug upstream).
453pub async fn prepare_op(
454    runtime: &KhiveRuntime,
455    token: &NamespaceToken,
456    tool: &str,
457    args: &Value,
458) -> RuntimeResult<AtomicOpPlan> {
459    match tool {
460        // `expected_kind: None` here — same reasoning as the `"delete"` arm
461        // below: callers that need `update(kind=...)` parity must resolve
462        // the kind spec themselves (it needs a `VerbRegistry`, unreachable
463        // from this crate: see `AtomicUpdateKind`'s doc comment) and call
464        // `prepare_update` directly with the resolved value; `kkernel`'s
465        // `--atomic` seam does exactly this and bypasses this dispatch arm.
466        // A caller reaching `prepare_op("update", ...)` without going
467        // through that seam gets kind-unchecked behavior.
468        "update" => prepare_update(runtime, token, args, None).await,
469        // `expected_kind: None` here — callers that need `delete(kind=...)`
470        // parity must resolve the kind spec themselves (it needs a
471        // `VerbRegistry`, unreachable from this crate: see
472        // `AtomicDeleteKind`'s doc comment) and call `prepare_delete`
473        // directly with the resolved value; `kkernel`'s `--atomic` seam does
474        // exactly this and bypasses this dispatch arm. A caller reaching
475        // `prepare_op("delete", ...)` without going through that seam gets
476        // kind-unchecked behavior.
477        "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
499// ---------------------------------------------------------------------------
500// create (AddEntity / AddNote)
501// ---------------------------------------------------------------------------
502
503/// Build the prepared plan for an `AddEntity` proposal change. The entity
504/// row and FTS document are committed together; vector indexing is deferred
505/// until after commit because embedding may suspend. `kind` must already be
506/// canonicalized by the caller because pack-aware resolution requires a
507/// `VerbRegistry`.
508pub 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    // Order-sensitive pair — see `insert_document_statements`'s adjacency
553    // contract: the map upsert's `last_insert_rowid()` must read back the
554    // FTS insert immediately before it.
555    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
571/// Build the prepared plan for an `AddNote` proposal change. Mirrors
572/// [`prepare_add_entity`]'s shape and the same
573/// `kind`-already-canonicalized split. `annotates` is out of scope: the
574/// proposal `NoteDraft` this backs carries no annotates targets, unlike
575/// `KhiveRuntime::create_note`'s general-purpose signature.
576pub 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    // Same note-write validator `create_note_inner` runs: this path builds its
588    // args itself and dispatches no pack hook, so without this call a proposal
589    // changeset would be the one note-write that stores caller-supplied owned
590    // identity properties verbatim. The token here is the applying caller's
591    // (threaded in by the apply worker), not the proposer's.
592    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(&note),
614        guard: Some(AffectedRowGuard::exactly(1)),
615    }];
616    // Order-sensitive pair — see `insert_document_statements`'s adjacency
617    // contract.
618    for statement in insert_document_statements("fts_notes", &note_fts_document(&note)) {
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
636// ---------------------------------------------------------------------------
637// update
638// ---------------------------------------------------------------------------
639
640/// Mirrors `khive-pack-kg::handlers::update::reject_inapplicable_fields`: a
641/// hard `InvalidInput` when a caller passes a field that does not apply to
642/// the resolved substrate (e.g. `salience` on an entity, or
643/// `description` on a note). That function has no dependency edge
644/// back to `khive-runtime`, so its exact field-applicability check list and
645/// error message shape are reimplemented here rather than imported: same
646/// pattern as `optional_string_patch` above. Presence is checked directly on
647/// the raw args object (this module has no `UpdateParams` struct); a JSON
648/// `null` value is treated as absent, matching `Option<T>` deserialization
649/// semantics.
650fn 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                // ADR-014 tri-state: a PRESENT key (including JSON `null`,
693                // the explicit clear) is inapplicable to notes.
694                Some("entity_type")
695            } else {
696                None
697            };
698            (
699                bad,
700                "name, content, salience, decay_factor, properties, tags",
701            )
702        }
703        // `update` admits `kind="edge"` per `ATOMIC_ADMISSIBLE_VERBS`, so
704        // this arm must reject entity/note-only fields (e.g. `name`) on an
705        // edge update rather than silently skip the guard, mirroring
706        // `khive-pack-kg::handlers::update::reject_inapplicable_fields`'s
707        // `KindSpec::Edge` arm.
708        "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                // ADR-014 tri-state: a PRESENT key (including JSON `null`,
723                // the explicit clear) is inapplicable to edges.
724                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
746/// Caller-supplied update-kind expectation, resolved via the canonical
747/// `resolve_kind_spec` at the kkernel `--atomic` seam: the same pattern
748/// [`AtomicDeleteKind`] uses. Without this check, `update(kind="document",
749/// id=<concept>)` would be canonically `NotFound` but the atomic path would
750/// ignore the explicit kind and mutate the resolved entity anyway.
751/// `khive-runtime` must not depend on `khive-pack-kg`, so this is a plain
752/// substrate-level shape rather than `khive_pack_kg::handlers::KindSpec`
753/// itself: the kkernel seam does the pack-aware resolution and passes down
754/// only what `prepare_update` needs to enforce the mismatch check.
755pub enum AtomicUpdateKind {
756    Entity { specific: Option<String> },
757    Note { specific: Option<String> },
758    Edge,
759}
760
761/// Enforce a caller's explicit update-kind discriminator against a resolved
762/// note before any pack hook can inspect or normalize the request. Canonical
763/// KG dispatch performs this mismatch check before its hook; the atomic
764/// adapter calls this same helper to preserve that error ordering.
765pub 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 &note.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(&note, 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, &note, &patch, &mut plan).await?;
842    }
843    Ok((updated, plan))
844}
845
846/// Build an atomic update plan from the exact note snapshot already supplied
847/// to a pack update hook. Persistence is guarded by that snapshot's revision
848/// and deletion marker, so hook normalization cannot race a second read. The
849/// caller must first run the registry normalizer/validator against this snapshot.
850/// This shared canonical/atomic seam prepares all patch fields with the owning
851/// kind's property policy (the value `VerbRegistry::prepare_note_update_policy`
852/// returned for this snapshot), then derives and attaches the owner's typed
853/// graph effects. It returns the projected note and one atomic plan, preserving
854/// the caller's operation index.
855pub 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
884/// The only companion attachment site. Canonical, CLI atomic, and stream note
885/// updates all enter through `prepare_update_from_note_snapshot` after running
886/// the registry's normalizer/validator on the same snapshot.
887async 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    // Refuse already-known stale input before asking the owner to derive effects.
902    // A change after this read remains protected by the note's first CAS statement.
903    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                // A live annotation inserted between the owner's read and this
959                // lookup belongs to that writer. Refuse instead of replacing its
960                // weight/metadata; a fresh preparation can preserve it explicitly.
961                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
1029/// `expected_kind`: `None` when the caller omitted `kind` (no check, parity
1030/// with canonical's own optional discriminator); `Some(_)` enforces an
1031/// exact-parity mismatch check against the resolved record's actual
1032/// substrate/specific kind, mirroring `handle_update`'s
1033/// `entity.kind != *k` / note kind checks (update.rs:200-201, :229-234).
1034/// Task notes must pass the GTD pack hook before this or
1035/// `prepare_op("update", ..)`; neither runs it.
1036pub 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    // Mirrors update.rs's entity_kind immutability guard: entity_kind is a
1045    // legacy top-level field, independent of the `kind` substrate
1046    // discriminator handled elsewhere.
1047    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            // Decide step lives in curation.rs's `prepare_update_entity` —
1073            // the SAME function canonical `update_entity` calls. Only the
1074            // arg-extraction (raw JSON -> `EntityPatch`) and the plan-shape
1075            // wiring are atomic-path-specific: the domain object becomes a
1076            // `PlanStatement` via `entity_replace_if_unchanged_statement`,
1077            // whose CAS predicate binds the snapshot's revision and deletion
1078            // marker, under an exactly-one affected-row guard.
1079            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            // Patch application lives in curation.rs's
1112            // `prepare_update_note_from_snapshot` — the same implementation
1113            // canonical guarded update calls, including salience/decay range
1114            // validation. The plan retains this exact snapshot's revision.
1115            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        // `Resolved` (khive-runtime::operations) has no `Edge` variant — an
1131        // id that isn't an entity/note/pack-private/event record is checked
1132        // against the graph store directly, mirroring
1133        // `khive-pack-kg::handlers::KgPack::infer_kind_from_uuid`'s own
1134        // entity/note-then-edge fallback order. `update` admits
1135        // `kind="edge"`, so this arm must be able to build a plan for one.
1136        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
1151/// Build an entity update plan from a typed patch. Proposal changesets use
1152/// this entry point so their explicit `description: null` clear operation is
1153/// preserved instead of being collapsed by raw verb deserialization.
1154pub 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
1218/// Edge branch of `prepare_update`. Mirrors `KhiveRuntime::update_edge`'s
1219/// patch semantics: `relation`/`weight`/`properties` are the only applicable
1220/// fields, a changed `relation` is endpoint-validated first, `weight` is
1221/// range-checked, and `properties` REPLACES `metadata` wholesale (no merge).
1222/// See `docs/api/atomic_prepare.md#prepare_update_edge` for the DML-shape parity
1223/// detail with `update_edge`.
1224///
1225/// Invariant (symmetric relations `competes_with`/`composed_with`): this
1226/// function must never branch on a prepare-time conflict probe — a different
1227/// op in the same atomic unit could change the conflict landscape between
1228/// probe and commit, making any such branch stale by construction. It always
1229/// emits BOTH statements from [`edge_symmetric_delete_if_conflict_statement`]
1230/// and [`edge_symmetric_absorb_or_update_inplace_statement`], each carrying
1231/// its own commit-time `WHERE`/`CASE WHEN` predicate that re-evaluates the
1232/// conflict condition fresh inside the transaction. This function reads no
1233/// state to guess a surviving id; the plan instead carries `edge_natural_key`
1234/// so a post-commit caller derives the actual surviving id from the
1235/// committed row, never from a value computed before the rest of this atomic
1236/// unit has even run.
1237async 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        // The write for a symmetric relation never branches on a
1296        // prepare-time probe result: it always carries both self-guarding,
1297        // commit-time-predicate statements (see their doc comment in
1298        // khive-db's graph.rs for the full rationale). This avoids the
1299        // staleness window a prepare-time probe would expose: an earlier op
1300        // in the same atomic unit could change the conflict landscape before
1301        // commit. Canonical's own probe-then-branch
1302        // `update_edge_symmetric_dml` has no such exposure (single
1303        // transaction, no interleaving) and is unaffected.
1304        let metadata_str = edge
1305            .metadata
1306            .as_ref()
1307            .map(|v| serde_json::to_string(v).unwrap_or_default());
1308
1309        // `updated_at` must strictly advance past the snapshot even when two
1310        // operations land inside one clock microsecond; saturating to
1311        // i64::MAX would let the CAS accept a write without advancing its
1312        // revision, so that is not a valid fallback (mirrors the note path).
1313        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        // No prepare-time read needed: the two statements above are
1357        // self-guarding at commit time (see their doc comment). Post-commit
1358        // result rendering derives the actual surviving id from THIS
1359        // natural key, never from a value computed here.
1360        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        // Non-symmetric: guarded replace of the read snapshot rather than
1368        // `graph.upsert_edge`'s unconditional natural-key upsert — a zero
1369        // affected-row result now means a concurrent writer moved this edge
1370        // between PREPARE and commit, and the atomic unit must roll back
1371        // instead of silently overwriting it.
1372        //
1373        // `updated_at` must strictly advance past the snapshot even when two
1374        // operations land inside one clock microsecond; saturating to
1375        // i64::MAX would let the CAS accept a write without advancing its
1376        // revision, so that is not a valid fallback (mirrors the note path).
1377        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    // Mirrors `update_edge`'s unconditional post-mutation `EdgeUpdated`
1402    // event append, keyed on the original `edge_id` the caller supplied:
1403    // canonical does the same (the event target is `edge_id`, not the
1404    // post-absorption surviving id).
1405    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
1429// ---------------------------------------------------------------------------
1430// delete
1431// ---------------------------------------------------------------------------
1432
1433/// Caller-supplied delete-kind expectation, resolved via the canonical
1434/// `resolve_kind_spec` at the kkernel `--atomic` seam. `khive-runtime` must
1435/// not depend on `khive-pack-kg` (packs depend on the runtime, not the other
1436/// way around), so this is a plain substrate-level shape rather than
1437/// `khive_pack_kg::handlers::KindSpec` itself: the kkernel seam does the
1438/// pack-aware `resolve_kind_spec` resolution (which needs a `VerbRegistry`,
1439/// unreachable from this crate) and passes down only what `prepare_delete`
1440/// needs to enforce the mismatch check.
1441///
1442/// `delete` admits `kind="edge"` per `ATOMIC_ADMISSIBLE_VERBS`, hence the
1443/// `Edge` variant. `Event`/`Proposal` remain rejected at the kkernel seam
1444/// (not v1-admissible for atomic delete at all).
1445pub enum AtomicDeleteKind {
1446    Entity { specific: Option<String> },
1447    Note { specific: Option<String> },
1448    Edge,
1449}
1450
1451/// `expected_kind`: `None` when the caller omitted `kind` (no check, parity
1452/// with canonical's own optional discriminator); `Some(_)` enforces an
1453/// exact-parity mismatch check against the resolved record's actual
1454/// substrate/specific kind, mirroring `handle_delete`'s
1455/// `entity.kind != *expected` / `note.kind != *expected` checks.
1456pub 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    // `delete(id, hard=true)` is the public purge route after a prior soft
1470    // delete, so it must resolve including already-tombstoned rows (a
1471    // live-only resolve would never find one). Soft delete keeps the
1472    // live-only resolve: a soft delete of an already-tombstoned row is a
1473    // no-op, matching non-atomic behavior.
1474    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            // Storage parity: `entity_soft_delete_statement`/
1499            // `entity_hard_delete_statement` are the SAME khive-db builders
1500            // khive-db's own `SqlEntityStore::delete_entity` calls — no DML
1501            // text is hand-duplicated here.
1502            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                // Same builder canonical `delete_entity`'s hard-delete
1538                // cascade calls (`graph.purge_incident_edges`).
1539                statements.push(PlanStatement {
1540                    statement: purge_incident_edges_statement(id),
1541                    guard: None,
1542                });
1543            }
1544            // FTS + vector index purge, matching operations.rs
1545            // `delete_entity`: both soft and hard delete clean indexes (a
1546            // hard delete of an already-tombstoned record must still purge
1547            // them); only hard additionally cascades edges above.
1548            push_index_purge_statements(
1549                runtime,
1550                &mut statements,
1551                "fts_entities",
1552                &namespace,
1553                id,
1554                "atomic-delete-entity",
1555            )
1556            .await?;
1557            // operations.rs's `delete_entity` appends an `EntityDeleted`
1558            // event after a successful row delete, on both soft and hard
1559            // delete. `apply_plan` never reaches this statement unless the
1560            // guarded row statement above affected a row, so no extra `if`
1561            // is needed here.
1562            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 &note.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(&note).await? {
1594                return Err(error);
1595            }
1596            let namespace = note.namespace.clone();
1597            // Storage parity: `note_soft_delete_statement`/
1598            // `note_hard_delete_statement` are the SAME khive-db builders
1599            // khive-db's own `SqlNoteStore::delete_note` calls.
1600            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            // FTS + vector index purge, matching operations.rs
1641            // `delete_note`: both soft and hard delete clean indexes (a hard
1642            // delete of an already-tombstoned record must still purge
1643            // them); only hard additionally cascades edges above.
1644            push_index_purge_statements(
1645                runtime,
1646                &mut statements,
1647                "fts_notes",
1648                &namespace,
1649                id,
1650                "atomic-delete-note",
1651            )
1652            .await?;
1653            // operations.rs's `delete_note` appends a `NoteDeleted` event
1654            // after a successful row delete, on both soft and hard delete:
1655            // same reasoning as the entity branch above.
1656            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                // A committed atomic note delete must fire the same
1669                // pack-installed note-mutation hook `operations.rs::
1670                // delete_note` fires, so a warm ANN cache sees the deletion
1671                // even when the mutation went through the atomic-plan path.
1672                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        // `Resolved` has no `Edge` variant (same reasoning as
1682        // `prepare_update`'s fallback above) — probe the graph store
1683        // directly.
1684        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
1706/// Edge branch of `prepare_delete`. Mirrors
1707/// `khive-runtime::operations::KhiveRuntime::delete_edge` exactly: hard
1708/// delete cascades `purge_incident_edges` (any `annotates` edge — or any
1709/// other edge — pointing AT this edge as a node) BEFORE deleting the edge
1710/// row itself, then a soft or hard delete statement, then an unconditional
1711/// `EdgeDeleted` event (edges are never FTS/vector-indexed, so unlike the
1712/// entity/note branches there is no index purge here — `delete_edge` has
1713/// none either).
1714async 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        // Mirrors `delete_edge`'s `graph.purge_incident_edges(edge_id)` —
1734        // unguarded: zero incident edges is a legitimate outcome, not a
1735        // failure (same reasoning as the entity/note cascade-edges
1736        // statements above).
1737        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
1770// ---------------------------------------------------------------------------
1771// link
1772// ---------------------------------------------------------------------------
1773
1774fn 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    // Top-level `dependency_kind` param merges into `metadata`: only fills
1800    // the key when metadata doesn't already carry one. Calls the same
1801    // `khive_runtime::merge_entry_metadata` `khive-pack-kg`'s canonical
1802    // `handle_link` calls, so both sides depend on one function instead of
1803    // each maintaining their own copy.
1804    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    // Endpoint-kind `dependency_kind` inference for `depends_on` edges,
1817    // matching operations.rs `link()`: only applies when both endpoints
1818    // resolve as entities and the key is still absent after the
1819    // top-level-param merge above. Runs against the canonical endpoints,
1820    // mirroring `KhiveRuntime::link`'s own ordering (canonicalize, then
1821    // infer).
1822    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    // The guarded mutation closes both atomic seams: endpoints are re-probed
1876    // inside the transaction, and the natural-key row must still match the
1877    // prepare snapshot. That makes the disposition used below truthful.
1878    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
1934// ---------------------------------------------------------------------------
1935// merge (entity-only)
1936// ---------------------------------------------------------------------------
1937
1938// Full atomic-merge parity (field folding, survivor FTS/vector reindex,
1939// loser index purge, merge provenance, same-kind rejection) is deferred:
1940// atomic `merge` is rejected entirely at the pre-runtime admissibility
1941// guard (`khive_types::pack::ATOMIC_KNOWN_UNIMPLEMENTED_VERBS`, alongside
1942// `propose`/`review`/`withdraw`). This function still produces a plan
1943// (kept for the existing direct-prepare test coverage below and as
1944// defense in depth), but the CLI's `--atomic` surface never reaches it,
1945// since `check_atomic_admissible` rejects `merge` before any runtime is
1946// built.
1947async 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// ---------------------------------------------------------------------------
2023// post-commit effects
2024// ---------------------------------------------------------------------------
2025
2026/// Embedding metadata produced by one successfully applied reindex effect.
2027///
2028/// The effect identity is retained so response builders can attach an advisory
2029/// to the exact atomic op that scheduled the reindex. Effects whose target is
2030/// no longer present are omitted from the returned outcome list.
2031#[derive(Clone, Debug, PartialEq, Eq)]
2032pub struct PostCommitEmbeddingOutcome {
2033    /// The committed reindex effect that produced this outcome.
2034    pub effect: PostCommitEffect,
2035    /// Actual input bounding observed while executing the effect.
2036    pub truncation: crate::retrieval::EmbeddingTruncationReport,
2037}
2038
2039/// Run every deferred [`PostCommitEffect`] after a committed atomic unit.
2040pub 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
2050/// Truncation-reporting form of [`apply_post_commit_effects`]. Re-fetches each
2051/// target's now-committed row outside any transaction and reuses the existing
2052/// `reindex_entity`/`reindex_note` (FTS + embedding, same as the non-atomic
2053/// path) for exact parity. Returns the typed embedding outcome for each reindex
2054/// effect so callers can preserve write-response advisories instead of
2055/// discarding them after commit.
2056pub 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, &note).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            // Atomic note updates bypass the regular update_note hook. Notify
2119            // in-process consumers only after the current version was indexed.
2120            runtime.fire_note_mutation_hook(&note.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            // The committed row may already be gone; use the captured kind.
2128            runtime.fire_note_mutation_hook(&kind, note_id).await;
2129            Ok(None)
2130        }
2131        PostCommitEffect::GtdAudit { .. } => {
2132            // The kkernel caller owns the GTD pack's separate audit side write.
2133            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    /// Owns a file-backed runtime and removes its database directory after shutdown.
2152    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    /// Atomic `update` must reject a field that does not apply to the
2223    /// resolved substrate: parity with
2224    /// `khive-pack-kg::handlers::update::reject_inapplicable_fields`.
2225    /// Without this check, atomic prepare would silently ignore `salience`
2226    /// on an entity: it would set every entity field to its current value,
2227    /// bump `updated_at`, satisfy the `exactly(1)` guard, and commit: a
2228    /// spurious no-op reported as success.
2229    #[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        // A valid entity update (name/description/tags) must still work.
2258        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    /// The `--atomic` seam intentionally requires a full UUID (never a short
2286    /// hex prefix) for its `id` fields — prefix resolution already happened
2287    /// upstream, at the kkernel CLI boundary that has namespace context
2288    /// (`resolve_kg_ids_in_args`). The rejection message must say *why* a
2289    /// prefix cannot be resolved at this stage, not just restate that a full
2290    /// UUID is required.
2291    #[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    /// Symmetric note-substrate case: `description` is entity-only and
2388    /// must be rejected the same way update.rs rejects it.
2389    #[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        // A valid note update (content) must still work.
2419        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    /// Updating a note's content inside an atomic unit must, after commit,
2621    /// leave the note recallable via FTS under its new content and its
2622    /// vector row refreshed: parity with the non-atomic
2623    /// `update_note` -> `reindex_note` path.
2624    #[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        // Sanity: no vector row yet for the stub model.
2643        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        // FTS: the note must be recallable under its NEW content.
2690        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        // Vector: a row must now exist for the registered stub model.
2704        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    /// The atomic-plan path must fire the pack-installed note-mutation hook
2712    /// for both an atomic note UPDATE (`PostCommitEffect::ReindexNote`'s
2713    /// handler fires it after its own reindex, mirroring `update_note()`
2714    /// on the non-atomic path) and an atomic note DELETE (`DeletePlan`
2715    /// carries a `PostCommitEffect::NoteDeleted` that
2716    /// `apply_post_commit_effects` dispatches directly, mirroring
2717    /// `operations.rs::delete_note`'s direct-fire, no-refetch shape: the
2718    /// row may already be gone by the time this runs, for a hard delete).
2719    /// A minimal counting hook proves both fire; no `khive-pack-memory`
2720    /// dependency is needed at this layer, since the hook itself is
2721    /// generic.
2722    #[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        // Update path.
2742        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        // Delete path (soft delete: the row still exists, but the hook
2772        // fires directly from the captured kind rather than refetching).
2773        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                    // Explicit embedding is required here: with embed=None
2855                    // and no existing vector row, atomic_runner legitimately
2856                    // replaces ReindexNote with NoteChanged (inheritance
2857                    // policy), leaving no reindex for the fault to fail.
2858                    &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        // The already-committed note rows remain intact, but both deferred
2912        // reindexes now fail. The following deletion hook must still fire.
2913        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    /// Atomic delete must purge the note's FTS row and vector row for both
2937    /// soft and hard delete: parity with `KhiveRuntime::delete_note`'s
2938    /// index-cleanup contract.
2939    #[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, &note)
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    /// Atomic delete must purge the entity's FTS row and vector row for
3020    /// both soft and hard delete: parity with
3021    /// `KhiveRuntime::delete_entity`'s index-cleanup contract.
3022    #[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, &note)
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    /// Atomic prepare validates the requested orientation before canonicalizing
3168    /// a symmetric edge for persistence. Fixed UUIDs force target < source so a
3169    /// regression would report the reverse ordered pair.
3170    #[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    /// Atomic link must persist an explicit top-level `dependency_kind`
3214    /// param into edge metadata, and must infer one for `depends_on` edges
3215    /// when absent: parity with
3216    /// `link.rs`'s `merge_entry_metadata` and `operations.rs`'s
3217    /// `infer_dependency_kind` table.
3218    #[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        // (a) explicit top-level `dependency_kind` param persists in metadata.
3238        {
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        // (b) `depends_on` with no explicit dependency_kind infers from
3265        // endpoint kinds: (service, service) -> "runtime".
3266        {
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    /// Raw natural-key probe of `graph_edges` (namespace, source_id,
3293    /// target_id, relation) — returns `(weight, metadata_json, deleted_at)`
3294    /// for exactly the ONE row a UNIQUE(namespace, source_id, target_id,
3295    /// relation) constraint permits. `None` means no row at all.
3296    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    /// Atomic `link` must be an upsert, exactly like canonical `link` ->
3340    /// `upsert_edge`'s natural-key `ON CONFLICT` arm: re-linking an
3341    /// already-linked triple must succeed and update weight/metadata, not
3342    /// hit the `UNIQUE(namespace, source_id, target_id, relation)`
3343    /// constraint and roll back the whole atomic unit.
3344    #[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        // (a) a fresh link still works.
3359        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        // (b) re-linking the SAME triple with a different weight/metadata
3388        // must SUCCEED (not a constraint-violation rollback) and UPDATE the
3389        // existing row in place — natural key stays unique.
3390        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    /// Atomic `link` refuses a soft-deleted triple by default and only
3430    /// resurrects it when the caller opts in explicitly.
3431    #[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        // Soft-delete the edge row directly (natural-key UPDATE — mirrors
3465        // what `delete_edge(hard=false)` does to this same row).
3466        {
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    /// Atomic delete of an entity and a note must succeed even when the
3549    /// registered embedding model's `vec_*` table has never been lazily
3550    /// created (a fresh DB registers models before any vector store is
3551    /// opened): the raw purge DML must skip tables that don't exist
3552    /// rather than hit `no such table` and roll back the whole atomic
3553    /// unit. FTS purge still fires (those tables always exist) and the
3554    /// delete itself is a clean commit.
3555    #[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        // Seed via raw upsert ONLY — never call reindex_entity/reindex_note
3564        // or vectors_for_model, so the stub model's `vec_*` table is never
3565        // lazily created (opening the vector store is what creates it).
3566        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    /// Atomic hard delete must be able to purge a record that was already
3626    /// soft-deleted: parity with `delete(id, hard=true)` being the public
3627    /// purge route after a prior soft delete (the non-atomic hard path
3628    /// resolves including deleted rows and its DML carries no `deleted_at`
3629    /// predicate).
3630    #[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, &note)
3664            .await
3665            .expect("seed index rows");
3666
3667        // First: SOFT delete both (via atomic prepare) so they're tombstoned
3668        // going into the hard-delete attempt below.
3669        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        // Now: HARD delete the already-tombstoned records.
3703        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    // ------------------------------------------------------------------
3776    // event-store append parity
3777    // ------------------------------------------------------------------
3778
3779    /// Fetch every event of `kind` targeting `target_id`, via the same
3780    /// `EventStore::query_events` surface `--atomic` callers would use to
3781    /// verify parity — not a raw SQL probe.
3782    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    /// Atomic `update(id=<entity>, name=...)` must append an
3804    /// `EntityUpdated` event, matching `curation::update_entity`: the
3805    /// event is appended unconditionally after a successful row update,
3806    /// not only on the reindex-triggering subset.
3807    #[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    /// Atomic soft and hard delete of an entity must each append an
3858    /// `EntityDeleted` event, matching `operations::delete_entity`, which
3859    /// fires on both delete modes.
3860    #[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    /// Atomic soft and hard delete of a note must each append a
3966    /// `NoteDeleted` event, matching `operations::delete_note`, which
3967    /// fires on both delete modes.
3968    #[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    /// `update` admits `kind="edge"` per `ATOMIC_ADMISSIBLE_VERBS`; this
4020    /// asserts `prepare_update` actually builds a plan for one, a
4021    /// non-symmetric relation (`extends`) exercises the
4022    /// `edge_upsert_statement` reuse branch, and that the committed row +
4023    /// `EdgeUpdated` event match canonical `update_edge`'s shape (weight
4024    /// persisted, relation unchanged, exactly one event).
4025    #[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    /// ADR-115 Amendment 1 §3: the reserved `khive:secret_gate` property key
4085    /// must be rejected on atomic edge-metadata updates the same way it is
4086    /// on canonical `update_edge`.
4087    #[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    /// The symmetric-relation conflict-absorption branch of
4134    /// `prepare_update_edge` — mirrors `update_edge_symmetric_dml`'s case
4135    /// (b): changing a non-symmetric edge's `relation` to a symmetric one
4136    /// whose canonical natural key collides with an ALREADY-EXISTING
4137    /// symmetric edge between the same two entities must delete the
4138    /// requested (non-canonical) row and leave the surviving canonical row
4139    /// untouched (ADR-039 ON CONFLICT DO NOTHING), rather than raising a
4140    /// uniqueness error OR overwriting the survivor with the discarded
4141    /// edge's attributes (khive#1213).
4142    #[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        // The non-canonical edge under test: A -> B, non-symmetric relation.
4156        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        // The pre-existing canonical row this update will collide with once
4163        // `relation` becomes `competes_with` (symmetric).
4164        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        // The plan does not compute a prepare-time advisory surviving id
4180        // (`target_id` is just the requested id): it carries
4181        // `edge_natural_key` so a post-commit caller can derive the real
4182        // surviving id itself. Assert the plan carries the right natural
4183        // key to look up; the actual surviving row's identity is verified
4184        // against the DB after commit, below.
4185        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        // The requested (non-canonical) row must be gone.
4213        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        // ADR-039 DO NOTHING: the surviving canonical row keeps its OWN
4223        // pre-existing attributes — the discarded edge's patched weight
4224        // (0.9) must never overwrite it.
4225        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        // Event target is the CALLER-supplied id, not the surviving id —
4237        // mirrors `update_edge`'s event using `edge_id` (the caller's
4238        // original argument), not the post-absorption id.
4239        let events =
4240            events_for_target(&runtime, &token, requested_id, EventKind::EdgeUpdated).await;
4241        assert_eq!(events.len(), 1);
4242    }
4243
4244    /// Regression for khive #1753: an entity atomic update plan is prepared
4245    /// from one revision, a concurrent (out-of-plan) writer advances the row
4246    /// past that revision before the plan commits, and the plan's guarded
4247    /// `entity_replace_if_unchanged_statement` must then affect zero rows —
4248    /// which `AffectedRowGuard::exactly(1)` must turn into a whole-unit
4249    /// rollback, not a silent overwrite of the concurrent writer's change.
4250    /// Before threading the expected revision into the plan's `WHERE`
4251    /// predicate, this plan used the unconditional `entity_upsert_statement`,
4252    /// which always affects exactly 1 row regardless of staleness — the
4253    /// guard could never fire and this test would redden (the outcome would
4254    /// be `Committed`, and the concurrent writer's `name` change would be
4255    /// lost under the stale plan's `description`-only patch).
4256    #[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        // PREPARE time: build the plan from the current (soon-to-be-stale)
4277        // revision.
4278        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        // Read the plan's OWN bound values, so the isolation below is
4288        // derived from the plan rather than assumed about it. Per
4289        // `entity_replace_if_unchanged_statement`, `?8` (index 7) is the
4290        // replacement revision, `?12` (index 11) the target id, `?13`
4291        // (index 12) the expected revision, and `?14` (index 13) the
4292        // expected deletion marker.
4293        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            // Locate by LABEL, never by the guard text. A locator keyed on
4299            // `?8 > updated_at` would stop finding the statement in exactly the
4300            // mutation run that deletes that conjunct, so the arm would report
4301            // on this fixture's locator instead of on the guard.
4302            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        // The identity conjunct, on the same footing as the revision and the
4322        // deletion marker. `id = ?12` is live in the same UPDATE, so a plan
4323        // that bound any other row's id would affect zero rows and roll the
4324        // unit back for a reason this test does not name — an outcome
4325        // indistinguishable from the one it does name.
4326        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        // A concurrent writer commits BEFORE the plan runs, advancing the
4335        // row's revision past what the plan's guard expects.
4336        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        // Pin the stored revision one microsecond BELOW the plan's replacement.
4349        // This is what makes the test specific to the guard it names. Both the
4350        // plan and the concurrent writer derive their revision from
4351        // `max(now, expected + 1)`, so left alone the concurrent write lands at
4352        // or past the plan's own replacement and BOTH conjuncts refuse — the
4353        // test would then stay green if either guard were deleted. Pinning
4354        // leaves `?8 > updated_at` satisfied, so it cannot be the refuser,
4355        // while `updated_at = ?13` is violated, which is the single condition
4356        // this test names. The shape is reachable in production whenever the
4357        // concurrent writer's clock trails the preparing writer's.
4358        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        // The remaining non-target conjunct. `deleted_at IS ?14` is live in the
4388        // same UPDATE, so without this read the fixture would stay green if the
4389        // deletion marker were the actual refuser — the refusal would look
4390        // identical. Read the row as it stands at DML time, after both the
4391        // concurrent writer and the pin.
4392        {
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            // Keep this legacy timestamp-conjunct oracle independent of the
4410            // newly added persisted-version predicate. The production plan
4411            // would also refuse on its old version; this fixture pins only
4412            // that additional predicate to the observed current value.
4413            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                // Attribution rather than inference. The premises above
4437                // establish that every non-target conjunct holds; this names
4438                // the statement that actually refused. Without it a unit that
4439                // rolled back at some other statement reads identically from
4440                // the outside, which is the whole failure mode those premises
4441                // were approximating.
4442                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    /// Edge counterpart of `atomic_entity_update_plan_stale_revision_rolls_back_unit`:
4470    /// a non-symmetric edge atomic update plan prepared from one revision
4471    /// must roll back the whole unit when a concurrent writer advances the
4472    /// row first, rather than silently overwriting that writer's change.
4473    #[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        // PREPARE time: build the plan from the current (soon-to-be-stale)
4493        // revision.
4494        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        // Read the plan's OWN bound values, so the isolation below is
4504        // derived from the plan rather than assumed about it. Per
4505        // `edge_replace_if_unchanged_statement`, `?6` (index 5) is the
4506        // replacement revision, `?10` (index 9) the target id, `?11`
4507        // (index 10) the expected revision, and `?12` (index 11) the expected
4508        // deletion marker.
4509        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            // Locate by LABEL, never by the guard text — see the entity sibling
4515            // above: a locator keyed on the conjunct disappears in exactly the
4516            // mutation run that deletes it.
4517            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        // The identity conjunct, for the same reason as the entity sibling:
4537        // `id = ?10` is live in the same UPDATE and a misbound target refuses
4538        // indistinguishably.
4539        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        // A concurrent writer commits BEFORE the plan runs, advancing the
4548        // row's revision past what the plan's guard expects.
4549        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        // Pin the stored revision one microsecond BELOW the plan's replacement,
4562        // for the same reason as the entity sibling above: both the plan and
4563        // the concurrent writer derive their revision from
4564        // `max(now, expected + 1)`, so left alone BOTH conjuncts refuse and the
4565        // test would stay green if either guard were deleted. Pinning leaves
4566        // `?6 > updated_at` satisfied, so it cannot be the refuser, while
4567        // `updated_at = ?11` is violated — the single condition this test names.
4568        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        // The remaining non-target conjunct, for the same reason as the entity
4597        // sibling: `deleted_at IS ?12` is live in the same UPDATE and would
4598        // produce an indistinguishable refusal.
4599        {
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                // Attribution rather than inference — see the entity sibling.
4633                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    /// A soft-deleted surviving canonical row must not be resurrected by a
4665    /// conflicting symmetric-relation update (ADR-039 DO NOTHING; khive#1213):
4666    /// the requested edge is still deleted (conflict absorbed), but the
4667    /// tombstoned survivor must stay tombstoned.
4668    #[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    /// The same-unit race: `[delete(X), update(X -> competes_with)]` where
4737    /// an already-existing canonical row sits at the post-update natural
4738    /// key. Both ops' async prepare passes run before either commits, so at
4739    /// prepare time `X` still exists and both plans build. At commit time
4740    /// `delete(X)` removes it first; `update(X -> competes_with)`'s own
4741    /// commit-time statements must then fail loud (its target no longer
4742    /// exists) rather than silently absorbing into the pre-existing
4743    /// canonical row it never causally touched. The whole atomic unit must
4744    /// roll back — parity with canonical `update_edge`'s `NotFound` for a
4745    /// missing edge, expressed here as the unit-level abort for any op
4746    /// whose commit-time guard fails.
4747    #[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        // The row op 1 will try to update — deleted by op 0 in the SAME unit.
4761        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        // The pre-existing canonical row the buggy `id = ?2 OR natural-key`
4768        // predicate used to silently absorb into.
4769        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        // Whole-unit rollback: op 0's delete must be undone too.
4811        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        // The pre-existing canonical row must be completely untouched.
4820        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    /// The symmetric absorption delete's atomic commit-time statement
4832    /// (`edge_symmetric_delete_if_conflict_statement`) must refuse a plan
4833    /// built from a since-changed snapshot, exactly like the non-symmetric
4834    /// `atomic_edge_update_plan_stale_revision_rolls_back_unit` case above —
4835    /// even though a genuine canonical survivor exists at the target natural
4836    /// key. `E`'s own revision changes (via an unrelated production
4837    /// `update_edge` call) AFTER `prepare_update` captured its snapshot but
4838    /// BEFORE the plan runs; the whole atomic unit must roll back rather than
4839    /// silently absorbing `E` into the survivor and discarding the
4840    /// concurrent write.
4841    #[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        // The pre-existing canonical row this plan will try to absorb into.
4855        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        // E: the requested edge under test.
4862        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        // PREPARE time: build the plan from the current (soon-to-be-stale)
4869        // revision.
4870        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        // A concurrent writer commits BEFORE the plan runs, advancing E's
4880        // revision past what the plan's guard expects. Relation stays
4881        // `extends` (non-symmetric), so this lands via the already-guarded
4882        // replace path — independent of the absorption bug under test.
4883        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        // Statement 2's in-place arm carries a strict-advance conjunct
4897        // (`?7 > updated_at`) beside the expected-revision conjunct
4898        // (`updated_at = ?10`) this test names. Both the plan and the
4899        // concurrent writer derive their revision from `max(now, expected + 1)`
4900        // and the concurrent writer ran LATER, so left alone the stored
4901        // revision sits at or past the plan's replacement and BOTH conjuncts
4902        // refuse — the test would stay green with either one deleted. Pin the
4903        // stored revision one microsecond below the plan's replacement, the
4904        // same treatment the entity and edge siblings above already carry.
4905        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        // A symmetric plan carries TWO guarded statements, and it matters
4962        // which one refuses.
4963        //
4964        //   1. `edge-symmetric-delete-if-conflict`, guarded
4965        //      `{expected_min: 0, expected_max: Some(1)}` — a zero-row delete
4966        //      SATISFIES this guard. It does not roll the unit back. Its
4967        //      effect here is to leave SQLite's `changes()` at 0.
4968        //   2. `edge-symmetric-absorb-or-update-inplace`, guarded
4969        //      `exactly(1)` — this is the statement that refuses and rolls the
4970        //      unit back, through its in-place arm
4971        //      `(id = ?2 AND changes() = 0 AND updated_at = ?10 AND
4972        //        deleted_at IS ?11 AND ?7 > updated_at)`,
4973        //      whose `updated_at = ?10` is the stale-revision guard under
4974        //      test. The absorbed arm needs `changes() = 1`, so the delete
4975        //      refusing is what selects the in-place arm.
4976        //
4977        // So the premises come in two layers. The delete's own predicates
4978        // (`?6`, `?7`, and the EXISTS survivor at `?3`/`?4`/`?5`) are premised
4979        // because they decide WHETHER the delete fires, and therefore which
4980        // arm of statement 2 is live. Statement 2's remaining conjuncts are
4981        // premised because each of them refuses identically to the revision
4982        // guard this test names. Read the plans' own bound values and check
4983        // the world against them, immediately before the DML.
4984        {
4985            let statements = match &plan {
4986                AtomicOpPlan::Update(p) => p.statements.clone(),
4987                other => panic!("expected an Update plan, got {other:?}"),
4988            };
4989            // By LABEL, for the same reason as the siblings above.
4990            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            // The delete's OWN identity conjuncts. Reading them into the
5009            // survivor-count panic message below is not asserting them: point
5010            // `?1` or `?2` at a row that does not exist and the delete still
5011            // affects zero rows, its `0..=1` guard still accepts that, and
5012            // statement 2 still refuses on the pinned stale revision — so the
5013            // test stays green while establishing nothing about which row the
5014            // named delete attempted.
5015            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            // The EXISTS arm. A missing survivor refuses the delete on its own
5048            // and looks identical from outside, so the premise has to be that
5049            // arm's OWN question. Do not reconstruct it from the seeded row's
5050            // endpoints: the natural key the plan binds is the canonical
5051            // ordering, which need not equal the order the survivor was stored
5052            // in, and a hand-built comparison would be asserting my model of
5053            // canonicalisation rather than the predicate. Run the subquery
5054            // verbatim against the plan's own bound parameters instead.
5055            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        // Statement 2's remaining conjuncts. Each of these refuses the in-place
5094        // arm identically to the revision guard this test names, so each has to
5095        // be shown satisfied for the attribution to hold.
5096        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                // The decisive attribution, and the reason the two-layer
5141                // premise stack above exists: `failed_op_index` cannot
5142                // distinguish the two statements, because both live in op 0.
5143                // Naming the label is what says the in-place absorb refused
5144                // rather than the delete — and the delete's `0..=1` guard
5145                // means it CANNOT be the refuser, so a run that reported it
5146                // here would be evidence the guards had been rewired.
5147                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    /// `update` rejects an entity/note-only field (`name`) on an edge
5193    /// target, mirroring
5194    /// `khive-pack-kg::handlers::update::reject_inapplicable_fields`'s
5195    /// `KindSpec::Edge` arm.
5196    #[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    /// `delete` admits `kind="edge"` per `ATOMIC_ADMISSIBLE_VERBS`; this
5230    /// asserts `prepare_delete` actually builds a plan for one on both soft
5231    /// and hard delete, matching `operations::delete_edge`'s row-mode DML
5232    /// and unconditional `EdgeDeleted` event.
5233    #[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    /// Parity boundary: an atomic `update` of a note appends exactly one
5296    /// `NoteUpdated` event, because this path and canonical `update_note` build
5297    /// their plan through the same `prepare_versioned_note_update`, which is
5298    /// where the event statements are added.
5299    ///
5300    /// This test used to assert the opposite. That was a faithful record of a
5301    /// gap rather than a contract: notes were the substrate that recorded no
5302    /// update at all, so the parity it certified was parity with nothing.
5303    #[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    /// Atomic `link` commits its mutation and event-plane observation in the
5361    /// same unit.
5362    #[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    // ------------------------------------------------------------------
5433    // AddEntity and AddNote plans alongside link
5434    // ------------------------------------------------------------------
5435
5436    #[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    /// ADR-115 Amendment 1 §3: proposal materialization (`prepare_add_entity`)
5497    /// must reject the reserved `khive:secret_gate` property key the same way
5498    /// canonical `create` does — proposal apply is reservation-only, never an
5499    /// exemption-consuming path.
5500    #[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    /// Note-substrate counterpart of
5525    /// `prepare_add_entity_rejects_reserved_secret_gate_property`.
5526    #[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        // Prepare sees x before the transaction; the guarded link must detect
5721        // that the preceding hard delete removed it inside the transaction.
5722        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    /// ADR-014 tri-state on the atomic path: `entity_type: null` must
5819    /// explicitly CLEAR a stored entity type (and reindex), not collapse to
5820    /// "unchanged" like the old `optional_create_string` did.
5821    #[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    /// ADR-014 tri-state on the atomic path: a PRESENT `entity_type` key,
5883    /// including JSON `null`, must be rejected on note and edge targets
5884    /// (parity with `khive-pack-kg`'s `reject_inapplicable_fields`).
5885    #[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}