1use crate::error::RuntimeResult;
8use crate::{KhiveRuntime, NamespaceToken};
9use khive_channel::VerifiedRecipientReceipt;
10use khive_db::stores::note::transport::SenderTransportStore;
11pub use khive_db::stores::note::transport::{
12 EnvelopeKey, FailureClass, HoldReason, PolicyMode, SenderAssurance, SenderEnvelope,
13 SenderRecord, TransportState,
14};
15
16#[derive(Clone, Copy, Debug, PartialEq, Eq, serde::Serialize)]
19#[serde(rename_all = "snake_case")]
20pub enum TransportStatus {
21 Pending,
22 RecipientStored,
23 RecipientQuarantined,
24 Failed,
25 Unknown,
26}
27
28impl KhiveRuntime {
29 fn sender_transport_store(&self) -> SenderTransportStore {
30 SenderTransportStore::new(self.backend().pool_arc())
31 }
32 pub async fn sender_transport_status(
35 &self,
36 token: &NamespaceToken,
37 outbound_note_id: uuid::Uuid,
38 ) -> RuntimeResult<TransportStatus> {
39 let Some(row) = self
40 .sender_transport_store()
41 .get_by_outbound_note_id(token.namespace().as_str(), outbound_note_id)
42 .await?
43 else {
44 return Ok(TransportStatus::Unknown);
45 };
46 if row.hold_reason.is_some() {
47 return Ok(TransportStatus::Pending);
48 }
49 Ok(match row.state {
50 TransportState::Pending => TransportStatus::Pending,
51 TransportState::RecipientStored => TransportStatus::RecipientStored,
52 TransportState::RecipientQuarantined => TransportStatus::RecipientQuarantined,
53 TransportState::Failed => TransportStatus::Failed,
54 })
55 }
56 pub async fn create_sender_transport(
61 &self,
62 token: &NamespaceToken,
63 mut envelope: SenderEnvelope,
64 ) -> RuntimeResult<SenderRecord> {
65 envelope.namespace = token.namespace().as_str().to_owned();
66 envelope.sender_assurance = SenderAssurance::Claimed;
67 Ok(self
68 .sender_transport_store()
69 .create(envelope, false)
70 .await?)
71 }
72 pub async fn reencrypt_sender_transport_after_confirmed_key_change(
77 &self,
78 token: &NamespaceToken,
79 mut envelope: SenderEnvelope,
80 ) -> RuntimeResult<SenderRecord> {
81 envelope.namespace = token.namespace().as_str().to_owned();
82 envelope.sender_assurance = SenderAssurance::Claimed;
83 Ok(self.sender_transport_store().create(envelope, true).await?)
84 }
85 pub async fn sender_transport(&self, key: EnvelopeKey) -> RuntimeResult<Option<SenderRecord>> {
87 Ok(self.sender_transport_store().get(key).await?)
88 }
89 pub async fn pending_sender_transports(
92 &self,
93 token: &NamespaceToken,
94 kind: &str,
95 slug: &str,
96 now: i64,
97 limit: u32,
98 ) -> RuntimeResult<Vec<SenderRecord>> {
99 Ok(self
100 .sender_transport_store()
101 .list_pending(token.namespace().as_str(), kind, slug, now, limit)
102 .await?)
103 }
104 pub async fn record_sender_transport_admission(
107 &self,
108 key: EnvelopeKey,
109 admitted_at: chrono::DateTime<chrono::Utc>,
110 ) -> RuntimeResult<()> {
111 Ok(self
112 .sender_transport_store()
113 .record_admission(key, admitted_at.timestamp_micros())
114 .await?)
115 }
116 pub async fn record_sender_transport_failure(
119 &self,
120 key: EnvelopeKey,
121 class: FailureClass,
122 next_retry_at: Option<i64>,
123 ) -> RuntimeResult<()> {
124 Ok(self
125 .sender_transport_store()
126 .record_failure(key, class, next_retry_at)
127 .await?)
128 }
129 pub async fn hold_sender_transport(
133 &self,
134 key: EnvelopeKey,
135 reason: Option<HoldReason>,
136 ) -> RuntimeResult<()> {
137 Ok(self.sender_transport_store().hold(key, reason).await?)
138 }
139 pub async fn accept_verified_recipient_receipt(
167 &self,
168 key: EnvelopeKey,
169 state: TransportState,
170 receipt: VerifiedRecipientReceipt,
171 ) -> RuntimeResult<()> {
172 let receipt = serde_json::to_value(receipt.receipt()).map_err(|error| {
173 crate::RuntimeError::InvalidInput(format!("invalid receipt: {error}"))
174 })?;
175 Ok(self
176 .sender_transport_store()
177 .accept_receipt(key, state, receipt)
178 .await?)
179 }
180}
181
182#[cfg(test)]
183mod tests {
184 use super::*;
185 use khive_channel::{
186 receipt_signing_input, DeliveryReceipt, DeliveryReceiptBinding, ReceiptDisposition,
187 };
188 use ring::signature::{Ed25519KeyPair, KeyPair};
189 use uuid::Uuid;
190
191 fn signed_receipt(
192 mut receipt: DeliveryReceipt,
193 seed: &[u8; 32],
194 ) -> (DeliveryReceipt, [u8; 32]) {
195 let key_pair = Ed25519KeyPair::from_seed_unchecked(seed).unwrap();
196 let input = receipt_signing_input(&receipt).unwrap();
197 receipt.signature = key_pair.sign(&input).as_ref().to_vec();
198 let pinned_key = key_pair.public_key().as_ref().try_into().unwrap();
199 (receipt, pinned_key)
200 }
201
202 fn envelope() -> SenderEnvelope {
203 SenderEnvelope {
204 namespace: "untrusted-envelope-attribution".into(),
205 logical_message_id: Uuid::new_v4(),
206 outbound_note_id: Uuid::new_v4(),
207 kind: "khive".into(),
208 slug: "device".into(),
209 credential_ref: "keys/device".into(),
210 recipient_address: format!("khive1:example/{}", Uuid::nil()),
211 protocol_version: 1,
212 sender_agent_id: Uuid::new_v4().to_string(),
213 sender_assurance: SenderAssurance::DaemonBearer,
214 recipient_agent_id: Uuid::nil().to_string(),
215 recipient_device_id: Uuid::new_v4(),
216 recipient_key_epoch: 1,
217 contact_generation: 1,
218 sender_key_epoch: 1,
219 recipient_key_fingerprint: "ab".repeat(32),
220 enc: vec![1; 32],
221 ciphertext: vec![2, 0, 255],
222 }
223 }
224
225 #[tokio::test]
226 async fn verified_receipt_uses_bound_backend_and_token_attribution() {
227 let runtime = KhiveRuntime::memory().unwrap();
228 let unrelated = KhiveRuntime::memory().unwrap();
229 let token = NamespaceToken::local();
230 let envelope = envelope();
231 let key = envelope.key();
232 let stored = runtime
233 .create_sender_transport(&token, envelope.clone())
234 .await
235 .unwrap();
236 assert_eq!(stored.envelope.namespace, token.namespace().as_str());
237 assert_eq!(stored.envelope.sender_assurance, SenderAssurance::Claimed);
238 assert_eq!(
239 runtime
240 .sender_transport(key)
241 .await
242 .unwrap()
243 .unwrap()
244 .envelope
245 .sender_assurance,
246 SenderAssurance::Claimed
247 );
248 assert!(
249 unrelated.sender_transport(key).await.unwrap().is_none(),
250 "transport must use bound backend"
251 );
252 let hold = HoldReason::PolicyDenied {
253 mode: PolicyMode::Enforce,
254 revision: 7,
255 };
256 let admitted_at = chrono::Utc::now();
257 runtime
258 .record_sender_transport_admission(key, admitted_at)
259 .await
260 .unwrap();
261 let admitted = runtime.sender_transport(key).await.unwrap().unwrap();
262 assert_eq!(admitted.admitted_at, Some(admitted_at.timestamp_micros()));
263 assert_eq!(
264 admitted.next_retry_at,
265 Some(admitted_at.timestamp_micros() + 600_000_000)
266 );
267 runtime
268 .hold_sender_transport(key, Some(hold))
269 .await
270 .unwrap();
271 assert_eq!(
272 runtime
273 .sender_transport(key)
274 .await
275 .unwrap()
276 .unwrap()
277 .hold_reason,
278 Some(hold)
279 );
280 let receipt = DeliveryReceipt {
281 binding: DeliveryReceiptBinding {
282 protocol_version: 1,
283 logical_message_id: envelope.logical_message_id,
284 sender_agent_id: envelope.sender_agent_id,
285 recipient_agent_id: envelope.recipient_agent_id,
286 recipient_device_id: envelope.recipient_device_id,
287 recipient_key_epoch: 1,
288 contact_generation: 1,
289 delivery_attempt_id: Uuid::new_v4(),
290 },
291 disposition: ReceiptDisposition::Stored,
292 signature: vec![7],
293 };
294 let (receipt, pinned_key) = signed_receipt(receipt, &[19; 32]);
295 assert!(runtime
296 .accept_verified_recipient_receipt(
297 key,
298 TransportState::RecipientQuarantined,
299 VerifiedRecipientReceipt::verify(receipt.clone(), &pinned_key).unwrap()
300 )
301 .await
302 .is_err());
303 runtime
304 .accept_verified_recipient_receipt(
305 key,
306 TransportState::RecipientStored,
307 VerifiedRecipientReceipt::verify(receipt, &pinned_key).unwrap(),
308 )
309 .await
310 .unwrap();
311 assert_eq!(
312 runtime
313 .sender_transport(key)
314 .await
315 .unwrap()
316 .unwrap()
317 .hold_reason,
318 None
319 );
320 assert_eq!(
321 runtime.sender_transport(key).await.unwrap().unwrap().state,
322 TransportState::RecipientStored
323 );
324 }
325 #[tokio::test]
326 async fn caller_sender_assurance_is_claimed() {
327 for assurance in [
328 SenderAssurance::DaemonBearer,
329 SenderAssurance::ActorSignature,
330 ] {
331 let runtime = KhiveRuntime::memory().unwrap();
332 let token = NamespaceToken::local();
333 let mut envelope = envelope();
334 envelope.sender_assurance = assurance;
335 let first = runtime
336 .create_sender_transport(&token, envelope.clone())
337 .await
338 .unwrap();
339 assert_eq!(
340 first.envelope.sender_assurance,
341 SenderAssurance::Claimed,
342 "create must normalize caller assurance"
343 );
344 assert_eq!(
345 runtime.sender_transport(envelope.key()).await.unwrap(),
346 Some(first)
347 );
348 runtime
349 .hold_sender_transport(envelope.key(), Some(HoldReason::RecipientKeyChanged))
350 .await
351 .unwrap();
352 envelope.recipient_key_epoch += 1;
353 let next = runtime
354 .reencrypt_sender_transport_after_confirmed_key_change(&token, envelope.clone())
355 .await
356 .expect("re-encryption must normalize caller assurance");
357 assert_eq!(next.envelope.sender_assurance, SenderAssurance::Claimed);
358 assert_eq!(
359 runtime.sender_transport(envelope.key()).await.unwrap(),
360 Some(next)
361 );
362 }
363 }
364
365 async fn accept_status_receipt(
366 runtime: &KhiveRuntime,
367 envelope: &SenderEnvelope,
368 disposition: ReceiptDisposition,
369 ) {
370 let state = match disposition {
371 ReceiptDisposition::Stored => TransportState::RecipientStored,
372 ReceiptDisposition::Quarantined => TransportState::RecipientQuarantined,
373 };
374 let receipt = DeliveryReceipt {
375 binding: DeliveryReceiptBinding {
376 protocol_version: envelope.protocol_version,
377 logical_message_id: envelope.logical_message_id,
378 sender_agent_id: envelope.sender_agent_id.clone(),
379 recipient_agent_id: envelope.recipient_agent_id.clone(),
380 recipient_device_id: envelope.recipient_device_id,
381 recipient_key_epoch: envelope.recipient_key_epoch,
382 contact_generation: envelope.contact_generation,
383 delivery_attempt_id: Uuid::new_v4(),
384 },
385 disposition,
386 signature: vec![],
387 };
388 let (receipt, pinned_key) = signed_receipt(receipt, &[23; 32]);
389 runtime
390 .accept_verified_recipient_receipt(
391 envelope.key(),
392 state,
393 VerifiedRecipientReceipt::verify(receipt, &pinned_key).unwrap(),
394 )
395 .await
396 .unwrap();
397 }
398
399 #[tokio::test]
400 async fn transport_status_pending_admission_and_hold() {
401 let runtime = KhiveRuntime::memory().unwrap();
402 let token = NamespaceToken::local();
403 let e = envelope();
404 runtime
405 .create_sender_transport(&token, e.clone())
406 .await
407 .unwrap();
408 assert_eq!(
409 runtime
410 .sender_transport_status(&token, e.outbound_note_id)
411 .await
412 .unwrap(),
413 TransportStatus::Pending,
414 "a newly persisted envelope is pending"
415 );
416 runtime
417 .record_sender_transport_admission(e.key(), chrono::Utc::now())
418 .await
419 .unwrap();
420 assert_eq!(
421 runtime
422 .sender_transport_status(&token, e.outbound_note_id)
423 .await
424 .unwrap(),
425 TransportStatus::Pending,
426 "service admission is not recipient storage"
427 );
428 runtime
429 .hold_sender_transport(e.key(), Some(HoldReason::InsufficientCredit))
430 .await
431 .unwrap();
432 assert_eq!(
433 runtime
434 .sender_transport_status(&token, e.outbound_note_id)
435 .await
436 .unwrap(),
437 TransportStatus::Pending,
438 "an insufficient-credit hold is pending, never a further status"
439 );
440 }
441
442 #[tokio::test]
443 async fn transport_status_failed_and_verified_receipts() {
444 let runtime = KhiveRuntime::memory().unwrap();
445 let token = NamespaceToken::local();
446 for (disposition, expected) in [
447 (ReceiptDisposition::Stored, TransportStatus::RecipientStored),
448 (
449 ReceiptDisposition::Quarantined,
450 TransportStatus::RecipientQuarantined,
451 ),
452 ] {
453 let e = envelope();
454 runtime
455 .create_sender_transport(&token, e.clone())
456 .await
457 .unwrap();
458 runtime
459 .record_sender_transport_failure(e.key(), FailureClass::Permanent, None)
460 .await
461 .unwrap();
462 assert_eq!(
463 runtime
464 .sender_transport_status(&token, e.outbound_note_id)
465 .await
466 .unwrap(),
467 TransportStatus::Failed,
468 "a permanent local refusal is failed"
469 );
470 accept_status_receipt(&runtime, &e, disposition).await;
471 assert_eq!(
472 runtime
473 .sender_transport_status(&token, e.outbound_note_id)
474 .await
475 .unwrap(),
476 expected,
477 "a verified receipt supersedes the local failure"
478 );
479 }
480 assert_eq!(
481 runtime
482 .sender_transport_status(&token, Uuid::new_v4())
483 .await
484 .unwrap(),
485 TransportStatus::Unknown,
486 "an absent outbound UUID is unknown"
487 );
488 }
489
490 #[tokio::test]
491 async fn transport_status_is_primary_namespace_and_backend_scoped() {
492 let runtime = KhiveRuntime::memory().unwrap();
493 let token_a = NamespaceToken::for_namespace(crate::Namespace::parse("sender-a").unwrap());
494 let token_b = NamespaceToken::mint_with_visibility(
495 crate::Namespace::parse("sender-b").unwrap(),
496 vec![crate::Namespace::parse("sender-a").unwrap()],
497 crate::ActorRef::anonymous(),
498 );
499 let e = envelope();
500 runtime
501 .create_sender_transport(&token_a, e.clone())
502 .await
503 .unwrap();
504 assert_eq!(
505 runtime
506 .sender_transport(e.key())
507 .await
508 .unwrap()
509 .unwrap()
510 .envelope
511 .namespace,
512 "sender-a"
513 );
514 assert_eq!(
515 runtime
516 .sender_transport_status(&token_a, e.outbound_note_id)
517 .await
518 .unwrap(),
519 TransportStatus::Pending
520 );
521 assert_eq!(
522 runtime
523 .sender_transport_status(&token_b, e.outbound_note_id)
524 .await
525 .unwrap(),
526 TransportStatus::Unknown,
527 "additional visible namespaces cannot reveal transport status"
528 );
529 let unrelated = KhiveRuntime::memory().unwrap();
530 assert_eq!(
531 unrelated
532 .sender_transport_status(&token_a, e.outbound_note_id)
533 .await
534 .unwrap(),
535 TransportStatus::Unknown,
536 "a different bound backend has no record"
537 );
538 }
539
540 #[tokio::test]
541 async fn transport_status_read_preserves_every_sender_field() {
542 let runtime = KhiveRuntime::memory().unwrap();
543 let token = NamespaceToken::local();
544 let e = envelope();
545 runtime
546 .create_sender_transport(&token, e.clone())
547 .await
548 .unwrap();
549 let before = runtime.sender_transport(e.key()).await.unwrap().unwrap();
550 assert_eq!(before.admitted_at, None);
551 for _ in 0..3 {
552 assert_eq!(
553 runtime
554 .sender_transport_status(&token, e.outbound_note_id)
555 .await
556 .unwrap(),
557 TransportStatus::Pending
558 );
559 }
560 let after = runtime.sender_transport(e.key()).await.unwrap().unwrap();
561 assert_eq!(
562 after, before,
563 "a status read preserves envelope bytes, holds, attempts and timestamps"
564 );
565 }
566}