Skip to main content

interlink/
delivery.rs

1//! Recoverable host-delivery failures, persisted before releasing a bus message.
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}
20
21pub struct FailedDeliveries {
22    path: PathBuf,
23}
24
25impl FailedDeliveries {
26    pub fn new(path: &Path) -> Self {
27        Self {
28            path: path.to_owned(),
29        }
30    }
31
32    pub fn list(&self) -> Result<Vec<FailedDelivery>> {
33        match std::fs::read(&self.path) {
34            Ok(bytes) => Ok(serde_json::from_slice(&bytes)?),
35            Err(e) if e.kind() == ErrorKind::NotFound => Ok(Vec::new()),
36            Err(e) => Err(e.into()),
37        }
38    }
39
40    pub fn retain(&self, msg_id: &str, sender: &str, text: &str, reason: &str) -> Result<()> {
41        let _lock = lock(&self.path.with_extension("lock"))?;
42        let mut records = self.list()?;
43        if records
44            .iter()
45            .any(|r| r.msg_id == msg_id && r.sender == sender)
46        {
47            return Ok(());
48        }
49        if records.len() >= 64 {
50            bail!("failed-delivery storage is full; read and discard recovered entries");
51        }
52        records.push(FailedDelivery {
53            id: mint_session_id()?,
54            msg_id: msg_id.into(),
55            sender: sender.into(),
56            text: text.into(),
57            reason: reason.into(),
58        });
59        atomic_write(&self.path, &serde_json::to_vec(&records)?)
60    }
61
62    pub fn remove(&self, id: &str) -> Result<()> {
63        let _lock = lock(&self.path.with_extension("lock"))?;
64        let mut records = self.list()?;
65        records.retain(|r| r.id != id);
66        atomic_write(&self.path, &serde_json::to_vec(&records)?)
67    }
68}