Skip to main content

interlink/
delivery.rs

1//! Recoverable notice diagnostics and legacy full-message delivery failures.
2
3use 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        // Repeated wake-up failures need one diagnostic slot, not an unbounded
68        // sequence of stale notices. Legacy full messages must remain recoverable.
69        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}