1use super::{map_err, SqlNoteStore};
3use crate::pool::ConnectionPool;
4use khive_storage::{StorageCapability, StorageError, StorageResult};
5use rusqlite::{params, OptionalExtension};
6use serde::{Deserialize, Serialize};
7use serde_json::Value;
8use std::sync::Arc;
9use uuid::Uuid;
10
11#[derive(Clone, Copy, Debug, PartialEq, Eq)]
12pub struct EnvelopeKey {
13 pub logical_message_id: Uuid,
14 pub recipient_device_id: Uuid,
15 pub recipient_key_epoch: u64,
16}
17
18#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
21pub struct SenderEnvelope {
22 pub namespace: String,
23 pub logical_message_id: Uuid,
24 pub outbound_note_id: Uuid,
25 pub kind: String,
26 pub slug: String,
27 pub credential_ref: String,
28 pub recipient_address: String,
29 pub protocol_version: u32,
30 pub sender_agent_id: String,
31 pub sender_assurance: SenderAssurance,
32 pub recipient_agent_id: String,
33 pub recipient_device_id: Uuid,
34 pub recipient_key_epoch: u64,
35 pub contact_generation: u64,
36 pub sender_key_epoch: u64,
37 pub recipient_key_fingerprint: String,
38 pub enc: Vec<u8>,
39 pub ciphertext: Vec<u8>,
40}
41impl SenderEnvelope {
42 pub fn key(&self) -> EnvelopeKey {
43 EnvelopeKey {
44 logical_message_id: self.logical_message_id,
45 recipient_device_id: self.recipient_device_id,
46 recipient_key_epoch: self.recipient_key_epoch,
47 }
48 }
49 pub fn validate(&self) -> StorageResult<()> {
50 for id in [&self.sender_agent_id, &self.recipient_agent_id] {
51 if Uuid::parse_str(id).ok().map(|id| id.to_string()).as_deref() != Some(id.as_str()) {
52 return Err(invalid("agent id must be a canonical UUID"));
53 }
54 }
55 if self.enc.len() != 32 || self.ciphertext.len() > 65_536 {
56 return Err(invalid("invalid envelope byte lengths"));
57 }
58 if [
59 self.recipient_key_epoch,
60 self.sender_key_epoch,
61 self.contact_generation,
62 ]
63 .iter()
64 .any(|n| *n == 0 || *n > u32::MAX as u64)
65 {
66 return Err(invalid("epoch/generation outside supported range"));
67 }
68 if self.protocol_version != 1
69 || self.kind.is_empty()
70 || self.slug.is_empty()
71 || self.credential_ref.is_empty()
72 {
73 return Err(invalid("invalid transport identity"));
74 }
75 if self.recipient_key_fingerprint.len() != 64
76 || !khive_types::is_lowercase_hex(&self.recipient_key_fingerprint)
77 {
78 return Err(invalid("fingerprint must be 32 lowercase hex bytes"));
79 }
80 Ok(())
81 }
82}
83
84#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
86#[serde(rename_all = "snake_case")]
87pub enum SenderAssurance {
88 Claimed,
89 DaemonBearer,
90 ActorSignature,
91}
92
93#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
94#[serde(rename_all = "snake_case")]
95pub enum TransportState {
96 Pending,
97 RecipientStored,
98 RecipientQuarantined,
99 Failed,
100}
101#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
102#[serde(rename_all = "snake_case")]
103pub enum FailureClass {
104 Transient,
105 Authentication,
106 Permanent,
107}
108#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
109#[serde(rename_all = "snake_case")]
110pub enum HoldReason {
111 InsufficientCredit,
112 RecipientKeyChanged,
113 PolicyDenied {
115 mode: PolicyMode,
116 revision: u64,
117 },
118}
119
120#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
122#[serde(rename_all = "snake_case")]
123pub enum PolicyMode {
124 Off,
125 Shadow,
126 Enforce,
127}
128#[derive(Clone, Debug, PartialEq, Eq)]
129pub struct SenderRecord {
130 pub envelope: SenderEnvelope,
131 pub state: TransportState,
132 pub attempt_count: u64,
133 pub envelope_seq: u64,
134 pub next_retry_at: Option<i64>,
135 pub last_failure_class: Option<FailureClass>,
136 pub hold_reason: Option<HoldReason>,
137 pub receipt: Option<Value>,
138 pub created_at: i64,
139 pub updated_at: i64,
140 pub admitted_at: Option<i64>,
141}
142
143fn invalid(message: &str) -> StorageError {
144 StorageError::InvalidInput {
145 capability: StorageCapability::Notes,
146 operation: "sender_transport".into(),
147 message: message.into(),
148 }
149}
150trait StorageSpelling {
151 fn storage_spelling(&self) -> &'static str;
152}
153impl StorageSpelling for SenderAssurance {
154 fn storage_spelling(&self) -> &'static str {
155 match self {
156 Self::Claimed => "claimed",
157 Self::DaemonBearer => "daemon_bearer",
158 Self::ActorSignature => "actor_signature",
159 }
160 }
161}
162impl StorageSpelling for TransportState {
163 fn storage_spelling(&self) -> &'static str {
164 match self {
165 Self::Pending => "pending",
166 Self::RecipientStored => "recipient_stored",
167 Self::RecipientQuarantined => "recipient_quarantined",
168 Self::Failed => "failed",
169 }
170 }
171}
172impl StorageSpelling for FailureClass {
173 fn storage_spelling(&self) -> &'static str {
174 match self {
175 Self::Transient => "transient",
176 Self::Authentication => "authentication",
177 Self::Permanent => "permanent",
178 }
179 }
180}
181impl StorageSpelling for HoldReason {
182 fn storage_spelling(&self) -> &'static str {
183 match self {
184 Self::InsufficientCredit => "insufficient_credit",
185 Self::RecipientKeyChanged => "recipient_key_changed",
186 Self::PolicyDenied { .. } => "policy_denied",
187 }
188 }
189}
190impl StorageSpelling for PolicyMode {
191 fn storage_spelling(&self) -> &'static str {
192 match self {
193 Self::Off => "off",
194 Self::Shadow => "shadow",
195 Self::Enforce => "enforce",
196 }
197 }
198}
199fn encode(value: &impl StorageSpelling) -> &'static str {
200 value.storage_spelling()
201}
202fn decode<T: serde::de::DeserializeOwned>(value: String) -> rusqlite::Result<T> {
203 serde_json::from_value(Value::String(value)).map_err(|e| {
204 rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(e))
205 })
206}
207fn uuid(row: &rusqlite::Row<'_>, index: usize) -> rusqlite::Result<Uuid> {
208 Uuid::parse_str(&row.get::<_, String>(index)?).map_err(|e| {
209 rusqlite::Error::FromSqlConversionFailure(index, rusqlite::types::Type::Text, Box::new(e))
210 })
211}
212const COLUMNS: &str = concat!(
213 "namespace, logical_message_id, outbound_note_id, kind, slug, credential_ref, ",
214 "recipient_address, protocol_version, sender_agent_id, recipient_agent_id, ",
215 "recipient_device_id, recipient_key_epoch, contact_generation, sender_key_epoch, ",
216 "recipient_key_fingerprint, enc, ciphertext, state, attempt_count, next_retry_at, ",
217 "last_failure_class, hold_reason, receipt, created_at, updated_at, envelope_seq, ",
218 "policy_mode, policy_revision, sender_assurance, admitted_at",
219);
220fn unsigned_column(value: i64, index: usize) -> rusqlite::Result<u64> {
221 u64::try_from(value).map_err(|error| {
222 rusqlite::Error::FromSqlConversionFailure(
223 index,
224 rusqlite::types::Type::Integer,
225 Box::new(error),
226 )
227 })
228}
229fn read_unsigned(row: &rusqlite::Row<'_>, index: usize) -> rusqlite::Result<u64> {
230 unsigned_column(row.get(index)?, index)
231}
232fn read_optional_unsigned(row: &rusqlite::Row<'_>, index: usize) -> rusqlite::Result<Option<u64>> {
233 row.get::<_, Option<i64>>(index)?
234 .map(|value| unsigned_column(value, index))
235 .transpose()
236}
237fn sql_integer(value: u64) -> StorageResult<i64> {
238 i64::try_from(value)
239 .map_err(|_| invalid("unsigned transport value exceeds SQLite integer range"))
240}
241fn read_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<SenderRecord> {
242 Ok(SenderRecord {
243 envelope: SenderEnvelope {
244 namespace: row.get(0)?,
245 logical_message_id: uuid(row, 1)?,
246 outbound_note_id: uuid(row, 2)?,
247 kind: row.get(3)?,
248 slug: row.get(4)?,
249 credential_ref: row.get(5)?,
250 recipient_address: row.get(6)?,
251 protocol_version: row.get(7)?,
252 sender_agent_id: row.get(8)?,
253 sender_assurance: decode(row.get(28)?)?,
254 recipient_agent_id: row.get(9)?,
255 recipient_device_id: uuid(row, 10)?,
256 recipient_key_epoch: read_unsigned(row, 11)?,
257 contact_generation: read_unsigned(row, 12)?,
258 sender_key_epoch: read_unsigned(row, 13)?,
259 recipient_key_fingerprint: row.get(14)?,
260 enc: row.get(15)?,
261 ciphertext: row.get(16)?,
262 },
263 state: decode(row.get(17)?)?,
264 attempt_count: read_unsigned(row, 18)?,
265 next_retry_at: row.get(19)?,
266 last_failure_class: row.get::<_, Option<String>>(20)?.map(decode).transpose()?,
267 hold_reason: match row.get::<_, Option<String>>(21)?.as_deref() {
268 Some("policy_denied") => Some(HoldReason::PolicyDenied {
269 mode: decode(row.get(26)?)?,
270 revision: read_unsigned(row, 27)?,
271 }),
272 reason => reason.map(|r| decode(r.to_owned())).transpose()?,
273 },
274 receipt: row
275 .get::<_, Option<String>>(22)?
276 .map(|s| {
277 serde_json::from_str(&s).map_err(|e| {
278 rusqlite::Error::FromSqlConversionFailure(
279 22,
280 rusqlite::types::Type::Text,
281 Box::new(e),
282 )
283 })
284 })
285 .transpose()?,
286 envelope_seq: read_unsigned(row, 25)?,
287 created_at: row.get(23)?,
288 updated_at: row.get(24)?,
289 admitted_at: row.get(29)?,
290 })
291}
292fn load(conn: &rusqlite::Connection, key: EnvelopeKey) -> rusqlite::Result<Option<SenderRecord>> {
293 let epoch = i64::try_from(key.recipient_key_epoch)
294 .map_err(|error| rusqlite::Error::ToSqlConversionFailure(Box::new(error)))?;
295 conn.query_row(
296 &LOAD_SQL.replace("{COLUMNS}", COLUMNS),
297 params![
298 key.logical_message_id.to_string(),
299 key.recipient_device_id.to_string(),
300 epoch
301 ],
302 read_row,
303 )
304 .optional()
305}
306
307const LOAD_SQL: &str = concat!(
308 "SELECT {COLUMNS} FROM comm_sender_transport WHERE logical_message_id=?1 AND ",
309 "recipient_device_id=?2 AND recipient_key_epoch=?3",
310);
311
312const TERMINAL_SQL: &str = include_str!("../../../../sql/comm-sender-receipt-exists.sql");
313
314const PRIOR_SQL: &str = concat!(
315 "SELECT {COLUMNS} FROM comm_sender_transport WHERE logical_message_id=?1 ORDER BY ",
316 "envelope_seq DESC LIMIT 1",
317);
318
319const DEVICE_EPOCH_SQL: &str =
320 include_str!("../../../../sql/comm-sender-device-key-epoch-max-select.sql");
321
322const INSERT_SQL: &str = include_str!("../../../../sql/comm-sender-envelope-insert.sql");
323
324const PENDING_SQL: &str = concat!(
325 "SELECT {COLUMNS} FROM comm_sender_transport AS t WHERE namespace=?1 AND kind=?2 ",
326 "AND slug=?3 AND state='pending' AND hold_reason IS NULL AND (next_retry_at IS ",
327 "NULL OR next_retry_at<=?4) AND NOT EXISTS(SELECT 1 FROM comm_sender_transport AS ",
328 "done WHERE done.logical_message_id=t.logical_message_id AND done.receipt IS NOT ",
329 "NULL) AND envelope_seq=(SELECT MAX(envelope_seq) FROM comm_sender_transport ",
330 "WHERE logical_message_id=t.logical_message_id) ORDER BY ",
331 "created_at,logical_message_id LIMIT ?5",
332);
333
334const FAILURE_SQL: &str = include_str!("../../../../sql/comm-sender-failure-update.sql");
335
336const ADMISSION_SQL: &str = include_str!("../../../../sql/comm-sender-admission-update.sql");
337
338const HOLD_SQL: &str = include_str!("../../../../sql/comm-sender-hold-update.sql");
339
340const RECEIPT_SQL: &str = include_str!("../../../../sql/comm-sender-receipt-update.sql");
341
342pub struct SenderTransportStore {
344 notes: SqlNoteStore,
345}
346impl SenderTransportStore {
347 pub fn new(pool: Arc<ConnectionPool>) -> Self {
348 Self {
349 notes: SqlNoteStore::new(pool, false),
350 }
351 }
352 pub async fn get(&self, key: EnvelopeKey) -> StorageResult<Option<SenderRecord>> {
353 self.notes
354 .with_reader("sender_transport_get", move |conn| load(conn, key))
355 .await
356 }
357 pub async fn get_by_outbound_note_id(
360 &self,
361 namespace: &str,
362 outbound_note_id: Uuid,
363 ) -> StorageResult<Option<SenderRecord>> {
364 let namespace = namespace.to_owned();
365 self.notes
366 .with_reader("sender_transport_status", move |conn| {
367 conn.query_row(
368 &format!(
369 "SELECT {COLUMNS} FROM comm_sender_transport \
370 WHERE namespace = ?1 AND outbound_note_id = ?2 \
371 ORDER BY (receipt IS NOT NULL) DESC, envelope_seq DESC LIMIT 1"
372 ),
373 params![namespace, outbound_note_id.to_string()],
374 read_row,
375 )
376 .optional()
377 })
378 .await
379 }
380 pub async fn create(
383 &self,
384 envelope: SenderEnvelope,
385 confirmed_key_change: bool,
386 ) -> StorageResult<SenderRecord> {
387 envelope.validate()?;
388 self.notes
389 .with_writer_tx_storage("sender_transport_create", move |conn| {
390 let op = "sender_transport_create";
391 if let Some(existing) = load(conn, envelope.key()).map_err(|e| map_err(e, op))? {
392 if confirmed_key_change {
393 return Err(invalid(
394 "confirmed re-encryption requires a new key identity",
395 ));
396 }
397 if existing.envelope != envelope {
398 return Err(invalid("envelope_conflict"));
399 }
400 return Ok(existing);
401 }
402 let terminal: bool = conn
403 .query_row(
404 TERMINAL_SQL,
405 [envelope.logical_message_id.to_string()],
406 |row| row.get(0),
407 )
408 .map_err(|e| map_err(e, op))?;
409 if terminal {
410 return Err(invalid("logical message already has a recipient receipt"));
411 }
412 let prior = conn
413 .query_row(
414 &PRIOR_SQL.replace("{COLUMNS}", COLUMNS),
415 [envelope.logical_message_id.to_string()],
416 read_row,
417 )
418 .optional()
419 .map_err(|e| map_err(e, op))?;
420 let envelope_seq = if let Some(prior) = prior {
421 if !confirmed_key_change
422 || prior.state != TransportState::Pending
423 || prior.hold_reason != Some(HoldReason::RecipientKeyChanged)
424 {
425 return Err(invalid(
426 "new envelope requires confirmed recipient key change",
427 ));
428 }
429 let previous_epoch: Option<u64> = conn
430 .query_row(
431 DEVICE_EPOCH_SQL,
432 params![
433 envelope.logical_message_id.to_string(),
434 envelope.recipient_device_id.to_string()
435 ],
436 |row| read_optional_unsigned(row, 0),
437 )
438 .map_err(|e| map_err(e, op))?;
439 if previous_epoch.is_some_and(|epoch| envelope.recipient_key_epoch <= epoch) {
440 return Err(invalid("same device key epoch must increase"));
441 }
442 let a = &prior.envelope;
443 let b = &envelope;
444 if a.sender_assurance != b.sender_assurance {
445 return Err(invalid("sender_assurance_conflict"));
446 }
447 if a.namespace != b.namespace
448 || a.outbound_note_id != b.outbound_note_id
449 || a.kind != b.kind
450 || a.slug != b.slug
451 || a.sender_agent_id != b.sender_agent_id
452 || a.sender_key_epoch != b.sender_key_epoch
453 || a.recipient_agent_id != b.recipient_agent_id
454 || a.recipient_address != b.recipient_address
455 {
456 return Err(invalid("logical message identity cannot change"));
457 }
458 prior
459 .envelope_seq
460 .checked_add(1)
461 .filter(|seq| *seq <= i64::MAX as u64)
462 .ok_or_else(|| invalid("envelope sequence exhausted"))?
463 } else {
464 if confirmed_key_change {
465 return Err(invalid("no prior envelope to re-encrypt"));
466 }
467 1
468 };
469 let now = chrono::Utc::now().timestamp_micros();
470 conn.execute(
471 INSERT_SQL,
472 params![
473 envelope.namespace,
474 envelope.logical_message_id.to_string(),
475 envelope.outbound_note_id.to_string(),
476 envelope.kind,
477 envelope.slug,
478 envelope.credential_ref,
479 envelope.recipient_address,
480 envelope.protocol_version,
481 envelope.sender_agent_id,
482 envelope.recipient_agent_id,
483 envelope.recipient_device_id.to_string(),
484 sql_integer(envelope.recipient_key_epoch)?,
485 sql_integer(envelope.contact_generation)?,
486 sql_integer(envelope.sender_key_epoch)?,
487 envelope.recipient_key_fingerprint,
488 envelope.enc,
489 envelope.ciphertext,
490 now,
491 sql_integer(envelope_seq)?,
492 encode(&envelope.sender_assurance)
493 ],
494 )
495 .map_err(|e| map_err(e, op))?;
496 load(conn, envelope.key())
497 .map_err(|e| map_err(e, op))?
498 .ok_or_else(|| invalid("inserted record disappeared"))
499 })
500 .await
501 }
502 pub async fn list_pending(
505 &self,
506 namespace: &str,
507 kind: &str,
508 slug: &str,
509 now: i64,
510 limit: u32,
511 ) -> StorageResult<Vec<SenderRecord>> {
512 let (namespace, kind, slug) = (namespace.to_owned(), kind.to_owned(), slug.to_owned());
513 self.notes
514 .with_reader("sender_transport_pending", move |conn| {
515 let mut stmt = conn.prepare(&PENDING_SQL.replace("{COLUMNS}", COLUMNS))?;
516 let rows = stmt
517 .query_map(
518 params![namespace, kind, slug, now, limit.min(1000)],
519 read_row,
520 )?
521 .collect();
522 rows
523 })
524 .await
525 }
526 pub async fn record_failure(
529 &self,
530 key: EnvelopeKey,
531 class: FailureClass,
532 next_retry_at: Option<i64>,
533 ) -> StorageResult<()> {
534 self.notes
535 .with_writer_tx_storage("sender_transport_failure", move |conn| {
536 let op = "sender_transport_failure";
537 let row = load(conn, key)
538 .map_err(|e| map_err(e, op))?
539 .ok_or_else(|| invalid("unknown sender record"))?;
540 if row.state != TransportState::Pending {
541 return Err(invalid("sender record is not pending"));
542 }
543 if row.hold_reason.is_some() {
544 return Err(invalid("sender record is held"));
545 }
546 let state = if class == FailureClass::Permanent {
547 TransportState::Failed
548 } else {
549 TransportState::Pending
550 };
551 let attempts = if class == FailureClass::Authentication {
552 row.attempt_count
553 } else {
554 row.attempt_count.saturating_add(1).min(i64::MAX as u64)
555 };
556 let retry = if class == FailureClass::Transient {
557 next_retry_at
558 } else {
559 None
560 };
561 conn.execute(
562 FAILURE_SQL,
563 params![
564 key.logical_message_id.to_string(),
565 key.recipient_device_id.to_string(),
566 sql_integer(key.recipient_key_epoch)?,
567 encode(&state),
568 sql_integer(attempts)?,
569 retry,
570 encode(&class),
571 chrono::Utc::now().timestamp_micros()
572 ],
573 )
574 .map_err(|e| map_err(e, op))?;
575 Ok(())
576 })
577 .await
578 }
579 pub async fn record_admission(&self, key: EnvelopeKey, admitted_at: i64) -> StorageResult<()> {
582 self.notes
583 .with_writer_tx_storage("sender_transport_admission", move |conn| {
584 let op = "sender_transport_admission";
585 let row = load(conn, key)
586 .map_err(|e| map_err(e, op))?
587 .ok_or_else(|| invalid("unknown sender record"))?;
588 if row.state != TransportState::Pending {
589 return Err(invalid("sender record is not pending"));
590 }
591 if row.hold_reason.is_some() {
592 return Err(invalid("sender record is held"));
593 }
594 let next_retry_at = admitted_at
595 .checked_add(600_000_000)
596 .ok_or_else(|| invalid("admission deadline exceeds SQLite integer range"))?;
597 conn.execute(
598 ADMISSION_SQL,
599 params![
600 key.logical_message_id.to_string(),
601 key.recipient_device_id.to_string(),
602 sql_integer(key.recipient_key_epoch)?,
603 admitted_at,
604 next_retry_at,
605 chrono::Utc::now().timestamp_micros()
606 ],
607 )
608 .map_err(|e| map_err(e, op))?;
609 Ok(())
610 })
611 .await
612 }
613 pub async fn hold(&self, key: EnvelopeKey, reason: Option<HoldReason>) -> StorageResult<()> {
617 self.notes
618 .with_writer_tx_storage("sender_transport_hold", move |conn| {
619 let op = "sender_transport_hold";
620 let row = load(conn, key)
621 .map_err(|e| map_err(e, op))?
622 .ok_or_else(|| invalid("unknown sender record"))?;
623 if row.state != TransportState::Pending {
624 return Err(invalid("sender record is not pending"));
625 }
626 if row.hold_reason == Some(HoldReason::RecipientKeyChanged)
627 && reason != Some(HoldReason::RecipientKeyChanged)
628 {
629 return Err(invalid("key change requires confirmed re-encryption"));
630 }
631 let (policy_mode, policy_revision) = match reason {
632 Some(HoldReason::PolicyDenied { mode, revision }) => {
633 (Some(encode(&mode)), Some(sql_integer(revision)?))
634 }
635 _ => (None, None),
636 };
637 conn.execute(
638 HOLD_SQL,
639 params![
640 key.logical_message_id.to_string(),
641 key.recipient_device_id.to_string(),
642 sql_integer(key.recipient_key_epoch)?,
643 reason.map(|r| encode(&r)),
644 chrono::Utc::now().timestamp_micros(),
645 policy_mode,
646 policy_revision
647 ],
648 )
649 .map_err(|e| map_err(e, op))?;
650 Ok(())
651 })
652 .await
653 }
654 pub async fn accept_receipt(
657 &self,
658 key: EnvelopeKey,
659 state: TransportState,
660 receipt: Value,
661 ) -> StorageResult<()> {
662 self.notes
663 .with_writer_tx_storage("sender_transport_receipt", move |conn| {
664 let op = "sender_transport_receipt";
665 let row = load(conn, key)
666 .map_err(|e| map_err(e, op))?
667 .ok_or_else(|| invalid("unknown sender record"))?;
668 let disposition = match state {
669 TransportState::RecipientStored => "stored",
670 TransportState::RecipientQuarantined => "quarantined",
671 _ => return Err(invalid("receipt target is not a recipient outcome")),
672 };
673 if receipt.get("disposition").and_then(Value::as_str) != Some(disposition) {
674 return Err(invalid("receipt disposition mismatch"));
675 }
676 let binding = receipt
677 .get("binding")
678 .ok_or_else(|| invalid("missing receipt binding"))?;
679 let e = &row.envelope;
680 let expected = [
681 ("protocol_version", serde_json::json!(e.protocol_version)),
682 (
683 "logical_message_id",
684 serde_json::json!(e.logical_message_id),
685 ),
686 ("sender_agent_id", serde_json::json!(e.sender_agent_id)),
687 (
688 "recipient_agent_id",
689 serde_json::json!(e.recipient_agent_id),
690 ),
691 (
692 "recipient_device_id",
693 serde_json::json!(e.recipient_device_id),
694 ),
695 (
696 "recipient_key_epoch",
697 serde_json::json!(e.recipient_key_epoch),
698 ),
699 (
700 "contact_generation",
701 serde_json::json!(e.contact_generation),
702 ),
703 ];
704 for (field, value) in expected {
705 if binding.get(field) != Some(&value) {
706 return Err(invalid("receipt binding mismatch"));
707 }
708 }
709 let attempt = binding
710 .get("delivery_attempt_id")
711 .and_then(Value::as_str)
712 .ok_or_else(|| invalid("missing delivery attempt"))?;
713 if Uuid::parse_str(attempt)
714 .ok()
715 .map(|id| id.to_string())
716 .as_deref()
717 != Some(attempt)
718 {
719 return Err(invalid("invalid delivery attempt"));
720 }
721 if let Some(accepted) = row.receipt {
722 if accepted != receipt || row.state != state {
723 return Err(invalid("receipt_conflict"));
724 }
725 return Ok(());
726 }
727 conn.execute(
728 RECEIPT_SQL,
729 params![
730 key.logical_message_id.to_string(),
731 key.recipient_device_id.to_string(),
732 sql_integer(key.recipient_key_epoch)?,
733 encode(&state),
734 receipt.to_string(),
735 chrono::Utc::now().timestamp_micros()
736 ],
737 )
738 .map_err(|e| map_err(e, op))?;
739 Ok(())
740 })
741 .await
742 }
743}
744#[cfg(test)]
745mod tests;