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::DeliveryReceipt;
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};
15use uuid::Uuid;
16
17#[derive(Debug, thiserror::Error, PartialEq, Eq)]
18pub enum ReceiptVerificationError {
19    #[error("receipt agent identifier is not a canonical UUID")]
20    InvalidAgentId,
21    #[error("receipt signature must contain exactly 64 bytes")]
22    InvalidSignatureLength,
23    #[error("recipient receipt signature is invalid")]
24    InvalidSignature,
25}
26
27fn canonical_agent_id(value: &str) -> Result<Uuid, ReceiptVerificationError> {
28    let id = Uuid::parse_str(value).map_err(|_| ReceiptVerificationError::InvalidAgentId)?;
29    if id.to_string() != value {
30        return Err(ReceiptVerificationError::InvalidAgentId);
31    }
32    Ok(id)
33}
34
35fn receipt_signing_input(receipt: &DeliveryReceipt) -> Result<Vec<u8>, ReceiptVerificationError> {
36    let binding = &receipt.binding;
37    let sender_agent_id = canonical_agent_id(&binding.sender_agent_id)?;
38    let recipient_agent_id = canonical_agent_id(&binding.recipient_agent_id)?;
39    let mut input = b"khive-node-v1/receipt\0".to_vec();
40    input.extend_from_slice(&binding.protocol_version.to_be_bytes());
41    input.extend_from_slice(binding.logical_message_id.as_bytes());
42    input.extend_from_slice(sender_agent_id.as_bytes());
43    input.extend_from_slice(recipient_agent_id.as_bytes());
44    input.extend_from_slice(binding.recipient_device_id.as_bytes());
45    input.extend_from_slice(&binding.recipient_key_epoch.to_be_bytes());
46    input.extend_from_slice(&binding.contact_generation.to_be_bytes());
47    input.extend_from_slice(binding.delivery_attempt_id.as_bytes());
48    input.push(match receipt.disposition {
49        khive_channel::ReceiptDisposition::Stored => 1,
50        khive_channel::ReceiptDisposition::Quarantined => 2,
51    });
52    Ok(input)
53}
54
55/// A receipt whose signature was checked with the owner's pinned recipient key.
56/// The private payload prevents callers from asserting verification themselves.
57#[derive(Debug)]
58pub struct VerifiedRecipientReceipt(DeliveryReceipt);
59impl VerifiedRecipientReceipt {
60    /// Verify the receipt with the signing public key pinned by the owner for
61    /// this contact at the receipt's recipient key epoch.
62    pub fn verify(
63        receipt: DeliveryReceipt,
64        pinned_signing_public_key: &[u8; 32],
65    ) -> Result<Self, ReceiptVerificationError> {
66        if receipt.signature.len() != 64 {
67            return Err(ReceiptVerificationError::InvalidSignatureLength);
68        }
69        let input = receipt_signing_input(&receipt)?;
70        ring::signature::UnparsedPublicKey::new(
71            &ring::signature::ED25519,
72            pinned_signing_public_key,
73        )
74        .verify(&input, &receipt.signature)
75        .map_err(|_| ReceiptVerificationError::InvalidSignature)?;
76        Ok(Self(receipt))
77    }
78}
79
80impl KhiveRuntime {
81    fn sender_transport_store(&self) -> SenderTransportStore {
82        SenderTransportStore::new(self.backend().pool_arc())
83    }
84    /// Persist immutable bytes before first submission. An exact retry is a no-op;
85    /// a different envelope with the same identity is refused. The namespace is
86    /// caller attribution and is always taken from the token. Sender assurance
87    /// is claimed until authenticated verification is available.
88    pub async fn create_sender_transport(
89        &self,
90        token: &NamespaceToken,
91        mut envelope: SenderEnvelope,
92    ) -> RuntimeResult<SenderRecord> {
93        envelope.namespace = token.namespace().as_str().to_owned();
94        envelope.sender_assurance = SenderAssurance::Claimed;
95        Ok(self
96            .sender_transport_store()
97            .create(envelope, false)
98            .await?)
99    }
100    /// Re-encrypt only after the caller confirms the recipient key change. The
101    /// old record must be on a key-change hold, and the new epoch must increase.
102    /// Cryptographic construction and directory confirmation belong to the caller.
103    /// Caller-supplied assurance is normalized to claimed as on initial creation.
104    pub async fn reencrypt_sender_transport_after_confirmed_key_change(
105        &self,
106        token: &NamespaceToken,
107        mut envelope: SenderEnvelope,
108    ) -> RuntimeResult<SenderRecord> {
109        envelope.namespace = token.namespace().as_str().to_owned();
110        envelope.sender_assurance = SenderAssurance::Claimed;
111        Ok(self.sender_transport_store().create(envelope, true).await?)
112    }
113    /// Read by durable identity. Absence means unknown; unknown is not persisted.
114    pub async fn sender_transport(&self, key: EnvelopeKey) -> RuntimeResult<Option<SenderRecord>> {
115        Ok(self.sender_transport_store().get(key).await?)
116    }
117    /// Due unheld submissions for the exact configured kind and slug. Times use
118    /// Unix microseconds; choosing the backoff and jitter is the caller's job.
119    pub async fn pending_sender_transports(
120        &self,
121        token: &NamespaceToken,
122        kind: &str,
123        slug: &str,
124        now: i64,
125        limit: u32,
126    ) -> RuntimeResult<Vec<SenderRecord>> {
127        Ok(self
128            .sender_transport_store()
129            .list_pending(token.namespace().as_str(), kind, slug, now, limit)
130            .await?)
131    }
132    /// Record service admission time and schedule resubmission 600 seconds later.
133    /// The timestamp is stored as Unix microseconds and survives restarts.
134    pub async fn record_sender_transport_admission(
135        &self,
136        key: EnvelopeKey,
137        admitted_at: chrono::DateTime<chrono::Utc>,
138    ) -> RuntimeResult<()> {
139        Ok(self
140            .sender_transport_store()
141            .record_admission(key, admitted_at.timestamp_micros())
142            .await?)
143    }
144    /// Record a failed attempt. Authentication errors remain pending and do not
145    /// advance the attempt count. This never modifies the envelope bytes.
146    pub async fn record_sender_transport_failure(
147        &self,
148        key: EnvelopeKey,
149        class: FailureClass,
150        next_retry_at: Option<i64>,
151    ) -> RuntimeResult<()> {
152        Ok(self
153            .sender_transport_store()
154            .record_failure(key, class, next_retry_at)
155            .await?)
156    }
157    /// Set/release credit or policy holds, or set a key-change hold. A policy hold
158    /// carries the evaluated mode and revision. Key-change holds can only be
159    /// superseded through explicit confirmed re-encryption.
160    pub async fn hold_sender_transport(
161        &self,
162        key: EnvelopeKey,
163        reason: Option<HoldReason>,
164    ) -> RuntimeResult<()> {
165        Ok(self.sender_transport_store().hold(key, reason).await?)
166    }
167    /// Accept an already signature-verified recipient receipt, including after a
168    /// local failure. Every known binding field and target disposition must match
169    /// durable data. This method performs no signature verification itself.
170    pub async fn accept_verified_recipient_receipt(
171        &self,
172        key: EnvelopeKey,
173        state: TransportState,
174        receipt: VerifiedRecipientReceipt,
175    ) -> RuntimeResult<()> {
176        let receipt = serde_json::to_value(receipt.0).map_err(|error| {
177            crate::RuntimeError::InvalidInput(format!("invalid receipt: {error}"))
178        })?;
179        Ok(self
180            .sender_transport_store()
181            .accept_receipt(key, state, receipt)
182            .await?)
183    }
184}
185
186#[cfg(test)]
187mod tests {
188    use super::*;
189    use khive_channel::{DeliveryReceiptBinding, ReceiptDisposition};
190    use ring::signature::{Ed25519KeyPair, KeyPair};
191    use uuid::Uuid;
192
193    fn uuid(value: &str) -> Uuid {
194        Uuid::parse_str(value).unwrap()
195    }
196
197    fn decode_hex(value: &str) -> Vec<u8> {
198        value
199            .as_bytes()
200            .chunks_exact(2)
201            .map(|pair| u8::from_str_radix(std::str::from_utf8(pair).unwrap(), 16).unwrap())
202            .collect()
203    }
204
205    fn vector_receipt(disposition: ReceiptDisposition) -> DeliveryReceipt {
206        let signature = match disposition {
207            ReceiptDisposition::Stored => {
208                "00de06d16479cd2a99fd6e66ae90ad66716c9af015fcba6b65db52029cbbe83bc9161dac398929cfb7955e3eb3cd8f4ae69e2cde7331d9aeb2a1f597a6e12f01"
209            }
210            ReceiptDisposition::Quarantined => {
211                "1925fea98935bc24b684df54d625b60748de38292a55a6215802607e51b44fda8190120b7b5460562e1ab3404c7096931144ab6ac2d3201e1b5543b96939aa0c"
212            }
213        };
214        DeliveryReceipt {
215            binding: DeliveryReceiptBinding {
216                protocol_version: 1,
217                logical_message_id: uuid("6f1c2d3e-4a5b-4c6d-8e7f-90a1b2c3d4e5"),
218                sender_agent_id: "01920000-0000-7000-8000-00000000a001".into(),
219                recipient_agent_id: "01920000-0000-7000-8000-00000000a002".into(),
220                recipient_device_id: uuid("01920000-0000-7000-8000-00000000d002"),
221                recipient_key_epoch: 2,
222                contact_generation: 3,
223                delivery_attempt_id: uuid("01920000-0000-7000-8000-0000000e0001"),
224            },
225            disposition,
226            signature: decode_hex(signature),
227        }
228    }
229
230    fn vector_key(value: &str) -> [u8; 32] {
231        decode_hex(value).try_into().unwrap()
232    }
233
234    fn signed_receipt(
235        mut receipt: DeliveryReceipt,
236        seed: &[u8; 32],
237    ) -> (DeliveryReceipt, [u8; 32]) {
238        let key_pair = Ed25519KeyPair::from_seed_unchecked(seed).unwrap();
239        let input = receipt_signing_input(&receipt).unwrap();
240        receipt.signature = key_pair.sign(&input).as_ref().to_vec();
241        let pinned_key = key_pair.public_key().as_ref().try_into().unwrap();
242        (receipt, pinned_key)
243    }
244
245    #[test]
246    fn recipient_receipt_vectors_match_signing_input_and_verify() {
247        let recipient_key =
248            vector_key("a914d2b78bbef06e728db06ad577d1c09d04dae4a078ab7b7574187d9dc5d032");
249        let vectors = [
250            (
251                vector_receipt(ReceiptDisposition::Stored),
252                "6b686976652d6e6f64652d76312f7265636569707400000000016f1c2d3e4a5b4c6d8e7f90a1b2c3d4e50192000000007000800000000000a0010192000000007000800000000000a0020192000000007000800000000000d00200000000000000020000000000000003019200000000700080000000000e000101",
253            ),
254            (
255                vector_receipt(ReceiptDisposition::Quarantined),
256                "6b686976652d6e6f64652d76312f7265636569707400000000016f1c2d3e4a5b4c6d8e7f90a1b2c3d4e50192000000007000800000000000a0010192000000007000800000000000a0020192000000007000800000000000d00200000000000000020000000000000003019200000000700080000000000e000102",
257            ),
258        ];
259        for (receipt, expected_input) in vectors {
260            assert_eq!(
261                receipt_signing_input(&receipt).unwrap(),
262                decode_hex(expected_input)
263            );
264            VerifiedRecipientReceipt::verify(receipt, &recipient_key)
265                .expect("recipient receipt vector must verify");
266        }
267    }
268
269    #[test]
270    fn invalid_recipient_receipt_signatures_are_refused() {
271        let recipient_key =
272            vector_key("a914d2b78bbef06e728db06ad577d1c09d04dae4a078ab7b7574187d9dc5d032");
273        let sender_key =
274            vector_key("9016672157bdb5b3529477312593f8e6fbf59641a52a374d50bd72fdf0f5d2af");
275        let stored = vector_receipt(ReceiptDisposition::Stored);
276
277        let mut quarantined_input = stored.clone();
278        quarantined_input.disposition = ReceiptDisposition::Quarantined;
279        assert!(VerifiedRecipientReceipt::verify(quarantined_input, &recipient_key).is_err());
280
281        let mut different_attempt = stored.clone();
282        different_attempt.binding.delivery_attempt_id =
283            uuid("01920000-0000-7000-8000-0000000e0002");
284        assert!(VerifiedRecipientReceipt::verify(different_attempt, &recipient_key).is_err());
285
286        let mut different_message = stored.clone();
287        different_message.binding.logical_message_id = uuid("6f1c2d3e-4a5b-4c6d-8e7f-90a1b2c3d4e6");
288        assert!(VerifiedRecipientReceipt::verify(different_message, &recipient_key).is_err());
289
290        let mut different_sender = stored.clone();
291        different_sender.binding.sender_agent_id = "01920000-0000-7000-8000-00000000a003".into();
292        assert!(VerifiedRecipientReceipt::verify(different_sender, &recipient_key).is_err());
293        assert!(VerifiedRecipientReceipt::verify(stored.clone(), &sender_key).is_err());
294
295        let mut invalid_agent_id = stored.clone();
296        invalid_agent_id.binding.sender_agent_id = "not-a-uuid".into();
297        assert!(VerifiedRecipientReceipt::verify(invalid_agent_id, &recipient_key).is_err());
298
299        let mut wrong_signature_length = stored;
300        wrong_signature_length.signature.pop();
301        assert!(VerifiedRecipientReceipt::verify(wrong_signature_length, &recipient_key).is_err());
302    }
303
304    fn envelope() -> SenderEnvelope {
305        SenderEnvelope {
306            namespace: "untrusted-envelope-attribution".into(),
307            logical_message_id: Uuid::new_v4(),
308            outbound_note_id: Uuid::new_v4(),
309            kind: "khive".into(),
310            slug: "device".into(),
311            credential_ref: "keys/device".into(),
312            recipient_address: format!("khive1:example/{}", Uuid::nil()),
313            protocol_version: 1,
314            sender_agent_id: Uuid::new_v4().to_string(),
315            sender_assurance: SenderAssurance::DaemonBearer,
316            recipient_agent_id: Uuid::nil().to_string(),
317            recipient_device_id: Uuid::new_v4(),
318            recipient_key_epoch: 1,
319            contact_generation: 1,
320            sender_key_epoch: 1,
321            recipient_key_fingerprint: "ab".repeat(32),
322            enc: vec![1; 32],
323            ciphertext: vec![2, 0, 255],
324        }
325    }
326
327    #[tokio::test]
328    async fn verified_receipt_uses_bound_backend_and_token_attribution() {
329        let runtime = KhiveRuntime::memory().unwrap();
330        let unrelated = KhiveRuntime::memory().unwrap();
331        let token = NamespaceToken::local();
332        let envelope = envelope();
333        let key = envelope.key();
334        let stored = runtime
335            .create_sender_transport(&token, envelope.clone())
336            .await
337            .unwrap();
338        assert_eq!(stored.envelope.namespace, token.namespace().as_str());
339        assert_eq!(stored.envelope.sender_assurance, SenderAssurance::Claimed);
340        assert_eq!(
341            runtime
342                .sender_transport(key)
343                .await
344                .unwrap()
345                .unwrap()
346                .envelope
347                .sender_assurance,
348            SenderAssurance::Claimed
349        );
350        assert!(
351            unrelated.sender_transport(key).await.unwrap().is_none(),
352            "transport must use bound backend"
353        );
354        let hold = HoldReason::PolicyDenied {
355            mode: PolicyMode::Enforce,
356            revision: 7,
357        };
358        let admitted_at = chrono::Utc::now();
359        runtime
360            .record_sender_transport_admission(key, admitted_at)
361            .await
362            .unwrap();
363        let admitted = runtime.sender_transport(key).await.unwrap().unwrap();
364        assert_eq!(admitted.admitted_at, Some(admitted_at.timestamp_micros()));
365        assert_eq!(
366            admitted.next_retry_at,
367            Some(admitted_at.timestamp_micros() + 600_000_000)
368        );
369        runtime
370            .hold_sender_transport(key, Some(hold))
371            .await
372            .unwrap();
373        assert_eq!(
374            runtime
375                .sender_transport(key)
376                .await
377                .unwrap()
378                .unwrap()
379                .hold_reason,
380            Some(hold)
381        );
382        let receipt = DeliveryReceipt {
383            binding: DeliveryReceiptBinding {
384                protocol_version: 1,
385                logical_message_id: envelope.logical_message_id,
386                sender_agent_id: envelope.sender_agent_id,
387                recipient_agent_id: envelope.recipient_agent_id,
388                recipient_device_id: envelope.recipient_device_id,
389                recipient_key_epoch: 1,
390                contact_generation: 1,
391                delivery_attempt_id: Uuid::new_v4(),
392            },
393            disposition: ReceiptDisposition::Stored,
394            signature: vec![7],
395        };
396        let (receipt, pinned_key) = signed_receipt(receipt, &[19; 32]);
397        assert!(runtime
398            .accept_verified_recipient_receipt(
399                key,
400                TransportState::RecipientQuarantined,
401                VerifiedRecipientReceipt::verify(receipt.clone(), &pinned_key).unwrap()
402            )
403            .await
404            .is_err());
405        runtime
406            .accept_verified_recipient_receipt(
407                key,
408                TransportState::RecipientStored,
409                VerifiedRecipientReceipt::verify(receipt, &pinned_key).unwrap(),
410            )
411            .await
412            .unwrap();
413        assert_eq!(
414            runtime
415                .sender_transport(key)
416                .await
417                .unwrap()
418                .unwrap()
419                .hold_reason,
420            None
421        );
422        assert_eq!(
423            runtime.sender_transport(key).await.unwrap().unwrap().state,
424            TransportState::RecipientStored
425        );
426    }
427    #[tokio::test]
428    async fn caller_sender_assurance_is_claimed() {
429        for assurance in [
430            SenderAssurance::DaemonBearer,
431            SenderAssurance::ActorSignature,
432        ] {
433            let runtime = KhiveRuntime::memory().unwrap();
434            let token = NamespaceToken::local();
435            let mut envelope = envelope();
436            envelope.sender_assurance = assurance;
437            let first = runtime
438                .create_sender_transport(&token, envelope.clone())
439                .await
440                .unwrap();
441            assert_eq!(
442                first.envelope.sender_assurance,
443                SenderAssurance::Claimed,
444                "create must normalize caller assurance"
445            );
446            assert_eq!(
447                runtime.sender_transport(envelope.key()).await.unwrap(),
448                Some(first)
449            );
450            runtime
451                .hold_sender_transport(envelope.key(), Some(HoldReason::RecipientKeyChanged))
452                .await
453                .unwrap();
454            envelope.recipient_key_epoch += 1;
455            let next = runtime
456                .reencrypt_sender_transport_after_confirmed_key_change(&token, envelope.clone())
457                .await
458                .expect("re-encryption must normalize caller assurance");
459            assert_eq!(next.envelope.sender_assurance, SenderAssurance::Claimed);
460            assert_eq!(
461                runtime.sender_transport(envelope.key()).await.unwrap(),
462                Some(next)
463            );
464        }
465    }
466}