Skip to main content

qefro_backend_sdk/
event_outbox.rs

1//! Durable local outbox for Business Events (`ctx.emit()` = persist, then deliver).
2
3use crate::business_events::EmittedBusinessEvent;
4use serde_json::Value;
5use sha2::{Digest, Sha256};
6use std::fs;
7use std::io::Write;
8use std::path::{Path, PathBuf};
9use std::time::SystemTime;
10
11const MAX_FILES: usize = 200;
12
13pub fn default_outbox_dir(signing_secret: &str) -> PathBuf {
14    let digest = Sha256::digest(signing_secret.as_bytes());
15    let shard = hex::encode(&digest[..6]);
16    std::env::temp_dir().join("qefro-event-outbox").join(shard)
17}
18
19fn file_name_for(event_id: &str) -> String {
20    let safe: String = event_id
21        .chars()
22        .map(|c| {
23            if c.is_ascii_alphanumeric() || c == '-' || c == '_' || c == '.' {
24                c
25            } else {
26                '_'
27            }
28        })
29        .take(180)
30        .collect();
31    format!("{safe}.json")
32}
33
34#[derive(Clone)]
35pub struct EventOutbox {
36    dir: PathBuf,
37}
38
39impl EventOutbox {
40    pub fn new(dir: impl AsRef<Path>) -> Self {
41        let dir = dir.as_ref().to_path_buf();
42        let _ = fs::create_dir_all(&dir);
43        Self { dir }
44    }
45
46    pub fn put(&self, event: EmittedBusinessEvent) -> anyhow::Result<EmittedBusinessEvent> {
47        let event_id = event
48            .event_id
49            .as_deref()
50            .map(str::trim)
51            .filter(|s| !s.is_empty())
52            .ok_or_else(|| anyhow::anyhow!("outbox.put requires event_id"))?;
53        fs::create_dir_all(&self.dir)?;
54        let dest = self.dir.join(file_name_for(event_id));
55        let tmp = dest.with_extension("json.tmp");
56        let mut body = serde_json::to_value(&event)?;
57        if let Some(obj) = body.as_object_mut() {
58            obj.insert(
59                "queued_at".into(),
60                Value::String(chrono::Utc::now().to_rfc3339()),
61            );
62        }
63        let bytes = serde_json::to_vec(&body)?;
64        {
65            let mut file = fs::File::create(&tmp)?;
66            file.write_all(&bytes)?;
67            file.sync_all()?;
68        }
69        fs::rename(&tmp, &dest)?;
70        self.prune();
71        Ok(event)
72    }
73
74    pub fn pending(&self) -> Vec<EmittedBusinessEvent> {
75        let mut out = Vec::new();
76        let Ok(entries) = fs::read_dir(&self.dir) else {
77            return out;
78        };
79        for entry in entries.flatten() {
80            let path = entry.path();
81            if path.extension().and_then(|e| e.to_str()) != Some("json") {
82                continue;
83            }
84            let Ok(raw) = fs::read_to_string(&path) else {
85                continue;
86            };
87            let Ok(mut value) = serde_json::from_str::<Value>(&raw) else {
88                continue;
89            };
90            if let Some(obj) = value.as_object_mut() {
91                obj.remove("queued_at");
92            }
93            if let Ok(event) = serde_json::from_value::<EmittedBusinessEvent>(value) {
94                if !event.event_type.is_empty() {
95                    out.push(event);
96                }
97            }
98        }
99        out
100    }
101
102    pub fn ack(&self, event_id: &str) {
103        let id = event_id.trim();
104        if id.is_empty() {
105            return;
106        }
107        let _ = fs::remove_file(self.dir.join(file_name_for(id)));
108    }
109
110    fn prune(&self) {
111        let Ok(entries) = fs::read_dir(&self.dir) else {
112            return;
113        };
114        let mut files: Vec<(SystemTime, PathBuf)> = entries
115            .flatten()
116            .map(|e| e.path())
117            .filter(|p| p.extension().and_then(|e| e.to_str()) == Some("json"))
118            .filter_map(|p| {
119                let mtime = fs::metadata(&p).ok()?.modified().ok()?;
120                Some((mtime, p))
121            })
122            .collect();
123        if files.len() <= MAX_FILES {
124            return;
125        }
126        files.sort_by_key(|(t, _)| *t);
127        let drop_n = files.len() - MAX_FILES;
128        for (_, path) in files.into_iter().take(drop_n) {
129            let _ = fs::remove_file(path);
130        }
131    }
132}