Skip to main content

sova_devtools/
hub.rs

1//! Process-wide ring buffer + SSE fan-out.
2
3use crate::collector::{LogLine, RequestMeta, RequestSnapshot, now_ms};
4use serde::Serialize;
5use serde_json::{json, Value};
6use sova_sse::{SseChannel, SseEvent};
7use std::collections::{HashMap, VecDeque};
8use std::sync::{
9    atomic::{AtomicU64, Ordering},
10    Arc, Mutex,
11};
12use std::time::Duration;
13
14static SEQ: AtomicU64 = AtomicU64::new(1);
15
16pub fn next_id() -> String {
17    format!("dt-{}", SEQ.fetch_add(1, Ordering::Relaxed))
18}
19
20#[derive(Clone, Debug, Serialize)]
21pub struct CustomEvent {
22    pub id: String,
23    pub kind: String,
24    pub payload: Value,
25    pub ts_ms: u64,
26}
27
28#[derive(Clone, Debug, Serialize)]
29pub struct MemorySample {
30    pub ts_ms: u64,
31    pub rss_bytes: Option<u64>,
32}
33
34struct HubInner {
35    requests: VecDeque<RequestSnapshot>,
36    by_id: HashMap<String, RequestSnapshot>,
37    logs: VecDeque<LogLine>,
38    custom: VecDeque<CustomEvent>,
39    memory: VecDeque<MemorySample>,
40    plugins: Vec<String>,
41    profile: String,
42    event_seq: u64,
43    custom_cap: usize,
44    memory_cap: usize,
45}
46
47/// Shared DevTools state installed on the app.
48#[derive(Clone)]
49pub struct DevToolsHub {
50    inner: Arc<Mutex<HubInner>>,
51    pub channel: SseChannel,
52    request_cap: usize,
53    log_cap: usize,
54}
55
56impl DevToolsHub {
57    pub fn new(request_cap: usize, log_cap: usize) -> Self {
58        let channel = SseChannel::new(256).history_cap(100);
59        Self {
60            inner: Arc::new(Mutex::new(HubInner {
61                requests: VecDeque::new(),
62                by_id: HashMap::new(),
63                logs: VecDeque::new(),
64                custom: VecDeque::new(),
65                memory: VecDeque::new(),
66                plugins: Vec::new(),
67                profile: String::new(),
68                event_seq: 0,
69                custom_cap: 100,
70                memory_cap: 120,
71            })),
72            channel,
73            request_cap: request_cap.max(10),
74            log_cap: log_cap.max(50),
75        }
76    }
77
78    pub fn set_config_info(&self, plugins: Vec<String>, profile: String) {
79        let mut g = self.inner.lock().unwrap();
80        g.plugins = plugins;
81        g.profile = profile;
82    }
83
84    fn next_eid(g: &mut HubInner) -> String {
85        g.event_seq += 1;
86        g.event_seq.to_string()
87    }
88
89    pub fn push_snapshot(&self, snap: RequestSnapshot) {
90        let meta = RequestMeta::from(&snap);
91        let mut g = self.inner.lock().unwrap();
92        let eid = Self::next_eid(&mut g);
93        g.by_id.insert(snap.id.clone(), snap.clone());
94        g.requests.push_back(snap);
95        while g.requests.len() > self.request_cap {
96            if let Some(old) = g.requests.pop_front() {
97                g.by_id.remove(&old.id);
98            }
99        }
100        drop(g);
101        let data = serde_json::to_string(&json!({
102            "type": "request.finished",
103            "meta": meta,
104        }))
105        .unwrap_or_else(|_| "{}".into());
106        self.channel.publish(
107            SseEvent::data(data)
108                .id(eid)
109                .event("request.finished"),
110        );
111    }
112
113    pub fn push_log(&self, line: LogLine) {
114        let mut g = self.inner.lock().unwrap();
115        let eid = Self::next_eid(&mut g);
116        g.logs.push_back(line.clone());
117        while g.logs.len() > self.log_cap {
118            g.logs.pop_front();
119        }
120        drop(g);
121        let data = serde_json::to_string(&json!({
122            "type": "log.line",
123            "line": line,
124        }))
125        .unwrap_or_else(|_| "{}".into());
126        self.channel
127            .publish(SseEvent::data(data).id(eid).event("log.line"));
128    }
129
130    /// Emit a custom application/plugin event onto the DevTools SSE feed.
131    pub fn emit(&self, kind: impl Into<String>, payload: Value) {
132        let ev = CustomEvent {
133            id: next_id(),
134            kind: kind.into(),
135            payload,
136            ts_ms: now_ms(),
137        };
138        let mut g = self.inner.lock().unwrap();
139        let eid = Self::next_eid(&mut g);
140        let cap = g.custom_cap;
141        g.custom.push_back(ev.clone());
142        while g.custom.len() > cap {
143            g.custom.pop_front();
144        }
145        drop(g);
146        let data = serde_json::to_string(&json!({
147            "type": "custom",
148            "event": ev,
149        }))
150        .unwrap_or_else(|_| "{}".into());
151        self.channel
152            .publish(SseEvent::data(data).id(eid).event("custom"));
153    }
154
155    pub fn push_memory_sample(&self, rss_bytes: Option<u64>) {
156        let sample = MemorySample {
157            ts_ms: now_ms(),
158            rss_bytes,
159        };
160        let mut g = self.inner.lock().unwrap();
161        let eid = Self::next_eid(&mut g);
162        let cap = g.memory_cap;
163        g.memory.push_back(sample.clone());
164        while g.memory.len() > cap {
165            g.memory.pop_front();
166        }
167        drop(g);
168        let data = serde_json::to_string(&json!({
169            "type": "memory.sample",
170            "sample": sample,
171        }))
172        .unwrap_or_else(|_| "{}".into());
173        self.channel
174            .publish(SseEvent::data(data).id(eid).event("memory.sample"));
175    }
176
177    pub fn get(&self, id: &str) -> Option<RequestSnapshot> {
178        self.inner.lock().unwrap().by_id.get(id).cloned()
179    }
180
181    pub fn list_meta(&self, limit: usize) -> Vec<RequestMeta> {
182        let g = self.inner.lock().unwrap();
183        g.requests
184            .iter()
185            .rev()
186            .take(limit)
187            .map(RequestMeta::from)
188            .collect()
189    }
190
191    pub fn recent_logs(&self, limit: usize) -> Vec<LogLine> {
192        let g = self.inner.lock().unwrap();
193        g.logs.iter().rev().take(limit).cloned().collect()
194    }
195
196    pub fn recent_custom(&self, limit: usize) -> Vec<CustomEvent> {
197        let g = self.inner.lock().unwrap();
198        g.custom.iter().rev().take(limit).cloned().collect()
199    }
200
201    pub fn recent_memory(&self, limit: usize) -> Vec<MemorySample> {
202        let g = self.inner.lock().unwrap();
203        g.memory.iter().rev().take(limit).cloned().collect()
204    }
205
206    pub fn config_json(&self) -> serde_json::Value {
207        let g = self.inner.lock().unwrap();
208        json!({
209            "profile": g.profile,
210            "plugins": g.plugins,
211            "features": compile_features(),
212        })
213    }
214}
215
216/// Best-effort process RSS (Linux `/proc/self/status`).
217pub fn process_rss_bytes() -> Option<u64> {
218    #[cfg(target_os = "linux")]
219    {
220        let s = std::fs::read_to_string("/proc/self/status").ok()?;
221        for line in s.lines() {
222            if let Some(rest) = line.strip_prefix("VmRSS:") {
223                let kb: u64 = rest.split_whitespace().next()?.parse().ok()?;
224                return Some(kb.saturating_mul(1024));
225            }
226        }
227        None
228    }
229    #[cfg(not(target_os = "linux"))]
230    {
231        None
232    }
233}
234
235/// Background RSS sampler → SSE `memory.sample`.
236pub fn spawn_memory_sampler(hub: DevToolsHub, interval: Duration) {
237    tokio::spawn(async move {
238        let mut tick = tokio::time::interval(interval);
239        tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
240        loop {
241            tick.tick().await;
242            hub.push_memory_sample(process_rss_bytes());
243        }
244    });
245}
246
247/// Forward auth/mail EventBus events into [`DevToolsHub::emit`].
248pub fn wire_event_bus(app: &mut sova_core::App, hub: DevToolsHub) {
249    let bus = app.events();
250
251    #[cfg(feature = "auth")]
252    {
253        let h = hub.clone();
254        bus.listen::<sova_auth::UserRegistered, _>(move |e| {
255            h.emit(
256                "auth.user_registered",
257                json!({ "user_id": e.user_id, "email": e.email }),
258            );
259        });
260        let h = hub.clone();
261        bus.listen::<sova_auth::UserLoggedIn, _>(move |e| {
262            h.emit(
263                "auth.user_logged_in",
264                json!({ "user_id": e.user_id, "email": e.email }),
265            );
266        });
267    }
268
269    #[cfg(feature = "mail")]
270    {
271        let h = hub.clone();
272        bus.listen::<sova_mail::MailSent, _>(move |e| {
273            h.emit(
274                "mail.sent",
275                json!({ "to": e.to, "subject": e.subject }),
276            );
277        });
278    }
279
280    let _ = bus;
281    let _ = hub;
282}
283
284fn compile_features() -> Vec<&'static str> {
285    #[allow(clippy::vec_init_then_push, unused_mut)]
286    {
287        let mut v = Vec::new();
288        #[cfg(feature = "session")]
289        v.push("session");
290        #[cfg(feature = "mail")]
291        v.push("mail");
292        #[cfg(feature = "http")]
293        v.push("http");
294        #[cfg(feature = "db")]
295        v.push("db");
296        #[cfg(feature = "tasks")]
297        v.push("tasks");
298        #[cfg(feature = "auth")]
299        v.push("auth");
300        #[cfg(feature = "i18n")]
301        v.push("i18n");
302        #[cfg(feature = "csrf")]
303        v.push("csrf");
304        #[cfg(feature = "passport")]
305        v.push("passport");
306        #[cfg(feature = "store")]
307        v.push("store");
308        #[cfg(feature = "redis")]
309        v.push("redis");
310        #[cfg(feature = "rate-limit")]
311        v.push("rate-limit");
312        v
313    }
314}