1use 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
71async 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
177pub 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 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
233pub 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
251pub 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
263pub 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 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 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 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 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}