1use 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#[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#[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
82pub const LOCAL_QUARANTINE_BOUND: usize = 512;
86
87pub 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 pub parsed_plaintext: Option<Value>,
100}
101#[derive(Clone, Debug)]
104pub struct RecipientCommit {
105 pub note: Option<Note>,
108 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#[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 pub note_id: Option<Uuid>,
131 pub disposition: RecipientDisposition,
132 pub created: bool,
133 pub note: Option<Note>,
135 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
148const 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
193fn 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 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 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 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 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 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 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 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;