Skip to main content

khive_db/stores/note/recipient/
mod.rs

1//! Atomic recipient message, replay, quarantine and acknowledgement persistence.
2use super::{assign_note_seq, map_err, SqlNoteStore};
3use crate::pool::ConnectionPool;
4use khive_storage::{Note, 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, Serialize, Deserialize)]
12#[serde(rename_all = "snake_case")]
13pub enum RecipientDisposition {
14    Stored,
15    Quarantined,
16}
17
18/// Closed class of terminal acknowledgement refusals from ADR-105 A.9.
19#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
20#[serde(rename_all = "snake_case")]
21pub enum AcknowledgementRetirementReason {
22    PermanentTransport,
23}
24impl AcknowledgementRetirementReason {
25    fn as_str(self) -> &'static str {
26        match self {
27            Self::PermanentTransport => "permanent_transport",
28        }
29    }
30}
31
32/// Unsigned binding and disposition read from the durable acknowledgement journal.
33#[derive(Clone, Debug, PartialEq, Eq)]
34pub struct AckJournalEntry {
35    pub delivery_attempt_id: Uuid,
36    pub binding: Value,
37    pub disposition: RecipientDisposition,
38    pub attempt_count: u64,
39    pub not_before: Option<i64>,
40    pub created_at: i64,
41    pub updated_at: i64,
42}
43impl RecipientDisposition {
44    fn as_str(self) -> &'static str {
45        match self {
46            Self::Stored => "stored",
47            Self::Quarantined => "quarantined",
48        }
49    }
50    fn parse(s: &str) -> StorageResult<Self> {
51        match s {
52            "stored" => Ok(Self::Stored),
53            "quarantined" => Ok(Self::Quarantined),
54            _ => Err(invalid("invalid stored disposition")),
55        }
56    }
57}
58#[derive(Clone, Copy, Debug, PartialEq, Eq)]
59pub enum QuarantineReason {
60    InvalidPlaintext,
61    InvalidMessage,
62    PolicyRejected,
63}
64impl QuarantineReason {
65    fn as_str(self) -> &'static str {
66        match self {
67            Self::InvalidPlaintext => "invalid_plaintext",
68            Self::InvalidMessage => "invalid_message",
69            Self::PolicyRejected => "policy_rejected",
70        }
71    }
72    fn parse(s: &str) -> StorageResult<Self> {
73        match s {
74            "invalid_plaintext" => Ok(Self::InvalidPlaintext),
75            "invalid_message" => Ok(Self::InvalidMessage),
76            "policy_rejected" => Ok(Self::PolicyRejected),
77            _ => Err(invalid("invalid stored quarantine reason")),
78        }
79    }
80}
81
82/// ADR-105 A.8 Receiving, step 3: the local quarantine bound this client sets,
83/// per local recipient agent, for quarantined items other than policy-refused
84/// ones. Past it the oldest such item is dropped and reported.
85pub const LOCAL_QUARANTINE_BOUND: usize = 512;
86
87/// ADR-105 A.8 Receiving, step 4: policy-refused items are not under the local
88/// quarantine bound. Each sender has its own bound, past which the oldest of
89/// that sender's policy-refused items is dropped and reported, so one sender's
90/// refusals never evict another's.
91pub const POLICY_REFUSED_PER_SENDER_BOUND: usize = 64;
92
93#[derive(Clone, Debug)]
94pub struct QuarantineRecord {
95    pub reason: QuarantineReason,
96    pub delivery_item: Vec<u8>,
97    /// The parsed plaintext object, present exactly for a policy refusal
98    /// (ADR-105 A.8 Receiving, step 4).
99    pub parsed_plaintext: Option<Value>,
100}
101/// Internal storage contract. Runtime supplies authenticated identity and a
102/// validated note; this structure is not a wire-ingest parameter.
103#[derive(Clone, Debug)]
104pub struct RecipientCommit {
105    /// The message note, present exactly for a stored disposition. A
106    /// quarantined delivery writes no message note.
107    pub note: Option<Note>,
108    /// Local recipient actor from the trusted enrollment binding.
109    pub recipient_actor: String,
110    pub binding: Value,
111    pub sender_agent_id: String,
112    pub logical_message_id: Uuid,
113    pub delivery_attempt_id: Uuid,
114    pub disposition: RecipientDisposition,
115    pub quarantine: Option<QuarantineRecord>,
116    pub in_reply_to: Option<Uuid>,
117    pub correlation: Option<String>,
118}
119/// A quarantined item dropped to keep a quarantine bound. Its replay identity
120/// stays claimed.
121#[derive(Clone, Debug, PartialEq, Eq)]
122pub struct EvictedQuarantine {
123    pub sender_agent_id: String,
124    pub logical_message_id: Uuid,
125    pub reason: QuarantineReason,
126}
127#[derive(Debug)]
128pub struct RecipientCommitResult {
129    /// The message note of a stored identity; a quarantined identity has none.
130    pub note_id: Option<Uuid>,
131    pub disposition: RecipientDisposition,
132    pub created: bool,
133    /// Present only for a newly committed note, for best-effort indexing.
134    pub note: Option<Note>,
135    /// Quarantined items this commit dropped to keep a quarantine bound.
136    pub evicted: Vec<EvictedQuarantine>,
137}
138fn invalid(message: &str) -> StorageError {
139    StorageError::InvalidInput {
140        capability: StorageCapability::Notes,
141        operation: "recipient_transport".into(),
142        message: message.into(),
143    }
144}
145
146const REPLAY_SQL: &str = include_str!("../../../../sql/comm-recipient-replay-select.sql");
147
148// An outbox transport row alone survives deletion of its message note. A
149// parent proves reply status only while that exact outbound note is live and
150// belongs to this recipient, addressed to the authenticated sender.
151const OUTBOUND_PARENT_SQL: &str =
152    include_str!("../../../../sql/comm-recipient-live-outbound-parent-exists.sql");
153
154const CORRELATION_SQL: &str =
155    include_str!("../../../../sql/comm-recipient-correlated-thread-select.sql");
156
157fn correlation_match_values(correlation: &str) -> ([Option<String>; 9], String) {
158    let raw = correlation.trim();
159    let mut spellings = std::array::from_fn(|_| None);
160    spellings[0] = Some(raw.to_owned());
161    let Ok(root) = Uuid::parse_str(raw) else {
162        return (spellings, correlation.to_owned());
163    };
164    spellings[1] = Some(root.as_hyphenated().to_string());
165    spellings[2] = Some(root.simple().to_string());
166    spellings[3] = Some(root.braced().to_string());
167    spellings[4] = Some(root.urn().to_string());
168    spellings[5] = Some(format!("{:X}", root.as_hyphenated()));
169    spellings[6] = Some(format!("{:X}", root.simple()));
170    spellings[7] = Some(format!("{:X}", root.braced()));
171    spellings[8] = Some(format!("{:X}", root.urn()));
172    (spellings, root.to_string())
173}
174
175const INSERT_NOTE_SQL: &str =
176    include_str!("../../../../sql/comm-recipient-message-note-insert.sql");
177
178const INSERT_REPLAY_SQL: &str = include_str!("../../../../sql/comm-recipient-replay-insert.sql");
179
180const QUARANTINE_SQL: &str = include_str!("../../../../sql/comm-recipient-quarantine-insert.sql");
181
182const LOCAL_BOUND_COUNT_SQL: &str =
183    include_str!("../../../../sql/comm-recipient-local-quarantine-count.sql");
184const LOCAL_BOUND_OLDEST_SQL: &str =
185    include_str!("../../../../sql/comm-recipient-local-quarantine-oldest-select.sql");
186const SENDER_BOUND_COUNT_SQL: &str =
187    include_str!("../../../../sql/comm-recipient-sender-quarantine-count.sql");
188const SENDER_BOUND_OLDEST_SQL: &str =
189    include_str!("../../../../sql/comm-recipient-sender-quarantine-oldest-select.sql");
190const DELETE_QUARANTINE_SQL: &str =
191    include_str!("../../../../sql/comm-recipient-quarantine-delete.sql");
192
193/// Make room for one more quarantined item of `reason` under its bound by
194/// dropping the oldest items in the same bound, oldest `created_at` first,
195/// then by key. Only the quarantine rows go: every replay identity stays claimed.
196fn evict_to_bound(
197    conn: &rusqlite::Connection,
198    recipient_agent_id: &str,
199    sender_agent_id: &str,
200    reason: QuarantineReason,
201    op: &'static str,
202) -> StorageResult<Vec<EvictedQuarantine>> {
203    let row = |r: &rusqlite::Row<'_>| {
204        Ok((
205            r.get::<_, String>(0)?,
206            r.get::<_, String>(1)?,
207            r.get::<_, String>(2)?,
208        ))
209    };
210    let oldest: Vec<(String, String, String)> = if reason == QuarantineReason::PolicyRejected {
211        let held: i64 = conn
212            .query_row(
213                SENDER_BOUND_COUNT_SQL,
214                params![recipient_agent_id, sender_agent_id],
215                |r| r.get(0),
216            )
217            .map_err(|e| map_err(e, op))?;
218        let excess = held + 1 - POLICY_REFUSED_PER_SENDER_BOUND as i64;
219        if excess <= 0 {
220            return Ok(Vec::new());
221        }
222        let mut stmt = conn
223            .prepare(SENDER_BOUND_OLDEST_SQL)
224            .map_err(|e| map_err(e, op))?;
225        let rows = stmt
226            .query_map(params![recipient_agent_id, sender_agent_id, excess], row)
227            .map_err(|e| map_err(e, op))?;
228        rows.collect::<Result<_, _>>().map_err(|e| map_err(e, op))?
229    } else {
230        let held: i64 = conn
231            .query_row(LOCAL_BOUND_COUNT_SQL, params![recipient_agent_id], |r| {
232                r.get(0)
233            })
234            .map_err(|e| map_err(e, op))?;
235        let excess = held + 1 - LOCAL_QUARANTINE_BOUND as i64;
236        if excess <= 0 {
237            return Ok(Vec::new());
238        }
239        let mut stmt = conn
240            .prepare(LOCAL_BOUND_OLDEST_SQL)
241            .map_err(|e| map_err(e, op))?;
242        let rows = stmt
243            .query_map(params![recipient_agent_id, excess], row)
244            .map_err(|e| map_err(e, op))?;
245        rows.collect::<Result<_, _>>().map_err(|e| map_err(e, op))?
246    };
247    oldest
248        .into_iter()
249        .map(|(sender, logical, reason)| {
250            conn.execute(DELETE_QUARANTINE_SQL, params![sender, logical])
251                .map_err(|e| map_err(e, op))?;
252            Ok(EvictedQuarantine {
253                sender_agent_id: sender,
254                logical_message_id: Uuid::parse_str(&logical)
255                    .map_err(|_| invalid("invalid quarantined logical message id"))?,
256                reason: QuarantineReason::parse(&reason)?,
257            })
258        })
259        .collect()
260}
261
262const ACK_LOOKUP_SQL: &str = include_str!("../../../../sql/comm-ack-binding-select.sql");
263
264const ACK_DUE_SQL: &str = include_str!("../../../../sql/comm-ack-due-select.sql");
265
266const ACK_FINISH_SQL: &str = include_str!("../../../../sql/comm-ack-finish-update.sql");
267
268const ACK_FAILED_TRY_SQL: &str = include_str!("../../../../sql/comm-ack-failed-try-update.sql");
269
270const ACK_RETIRE_SQL: &str = include_str!("../../../../sql/comm-ack-retire-update.sql");
271
272fn read_ack_entry(row: &rusqlite::Row<'_>) -> rusqlite::Result<AckJournalEntry> {
273    let attempt: String = row.get(0)?;
274    let delivery_attempt_id = Uuid::parse_str(&attempt).map_err(|error| {
275        rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(error))
276    })?;
277    let binding: String = row.get(1)?;
278    let binding = serde_json::from_str(&binding).map_err(|error| {
279        rusqlite::Error::FromSqlConversionFailure(1, rusqlite::types::Type::Text, Box::new(error))
280    })?;
281    let disposition: String = row.get(2)?;
282    let disposition = serde_json::from_value(Value::String(disposition)).map_err(|error| {
283        rusqlite::Error::FromSqlConversionFailure(2, rusqlite::types::Type::Text, Box::new(error))
284    })?;
285    let attempt_count = u64::try_from(row.get::<_, i64>(3)?).map_err(|error| {
286        rusqlite::Error::FromSqlConversionFailure(
287            3,
288            rusqlite::types::Type::Integer,
289            Box::new(error),
290        )
291    })?;
292    Ok(AckJournalEntry {
293        delivery_attempt_id,
294        binding,
295        disposition,
296        attempt_count,
297        not_before: row.get(4)?,
298        created_at: row.get(5)?,
299        updated_at: row.get(6)?,
300    })
301}
302
303const INSERT_ACK_SQL: &str = include_str!("../../../../sql/comm-ack-insert.sql");
304
305struct AckIdentity<'a> {
306    binding: &'a Value,
307    sender_agent_id: &'a str,
308    logical_message_id: Uuid,
309    delivery_attempt_id: Uuid,
310}
311
312fn record_ack(
313    conn: &rusqlite::Connection,
314    identity: AckIdentity<'_>,
315    disposition: RecipientDisposition,
316    now: i64,
317    op: &'static str,
318) -> StorageResult<()> {
319    let prior_ack: Option<(String, String)> = conn
320        .query_row(
321            ACK_LOOKUP_SQL,
322            [identity.delivery_attempt_id.to_string()],
323            |r| Ok((r.get(0)?, r.get(1)?)),
324        )
325        .optional()
326        .map_err(|error| map_err(error, op))?;
327    if let Some((prior_binding, prior_disposition)) = prior_ack {
328        let prior_binding: Value = serde_json::from_str(&prior_binding)
329            .map_err(|_| invalid("invalid stored ack binding"))?;
330        if prior_binding != *identity.binding || prior_disposition != disposition.as_str() {
331            return Err(invalid("ack attempt binding conflict"));
332        }
333    } else {
334        conn.execute(
335            INSERT_ACK_SQL,
336            params![
337                identity.delivery_attempt_id.to_string(),
338                identity.sender_agent_id,
339                identity.logical_message_id.to_string(),
340                identity.binding.to_string(),
341                disposition.as_str(),
342                now
343            ],
344        )
345        .map_err(|error| map_err(error, op))?;
346    }
347    Ok(())
348}
349
350pub struct RecipientTransportStore {
351    notes: SqlNoteStore,
352}
353impl RecipientTransportStore {
354    pub fn new(pool: Arc<ConnectionPool>) -> Self {
355        Self {
356            notes: SqlNoteStore::new(pool, false),
357        }
358    }
359
360    /// Pending entries at or before `now` (UTC microseconds), oldest first.
361    /// Pages are bounded to 1,000 entries; a zero limit returns no entries.
362    pub async fn list_due_acknowledgements(
363        &self,
364        now: i64,
365        limit: usize,
366    ) -> StorageResult<Vec<AckJournalEntry>> {
367        let limit = limit.min(1000) as i64;
368        self.notes
369            .with_reader("recipient_acknowledgement_due", move |conn| {
370                let mut statement = conn.prepare(ACK_DUE_SQL)?;
371                let entries = statement
372                    .query_map(params![now, limit], read_ack_entry)?
373                    .collect();
374                entries
375            })
376            .await
377    }
378
379    /// Finish a pending entry; terminal or absent entries return `false` unchanged.
380    pub async fn finish_acknowledgement(&self, delivery_attempt_id: Uuid) -> StorageResult<bool> {
381        self.notes
382            .with_writer_tx_storage("recipient_acknowledgement_finish", move |conn| {
383                let changed = conn
384                    .execute(
385                        ACK_FINISH_SQL,
386                        params![
387                            delivery_attempt_id.to_string(),
388                            chrono::Utc::now().timestamp_micros()
389                        ],
390                    )
391                    .map_err(|error| map_err(error, "recipient_acknowledgement_finish"))?;
392                Ok(changed != 0)
393            })
394            .await
395    }
396
397    /// Persist one failed try and its next eligible UTC-microsecond deadline.
398    /// The integer constraint aborts a counter overflow without changing the row.
399    pub async fn record_acknowledgement_failed_try(
400        &self,
401        delivery_attempt_id: Uuid,
402        not_before: i64,
403    ) -> StorageResult<bool> {
404        self.notes
405            .with_writer_tx_storage("recipient_acknowledgement_failed_try", move |conn| {
406                let changed = conn
407                    .execute(
408                        ACK_FAILED_TRY_SQL,
409                        params![
410                            delivery_attempt_id.to_string(),
411                            not_before,
412                            chrono::Utc::now().timestamp_micros()
413                        ],
414                    )
415                    .map_err(|error| map_err(error, "recipient_acknowledgement_failed_try"))?;
416                Ok(changed != 0)
417            })
418            .await
419    }
420
421    /// Retire a pending acknowledgement without touching its message or replay claim.
422    pub async fn retire_acknowledgement(
423        &self,
424        delivery_attempt_id: Uuid,
425        reason: AcknowledgementRetirementReason,
426    ) -> StorageResult<bool> {
427        self.notes
428            .with_writer_tx_storage("recipient_acknowledgement_retire", move |conn| {
429                let changed = conn
430                    .execute(
431                        ACK_RETIRE_SQL,
432                        params![
433                            delivery_attempt_id.to_string(),
434                            reason.as_str(),
435                            chrono::Utc::now().timestamp_micros()
436                        ],
437                    )
438                    .map_err(|error| map_err(error, "recipient_acknowledgement_retire"))?;
439                Ok(changed != 0)
440            })
441            .await
442    }
443    /// Answer an already committed logical message before its plaintext is
444    /// interpreted again. The authenticated caller supplies the current local
445    /// actor and binding; a hit writes only the new attempt's ack journal row.
446    /// A miss makes no claim: `commit` rechecks replay under its own writer
447    /// transaction after the new delivery has been validated.
448    pub async fn ack_if_replayed(
449        &self,
450        binding: Value,
451        recipient_actor: &str,
452    ) -> StorageResult<Option<RecipientCommitResult>> {
453        let (sender_agent_id, recipient_agent_id, logical_message_id, delivery_attempt_id) = {
454            let string_field = |name: &str| {
455                binding
456                    .get(name)
457                    .and_then(Value::as_str)
458                    .filter(|value| !value.is_empty())
459                    .map(str::to_owned)
460                    .ok_or_else(|| invalid("invalid replay receipt binding"))
461            };
462            let logical_message_id = Uuid::parse_str(&string_field("logical_message_id")?)
463                .map_err(|_| invalid("invalid replay logical message id"))?;
464            let delivery_attempt_id = Uuid::parse_str(&string_field("delivery_attempt_id")?)
465                .map_err(|_| invalid("invalid replay delivery attempt id"))?;
466            (
467                string_field("sender_agent_id")?,
468                string_field("recipient_agent_id")?,
469                logical_message_id,
470                delivery_attempt_id,
471            )
472        };
473        let recipient_actor = recipient_actor.to_owned();
474        self.notes
475            .with_writer_tx_storage("recipient_transport_replay", move |conn| {
476                let op = "recipient_transport_replay";
477                let prior: Option<(Option<String>, String, String, String)> = conn
478                    .query_row(
479                        REPLAY_SQL,
480                        params![sender_agent_id, logical_message_id.to_string()],
481                        |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
482                    )
483                    .optional()
484                    .map_err(|error| map_err(error, op))?;
485                let Some((id, disposition, prior_recipient, prior_actor)) = prior else {
486                    return Ok(None);
487                };
488                if prior_recipient != recipient_agent_id || prior_actor != recipient_actor {
489                    return Err(invalid("replay recipient binding changed"));
490                }
491                let note_id = id
492                    .map(|id| Uuid::parse_str(&id).map_err(|_| invalid("invalid replay note id")))
493                    .transpose()?;
494                let disposition = RecipientDisposition::parse(&disposition)?;
495                record_ack(
496                    conn,
497                    AckIdentity {
498                        binding: &binding,
499                        sender_agent_id: &sender_agent_id,
500                        logical_message_id,
501                        delivery_attempt_id,
502                    },
503                    disposition,
504                    chrono::Utc::now().timestamp_micros(),
505                    op,
506                )?;
507                Ok(Some(RecipientCommitResult {
508                    note_id,
509                    disposition,
510                    created: false,
511                    note: None,
512                    evicted: Vec::new(),
513                }))
514            })
515            .await
516    }
517
518    /// Commit all recipient effects together, or none. A replay keeps its first
519    /// disposition and creates only the acknowledgement for a new attempt.
520    ///
521    /// A stored disposition writes the message note. A quarantined one writes
522    /// the quarantine record and no message note (ADR-105 A.8 Receiving, step
523    /// 4), after dropping the oldest item of its quarantine bound when that
524    /// bound is full. Dropped items are returned in `evicted` and logged; their
525    /// replay identities stay claimed.
526    pub async fn commit(&self, mut input: RecipientCommit) -> StorageResult<RecipientCommitResult> {
527        match input.disposition {
528            RecipientDisposition::Quarantined => {
529                if input.quarantine.is_none() {
530                    return Err(invalid("quarantine requires replay bytes"));
531                }
532                if input.note.is_some() {
533                    return Err(invalid("quarantined delivery cannot carry a message note"));
534                }
535            }
536            RecipientDisposition::Stored => {
537                if input.quarantine.is_some() {
538                    return Err(invalid("stored message cannot carry quarantine bytes"));
539                }
540                if input.note.is_none() {
541                    return Err(invalid("stored message requires a message note"));
542                }
543            }
544        }
545        if let Some(q) = &input.quarantine {
546            if q.delivery_item.len() > 98_304
547                || !serde_json::from_slice::<Value>(&q.delivery_item).is_ok_and(|v| v.is_object())
548            {
549                return Err(invalid(
550                    "quarantine delivery item must be one JSON object within 98304 bytes",
551                ));
552            }
553            match (q.reason, &q.parsed_plaintext) {
554                (QuarantineReason::PolicyRejected, Some(plaintext)) if plaintext.is_object() => {}
555                (QuarantineReason::PolicyRejected, _) => {
556                    return Err(invalid(
557                        "policy-refused quarantine requires its parsed plaintext object",
558                    ));
559                }
560                (_, Some(_)) => {
561                    return Err(invalid(
562                        "only a policy-refused quarantine keeps parsed plaintext",
563                    ));
564                }
565                (_, None) => {}
566            }
567        }
568        if input.binding.get("sender_agent_id").and_then(Value::as_str)
569            != Some(input.sender_agent_id.as_str())
570            || input.binding.get("logical_message_id")
571                != Some(&serde_json::json!(input.logical_message_id))
572            || input.binding.get("delivery_attempt_id")
573                != Some(&serde_json::json!(input.delivery_attempt_id))
574        {
575            return Err(invalid("receipt key mismatch"));
576        }
577        let recipient = input
578            .binding
579            .get("recipient_agent_id")
580            .and_then(Value::as_str)
581            .ok_or_else(|| invalid("missing recipient agent"))?
582            .to_owned();
583        if input.recipient_actor.trim().is_empty() {
584            return Err(invalid("missing recipient actor"));
585        }
586        let actor = input.recipient_actor.clone();
587        if let Some(note) = &input.note {
588            if note
589                .properties
590                .as_ref()
591                .and_then(|p| p.get("to_actor"))
592                .and_then(Value::as_str)
593                != Some(actor.as_str())
594            {
595                return Err(invalid(
596                    "message note is not addressed to the recipient actor",
597                ));
598            }
599        }
600        let correlated_thread = if let (Some(correlation), Some(note)) =
601            (input.correlation.as_deref(), input.note.as_ref())
602        {
603            if let Some(sender) = note
604                .properties
605                .as_ref()
606                .and_then(|p| p.get("from_actor"))
607                .and_then(Value::as_str)
608            {
609                let (spellings, id) = correlation_match_values(correlation);
610                let namespace = note.namespace.clone();
611                let sender = sender.to_owned();
612                let actor_for_query = actor.clone();
613                let correlation = correlation.to_owned();
614                let matched: Option<(String, Option<String>)> = self
615                    .notes
616                    .with_reader("recipient_transport_correlation", move |conn| {
617                        conn.query_row(
618                            CORRELATION_SQL,
619                            params![
620                                namespace,
621                                sender,
622                                actor_for_query,
623                                correlation,
624                                &spellings[0],
625                                &spellings[1],
626                                &spellings[2],
627                                &spellings[3],
628                                &spellings[4],
629                                &spellings[5],
630                                &spellings[6],
631                                &spellings[7],
632                                &spellings[8],
633                                id
634                            ],
635                            |r| Ok((r.get(0)?, r.get(1)?)),
636                        )
637                        .optional()
638                    })
639                    .await?;
640                matched
641                    .map(|(id, thread)| {
642                        let matched_id = Uuid::parse_str(&id)
643                            .map_err(|_| invalid("invalid correlated note id"))?;
644                        Ok(thread
645                            .and_then(|s| Uuid::parse_str(&s).ok())
646                            .unwrap_or(matched_id))
647                    })
648                    .transpose()?
649            } else {
650                None
651            }
652        } else {
653            None
654        };
655        let result = self
656            .notes
657            .with_writer_tx_storage("recipient_transport_commit", move |conn| {
658                let op = "recipient_transport_commit";
659                let prior: Option<(Option<String>, String, String, String)> = conn
660                    .query_row(
661                        REPLAY_SQL,
662                        params![input.sender_agent_id, input.logical_message_id.to_string()],
663                        |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)),
664                    )
665                    .optional()
666                    .map_err(|e| map_err(e, op))?;
667                let now = chrono::Utc::now().timestamp_micros();
668                let mut evicted = Vec::new();
669                let (note_id, disposition, created) =
670                    if let Some((id, disposition, prior_recipient, prior_actor)) = prior {
671                        if prior_recipient != recipient || prior_actor != actor {
672                            return Err(invalid("replay recipient binding changed"));
673                        }
674                        (
675                            id.map(|id| {
676                                Uuid::parse_str(&id).map_err(|_| invalid("invalid replay note id"))
677                            })
678                            .transpose()?,
679                            RecipientDisposition::parse(&disposition)?,
680                            false,
681                        )
682                    } else {
683                        let note_id = if let Some(note) = input.note.as_mut() {
684                            let props = note
685                                .properties
686                                .as_mut()
687                                .and_then(Value::as_object_mut)
688                                .ok_or_else(|| invalid("missing message properties"))?;
689                            props
690                                .get("from_actor")
691                                .and_then(Value::as_str)
692                                .ok_or_else(|| invalid("missing sender actor"))?;
693                            if let Some(parent) = input.in_reply_to {
694                                let is_reply: bool = conn
695                                    .query_row(
696                                        OUTBOUND_PARENT_SQL,
697                                        params![
698                                            note.namespace.as_str(),
699                                            parent.to_string(),
700                                            recipient.as_str(),
701                                            input.sender_agent_id.as_str()
702                                        ],
703                                        |row| row.get(0),
704                                    )
705                                    .map_err(|e| map_err(e, op))?;
706                                if is_reply {
707                                    props.insert(
708                                        "message_kind".into(),
709                                        Value::String("reply".into()),
710                                    );
711                                }
712                            }
713                            let thread = correlated_thread.unwrap_or(note.id);
714                            props.insert("thread_id".into(), serde_json::json!(thread));
715                            let n = &*note;
716                            conn.execute(
717                                INSERT_NOTE_SQL,
718                                params![
719                                    n.id.to_string(),
720                                    n.namespace,
721                                    n.kind,
722                                    n.status,
723                                    n.name,
724                                    n.content,
725                                    n.salience,
726                                    n.decay_factor,
727                                    n.expires_at,
728                                    n.properties.as_ref().map(Value::to_string),
729                                    n.created_at,
730                                    n.updated_at,
731                                    n.deleted_at,
732                                    n.key
733                                ],
734                            )
735                            .map_err(|e| map_err(e, op))?;
736                            assign_note_seq(conn, &n.id.to_string()).map_err(|e| map_err(e, op))?;
737                            Some(n.id)
738                        } else {
739                            None
740                        };
741                        conn.execute(
742                            INSERT_REPLAY_SQL,
743                            params![
744                                input.sender_agent_id,
745                                input.logical_message_id.to_string(),
746                                recipient,
747                                actor,
748                                note_id.map(|id| id.to_string()),
749                                input.disposition.as_str(),
750                                now
751                            ],
752                        )
753                        .map_err(|e| map_err(e, op))?;
754                        if let Some(q) = &input.quarantine {
755                            evicted = evict_to_bound(
756                                conn,
757                                &recipient,
758                                &input.sender_agent_id,
759                                q.reason,
760                                op,
761                            )?;
762                            conn.execute(
763                                QUARANTINE_SQL,
764                                params![
765                                    input.sender_agent_id,
766                                    input.logical_message_id.to_string(),
767                                    recipient,
768                                    q.delivery_item,
769                                    q.reason.as_str(),
770                                    q.parsed_plaintext.as_ref().map(Value::to_string),
771                                    now
772                                ],
773                            )
774                            .map_err(|e| map_err(e, op))?;
775                        }
776                        (note_id, input.disposition, true)
777                    };
778                // This insertion deliberately shares the message/replay transaction.
779                record_ack(
780                    conn,
781                    AckIdentity {
782                        binding: &input.binding,
783                        sender_agent_id: &input.sender_agent_id,
784                        logical_message_id: input.logical_message_id,
785                        delivery_attempt_id: input.delivery_attempt_id,
786                    },
787                    disposition,
788                    now,
789                    op,
790                )?;
791                Ok(RecipientCommitResult {
792                    note_id,
793                    disposition,
794                    created,
795                    note: if created { input.note } else { None },
796                    evicted,
797                })
798            })
799            .await?;
800        for item in &result.evicted {
801            tracing::warn!(
802                sender_agent_id = %item.sender_agent_id,
803                logical_message_id = %item.logical_message_id,
804                reason = item.reason.as_str(),
805                "quarantine bound reached; dropped the oldest quarantined item"
806            );
807        }
808        Ok(result)
809    }
810}
811#[cfg(test)]
812mod tests;