Skip to main content

khive_runtime/
comm_recipient.rs

1//! Trusted in-process recipient ingest. No registry verb dispatches this API.
2//!
3//! The authenticated node-loop issuer is not yet present, so this entry point
4//! is unavailable to callers outside the runtime crate.
5//!
6//! ```compile_fail
7//! use khive_runtime::KhiveRuntime;
8//! fn main() { let _ = KhiveRuntime::ingest_verified_recipient; }
9//! ```
10//!
11//! ```compile_fail
12//! use khive_runtime::comm_recipient::LocalRecipientBinding;
13//! fn main() { let _ = std::mem::size_of::<LocalRecipientBinding>(); }
14//! ```
15//!
16//! ```compile_fail
17//! use khive_runtime::comm_recipient::VerifiedInboundContent;
18//! fn main() { let _ = std::mem::size_of::<VerifiedInboundContent>(); }
19//! ```
20#![allow(dead_code)] // The verified node-loop caller is supplied by #3537.
21use crate::{KhiveRuntime, NamespaceToken, RuntimeError, RuntimeResult};
22use khive_channel::InboundReceiptTicket;
23use khive_db::stores::note::recipient::{
24    AckJournalEntry, QuarantineRecord, RecipientCommit, RecipientTransportStore,
25};
26pub use khive_db::stores::note::recipient::{
27    AcknowledgementRetirementReason, QuarantineReason, RecipientCommitResult, RecipientDisposition,
28};
29use khive_storage::Note;
30use serde_json::json;
31use uuid::Uuid;
32
33/// An acknowledgement returned by the runtime's durable journal read.
34///
35/// Callers can inspect the binding and retry bookkeeping, but cannot manufacture
36/// an entry from uncommitted delivery data. This type stores no signed bytes.
37///
38/// ```compile_fail
39/// use khive_db::stores::note::recipient::AckJournalEntry;
40/// use khive_runtime::comm_recipient::DueAcknowledgementEntry;
41/// fn forge(journal: AckJournalEntry) {
42///     let _ = DueAcknowledgementEntry { journal };
43/// }
44/// ```
45#[derive(Debug)]
46pub struct DueAcknowledgementEntry {
47    journal: AckJournalEntry,
48}
49
50impl DueAcknowledgementEntry {
51    /// The binding recorded when the delivery was committed or replayed.
52    pub fn binding(&self) -> &serde_json::Value {
53        &self.journal.binding
54    }
55
56    /// The disposition first committed for this logical message.
57    pub fn disposition(&self) -> RecipientDisposition {
58        self.journal.disposition
59    }
60
61    /// The exact delivery attempt this acknowledgement finishes.
62    pub fn delivery_attempt_id(&self) -> Uuid {
63        self.journal.delivery_attempt_id
64    }
65
66    /// The number of failed acknowledgement tries recorded durably.
67    pub fn attempt_count(&self) -> u64 {
68        self.journal.attempt_count
69    }
70
71    /// The earliest retry time, in microseconds since the Unix epoch.
72    pub fn not_before(&self) -> Option<i64> {
73        self.journal.not_before
74    }
75}
76
77impl KhiveRuntime {
78    /// Read due acknowledgements from comm's assigned runtime, oldest first.
79    /// This is the only construction path for [`DueAcknowledgementEntry`].
80    pub async fn due_acknowledgements(
81        &self,
82        now: i64,
83        limit: usize,
84    ) -> RuntimeResult<Vec<DueAcknowledgementEntry>> {
85        let store = RecipientTransportStore::new(self.backend().pool_arc());
86        Ok(store
87            .list_due_acknowledgements(now, limit)
88            .await?
89            .into_iter()
90            .map(|journal| DueAcknowledgementEntry { journal })
91            .collect())
92    }
93
94    /// Finish a pending journal entry. Returns false for an absent or terminal
95    /// attempt, including an entry already finished by an earlier call.
96    pub async fn finish_acknowledgement(&self, delivery_attempt_id: Uuid) -> RuntimeResult<bool> {
97        let store = RecipientTransportStore::new(self.backend().pool_arc());
98        Ok(store.finish_acknowledgement(delivery_attempt_id).await?)
99    }
100
101    /// Record one failed try and its earliest retry time durably. An absent or
102    /// terminal attempt is unchanged and returns false.
103    pub async fn record_acknowledgement_failed_try(
104        &self,
105        delivery_attempt_id: Uuid,
106        not_before: i64,
107    ) -> RuntimeResult<bool> {
108        let store = RecipientTransportStore::new(self.backend().pool_arc());
109        Ok(store
110            .record_acknowledgement_failed_try(delivery_attempt_id, not_before)
111            .await?)
112    }
113
114    /// Retire a pending entry after a permanent transport refusal. The message
115    /// and its committed receipt remain intact. An absent or terminal attempt
116    /// is unchanged and returns false.
117    pub async fn retire_acknowledgement(
118        &self,
119        delivery_attempt_id: Uuid,
120        reason: AcknowledgementRetirementReason,
121    ) -> RuntimeResult<bool> {
122        let store = RecipientTransportStore::new(self.backend().pool_arc());
123        Ok(store
124            .retire_acknowledgement(delivery_attempt_id, reason)
125            .await?)
126    }
127}
128
129/// Local enrollment authority, supplied by the trusted slug owner. Never fill
130/// this from plaintext or a wire request. Actor labels are local, not agent UUIDs.
131pub(crate) struct LocalRecipientBinding {
132    pub actor: String,
133    pub realm: String,
134    pub slug: String,
135    pub agent_id: String,
136    pub device_id: Uuid,
137    pub key_epoch: u64,
138}
139/// The only purpose values a protocol-v1 sender may declare.
140#[derive(Clone, Copy, Debug, PartialEq, Eq)]
141pub(crate) enum DeclaredMessageKind {
142    Announce,
143    Report,
144    Ask,
145}
146impl DeclaredMessageKind {
147    fn as_str(self) -> &'static str {
148        match self {
149            Self::Announce => "announce",
150            Self::Report => "report",
151            Self::Ask => "ask",
152        }
153    }
154}
155
156/// Content has deliberately no sender, recipient, namespace or actor field.
157/// Identity comes exclusively from local binding and an authenticated ticket.
158pub(crate) enum VerifiedInboundContent {
159    Message {
160        content: String,
161        subject: Option<String>,
162        kind: Option<DeclaredMessageKind>,
163        in_reply_to: Option<Uuid>,
164        correlation: Option<String>,
165        sent_at: String,
166    },
167    Quarantine {
168        reason: QuarantineReason,
169        /// The parsed plaintext object, present exactly when the pair policy
170        /// refused a valid message (ADR-105 A.8 Receiving, step 4).
171        parsed_plaintext: Option<serde_json::Value>,
172    },
173}
174fn invalid(message: &str) -> RuntimeError {
175    RuntimeError::InvalidInput(message.into())
176}
177fn canonical_id(id: &str) -> bool {
178    Uuid::parse_str(id)
179        .ok()
180        .is_some_and(|u| u.to_string() == id)
181}
182impl KhiveRuntime {
183    /// Commit an authenticated delivery on comm's assigned runtime. The caller
184    /// MUST authenticate/decrypt and verify enrollment before creating the ticket.
185    /// This method performs no cryptography; it consumes the in-memory ticket,
186    /// checks local binding, and commits note/replay/ack/quarantine atomically.
187    ///
188    /// For quarantine, delivery_item is the exact original JSON object, bounded
189    /// at 98,304 bytes by the protocol request limit. No blob store is involved.
190    /// A failure commits nothing and returns no recipient disposition. Retryable
191    /// storage errors must not be converted into quarantine by the caller.
192    /// A SecretDetected refusal is deterministic: redelivering the same payload
193    /// is refused again.
194    pub(crate) async fn ingest_verified_recipient(
195        &self,
196        token: &NamespaceToken,
197        local: &LocalRecipientBinding,
198        ticket: InboundReceiptTicket,
199        payload: VerifiedInboundContent,
200        delivery_item: Vec<u8>,
201    ) -> RuntimeResult<RecipientCommitResult> {
202        let binding = ticket.binding();
203        if !canonical_id(&binding.sender_agent_id)
204            || !canonical_id(&binding.recipient_agent_id)
205            || !canonical_id(&local.agent_id)
206        {
207            return Err(invalid("verified receipt agents must be canonical UUIDs"));
208        }
209        if binding.protocol_version != 1
210            || binding.recipient_agent_id != local.agent_id
211            || binding.recipient_device_id != local.device_id
212            || binding.recipient_key_epoch != local.key_epoch
213        {
214            return Err(invalid(
215                "verified delivery does not match local recipient binding",
216            ));
217        }
218        if [
219            binding.recipient_key_epoch,
220            binding.contact_generation,
221            ticket.sender_key_epoch(),
222        ]
223        .iter()
224        .any(|n| *n == 0 || *n > u32::MAX as u64)
225        {
226            return Err(invalid("invalid ticket epoch or generation"));
227        }
228        if local.actor.trim().is_empty()
229            || local.slug.is_empty()
230            || local.realm.is_empty()
231            || local.realm.len() > 64
232            || !local
233                .realm
234                .bytes()
235                .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b"._-".contains(&b))
236        {
237            return Err(invalid("invalid local recipient route"));
238        }
239        // ADR-105 A.8 step 4: a successfully opened replay is answered from
240        // its durable claim before interpreting the new attempt's plaintext.
241        let binding_value = serde_json::to_value(binding)
242            .map_err(|error| invalid(&format!("invalid receipt binding: {error}")))?;
243        let recipient_store = RecipientTransportStore::new(self.backend().pool_arc());
244        if let Some(replay) = recipient_store
245            .ack_if_replayed(binding_value.clone(), &local.actor)
246            .await?
247        {
248            return Ok(replay);
249        }
250        self.validate_note_kind("message")?;
251        let from = format!("khive1:{}/{}", local.realm, binding.sender_agent_id);
252        let received_at = chrono::Utc::now().to_rfc3339();
253        let mut in_reply_to = None;
254        let mut correlation = None;
255        let (note, disposition, quarantine) = match payload {
256            VerifiedInboundContent::Message {
257                content,
258                subject,
259                kind,
260                in_reply_to: parent,
261                correlation: message_correlation,
262                sent_at,
263            } => {
264                let parsed_sent_at = chrono::DateTime::parse_from_rfc3339(&sent_at)
265                    .ok()
266                    .filter(|stamp| stamp.offset().local_minus_utc() == 0);
267                if let Some(stamp) = parsed_sent_at {
268                    // An empty body is a valid A.5 message and is stored like any other.
269                    crate::secret_gate::check_at(&content, "note", "content")?;
270                    if let Some(subject) = &subject {
271                        crate::secret_gate::check_at(subject, "note", "name")?;
272                    }
273                    in_reply_to = parent;
274                    correlation = message_correlation;
275                    let mut note = Note::new(token.namespace().as_str(), "message", content);
276                    note.name = subject.clone();
277                    let mut props = json!({
278                        "comm_schema_version": 1,
279                        "from": from,
280                        "from_actor": from,
281                        "to": local.actor,
282                        "to_actor": local.actor,
283                        "direction": "inbound",
284                        "read": false,
285                        "received_at": received_at,
286                        "channel_kind": "khive",
287                        "channel_slug": local.slug,
288                        "logical_message_id": binding.logical_message_id,
289                        "sent_at": stamp
290                            .with_timezone(&chrono::Utc)
291                            .to_rfc3339_opts(chrono::SecondsFormat::AutoSi, true),
292                        "message_kind": kind.map_or("unspecified", DeclaredMessageKind::as_str),
293                    });
294                    if let Some(subject) = subject {
295                        props["subject"] = json!(subject);
296                    }
297                    if let Some(parent) = in_reply_to {
298                        props["in_reply_to"] = json!(parent);
299                    }
300                    crate::secret_gate::check_json_at(&props, "note", "properties")?;
301                    note.properties = Some(props);
302                    (Some(note), RecipientDisposition::Stored, None)
303                } else {
304                    (
305                        None,
306                        RecipientDisposition::Quarantined,
307                        Some(QuarantineRecord {
308                            reason: QuarantineReason::InvalidPlaintext,
309                            delivery_item,
310                            parsed_plaintext: None,
311                        }),
312                    )
313                }
314            }
315            VerifiedInboundContent::Quarantine {
316                reason,
317                parsed_plaintext,
318            } => {
319                // Retained plaintext passes the same write-time secret gate as a stored message.
320                // Refuse before the quarantine/replay/ack transaction can write.
321                if let Some(plaintext) = &parsed_plaintext {
322                    crate::secret_gate::check_json_at(plaintext, "quarantine", "parsed_plaintext")?;
323                }
324                (
325                    None,
326                    RecipientDisposition::Quarantined,
327                    Some(QuarantineRecord {
328                        reason,
329                        delivery_item,
330                        parsed_plaintext,
331                    }),
332                )
333            }
334        };
335        let result = recipient_store
336            .commit(RecipientCommit {
337                note,
338                recipient_actor: local.actor.clone(),
339                binding: binding_value,
340                sender_agent_id: binding.sender_agent_id.clone(),
341                logical_message_id: binding.logical_message_id,
342                delivery_attempt_id: binding.delivery_attempt_id,
343                disposition,
344                quarantine,
345                in_reply_to,
346                correlation,
347            })
348            .await?;
349        if let Some(note) = &result.note {
350            // Like ordinary ingest, indexing is best-effort after the durable commit.
351            if let Ok(fts) = self.text_for_notes(token) {
352                if let Err(error) = fts
353                    .upsert_document(crate::curation::note_fts_document(note))
354                    .await
355                {
356                    tracing::warn!(note_id=%note.id,error=%error,"verified ingest FTS indexing failed");
357                }
358            }
359            for model in self.registered_embedding_model_names() {
360                match self
361                    .embed_document_with_model_outcome_for_token(
362                        token,
363                        &model,
364                        crate::curation::note_embedding_text_ref(note),
365                    )
366                    .await
367                {
368                    Ok(outcome) => {
369                        if let Err(error) = self
370                            .publish_note_vector_revision(token, note, &model, &outcome.vector)
371                            .await
372                        {
373                            tracing::warn!(note_id=%note.id,error=%error,"verified ingest vector indexing failed");
374                        }
375                    }
376                    Err(error) => {
377                        tracing::warn!(note_id=%note.id,error=%error,"verified ingest embedding failed")
378                    }
379                }
380            }
381        }
382        Ok(result)
383    }
384}
385#[cfg(test)]
386mod tests;
387
388#[cfg(test)]
389mod acknowledgement_journal_tests;