1use 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#[derive(Debug)]
58pub struct VerifiedRecipientReceipt(DeliveryReceipt);
59impl VerifiedRecipientReceipt {
60 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 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 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 pub async fn sender_transport(&self, key: EnvelopeKey) -> RuntimeResult<Option<SenderRecord>> {
115 Ok(self.sender_transport_store().get(key).await?)
116 }
117 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 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 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 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 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}