Skip to main content

khive_db/stores/note/transport/
mod.rs

1//! Durable sender-side records. No transport adapter performs SQL.
2use 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/// Immutable submission data. Credential references name a key-facility entry; never supply key
19/// material.
20#[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/// Sender identity assurance recorded at send time and preserved across transport attempts.
85#[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    /// Policy state evaluated at the refused transport attempt.
114    PolicyDenied {
115        mode: PolicyMode,
116        revision: u64,
117    },
118}
119
120/// Recorded evaluation mode; this does not enable or change runtime policy.
121#[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
342/// Store uses the note writer's transaction/queue routing and bounded pooled readers.
343pub 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    /// Read the caller's outbound transport record. A recipient receipt wins over
358    /// a newer envelope; without a receipt, use the latest local envelope sequence.
359    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    /// Exact retries reuse the record. Only confirmed key-change operations
381    /// may create another envelope.
382    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    /// Only due, unheld rows of the exact route and highest local envelope
503    /// sequence are submitted automatically.
504    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    /// Record a failure only for an unheld pending envelope. A hold suspends failure updates as
527    /// well as automatic retry until the hold is resolved.
528    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    /// Record a successful service admission and schedule resubmission 600 seconds later.
580    /// Re-admission of the same pending envelope replaces both timestamps.
581    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    /// Release credit or policy holds explicitly. Key-change holds release by creating a confirmed
614    /// new envelope. Policy holds require the evaluated mode and revision in the reason;
615    /// every other reason, including release (`None`), clears both policy columns.
616    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    /// Caller verifies the signature before reaching this store. Binding checks repeat inside
655    /// the transaction.
656    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;