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