Skip to main content

khive_runtime/curation/
note_curation.rs

1use super::{
2    kind_owned_properties, merge_properties, note_update_values_equal,
3    owner_established_property_named_in, reject_pack_managed_schedule_mutation,
4    stale_note_snapshot_error, EntityDedupMergePolicy, KhiveRuntime, NamespaceToken, NotePatch,
5    RuntimeError, RuntimeResult, Uuid, Value,
6};
7
8impl KhiveRuntime {
9    /// Apply a note patch to exactly the supplied read snapshot without
10    /// fetching the row again. The caller must persist it through
11    /// [`Self::update_note_from_snapshot_with_embedding_report`] or a write
12    /// plan guarded by the snapshot's `updated_at`/`deleted_at` values.
13    pub(crate) async fn prepare_update_note_from_snapshot(
14        &self,
15        _token: &NamespaceToken,
16        mut note: khive_storage::note::Note,
17        patch: NotePatch,
18    ) -> RuntimeResult<(khive_storage::note::Note, bool, bool)> {
19        if note.properties.as_ref().is_some_and(|properties| {
20            properties
21                .as_object()
22                .is_some_and(|map| map.contains_key(crate::secret_gate::RESERVED_WEB_RECEIPT_KEY))
23        }) {
24            return Err(RuntimeError::InvalidInput(
25                "web receipt notes are immutable through generic update".into(),
26            ));
27        }
28        // The stored row as read. A no-op answers with this, not with the
29        // patched snapshot: the patch may differ from the row in ways the
30        // no-op decision ignores (tag order), and nothing was written.
31        let stored = note.clone();
32        let original_name = note.name.clone();
33        let original_content = note.content.clone();
34        let original_salience = note.salience;
35        let original_decay_factor = note.decay_factor;
36        let original_properties = note.properties.clone();
37        let original_status = note.status.clone();
38        if patch
39            .update_policy
40            .kind
41            .as_deref()
42            .is_some_and(|kind| kind != note.kind)
43        {
44            return Err(RuntimeError::InvalidInput(
45                "note update policy does not match the stored note kind".into(),
46            ));
47        }
48        if patch.content.is_some() || patch.properties.is_some() {
49            if let Some(error) = self.stream_member_error(&note).await? {
50                return Err(error);
51            }
52        }
53        crate::secret_gate::reject_reserved_secret_gate_property(patch.properties.as_ref())?;
54        if let Some(ref content) = patch.content {
55            crate::secret_gate::check_at(content, "note", "content")?;
56        }
57        if let Some(Some(ref name)) = patch.name {
58            crate::secret_gate::check_at(name, "note", "name")?;
59        }
60        if let Some(ref props) = patch.properties {
61            crate::secret_gate::check_json_at(props, "note", "properties")?;
62        }
63
64        reject_pack_managed_schedule_mutation(&note, "update")?;
65
66        let mut text_changed = false;
67
68        if let Some(name_patch) = patch.name {
69            text_changed |= note.name != name_patch;
70            note.name = name_patch;
71        }
72        if let Some(content) = patch.content {
73            text_changed |= note.content != content;
74            note.content = content;
75        }
76        if let Some(salience_patch) = patch.salience {
77            // Reject invalid salience rather than silently clamping caller input.
78            if let Some(s) = salience_patch {
79                if !s.is_finite() || !(0.0..=1.0).contains(&s) {
80                    return Err(crate::RuntimeError::InvalidInput(format!(
81                        "salience must be a finite value in [0.0, 1.0]; got {s}"
82                    )));
83                }
84            }
85            note.salience = salience_patch;
86        }
87        if let Some(decay_patch) = patch.decay_factor {
88            // Reject invalid decay_factor rather than silently clamping caller input.
89            if let Some(d) = decay_patch {
90                if !d.is_finite() || d < 0.0 {
91                    return Err(crate::RuntimeError::InvalidInput(format!(
92                        "decay_factor must be a finite value >= 0.0; got {d}"
93                    )));
94                }
95            }
96            note.decay_factor = decay_patch;
97        }
98        if let Some(props) = patch.properties {
99            // Kind-owned identity is protected below pack hooks, including
100            // direct runtime and atomic/proposal update preparation. The merge
101            // path restores these same keys on its surviving row.
102            let owned_keys = kind_owned_properties(&note.kind);
103            if !owned_keys.is_empty() {
104                let object = props.as_object().ok_or_else(|| {
105                    if note.kind == "message" {
106                        RuntimeError::InvalidInput(
107                            "properties on a `message` note must be patched with an object: a \
108                             non-object patch would replace the transport-owned quarantine and \
109                             channel provenance established by `comm.ingest`"
110                                .into(),
111                        )
112                    } else {
113                        RuntimeError::InvalidInput(format!(
114                            "properties on a `{}` note must be patched with an object; \
115                             a non-object patch would erase its owner-established identity",
116                            note.kind
117                        ))
118                    }
119                })?;
120                if let Some(named) = owned_keys.iter().find(|key| object.contains_key(**key)) {
121                    if note.kind == "message" {
122                        return Err(RuntimeError::InvalidInput(format!(
123                            "`{named}` is transport-owned on a `message` note and cannot be patched; \
124                             only `comm.ingest` may establish quarantine disposition and channel \
125                             provenance"
126                        )));
127                    }
128                    return Err(RuntimeError::InvalidInput(format!(
129                        "`{named}` is not patchable on a `{}` note; \
130                         use `comm.heartbeat` to report health without changing the row's identity",
131                        note.kind
132                    )));
133                }
134            }
135            // On a pack-owned note kind, the properties in
136            // `OWNER_ESTABLISHED_PROPERTIES` are established by the owning pack
137            // and read back by it to decide something structural — who wrote
138            // the record and when, which author-side record it copies, which
139            // conversation it belongs to. A caller cannot patch them here.
140            // Only a patch that *names* one of them is refused, and naming is
141            // the exact test: the merge below is `PreferFrom`, so a patch that
142            // names an owned key would overwrite it while a patch that does
143            // not name it leaves it intact. Every other key still merges
144            // normally — arbitrary metadata on a pack-owned record (a
145            // `blocked_on` note on a `task`) has no other write path and must
146            // keep working.
147            if self.is_pack_owned_note_kind(&note.kind) {
148                // A non-object patch names nothing, so it slips past the
149                // named-key check below and then takes `merge_json`'s
150                // non-object `PreferFrom` arm, which replaces the whole
151                // property object rather than merging into it — erasing
152                // every owned key. Refused on every pack-owned kind, not only
153                // rows that currently carry an owned key, so an identical
154                // call cannot succeed or fail on state the caller cannot see.
155                if !props.is_object() {
156                    return Err(RuntimeError::InvalidInput(format!(
157                        "properties on a `{}` note must be patched with an object: a non-object \
158                         patch names no key, so it would replace the whole property object rather \
159                         than merging into it. Pass an object containing the keys you intend to \
160                         set.",
161                        note.kind
162                    )));
163                }
164                if let Some(named) = owner_established_property_named_in(&props) {
165                    return Err(RuntimeError::InvalidInput(format!(
166                        "`{named}` is not patchable on a `{}` note: the pack that owns this \
167                         kind establishes it and reads it back — to decide how the record is \
168                         attributed and grouped, or to reproduce it verbatim when the record \
169                         is re-emitted — so it is written by the owner and immutable to a \
170                         caller patch. Patch any other property key here, or omit \
171                         `{named}` from this patch.",
172                        note.kind
173                    )));
174                }
175            }
176            let incoming_properties = Some(props);
177            let (mut merged, _) = merge_properties(
178                &note.properties,
179                &incoming_properties,
180                EntityDedupMergePolicy::PreferFrom,
181            );
182            if let Some(properties) = merged.as_mut().and_then(Value::as_object_mut) {
183                for key in patch.update_policy.null_clearing_properties {
184                    if incoming_properties
185                        .as_ref()
186                        .and_then(|incoming| incoming.get(*key))
187                        .is_some_and(Value::is_null)
188                    {
189                        properties.remove(*key);
190                    }
191                }
192            }
193            note.properties = merged;
194        }
195        if let Some(status) = patch.kind_status {
196            note.status = status;
197        }
198
199        // The whole-note CAS persists the merged properties, including keys
200        // carried from the snapshot when the patch changes another field.
201        crate::secret_gate::reject_reserved_secret_gate_property(note.properties.as_ref())?;
202
203        // JSON object key order is not meaningful to callers. Tags are also
204        // set-like in every existing note reader, so their order is ignored
205        // for the no-op decision while duplicate entries remain meaningful.
206        // All other arrays retain ordinary JSON ordering semantics.
207        let changed = original_name != note.name
208            || original_content != note.content
209            || original_salience != note.salience
210            || original_decay_factor != note.decay_factor
211            || !note_update_values_equal(&original_properties, &note.properties)
212            || original_status != note.status;
213        if !changed {
214            return Ok((stored, text_changed, false));
215        }
216
217        // `updated_at` is also the optimistic-concurrency revision for
218        // full-note replacement. Make it strictly advance even when two
219        // operations land inside one clock microsecond. Saturation is not a
220        // valid fallback: reusing i64::MAX would make the CAS accept a write
221        // without advancing its revision.
222        let minimum_updated_at = note.updated_at.checked_add(1).ok_or_else(|| {
223            RuntimeError::Internal(format!(
224                "note {} updated_at is already at i64::MAX and cannot advance",
225                note.id
226            ))
227        })?;
228        note.updated_at = chrono::Utc::now()
229            .timestamp_micros()
230            .max(minimum_updated_at);
231        Ok((note, text_changed, true))
232    }
233
234    /// Patch-style note update.
235    #[cfg(test)]
236    pub(crate) async fn update_note(
237        &self,
238        token: &NamespaceToken,
239        id: Uuid,
240        patch: NotePatch,
241    ) -> RuntimeResult<khive_storage::note::Note> {
242        Ok(self
243            .update_note_with_embedding_report(token, id, patch)
244            .await?
245            .0)
246    }
247
248    pub async fn update_note_with_embedding_report(
249        &self,
250        token: &NamespaceToken,
251        id: Uuid,
252        patch: NotePatch,
253    ) -> RuntimeResult<(
254        khive_storage::note::Note,
255        crate::retrieval::EmbeddingTruncationReport,
256    )> {
257        let snapshot = self
258            .notes(token)?
259            .get_note(id)
260            .await?
261            .ok_or_else(|| RuntimeError::NotFound(format!("note {id}")))?;
262        self.update_note_from_snapshot_with_embedding_report(token, snapshot, patch)
263            .await
264    }
265
266    /// Patch and persist one note from a caller-owned read snapshot.
267    ///
268    /// This is the canonical seam for kind hooks that normalize coupled
269    /// fields from the current note. The same snapshot feeds normalization,
270    /// patch application, and the compare-and-swap write; a concurrent note
271    /// change therefore refuses the write instead of persisting derivations
272    /// computed from stale state.
273    pub async fn update_note_from_snapshot_with_embedding_report(
274        &self,
275        token: &NamespaceToken,
276        snapshot: khive_storage::note::Note,
277        patch: NotePatch,
278    ) -> RuntimeResult<(
279        khive_storage::note::Note,
280        crate::retrieval::EmbeddingTruncationReport,
281    )> {
282        let (note, plan) = self
283            .prepare_versioned_note_update(token, snapshot, patch)
284            .await?;
285        self.commit_prepared_note_update(token, note, crate::AtomicOpPlan::Update(Box::new(plan)))
286            .await
287    }
288
289    /// Commit a normalized and validated kind-owned update, including its typed
290    /// graph companions. Callers must first run `prepare_note_update_policy`
291    /// against this exact snapshot and pass the policy it returned; the shared
292    /// atomic prepare seam checks all patch fields before asking the kind hook
293    /// to derive any graph effects.
294    pub async fn update_note_from_snapshot_with_kind_effects(
295        &self,
296        token: &NamespaceToken,
297        snapshot: khive_storage::Note,
298        args: &Value,
299        policy: crate::NoteUpdatePolicy,
300        registry: &crate::VerbRegistry,
301    ) -> RuntimeResult<(
302        khive_storage::Note,
303        crate::retrieval::EmbeddingTruncationReport,
304    )> {
305        let (note, plan) = crate::atomic_prepare::prepare_update_from_note_snapshot(
306            self, token, args, None, snapshot, policy, registry,
307        )
308        .await?;
309        self.commit_prepared_note_update(token, note, plan).await
310    }
311
312    async fn commit_prepared_note_update(
313        &self,
314        token: &NamespaceToken,
315        note: khive_storage::Note,
316        plan: crate::AtomicOpPlan,
317    ) -> RuntimeResult<(
318        khive_storage::Note,
319        crate::retrieval::EmbeddingTruncationReport,
320    )> {
321        let id = note.id;
322        use crate::atomic_runner::{run_atomic_unit, AtomicOpFailure, AtomicRunOutcome};
323        match run_atomic_unit(self.sql().as_ref(), vec![plan]).await {
324            Ok(AtomicRunOutcome::Committed { post_commit }) => {
325                let outcomes = crate::atomic_prepare::apply_post_commit_effects_with_report(
326                    self,
327                    token,
328                    post_commit,
329                )
330                .await?;
331                let report = outcomes
332                    .into_iter()
333                    .next()
334                    .map(|outcome| outcome.truncation)
335                    .unwrap_or_default();
336                Ok((note, report))
337            }
338            Ok(AtomicRunOutcome::RolledBack {
339                failure: AtomicOpFailure::NoteConflict(conflict),
340                ..
341            }) => Err(conflict.into_error().into()),
342            Ok(AtomicRunOutcome::RolledBack {
343                failure: AtomicOpFailure::GuardFailed { .. },
344                ..
345            }) => Err(stale_note_snapshot_error(id)),
346            Ok(AtomicRunOutcome::RolledBack { failure, .. }) => Err(RuntimeError::Internal(
347                format!("note update rolled back: {failure:?}"),
348            )),
349            Err(error) => Err(RuntimeError::Storage(error.0)),
350        }
351    }
352}