1use std::io::ErrorKind;
4use std::path::{Path, PathBuf};
5
6use anyhow::{Result, bail};
7use serde::{Deserialize, Serialize};
8
9use crate::identity::mint_session_id;
10use crate::state::{atomic_write, lock};
11
12#[derive(Clone, Serialize, Deserialize)]
13pub struct FailedDelivery {
14 pub id: String,
15 pub msg_id: String,
16 pub sender: String,
17 pub text: String,
18 pub reason: String,
19 #[serde(default)]
20 pub notice: bool,
21}
22
23pub struct FailedDeliveries {
24 path: PathBuf,
25}
26
27impl FailedDeliveries {
28 pub fn new(path: &Path) -> Self {
29 Self {
30 path: path.to_owned(),
31 }
32 }
33
34 pub fn list(&self) -> Result<Vec<FailedDelivery>> {
35 match std::fs::read(&self.path) {
36 Ok(bytes) => Ok(serde_json::from_slice(&bytes)?),
37 Err(e) if e.kind() == ErrorKind::NotFound => Ok(Vec::new()),
38 Err(e) => Err(e.into()),
39 }
40 }
41
42 pub fn retain(&self, msg_id: &str, sender: &str, text: &str, reason: &str) -> Result<()> {
43 self.retain_entry(msg_id, sender, text, reason, false)
44 }
45
46 pub fn retain_notice(&self, msg_id: &str, text: &str, reason: &str) -> Result<()> {
47 self.retain_entry(msg_id, "Interlink", text, reason, true)
48 }
49
50 pub fn clear_notices(&self) -> Result<()> {
51 let _lock = lock(&self.path.with_extension("lock"))?;
52 let mut records = self.list()?;
53 records.retain(|r| !r.notice);
54 atomic_write(&self.path, &serde_json::to_vec(&records)?)
55 }
56
57 fn retain_entry(
58 &self,
59 msg_id: &str,
60 sender: &str,
61 text: &str,
62 reason: &str,
63 notice: bool,
64 ) -> Result<()> {
65 let _lock = lock(&self.path.with_extension("lock"))?;
66 let mut records = self.list()?;
67 if notice {
70 records.retain(|r| !r.notice);
71 }
72 if records
73 .iter()
74 .any(|r| r.msg_id == msg_id && r.sender == sender)
75 {
76 return Ok(());
77 }
78 if records.len() >= 64 {
79 bail!("failed-delivery storage is full; read and discard recovered entries");
80 }
81 records.push(FailedDelivery {
82 id: mint_session_id()?,
83 msg_id: msg_id.into(),
84 sender: sender.into(),
85 text: text.into(),
86 reason: reason.into(),
87 notice,
88 });
89 atomic_write(&self.path, &serde_json::to_vec(&records)?)
90 }
91
92 pub fn remove(&self, id: &str) -> Result<()> {
93 let _lock = lock(&self.path.with_extension("lock"))?;
94 let mut records = self.list()?;
95 records.retain(|r| r.id != id);
96 atomic_write(&self.path, &serde_json::to_vec(&records)?)
97 }
98}
99
100#[cfg(test)]
101mod tests {
102 use super::*;
103
104 #[test]
105 fn notice_retries_use_one_slot_and_preserve_legacy_messages() {
106 let dir = tempfile::tempdir().unwrap();
107 let store = FailedDeliveries::new(&dir.path().join("failed.json"));
108 store
109 .retain("legacy", "peer", "full body", "offline")
110 .unwrap();
111 for index in 0..70 {
112 store
113 .retain_notice(&index.to_string(), "fetch mailbox", "offline")
114 .unwrap();
115 }
116 let records = store.list().unwrap();
117 assert_eq!(records.len(), 2);
118 assert_eq!(records[0].text, "full body");
119 assert_eq!(records[1].msg_id, "69");
120 store.clear_notices().unwrap();
121 assert_eq!(store.list().unwrap().len(), 1);
122 }
123}