qefro_backend_sdk/
event_outbox.rs1use 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}