Skip to main content

khive_runtime/
keyed_memory.rs

1//! Memory identity is published only by the final DML of its atomic create.
2
3use khive_storage::note::Note;
4use khive_storage::{SqlStatement, SqlValue};
5use khive_types::{Details, KhiveError};
6use serde_json::Value;
7use uuid::Uuid;
8
9use crate::atomic_message::{AtomicNoteOptions, AtomicNoteSpec};
10use crate::atomic_runner::{run_atomic_unit, AtomicOpFailure, AtomicRunOutcome};
11use crate::note_create::{prepare_note_create, KeyPublication, KEY_CLAIM};
12use crate::{KhiveRuntime, NamespaceToken, RuntimeError, RuntimeResult};
13
14pub struct KeyedMemorySpec<'a> {
15    pub content: &'a str,
16    pub key: &'a str,
17    pub salience: f64,
18    pub decay_factor: f64,
19    pub properties: Value,
20    pub source_id: Option<Uuid>,
21    pub embedding_model: Option<&'a str>,
22}
23
24pub fn validate_memory_key(key: &str) -> RuntimeResult<()> {
25    if key.len() > 512 || key.contains('\0') {
26        return Err(RuntimeError::InvalidInput(
27            "key must be at most 512 UTF-8 bytes and must not contain U+0000".into(),
28        ));
29    }
30    Ok(())
31}
32
33fn idempotency_conflict(key: &str, existing: &Note) -> RuntimeError {
34    KhiveError::conflict(format!(
35        "idempotency_key_conflict: key {key:?} already exists; stored content differs (existing memory {})",
36        existing.id
37    ))
38        .with_details(Details::new_owned([
39            ("reason", "idempotency_key_conflict".into()),
40            ("key", key.to_owned()),
41            ("existing_id", existing.id.to_string()),
42        ]))
43        .into()
44}
45
46async fn resolve_holder(
47    runtime: &KhiveRuntime,
48    token: &NamespaceToken,
49    key: &str,
50) -> RuntimeResult<Option<Note>> {
51    let mut matches = runtime
52        .notes(token)?
53        .get_live_notes_by_key(token.namespace().as_str(), key, Some("memory"))
54        .await?;
55    match matches.len() {
56        0 => Ok(None),
57        1 => Ok(matches.pop()),
58        _ => Err(RuntimeError::Internal(
59            "memory key lookup returned multiple live holders".into(),
60        )),
61    }
62}
63
64struct VisibilityReceiptState {
65    note_present: bool,
66    epoch: Option<String>,
67    receipt_present: bool,
68    fences: Option<Vec<(String, u64)>>,
69}
70
71/// Join identity, independent provenance and original receipt in one snapshot.
72/// Joins use note identity first so a conflicting namespace cannot disappear
73/// behind a filter and masquerade as a genuinely absent receipt or marker.
74async fn visibility_receipt_state(
75    runtime: &KhiveRuntime,
76    token: &NamespaceToken,
77    note_id: Uuid,
78) -> RuntimeResult<VisibilityReceiptState> {
79    let unavailable = || {
80        crate::visibility_receipts::receipt_failure(
81            "receipt_store_unavailable",
82            Some(note_id),
83            false,
84        )
85    };
86    runtime
87        .require_visibility_cutover()
88        .map_err(|_| unavailable())?;
89    let mut reader = runtime.sql().reader().await.map_err(|_| unavailable())?;
90    let rows = reader
91        .query_all(SqlStatement {
92            sql: "SELECT n.kind AS note_kind, n.namespace AS note_namespace, \
93              e.namespace AS epoch_namespace, e.epoch, \
94              r.note_id AS receipt_id, r.namespace AS receipt_namespace, r.model_count, \
95              f.namespace AS fence_namespace, f.model, f.ann_write_log_seq \
96              FROM notes n LEFT JOIN memory_visibility_epochs e ON e.note_id = n.id \
97              LEFT JOIN memory_visibility_receipts r ON r.note_id = n.id \
98              LEFT JOIN memory_visibility_fences f ON f.note_id = n.id \
99              WHERE n.id = ?1 ORDER BY r.namespace, f.namespace, f.model"
100                .into(),
101            params: vec![SqlValue::Text(note_id.to_string())],
102            label: Some("memory-visibility-receipt-read".into()),
103        })
104        .await
105        .map_err(|_| unavailable())?;
106    let mut state = VisibilityReceiptState {
107        note_present: !rows.is_empty(),
108        epoch: None,
109        receipt_present: false,
110        fences: None,
111    };
112    let Some(first) = rows.first() else {
113        return Ok(state);
114    };
115    let namespace = token.namespace().as_str();
116    let text = |value: Option<&SqlValue>, expected: &str| match value {
117        Some(SqlValue::Text(actual)) => actual == expected,
118        _ => false,
119    };
120    let identity_matches = text(first.get("note_kind"), "memory")
121        && text(first.get("note_namespace"), namespace)
122        && text(first.get("epoch_namespace"), namespace);
123    if identity_matches {
124        if let Some(SqlValue::Text(epoch)) = first.get("epoch") {
125            state.epoch = Some(epoch.clone());
126        }
127    }
128    state.receipt_present = rows
129        .iter()
130        .any(|row| matches!(row.get("receipt_id"), Some(SqlValue::Text(_))));
131    let expected_count = match first.get("model_count") {
132        Some(SqlValue::Integer(count)) if *count >= 0 => usize::try_from(*count).ok(),
133        _ => None,
134    };
135    let mut valid = text(first.get("note_kind"), "memory")
136        && text(first.get("note_namespace"), namespace)
137        && state.receipt_present
138        && expected_count.is_some();
139    let mut fences = std::collections::BTreeMap::new();
140    for row in &rows {
141        if (matches!(row.get("receipt_id"), Some(SqlValue::Text(_)))
142            && !text(row.get("receipt_namespace"), namespace))
143            || matches!(row.get("fence_namespace"), Some(SqlValue::Text(ns))
144                if ns != namespace || !text(row.get("receipt_namespace"), ns))
145        {
146            state.epoch = None;
147        }
148        valid &= text(row.get("receipt_namespace"), namespace)
149            && matches!(row.get("model_count"), Some(SqlValue::Integer(count))
150                if usize::try_from(*count).ok() == expected_count);
151        match (
152            row.get("fence_namespace"),
153            row.get("model"),
154            row.get("ann_write_log_seq"),
155        ) {
156            (
157                Some(SqlValue::Text(ns)),
158                Some(SqlValue::Text(model)),
159                Some(SqlValue::Integer(seq)),
160            ) if ns == namespace && !model.is_empty() && *seq > 0 => {
161                valid &= fences.insert(model.clone(), *seq as u64).is_none();
162            }
163            (
164                Some(SqlValue::Null) | None,
165                Some(SqlValue::Null) | None,
166                Some(SqlValue::Null) | None,
167            ) => {}
168            _ => valid = false,
169        }
170    }
171    if valid && expected_count == Some(fences.len()) {
172        state.fences = Some(fences.into_iter().collect());
173    }
174    Ok(state)
175}
176
177/// Read original fences without issuing a token. Callers needing a replay token
178/// must additionally apply the independent provenance classification below.
179pub async fn memory_visibility_receipt(
180    runtime: &KhiveRuntime,
181    token: &NamespaceToken,
182    note_id: Uuid,
183) -> RuntimeResult<Option<Vec<(String, u64)>>> {
184    Ok(visibility_receipt_state(runtime, token, note_id)
185        .await?
186        .fences)
187}
188
189async fn classified_visibility_receipt(
190    runtime: &KhiveRuntime,
191    token: &NamespaceToken,
192    note_id: Uuid,
193    replay: bool,
194) -> RuntimeResult<Vec<(String, u64)>> {
195    let state = visibility_receipt_state(runtime, token, note_id).await?;
196    if replay && !state.note_present {
197        // The attempt's own unit rolled back on the key claim, so this request
198        // wrote nothing; the reason carries that to the error projection.
199        return Err(KhiveError::not_found("memory", "keyed replay holder")
200            .with_details(Details::new_owned([(
201                "reason",
202                "keyed_replay_holder_missing".into(),
203            )]))
204            .into());
205    }
206    let reason = match state.epoch.as_deref() {
207        Some("modern") => {
208            if let Some(fences) = state.fences {
209                return Ok(fences);
210            }
211            "receipt_temporarily_unavailable"
212        }
213        Some("legacy") if !state.receipt_present => "legacy_receipt_absent",
214        _ => "receipt_epoch_unknown",
215    };
216    Err(crate::visibility_receipts::receipt_failure(
217        reason,
218        Some(note_id),
219        replay,
220    ))
221}
222
223pub async fn create_keyed_memory(
224    runtime: &KhiveRuntime,
225    token: &NamespaceToken,
226    spec: KeyedMemorySpec<'_>,
227) -> RuntimeResult<(Note, Option<Uuid>, bool)> {
228    let (note, edge_id, replayed, _) =
229        create_keyed_memory_with_report(runtime, token, spec).await?;
230    Ok((note, edge_id, replayed))
231}
232
233/// Same as [`create_keyed_memory`], also returning the embedding-input
234/// truncation report computed for this call. A replay stores nothing new but
235/// retains the report from preparing this call's identical content.
236pub async fn create_keyed_memory_with_report(
237    runtime: &KhiveRuntime,
238    token: &NamespaceToken,
239    spec: KeyedMemorySpec<'_>,
240) -> RuntimeResult<(
241    Note,
242    Option<Uuid>,
243    bool,
244    crate::retrieval::EmbeddingTruncationReport,
245)> {
246    let (note, edge_id, replayed, _, report) =
247        create_keyed_memory_with_receipt_and_report(runtime, token, spec).await?;
248    Ok((note, edge_id, replayed, report))
249}
250
251/// Receipt-bearing keyed create. The original per-model receipt is retained
252/// on replay, including an explicit receipt for a write with zero models.
253pub async fn create_keyed_memory_with_receipt(
254    runtime: &KhiveRuntime,
255    token: &NamespaceToken,
256    spec: KeyedMemorySpec<'_>,
257) -> RuntimeResult<(Note, Option<Uuid>, bool, Vec<(String, u64)>)> {
258    let (note, edge_id, replayed, fences, _) =
259        create_keyed_memory_with_receipt_and_report(runtime, token, spec).await?;
260    Ok((note, edge_id, replayed, fences))
261}
262
263/// Return both the original visibility receipt and the embedding-input report.
264/// Replays retain the stored fences and the report computed while preparing
265/// this call's embedding input; the report does not describe the original write.
266pub async fn create_keyed_memory_with_receipt_and_report(
267    runtime: &KhiveRuntime,
268    token: &NamespaceToken,
269    spec: KeyedMemorySpec<'_>,
270) -> RuntimeResult<(
271    Note,
272    Option<Uuid>,
273    bool,
274    Vec<(String, u64)>,
275    crate::retrieval::EmbeddingTruncationReport,
276)> {
277    validate_memory_key(spec.key)?;
278    if spec.content.trim().is_empty() {
279        return Err(RuntimeError::InvalidInput(
280            "content must not be empty".into(),
281        ));
282    }
283    let (mut prepared, annotation_ids) = prepare_note_create(
284        runtime,
285        AtomicNoteSpec {
286            token,
287            id: None,
288            kind: "memory",
289            name: None,
290            content: spec.content,
291            properties: Some(spec.properties),
292        },
293        AtomicNoteOptions {
294            salience: Some(spec.salience),
295            decay_factor: Some(spec.decay_factor),
296            embedding_model: spec.embedding_model,
297            key: Some(spec.key),
298            memory_visibility_receipt: true,
299            ..Default::default()
300        },
301        &spec.source_id.into_iter().collect::<Vec<_>>(),
302        KeyPublication::AfterDependents,
303    )
304    .await?;
305    let mut note = prepared.notes.remove(0);
306    let edge_id = annotation_ids.first().copied();
307
308    for _attempt in 0..2 {
309        #[cfg(test)]
310        crate::keyed_memory_tests::checkpoint(token.namespace().as_str(), _attempt, false).await;
311        match run_atomic_unit(runtime.sql().as_ref(), prepared.plans.clone()).await {
312            Ok(AtomicRunOutcome::Committed { .. }) => {
313                note.key = Some(spec.key.to_owned());
314                note.version = 2;
315                let fences = classified_visibility_receipt(runtime, token, note.id, false).await?;
316                return Ok((note, edge_id, false, fences, prepared.embedding_truncation));
317            }
318            Ok(AtomicRunOutcome::RolledBack {
319                failure:
320                    AtomicOpFailure::GuardFailed {
321                        statement_label,
322                        observed: 0,
323                        ..
324                    },
325                ..
326            }) if statement_label.as_deref() == Some(KEY_CLAIM) => {
327                #[cfg(test)]
328                crate::keyed_memory_tests::checkpoint(token.namespace().as_str(), _attempt, true)
329                    .await;
330                if let Some(holder) = resolve_holder(runtime, token, spec.key).await? {
331                    if holder.content == spec.content {
332                        let fences =
333                            classified_visibility_receipt(runtime, token, holder.id, true).await?;
334                        return Ok((holder, None, true, fences, prepared.embedding_truncation));
335                    }
336                    return Err(idempotency_conflict(spec.key, &holder));
337                }
338            }
339            Ok(AtomicRunOutcome::RolledBack {
340                failed_op_index,
341                failure,
342            }) => {
343                // Admission is cached after the first successful cutover check,
344                // so a receipt store that broke since then surfaces here.
345                runtime.recheck_visibility_cutover()?;
346                return Err(RuntimeError::Internal(format!(
347                    "atomic memory write rolled back at op {failed_op_index}: {failure:?}"
348                )));
349            }
350            Err(error) => return Err(RuntimeError::Storage(error.0)),
351        }
352    }
353    Err(
354        KhiveError::unavailable("memory key holder disappeared during reconciliation")
355            .with_details(Details::new_owned([
356                ("reason", "key_holder_unresolved".into()),
357                ("key", spec.key.to_owned()),
358            ]))
359            .into(),
360    )
361}
362
363#[cfg(test)]
364mod receipt_read_tests {
365    use super::*;
366    use crate::DomainDisposition;
367
368    const KEY: &str = "receipt-race-key";
369    const CONTENT: &str = "private receipt race content";
370
371    async fn captured_holder() -> (KhiveRuntime, NamespaceToken, Note) {
372        let runtime = KhiveRuntime::memory().unwrap();
373        runtime.install_kind_registry(vec![], vec!["memory".into()]);
374        let token = runtime
375            .authorize(khive_types::Namespace::parse("receipt-read-race").unwrap())
376            .unwrap();
377        let (note, _, replayed, fences) = create_keyed_memory_with_receipt(
378            &runtime,
379            &token,
380            KeyedMemorySpec {
381                content: CONTENT,
382                key: KEY,
383                salience: 0.7,
384                decay_factor: 0.95,
385                properties: serde_json::json!({}),
386                source_id: None,
387                embedding_model: None,
388            },
389        )
390        .await
391        .unwrap();
392        assert!(!replayed);
393        assert!(fences.is_empty());
394        let holder = resolve_holder(&runtime, &token, KEY)
395            .await
396            .unwrap()
397            .unwrap();
398        assert_eq!(holder, note);
399        assert!(
400            classified_visibility_receipt(&runtime, &token, holder.id, true)
401                .await
402                .unwrap()
403                .is_empty()
404        );
405        (runtime, token, holder)
406    }
407
408    #[tokio::test]
409    async fn replay_receipt_reports_missing_after_captured_holder_is_hard_deleted() {
410        let (runtime, token, holder) = captured_holder().await;
411        // This is the exact boundary in the replay path: holder lookup has
412        // completed, but its joined receipt/epoch read has not started.
413        assert!(runtime.delete_note(&token, holder.id, true).await.unwrap());
414        assert!(runtime
415            .notes(&token)
416            .unwrap()
417            .get_note_including_deleted(holder.id)
418            .await
419            .unwrap()
420            .is_none());
421        let error = classified_visibility_receipt(&runtime, &token, holder.id, true)
422            .await
423            .unwrap_err();
424        assert!(
425            matches!(&error, RuntimeError::Khive(e) if e.kind() == khive_types::ErrorKind::NotFound)
426        );
427        let value = crate::error_projection::runtime_error_value(error, DomainDisposition::Unknown);
428        assert_eq!(value["kind"], "not_found");
429        assert_eq!(value["details"]["reason"], "keyed_replay_holder_missing");
430        assert_eq!(value["domain_disposition"], "not_committed");
431        assert!(value.get("retryable").is_none());
432        let encoded = value.to_string();
433        let id = holder.id.to_string();
434        for private in [KEY, CONTENT, id.as_str()] {
435            assert!(!encoded.contains(private));
436        }
437        assert!(memory_visibility_receipt(&runtime, &token, holder.id)
438            .await
439            .unwrap()
440            .is_none());
441        assert!(resolve_holder(&runtime, &token, KEY)
442            .await
443            .unwrap()
444            .is_none());
445
446        // The post-commit caller retains its existing uncertainty semantics.
447        let error = classified_visibility_receipt(&runtime, &token, holder.id, false)
448            .await
449            .unwrap_err();
450        let value = crate::error_projection::runtime_error_value(error, DomainDisposition::Unknown);
451        assert_eq!(value["details"]["reason"], "receipt_epoch_unknown");
452        assert_eq!(value["domain_disposition"], "unknown");
453    }
454
455    #[tokio::test]
456    async fn present_holder_with_missing_or_invalid_epoch_is_not_missing() {
457        for mutation in [
458            "DELETE FROM memory_visibility_epochs WHERE note_id = ?1",
459            "UPDATE memory_visibility_epochs SET epoch = 'unknown' WHERE note_id = ?1",
460            "UPDATE memory_visibility_epochs SET epoch = 'malformed' WHERE note_id = ?1",
461            "UPDATE memory_visibility_epochs SET epoch = 'legacy' WHERE note_id = ?1",
462            "UPDATE memory_visibility_epochs SET namespace = 'foreign-private-namespace' WHERE note_id = ?1",
463            "UPDATE notes SET namespace = 'foreign-private-namespace' WHERE id = ?1",
464        ] {
465            let (runtime, token, holder) = captured_holder().await;
466            {
467                // Model a damaged epoch value as well as valid but incomplete
468                // provenance. Restore constraint checks before reading it.
469                let writer = runtime.backend().pool().try_writer().unwrap();
470                writer
471                    .conn()
472                    .pragma_update(None, "ignore_check_constraints", true)
473                    .unwrap();
474                writer.conn().execute(mutation, [holder.id.to_string()]).unwrap();
475                writer
476                    .conn()
477                    .pragma_update(None, "ignore_check_constraints", false)
478                    .unwrap();
479            }
480            assert!(runtime
481                .notes(&token)
482                .unwrap()
483                .get_note_including_deleted(holder.id)
484                .await
485                .unwrap()
486                .is_some());
487            let error = classified_visibility_receipt(&runtime, &token, holder.id, true)
488                .await
489                .unwrap_err();
490            let value = crate::error_projection::runtime_error_value(error, DomainDisposition::Unknown);
491            assert_eq!(value["details"]["reason"], "receipt_epoch_unknown");
492            assert_eq!(value["details"]["memory_id"], holder.id.to_string());
493            assert_eq!(value["domain_disposition"], "not_committed");
494            assert_eq!(value["retryable"], false);
495            for private in [KEY, CONTENT, "foreign-private-namespace"] {
496                assert!(!value.to_string().contains(private));
497            }
498        }
499    }
500
501    #[tokio::test]
502    async fn unreadable_note_store_is_unavailable_not_missing() {
503        let (runtime, token, holder) = captured_holder().await;
504        runtime
505            .sql()
506            .writer()
507            .await
508            .unwrap()
509            .execute(SqlStatement {
510                sql: "ALTER TABLE notes RENAME TO private_unreadable_notes".into(),
511                params: vec![],
512                label: Some("receipt-unreadable-note-control".into()),
513            })
514            .await
515            .unwrap();
516        let error = classified_visibility_receipt(&runtime, &token, holder.id, true)
517            .await
518            .unwrap_err();
519        let value = crate::error_projection::runtime_error_value(error, DomainDisposition::Unknown);
520        assert_eq!(value["details"]["reason"], "receipt_store_unavailable");
521        assert_eq!(value["details"]["memory_id"], holder.id.to_string());
522        assert_eq!(value["domain_disposition"], "unknown");
523        assert_eq!(value["retryable"], true);
524        for private in [KEY, CONTENT, "private_unreadable_notes"] {
525            assert!(!value.to_string().contains(private));
526        }
527    }
528}