Skip to main content

khive_runtime/curation/
note_merge.rs

1//! Asynchronous note merge, with or without a caller-supplied guard.
2
3use super::*;
4
5impl KhiveRuntime {
6    /// Merge `from_id` note into `into_id` note.
7    ///
8    /// Both notes must exist in the namespace and have the same `kind`. Content is merged
9    /// per `content_strategy`. Properties are merged per `strategy`. `from_id` is
10    /// tombstoned (status='deleted', deleted_at set). Returns a summary.
11    ///
12    /// If `dry_run` is true, computes and returns the planned summary without mutating
13    /// any rows, edges, or indexes.
14    /// The NoteMerged event, including destructive edge preimages, commits in
15    /// the same SQL transaction as the note and edge changes.
16    pub async fn merge_note(
17        &self,
18        token: &NamespaceToken,
19        into_id: Uuid,
20        from_id: Uuid,
21        strategy: EntityDedupMergePolicy,
22        content_strategy: ContentMergeStrategy,
23        dry_run: bool,
24    ) -> RuntimeResult<MergeSummary> {
25        self.merge_note_with_reason(
26            token,
27            into_id,
28            from_id,
29            strategy,
30            content_strategy,
31            dry_run,
32            None,
33        )
34        .await
35    }
36
37    /// Merge `from_id` note into `into_id` note and include an optional audit reason.
38    // REASON: these arguments mirror the merge verb's policy, content strategy,
39    // dry-run, and audit-reason fields; a builder would only move that surface.
40    #[allow(clippy::too_many_arguments)]
41    pub async fn merge_note_with_reason(
42        &self,
43        token: &NamespaceToken,
44        into_id: Uuid,
45        from_id: Uuid,
46        strategy: EntityDedupMergePolicy,
47        content_strategy: ContentMergeStrategy,
48        dry_run: bool,
49        reason: Option<String>,
50    ) -> RuntimeResult<MergeSummary> {
51        self.merge_note_with_guard(
52            token,
53            into_id,
54            from_id,
55            strategy,
56            content_strategy,
57            dry_run,
58            reason,
59            None,
60        )
61        .await
62        .map(|(summary, _)| summary)
63    }
64
65    /// Merge `from_id` note into `into_id` note, refusing unless the caller's
66    /// guard still holds.
67    ///
68    /// The guard names the version of each note the caller read, facts it relied
69    /// on (see [`MergeAssertion`]), property values that replace the survivor's,
70    /// and one annotation for the provenance entry. All of it is checked inside
71    /// the merge transaction, on the writer connection, after both notes are read
72    /// there and before the first write; a refusal writes nothing. The returned
73    /// `kept_version` is the survivor's committed version, so the caller can
74    /// merge the next duplicate into the same survivor without reading it again.
75    ///
76    /// With no guard this is [`Self::merge_note_with_reason`]: the same
77    /// statements in the same order, with the same errors and effects.
78    // REASON: these arguments mirror the merge verb's policy, content strategy,
79    // dry-run, and audit-reason fields plus the guard; a builder would only move
80    // that surface.
81    #[allow(clippy::too_many_arguments)]
82    pub async fn merge_note_guarded(
83        &self,
84        token: &NamespaceToken,
85        into_id: Uuid,
86        from_id: Uuid,
87        strategy: EntityDedupMergePolicy,
88        content_strategy: ContentMergeStrategy,
89        dry_run: bool,
90        reason: Option<String>,
91        guard: NoteMergeGuard,
92    ) -> RuntimeResult<GuardedNoteMerge> {
93        let (summary, kept_version) = self
94            .merge_note_with_guard(
95                token,
96                into_id,
97                from_id,
98                strategy,
99                content_strategy,
100                dry_run,
101                reason,
102                Some(guard),
103            )
104            .await?;
105        Ok(GuardedNoteMerge {
106            summary,
107            kept_version,
108        })
109    }
110
111    /// Merge `from_id` note into `into_id` note and include an optional audit reason.
112    // REASON: these arguments mirror the merge verb's policy, content strategy,
113    // dry-run, and audit-reason fields; a builder would only move that surface.
114    #[allow(clippy::too_many_arguments)]
115    pub(super) async fn merge_note_with_guard(
116        &self,
117        token: &NamespaceToken,
118        into_id: Uuid,
119        from_id: Uuid,
120        strategy: EntityDedupMergePolicy,
121        content_strategy: ContentMergeStrategy,
122        dry_run: bool,
123        reason: Option<String>,
124        merge_guard: Option<NoteMergeGuard>,
125    ) -> RuntimeResult<(MergeSummary, i64)> {
126        if let Some(reason) = reason.as_deref() {
127            crate::secret_gate::check_at(reason, "merge", "reason")?;
128        }
129        if let Some(merge_guard) = merge_guard.as_ref() {
130            merge_guard.check_values()?;
131        }
132        if into_id == from_id {
133            return Err(RuntimeError::InvalidInput(
134                "cannot merge a note into itself".into(),
135            ));
136        }
137        let ns = token.namespace().as_str().to_string();
138        let fts_table = "fts_notes".to_string();
139        // Keep deletion, table preparation, and survivor reindex on the same
140        // immutable registry view; see the entity merge path above.
141        let embedding_plan = EmbeddingModelPlan::capture(self);
142        let vec_tables = embedding_plan.vector_tables();
143        let pack_rules = self.pack_edge_rules();
144
145        let note_store = self.notes(token)?;
146        let into_note = note_store
147            .get_note(into_id)
148            .await?
149            .ok_or_else(|| RuntimeError::NotFound("not found in this namespace".into()))?;
150        Self::ensure_namespace(&into_note.namespace, &ns)?;
151
152        let from_note = note_store
153            .get_note(from_id)
154            .await?
155            .ok_or_else(|| RuntimeError::NotFound("not found in this namespace".into()))?;
156        Self::ensure_namespace(&from_note.namespace, &ns)?;
157
158        if !dry_run {
159            for note in [&into_note, &from_note] {
160                if let Some(error) = self.stream_member_error(note).await? {
161                    return Err(error);
162                }
163            }
164        }
165        reject_pack_managed_schedule_mutation(&into_note, "merge")?;
166        reject_pack_managed_schedule_mutation(&from_note, "merge")?;
167
168        let _ = self.graph(token)?;
169        let _ = self.text_for_notes(token)?;
170        let _ = self.events(token)?;
171        for model_name in embedding_plan.model_names() {
172            let _ = self.vectors_for_model(token, model_name)?;
173        }
174
175        // Resolved here, where the runtime's installed pack-kind list is in
176        // reach; `merge_note_sql` runs on the writer connection with no runtime
177        // handle. Both notes share a kind (checked inside), so the into-note's
178        // kind decides for the merge.
179        let preserve_owner_established = self.is_pack_owned_note_kind(&into_note.kind);
180
181        let pool = self.backend().pool_arc();
182        let writer_task = pool
183            .writer_task_for_runtime_write(RuntimeWriteOperation::MergeNote)
184            .map_err(RuntimeError::Storage)?;
185        let event_context = MergeEventContext {
186            attribution: EventAttribution::from_token(token),
187            reason,
188            force: false,
189            strategy,
190            content_strategy,
191            kind: EventKind::NoteMerged,
192            substrate: SubstrateKind::Note,
193            event_id: None,
194        };
195
196        let (mut summary, updated_note) = if let Some(writer_task) = writer_task {
197            writer_task
198                .send(move |conn| {
199                    merge_note_sql(
200                        conn,
201                        ns,
202                        fts_table,
203                        vec_tables,
204                        into_id,
205                        from_id,
206                        strategy,
207                        content_strategy,
208                        dry_run,
209                        pack_rules,
210                        preserve_owner_established,
211                        MergeTxLimits::default(),
212                        Some(event_context),
213                        merge_guard,
214                    )
215                    .map_err(|e| {
216                        khive_storage::StorageError::driver(
217                            khive_storage::StorageCapability::Notes,
218                            "merge_note",
219                            e,
220                        )
221                    })
222                })
223                .await
224                .inspect_err(|error| khive_storage::usage::account_event_write(Err(error)))
225                .map_err(map_merge_note_storage_error)?
226        } else {
227            tokio::task::spawn_blocking(move || {
228                let guard = pool.writer()?;
229                let mut refusal = None;
230                let result = guard.transaction(|conn| {
231                    merge_note_sql(
232                        conn,
233                        ns,
234                        fts_table,
235                        vec_tables,
236                        into_id,
237                        from_id,
238                        strategy,
239                        content_strategy,
240                        dry_run,
241                        pack_rules,
242                        preserve_owner_established,
243                        MergeTxLimits::default(),
244                        Some(event_context),
245                        merge_guard,
246                    )
247                    .map_err(|error| match error {
248                        MergeSqlError::Sqlite(error) => error,
249                        MergeSqlError::Refusal(error) => {
250                            refusal = Some(error);
251                            SqliteError::InvalidData(
252                                "note merge refused by transactional policy".to_string(),
253                            )
254                        }
255                    })
256                });
257                match refusal {
258                    Some(error) => Err(error),
259                    None => result.map_err(RuntimeError::from),
260                }
261            })
262            .await
263            .map_err(|e| RuntimeError::Internal(e.to_string()))??
264        };
265
266        // Count only committed event rows; dry-run never inserts an event.
267        if !dry_run {
268            khive_storage::usage::account_event_write(Ok(1));
269            tracing::info!(
270                into_id = %summary.kept_id,
271                from_id = %summary.removed_id,
272                budget_rows = summary.tx_budget.rows_charged,
273                budget_bytes = summary.tx_budget.bytes_charged,
274                budget_max_rows = summary.tx_budget.max_rows,
275                budget_max_bytes = summary.tx_budget.max_bytes,
276                "merge_note: transaction materialization budget"
277            );
278        }
279
280        if !dry_run {
281            if !embedding_plan.is_empty() {
282                #[cfg(any(test, feature = "fault-injection"))]
283                let reindex_result =
284                    if crate::operations::consume_fts_fail_fault(&updated_note.namespace) {
285                        Err(RuntimeError::Internal("injected FTS failure".to_string()))
286                    } else {
287                        self.reindex_note_report_with_plan(token, &updated_note, &embedding_plan)
288                            .await
289                    };
290                #[cfg(not(any(test, feature = "fault-injection")))]
291                let reindex_result = self
292                    .reindex_note_report_with_plan(token, &updated_note, &embedding_plan)
293                    .await;
294
295                match reindex_result {
296                    Ok(report) => {
297                        summary.embedding_truncation = report.truncation;
298                        if !report.failures.is_empty() {
299                            summary.post_commit_reindex_error = Some(
300                                report
301                                    .failures
302                                    .iter()
303                                    .map(|failure| {
304                                        format!(
305                                            "model {} {}: {}",
306                                            failure.model,
307                                            failure.stage.as_str(),
308                                            failure.error
309                                        )
310                                    })
311                                    .collect::<Vec<_>>()
312                                    .join("; "),
313                            );
314                        }
315                    }
316                    Err(error) => {
317                        tracing::warn!(
318                            into_id = %summary.kept_id,
319                            from_id = %summary.removed_id,
320                            error = %error,
321                            "merge_note: committed merge but survivor reindex failed"
322                        );
323                        summary.post_commit_reindex_error = Some(error.to_string());
324                    }
325                }
326            }
327            // The note row is committed even when no embedding model is
328            // registered or the post-commit reindex reports an error.
329            self.fire_note_mutation_hook(&updated_note.kind, updated_note.id)
330                .await;
331        }
332
333        Ok((summary, updated_note.version))
334    }
335}