Skip to main content

entracte_lib/scheduler/
exports.rs

1//! Declarative export-adapter delivery (#156, slice 6b).
2//!
3//! On a scheduler event, every installed export adapter subscribed to it has
4//! the host render its *own* break stats (CSV/JSON) and deliver them to the
5//! adapter's consent-fixed destination: a local file, or an HTTP POST — the
6//! only path in the app that sends data off the machine. The plugin runs no
7//! code and cannot influence the destination (fixed in the signed manifest,
8//! shown in full in the consent dialog).
9//!
10//! Delivery is fire-and-forget on a spawned task so it never blocks the
11//! scheduler tick, and bounded (payload cap + HTTP timeout + no redirects) so
12//! a slow or hostile endpoint can't stall or redirect it. Any failure is
13//! logged, never surfaced — a broken sink must not break breaks.
14
15use std::time::Duration;
16
17use crate::hooks::HookEvent;
18use crate::plugins::{ExportConfig, ExportFormat, ExportSink};
19use crate::stats::{self, LoggedEvent};
20
21use super::Scheduler;
22
23/// Hard cap on a rendered payload. Generous for a break-stats log; bounds the
24/// write/POST so an ever-growing event history can't produce an unbounded
25/// request.
26const MAX_EXPORT_BYTES: usize = 5 * 1024 * 1024;
27
28/// Timeout for the whole HTTP delivery, so a hung endpoint can't pin the task.
29const HTTP_TIMEOUT: Duration = Duration::from_secs(10);
30
31/// Render the logged events in the requested format. Pure.
32fn 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
39/// Deliver one rendered payload to its configured sink. Over-cap payloads are
40/// dropped; all failures are logged, never propagated.
41async 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
60/// POST `payload` to the adapter's URL. A fresh client per call with a hard
61/// timeout and **redirects disabled** — a redirect could bounce the data to a
62/// different host than the one the user consented to.
63async 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
90/// Fire-and-forget: deliver the current break stats to every export adapter
91/// subscribed to `event`. Snapshots the configs under the registry lock, then
92/// renders + delivers off the lock on a spawned task. No subscribers → no
93/// work (and no file read).
94pub 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
100/// The delivery body, split from the spawn wrapper so it's directly awaitable
101/// in tests. Snapshots subscribers under the lock, then renders + delivers off
102/// it.
103async 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        // Destination is a directory → write fails; must not panic.
177        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        // Port 1 on loopback refuses/!listens → send() errors; must not panic.
190        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        // A non-subscribed event delivers nothing.
246        run_delivery(reg.clone(), events.clone(), HookEvent::PauseStart).await;
247        assert!(!dest.exists());
248
249        // The subscribed event renders the stats to the file.
250        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        // A one-shot raw HTTP listener: accept one connection, read the
258        // request, reply 200. Proves the POST reaches the consented address
259        // with the body.
260        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        // Route through `deliver_one` so its Http arm is exercised too.
280        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}