1use std::time::Duration;
16
17use crate::hooks::HookEvent;
18use crate::plugins::{ExportConfig, ExportFormat, ExportSink};
19use crate::stats::{self, LoggedEvent};
20
21use super::Scheduler;
22
23const MAX_EXPORT_BYTES: usize = 5 * 1024 * 1024;
27
28const HTTP_TIMEOUT: Duration = Duration::from_secs(10);
30
31fn render_stats(events: &[LoggedEvent], format: ExportFormat) -> String {
33 match format {
34 ExportFormat::Csv => stats::export_csv(events),
35 ExportFormat::Json => serde_json::to_string(events).unwrap_or_else(|_| "[]".to_string()),
36 }
37}
38
39async fn deliver_one(cfg: &ExportConfig, payload: String) {
42 if payload.len() > MAX_EXPORT_BYTES {
43 log::warn!(
44 "export: skipping delivery to {} — payload {} bytes exceeds the {MAX_EXPORT_BYTES}-byte cap",
45 cfg.destination,
46 payload.len()
47 );
48 return;
49 }
50 match cfg.sink {
51 ExportSink::File => {
52 if let Err(e) = std::fs::write(&cfg.destination, payload.as_bytes()) {
53 log::warn!("export: write to {} failed: {e}", cfg.destination);
54 }
55 }
56 ExportSink::Http => post(cfg, payload).await,
57 }
58}
59
60async fn post(cfg: &ExportConfig, payload: String) {
64 let content_type = match cfg.format {
65 ExportFormat::Csv => "text/csv",
66 ExportFormat::Json => "application/json",
67 };
68 let client = match reqwest::Client::builder()
69 .timeout(HTTP_TIMEOUT)
70 .redirect(reqwest::redirect::Policy::none())
71 .build()
72 {
73 Ok(c) => c,
74 Err(e) => {
75 log::warn!("export: HTTP client build failed: {e}");
76 return;
77 }
78 };
79 if let Err(e) = client
80 .post(&cfg.destination)
81 .header(reqwest::header::CONTENT_TYPE, content_type)
82 .body(payload)
83 .send()
84 .await
85 {
86 log::warn!("export: POST to {} failed: {e}", cfg.destination);
87 }
88}
89
90pub fn deliver_on_event(scheduler: &Scheduler, event: HookEvent) {
95 let registry = scheduler.plugins.clone();
96 let events_path = scheduler.events_path.clone();
97 tauri::async_runtime::spawn(run_delivery(registry, events_path, event));
98}
99
100async fn run_delivery(
104 registry: std::sync::Arc<tokio::sync::Mutex<crate::plugins::PluginRegistry>>,
105 events_path: std::path::PathBuf,
106 event: HookEvent,
107) {
108 let configs = { registry.lock().await.export_configs_for(event) };
109 if configs.is_empty() {
110 return;
111 }
112 let events = stats::read_all(&events_path);
113 for cfg in configs {
114 let payload = render_stats(&events, cfg.format);
115 deliver_one(&cfg, payload).await;
116 }
117}
118
119#[cfg(test)]
120mod tests {
121 use super::*;
122 use crate::stats::{EventPayload, LoggedEvent};
123 use std::io::{Read, Write};
124
125 fn sample_events() -> Vec<LoggedEvent> {
126 vec![LoggedEvent {
127 t: "2026-06-10T00:00:00Z"
128 .parse::<chrono::DateTime<chrono::Utc>>()
129 .unwrap(),
130 event: EventPayload::BreakResumed {
131 kind: crate::scheduler::BreakKind::Micro,
132 },
133 }]
134 }
135
136 #[test]
137 fn render_stats_csv_and_json_differ_and_carry_the_event() {
138 let events = sample_events();
139 let csv = render_stats(&events, ExportFormat::Csv);
140 let json = render_stats(&events, ExportFormat::Json);
141 assert!(csv.contains("break_resumed") || csv.contains("micro"));
142 assert!(json.starts_with('[') && json.contains("break_resumed"));
143 assert_ne!(csv, json);
144 }
145
146 #[tokio::test]
147 async fn file_sink_writes_the_payload() {
148 let dir = crate::test_support::temp_dir();
149 let dest = dir.path().join("breaks.json");
150 let cfg = ExportConfig {
151 sink: ExportSink::File,
152 format: ExportFormat::Json,
153 destination: dest.display().to_string(),
154 on: vec![HookEvent::BreakEnd],
155 };
156 deliver_one(&cfg, "[{\"x\":1}]".to_string()).await;
157 assert_eq!(std::fs::read_to_string(&dest).unwrap(), "[{\"x\":1}]");
158 }
159
160 #[tokio::test]
161 async fn oversized_payload_is_dropped_not_written() {
162 let dir = crate::test_support::temp_dir();
163 let dest = dir.path().join("breaks.csv");
164 let cfg = ExportConfig {
165 sink: ExportSink::File,
166 format: ExportFormat::Csv,
167 destination: dest.display().to_string(),
168 on: vec![HookEvent::BreakEnd],
169 };
170 deliver_one(&cfg, "x".repeat(MAX_EXPORT_BYTES + 1)).await;
171 assert!(!dest.exists(), "over-cap payload must not be written");
172 }
173
174 #[tokio::test]
175 async fn file_sink_write_failure_is_logged_not_panicked() {
176 let dir = crate::test_support::temp_dir();
178 let cfg = ExportConfig {
179 sink: ExportSink::File,
180 format: ExportFormat::Json,
181 destination: dir.path().display().to_string(),
182 on: vec![HookEvent::BreakEnd],
183 };
184 deliver_one(&cfg, "[]".to_string()).await;
185 }
186
187 #[tokio::test]
188 async fn http_post_failure_is_logged_not_panicked() {
189 let cfg = ExportConfig {
191 sink: ExportSink::Http,
192 format: ExportFormat::Csv,
193 destination: "http://127.0.0.1:1/ingest".to_string(),
194 on: vec![HookEvent::BreakEnd],
195 };
196 post(&cfg, "ts\n".to_string()).await;
197 }
198
199 #[tokio::test]
200 async fn run_delivery_renders_and_delivers_to_subscribers_only() {
201 use crate::plugins::{InstalledPlugin, Manifest, PluginKind, Signature, MANIFEST_VERSION};
202 use std::sync::Arc;
203 use tokio::sync::Mutex;
204
205 let dir = crate::test_support::temp_dir();
206 let events = dir.path().join("events.jsonl");
207 std::fs::write(
208 &events,
209 b"{\"t\":\"2026-06-10T00:00:00Z\",\"type\":\"break_resumed\",\"kind\":\"micro\"}\n",
210 )
211 .unwrap();
212 let dest = dir.path().join("out.json");
213
214 let manifest = Manifest {
215 manifest_version: MANIFEST_VERSION,
216 id: "com.x.exp".to_string(),
217 name: "E".to_string(),
218 version: "1.0.0".to_string(),
219 author: String::new(),
220 description: String::new(),
221 kind: PluginKind::Export,
222 module: None,
223 module_base64: None,
224 abi_version: None,
225 imports: vec![],
226 detect: None,
227 export: Some(ExportConfig {
228 sink: ExportSink::File,
229 format: ExportFormat::Json,
230 destination: dest.display().to_string(),
231 on: vec![HookEvent::BreakEnd],
232 }),
233 content: None,
234 assets: Vec::new(),
235 signature: Signature {
236 alg: "ed25519".to_string(),
237 public_key: String::new(),
238 sig: String::new(),
239 },
240 };
241 let mut reg = crate::plugins::PluginRegistry::default();
242 reg.insert(InstalledPlugin::from_export(&manifest));
243 let reg = Arc::new(Mutex::new(reg));
244
245 run_delivery(reg.clone(), events.clone(), HookEvent::PauseStart).await;
247 assert!(!dest.exists());
248
249 run_delivery(reg, events, HookEvent::BreakEnd).await;
251 let written = std::fs::read_to_string(&dest).unwrap();
252 assert!(written.contains("break_resumed"));
253 }
254
255 #[tokio::test]
256 async fn http_sink_posts_the_payload_to_the_destination() {
257 let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
261 let addr = listener.local_addr().unwrap();
262 let handle = std::thread::spawn(move || {
263 let (mut stream, _) = listener.accept().unwrap();
264 let mut buf = [0u8; 2048];
265 let n = stream.read(&mut buf).unwrap();
266 let req = String::from_utf8_lossy(&buf[..n]).to_string();
267 stream
268 .write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 0\r\n\r\n")
269 .unwrap();
270 req
271 });
272
273 let cfg = ExportConfig {
274 sink: ExportSink::Http,
275 format: ExportFormat::Json,
276 destination: format!("http://{addr}/ingest"),
277 on: vec![HookEvent::BreakEnd],
278 };
279 deliver_one(&cfg, "[{\"break\":1}]".to_string()).await;
281
282 let req = handle.join().unwrap();
283 assert!(req.starts_with("POST /ingest "), "got: {req}");
284 assert!(req.contains("[{\"break\":1}]"), "body delivered: {req}");
285 assert!(req
286 .to_lowercase()
287 .contains("content-type: application/json"));
288 }
289}