1#![allow(dead_code)] use 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#[derive(Debug)]
46pub struct DueAcknowledgementEntry {
47 journal: AckJournalEntry,
48}
49
50impl DueAcknowledgementEntry {
51 pub fn binding(&self) -> &serde_json::Value {
53 &self.journal.binding
54 }
55
56 pub fn disposition(&self) -> RecipientDisposition {
58 self.journal.disposition
59 }
60
61 pub fn delivery_attempt_id(&self) -> Uuid {
63 self.journal.delivery_attempt_id
64 }
65
66 pub fn attempt_count(&self) -> u64 {
68 self.journal.attempt_count
69 }
70
71 pub fn not_before(&self) -> Option<i64> {
73 self.journal.not_before
74 }
75}
76
77impl KhiveRuntime {
78 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 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 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 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
129pub(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#[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
156pub(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 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 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 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 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 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 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;