Skip to main content

khive_runtime/
keyed_message.rs

1//! A message pair publishes its outbound identity in the final atomic statement.
2
3use 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
22/// Create exactly two caller-namespaced message notes or return the outbound
23/// holder after rolling back a competing pair. The comm layer verifies payload
24/// and sibling integrity before treating `Existing` as a successful replay.
25pub 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
33/// Include every copy's attachments before the final key claim in the same unit.
34pub 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        // A caller-supplied UUID must never turn pair creation into an upsert.
67        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
129/// Look up the live outbound holder without preparing or attempting a write.
130/// The comm layer must validate the request and intact pair before returning it.
131pub 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}