Skip to main content

khive_runtime/
comm_transport.rs

1//! Runtime-owned sender transport persistence (ADR-105).
2//!
3//! Call these methods on the runtime bound to comm's assigned backend, exactly
4//! as for ordinary comm note operations. Routing happens before this boundary;
5//! the adapter never opens a database or supplies SQL. Receipt signatures are
6//! checked against the owner's pinned recipient key before durable state changes.
7use 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/// Sender-local delivery status, including the absence of a transport record.
17/// Holds retain `Pending`; `Unknown` is never persisted as a transport state.
18#[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    /// Read transport status in the token's primary namespace on the bound backend.
33    /// This does not infer the internal inbound sibling or mutate retry metadata.
34    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    /// Persist immutable bytes before first submission. An exact retry is a no-op;
57    /// a different envelope with the same identity is refused. The namespace is
58    /// caller attribution and is always taken from the token. Sender assurance
59    /// is claimed until authenticated verification is available.
60    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    /// Re-encrypt only after the caller confirms the recipient key change. The
73    /// old record must be on a key-change hold, and the new epoch must increase.
74    /// Cryptographic construction and directory confirmation belong to the caller.
75    /// Caller-supplied assurance is normalized to claimed as on initial creation.
76    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    /// Read by durable identity. Absence means unknown; unknown is not persisted.
86    pub async fn sender_transport(&self, key: EnvelopeKey) -> RuntimeResult<Option<SenderRecord>> {
87        Ok(self.sender_transport_store().get(key).await?)
88    }
89    /// Due unheld submissions for the exact configured kind and slug. Times use
90    /// Unix microseconds; choosing the backoff and jitter is the caller's job.
91    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    /// Record service admission time and schedule resubmission 600 seconds later.
105    /// The timestamp is stored as Unix microseconds and survives restarts.
106    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    /// Record a failed attempt. Authentication errors remain pending and do not
117    /// advance the attempt count. This never modifies the envelope bytes.
118    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    /// Set/release credit or policy holds, or set a key-change hold. A policy hold
130    /// carries the evaluated mode and revision. Key-change holds can only be
131    /// superseded through explicit confirmed re-encryption.
132    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    /// Accept an already signature-verified recipient receipt, including after a
140    /// local failure. Every known binding field and target disposition must match
141    /// durable data. This method performs no signature verification itself.
142    ///
143    /// An unverified payload cannot cross this boundary:
144    ///
145    /// ```compile_fail
146    /// use khive_channel::DeliveryReceipt;
147    /// use khive_runtime::{KhiveRuntime, comm_transport::{EnvelopeKey, TransportState}};
148    /// async fn accept(runtime: &KhiveRuntime, key: EnvelopeKey, receipt: DeliveryReceipt) {
149    ///     runtime.accept_verified_recipient_receipt(key, TransportState::RecipientStored, receipt)
150    ///         .await.unwrap();
151    /// }
152    /// ```
153    ///
154    /// Its positive twin uses the channel constructor and the same runtime method:
155    ///
156    /// ```
157    /// use khive_channel::{DeliveryReceipt, VerifiedRecipientReceipt};
158    /// use khive_runtime::{KhiveRuntime, comm_transport::{EnvelopeKey, TransportState}};
159    /// async fn accept(runtime: &KhiveRuntime, key: EnvelopeKey, receipt: DeliveryReceipt,
160    ///     pinned_key: &[u8; 32]) {
161    ///     let verified = VerifiedRecipientReceipt::verify(receipt, pinned_key).unwrap();
162    ///     runtime.accept_verified_recipient_receipt(key, TransportState::RecipientStored, verified)
163    ///         .await.unwrap();
164    /// }
165    /// ```
166    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}