khive_runtime/
keyed_message.rs1use khive_storage::{Note, SqlStatement, SqlValue};
4use khive_types::{Details, KhiveError};
5use uuid::Uuid;
6
7use crate::atomic_message::{prepare_atomic_notes, AtomicNoteOptions, AtomicNoteSpec};
8use crate::atomic_plan::{AffectedRowGuard, PlanStatement};
9use crate::atomic_runner::{run_atomic_unit, AtomicOpFailure, AtomicOpPlan, AtomicRunOutcome};
10use crate::{KhiveRuntime, RuntimeError, RuntimeResult};
11
12const KEY_CLAIM: &str = "message-key-claim";
13
14pub enum KeyedMessageWrite {
15 Created {
16 notes: Vec<Note>,
17 embedding_truncation: crate::retrieval::EmbeddingTruncationReport,
18 },
19 Existing(Uuid),
20}
21
22pub async fn create_keyed_message_pair(
26 runtime: &KhiveRuntime,
27 specs: [AtomicNoteSpec<'_>; 2],
28 physical_key: &str,
29) -> RuntimeResult<KeyedMessageWrite> {
30 create_keyed_message_pair_with_attachments(runtime, specs, physical_key, &[]).await
31}
32
33pub async fn create_keyed_message_pair_with_attachments(
35 runtime: &KhiveRuntime,
36 specs: [AtomicNoteSpec<'_>; 2],
37 physical_key: &str,
38 attachments: &[khive_storage::NewAttachment],
39) -> RuntimeResult<KeyedMessageWrite> {
40 crate::atomic_message::validate_note_attachments(runtime, attachments)?;
41 let namespace = specs[0].token.namespace().as_str().to_owned();
42 if specs.iter().any(|spec| spec.kind != "message")
43 || specs[1].token.namespace().as_str() != namespace
44 || specs[0].token.actor().id != specs[1].token.actor().id
45 || specs[0].id.is_none()
46 || specs[1].id.is_none()
47 || specs[0].id == specs[1].id
48 {
49 return Err(RuntimeError::InvalidInput(
50 "a keyed message pair requires distinct IDs in one actor namespace".into(),
51 ));
52 }
53 let mut prepared =
54 prepare_atomic_notes(runtime, specs.into(), AtomicNoteOptions::default()).await?;
55 crate::atomic_message::append_note_attachments(&mut prepared, attachments)?;
56 for note in &prepared.notes {
57 crate::secret_gate::reject_reserved_secret_gate_property(note.properties.as_ref())?;
58 }
59 let outbound_id = prepared.notes[0].id;
60 for (plan, note) in prepared.plans.iter_mut().zip(&prepared.notes) {
61 let AtomicOpPlan::AddNote(plan) = plan else {
62 return Err(RuntimeError::Internal(
63 "expected prepared message note".into(),
64 ));
65 };
66 plan.statements[0].statement =
68 khive_db::stores::note::note_insert_if_absent_statement(note);
69 }
70 let Some(AtomicOpPlan::AddNote(last)) = prepared.plans.last_mut() else {
71 return Err(RuntimeError::Internal(
72 "expected recipient message plan".into(),
73 ));
74 };
75 last.statements.push(PlanStatement {
76 statement: SqlStatement {
77 sql: "UPDATE OR IGNORE notes SET key = ?1 WHERE id = ?2 AND namespace = ?3 \
78 AND kind = 'message' AND key IS NULL AND deleted_at IS NULL"
79 .into(),
80 params: vec![
81 SqlValue::Text(physical_key.to_owned()),
82 SqlValue::Text(outbound_id.to_string()),
83 SqlValue::Text(namespace.clone()),
84 ],
85 label: Some(KEY_CLAIM.into()),
86 },
87 guard: Some(AffectedRowGuard::exactly(1)),
88 });
89
90 match run_atomic_unit(runtime.sql().as_ref(), prepared.plans).await {
91 Ok(AtomicRunOutcome::Committed { .. }) => {
92 prepared.notes[0].key = Some(physical_key.to_owned());
93 Ok(KeyedMessageWrite::Created {
94 notes: prepared.notes,
95 embedding_truncation: prepared.embedding_truncation,
96 })
97 }
98 Ok(AtomicRunOutcome::RolledBack {
99 failure:
100 AtomicOpFailure::GuardFailed {
101 statement_label,
102 observed: 0,
103 ..
104 },
105 ..
106 }) if statement_label.as_deref() == Some(KEY_CLAIM) => {
107 find_keyed_message_holder(runtime, &namespace, physical_key)
108 .await?
109 .map(KeyedMessageWrite::Existing)
110 .ok_or_else(|| {
111 KhiveError::unavailable("message key holder disappeared during reconciliation")
112 .with_details(Details::new_owned([(
113 "reason",
114 "key_holder_unresolved".into(),
115 )]))
116 .into()
117 })
118 }
119 Ok(AtomicRunOutcome::RolledBack {
120 failed_op_index,
121 failure,
122 }) => Err(RuntimeError::Internal(format!(
123 "atomic message pair rolled back at op {failed_op_index}: {failure:?}"
124 ))),
125 Err(error) => Err(RuntimeError::Storage(error.0)),
126 }
127}
128
129pub async fn find_keyed_message_holder(
132 runtime: &KhiveRuntime,
133 namespace: &str,
134 physical_key: &str,
135) -> RuntimeResult<Option<Uuid>> {
136 let holder = runtime
137 .sql()
138 .reader()
139 .await?
140 .query_scalar(SqlStatement {
141 sql: "SELECT id FROM notes WHERE namespace = ?1 AND kind = 'message' \
142 AND key = ?2 AND deleted_at IS NULL LIMIT 1"
143 .into(),
144 params: vec![
145 SqlValue::Text(namespace.to_owned()),
146 SqlValue::Text(physical_key.to_owned()),
147 ],
148 label: Some("message-key-holder".into()),
149 })
150 .await?;
151 match holder {
152 Some(SqlValue::Text(id)) => Uuid::parse_str(&id)
153 .map(Some)
154 .map_err(|error| RuntimeError::Internal(format!("invalid message holder id: {error}"))),
155 None => Ok(None),
156 Some(_) => Err(RuntimeError::Internal(
157 "message holder id is not text".into(),
158 )),
159 }
160}