Skip to main content

sova_devtools/
hub.rs

1//! Process-wide ring buffer + SSE fan-out.
2
3use crate::collector::{now_ms, LogLine, RequestMeta, RequestSnapshot};
4use serde::Serialize;
5use serde_json::{json, Map, Value};
6use sova_core::DevToolsConfigRegistry;
7use sova_sse::{SseChannel, SseEvent};
8use std::collections::{HashMap, VecDeque};
9use std::sync::{
10    atomic::{AtomicU64, Ordering},
11    Arc, Mutex,
12};
13use std::time::Duration;
14
15static SEQ: AtomicU64 = AtomicU64::new(1);
16
17pub fn next_id() -> String {
18    format!("dt-{}", SEQ.fetch_add(1, Ordering::Relaxed))
19}
20
21#[derive(Clone, Debug, Serialize)]
22pub struct CustomEvent {
23    pub id: String,
24    pub kind: String,
25    pub payload: Value,
26    pub ts_ms: u64,
27}
28
29#[derive(Clone, Debug, Serialize)]
30pub struct MemorySample {
31    pub ts_ms: u64,
32    pub rss_bytes: Option<u64>,
33    pub rss_peak_bytes: Option<u64>,
34    #[serde(skip_serializing_if = "Option::is_none")]
35    pub available_bytes: Option<u64>,
36}
37
38#[derive(Clone, Debug, Serialize)]
39pub struct MemorySummary {
40    pub samples: Vec<MemorySample>,
41    pub current: Option<u64>,
42    pub peak: Option<u64>,
43    pub min: Option<u64>,
44}
45
46struct HubInner {
47    requests: VecDeque<RequestSnapshot>,
48    by_id: HashMap<String, RequestSnapshot>,
49    logs: VecDeque<LogLine>,
50    custom: VecDeque<CustomEvent>,
51    memory: VecDeque<MemorySample>,
52    rss_peak: Option<u64>,
53    plugins: Vec<String>,
54    profile: String,
55    config_registry: Option<DevToolsConfigRegistry>,
56    event_seq: u64,
57    custom_cap: usize,
58    memory_cap: usize,
59}
60
61/// Shared DevTools state installed on the app.
62#[derive(Clone)]
63pub struct DevToolsHub {
64    inner: Arc<Mutex<HubInner>>,
65    pub channel: SseChannel,
66    request_cap: usize,
67    log_cap: usize,
68}
69
70impl DevToolsHub {
71    pub fn new(request_cap: usize, log_cap: usize) -> Self {
72        let channel = SseChannel::new(256).history_cap(100);
73        Self {
74            inner: Arc::new(Mutex::new(HubInner {
75                requests: VecDeque::new(),
76                by_id: HashMap::new(),
77                logs: VecDeque::new(),
78                custom: VecDeque::new(),
79                memory: VecDeque::new(),
80                rss_peak: None,
81                plugins: Vec::new(),
82                profile: String::new(),
83                config_registry: None,
84                event_seq: 0,
85                custom_cap: 100,
86                memory_cap: 120,
87            })),
88            channel,
89            request_cap: request_cap.max(10),
90            log_cap: log_cap.max(50),
91        }
92    }
93
94    pub fn set_config_info(&self, plugins: Vec<String>, profile: String) {
95        let mut g = self.inner.lock().unwrap();
96        g.plugins = plugins;
97        g.profile = profile;
98    }
99
100    pub fn set_config_registry(&self, registry: DevToolsConfigRegistry) {
101        self.inner.lock().unwrap().config_registry = Some(registry);
102    }
103
104    fn next_eid(g: &mut HubInner) -> String {
105        g.event_seq += 1;
106        g.event_seq.to_string()
107    }
108
109    pub fn push_snapshot(&self, snap: RequestSnapshot) {
110        let meta = RequestMeta::from(&snap);
111        let mut g = self.inner.lock().unwrap();
112        let eid = Self::next_eid(&mut g);
113        g.by_id.insert(snap.id.clone(), snap.clone());
114        g.requests.push_back(snap);
115        while g.requests.len() > self.request_cap {
116            if let Some(old) = g.requests.pop_front() {
117                g.by_id.remove(&old.id);
118            }
119        }
120        drop(g);
121        let data = serde_json::to_string(&json!({
122            "type": "request.finished",
123            "meta": meta,
124        }))
125        .unwrap_or_else(|_| "{}".into());
126        self.channel
127            .publish(SseEvent::data(data).id(eid).event("request.finished"));
128    }
129
130    pub fn push_log(&self, line: LogLine) {
131        let mut g = self.inner.lock().unwrap();
132        let eid = Self::next_eid(&mut g);
133        g.logs.push_back(line.clone());
134        while g.logs.len() > self.log_cap {
135            g.logs.pop_front();
136        }
137        drop(g);
138        let data = serde_json::to_string(&json!({
139            "type": "log.line",
140            "line": line,
141        }))
142        .unwrap_or_else(|_| "{}".into());
143        self.channel
144            .publish(SseEvent::data(data).id(eid).event("log.line"));
145    }
146
147    /// Emit a custom application/plugin event onto the DevTools SSE feed.
148    pub fn emit(&self, kind: impl Into<String>, payload: Value) {
149        let ev = CustomEvent {
150            id: next_id(),
151            kind: kind.into(),
152            payload,
153            ts_ms: now_ms(),
154        };
155        let mut g = self.inner.lock().unwrap();
156        let eid = Self::next_eid(&mut g);
157        let cap = g.custom_cap;
158        g.custom.push_back(ev.clone());
159        while g.custom.len() > cap {
160            g.custom.pop_front();
161        }
162        drop(g);
163        let data = serde_json::to_string(&json!({
164            "type": "custom",
165            "event": ev,
166        }))
167        .unwrap_or_else(|_| "{}".into());
168        self.channel
169            .publish(SseEvent::data(data).id(eid).event("custom"));
170    }
171
172    pub fn push_memory_sample(&self, rss_bytes: Option<u64>) {
173        let available_bytes = process_mem_available_bytes();
174        let mut g = self.inner.lock().unwrap();
175        if let Some(rss) = rss_bytes {
176            g.rss_peak = Some(match g.rss_peak {
177                Some(p) => p.max(rss),
178                None => rss,
179            });
180        }
181        let sample = MemorySample {
182            ts_ms: now_ms(),
183            rss_bytes,
184            rss_peak_bytes: g.rss_peak,
185            available_bytes,
186        };
187        let eid = Self::next_eid(&mut g);
188        let cap = g.memory_cap;
189        g.memory.push_back(sample.clone());
190        while g.memory.len() > cap {
191            g.memory.pop_front();
192        }
193        drop(g);
194        let data = serde_json::to_string(&json!({
195            "type": "memory.sample",
196            "sample": sample,
197        }))
198        .unwrap_or_else(|_| "{}".into());
199        self.channel
200            .publish(SseEvent::data(data).id(eid).event("memory.sample"));
201    }
202
203    pub fn get(&self, id: &str) -> Option<RequestSnapshot> {
204        self.inner.lock().unwrap().by_id.get(id).cloned()
205    }
206
207    pub fn list_meta(&self, limit: usize) -> Vec<RequestMeta> {
208        let g = self.inner.lock().unwrap();
209        g.requests
210            .iter()
211            .rev()
212            .take(limit)
213            .map(RequestMeta::from)
214            .collect()
215    }
216
217    pub fn recent_logs(&self, limit: usize) -> Vec<LogLine> {
218        let g = self.inner.lock().unwrap();
219        g.logs.iter().rev().take(limit).cloned().collect()
220    }
221
222    pub fn recent_custom(&self, limit: usize) -> Vec<CustomEvent> {
223        let g = self.inner.lock().unwrap();
224        g.custom.iter().rev().take(limit).cloned().collect()
225    }
226
227    pub fn recent_memory(&self, limit: usize) -> MemorySummary {
228        let g = self.inner.lock().unwrap();
229        let samples: Vec<MemorySample> = g.memory.iter().rev().take(limit).cloned().collect();
230        let mut current: Option<u64> = None;
231        let mut peak: Option<u64> = g.rss_peak;
232        let mut min: Option<u64> = None;
233        for s in &samples {
234            if let Some(rss) = s.rss_bytes {
235                if current.is_none() {
236                    current = Some(rss);
237                }
238                peak = Some(peak.map_or(rss, |p: u64| p.max(rss)));
239                min = Some(min.map_or(rss, |m: u64| m.min(rss)));
240            }
241        }
242        MemorySummary {
243            samples,
244            current,
245            peak,
246            min,
247        }
248    }
249
250    pub fn config_json(&self) -> serde_json::Value {
251        let g = self.inner.lock().unwrap();
252        let mounts = g
253            .config_registry
254            .as_ref()
255            .map(|r| Value::Object(r.snapshot()))
256            .unwrap_or_else(|| Value::Object(Map::new()));
257        json!({
258            "profile": g.profile,
259            "plugins": g.plugins,
260            "features": compile_features(),
261            "mounts": mounts,
262        })
263    }
264}
265
266/// Best-effort process RSS (Linux `/proc`, macOS `task_info`).
267pub fn process_rss_bytes() -> Option<u64> {
268    #[cfg(target_os = "linux")]
269    {
270        let s = std::fs::read_to_string("/proc/self/status").ok()?;
271        for line in s.lines() {
272            if let Some(rest) = line.strip_prefix("VmRSS:") {
273                let kb: u64 = rest.split_whitespace().next()?.parse().ok()?;
274                return Some(kb.saturating_mul(1024));
275            }
276        }
277        None
278    }
279    #[cfg(target_os = "macos")]
280    {
281        macos_rss_bytes()
282    }
283    #[cfg(not(any(target_os = "linux", target_os = "macos")))]
284    {
285        None
286    }
287}
288
289fn process_mem_available_bytes() -> Option<u64> {
290    #[cfg(target_os = "linux")]
291    {
292        let s = std::fs::read_to_string("/proc/meminfo").ok()?;
293        for line in s.lines() {
294            if let Some(rest) = line.strip_prefix("MemAvailable:") {
295                let kb: u64 = rest.split_whitespace().next()?.parse().ok()?;
296                return Some(kb.saturating_mul(1024));
297            }
298        }
299        None
300    }
301    #[cfg(not(target_os = "linux"))]
302    {
303        None
304    }
305}
306
307#[cfg(target_os = "macos")]
308fn macos_rss_bytes() -> Option<u64> {
309    // MACH_TASK_BASIC_INFO → resident_size (bytes).
310    #[repr(C)]
311    struct TaskBasicInfo {
312        suspend_count: u32,
313        virtual_size: u64,
314        resident_size: u64,
315        user_time: [u32; 2],
316        system_time: [u32; 2],
317        policy: i32,
318    }
319    const TASK_BASIC_INFO: u32 = 5;
320    const TASK_BASIC_INFO_COUNT: u32 =
321        (std::mem::size_of::<TaskBasicInfo>() / std::mem::size_of::<u32>()) as u32;
322
323    extern "C" {
324        fn mach_task_self() -> u32;
325        fn task_info(
326            target_task: u32,
327            flavor: u32,
328            task_info_out: *mut TaskBasicInfo,
329            task_info_count: *mut u32,
330        ) -> i32;
331    }
332
333    let mut info = unsafe { std::mem::zeroed::<TaskBasicInfo>() };
334    let mut count = TASK_BASIC_INFO_COUNT;
335    let kr = unsafe { task_info(mach_task_self(), TASK_BASIC_INFO, &mut info, &mut count) };
336    if kr == 0 {
337        Some(info.resident_size)
338    } else {
339        None
340    }
341}
342
343/// Background RSS sampler → SSE `memory.sample`.
344pub fn spawn_memory_sampler(hub: DevToolsHub, interval: Duration) {
345    tokio::spawn(async move {
346        let mut tick = tokio::time::interval(interval);
347        tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
348        loop {
349            tick.tick().await;
350            hub.push_memory_sample(process_rss_bytes());
351        }
352    });
353}
354
355/// Forward domain EventBus events into [`DevToolsHub::emit`].
356pub fn wire_event_bus(app: &mut sova_core::App, hub: DevToolsHub) {
357    let bus = app.events();
358
359    #[cfg(feature = "auth")]
360    {
361        let h = hub.clone();
362        bus.listen::<sova_auth::UserRegistered, _>(move |e| {
363            h.emit(
364                "auth.user_registered",
365                json!({ "user_id": e.user_id, "email": e.email }),
366            );
367        });
368        let h = hub.clone();
369        bus.listen::<sova_auth::UserLoggedIn, _>(move |e| {
370            h.emit(
371                "auth.user_logged_in",
372                json!({ "user_id": e.user_id, "email": e.email }),
373            );
374        });
375    }
376
377    #[cfg(feature = "mail")]
378    {
379        let h = hub.clone();
380        bus.listen::<sova_mail::MailSent, _>(move |e| {
381            h.emit("mail.sent", json!({ "to": e.to, "subject": e.subject }));
382        });
383    }
384
385    #[cfg(feature = "fs")]
386    {
387        let h = hub.clone();
388        bus.listen::<sova_fs::FileWritten, _>(move |e| {
389            h.emit("fs.file_written", json!({ "path": e.path }));
390        });
391        let h = hub.clone();
392        bus.listen::<sova_fs::FileRemoved, _>(move |e| {
393            h.emit("fs.file_removed", json!({ "path": e.path }));
394        });
395        let h = hub.clone();
396        bus.listen::<sova_fs::DirCreated, _>(move |e| {
397            h.emit("fs.dir_created", json!({ "path": e.path }));
398        });
399    }
400
401    #[cfg(feature = "csrf")]
402    {
403        let h = hub.clone();
404        bus.listen::<sova_csrf::CsrfMismatch, _>(move |e| {
405            h.emit(
406                "csrf.mismatch",
407                json!({ "method": e.method, "path": e.path }),
408            );
409        });
410    }
411
412    #[cfg(feature = "rate-limit")]
413    {
414        let h = hub.clone();
415        bus.listen::<sova_rate_limit::RateLimitExceeded, _>(move |e| {
416            h.emit(
417                "rate_limit.exceeded",
418                json!({
419                    "key": e.key,
420                    "limit": e.limit,
421                    "retry_after": e.retry_after,
422                }),
423            );
424        });
425    }
426
427    #[cfg(feature = "session")]
428    {
429        let h = hub.clone();
430        bus.listen::<sova_session::SessionRegenerated, _>(move |e| {
431            h.emit("session.regenerated", json!({ "had_user": e.had_user }));
432        });
433        let h = hub.clone();
434        bus.listen::<sova_session::SessionLogoutAll, _>(move |e| {
435            h.emit(
436                "session.logout_all",
437                json!({ "user_id": e.user_id, "count": e.count }),
438            );
439        });
440    }
441
442    #[cfg(feature = "tasks")]
443    {
444        let h = hub.clone();
445        bus.listen::<sova_tasks::TaskDispatched, _>(move |e| {
446            h.emit(
447                "tasks.dispatched",
448                json!({ "id": e.id, "name": e.name, "queue": e.queue }),
449            );
450        });
451        let h = hub.clone();
452        bus.listen::<sova_tasks::TaskFailed, _>(move |e| {
453            h.emit(
454                "tasks.failed",
455                json!({ "id": e.id, "name": e.name, "attempts": e.attempts }),
456            );
457        });
458    }
459
460    #[cfg(feature = "notifications")]
461    {
462        let h = hub.clone();
463        bus.listen::<sova_notifications::NotificationSent, _>(move |e| {
464            h.emit(
465                "notifications.sent",
466                json!({
467                    "channel": e.channel,
468                    "event": e.event,
469                    "recipients": e.recipients,
470                }),
471            );
472        });
473    }
474
475    #[cfg(feature = "passport")]
476    {
477        let h = hub.clone();
478        bus.listen::<sova_passport::ApiTokenRevoked, _>(move |e| {
479            h.emit(
480                "passport.api_token_revoked",
481                json!({ "user_id": e.user_id, "token_id": e.token_id }),
482            );
483        });
484    }
485
486    #[cfg(feature = "acme")]
487    {
488        let h = hub.clone();
489        bus.listen::<sova_acme::CertificateIssued, _>(move |e| {
490            h.emit(
491                "acme.certificate_issued",
492                json!({
493                    "domains": e.domains,
494                    "not_after_unix": e.not_after_unix,
495                }),
496            );
497        });
498        let h = hub.clone();
499        bus.listen::<sova_acme::CertificateRenewed, _>(move |e| {
500            h.emit(
501                "acme.certificate_renewed",
502                json!({
503                    "domains": e.domains,
504                    "not_after_unix": e.not_after_unix,
505                }),
506            );
507        });
508        let h = hub.clone();
509        bus.listen::<sova_acme::AcmeFailed, _>(move |e| {
510            h.emit(
511                "acme.failed",
512                json!({ "domains": e.domains, "error": e.error }),
513            );
514        });
515    }
516
517    let _ = bus;
518    let _ = hub;
519}
520
521fn compile_features() -> Vec<&'static str> {
522    #[allow(clippy::vec_init_then_push, unused_mut)]
523    {
524        let mut v = Vec::new();
525        #[cfg(feature = "session")]
526        v.push("session");
527        #[cfg(feature = "mail")]
528        v.push("mail");
529        #[cfg(feature = "http")]
530        v.push("http");
531        #[cfg(feature = "db")]
532        v.push("db");
533        #[cfg(feature = "tasks")]
534        v.push("tasks");
535        #[cfg(feature = "auth")]
536        v.push("auth");
537        #[cfg(feature = "i18n")]
538        v.push("i18n");
539        #[cfg(feature = "csrf")]
540        v.push("csrf");
541        #[cfg(feature = "passport")]
542        v.push("passport");
543        #[cfg(feature = "store")]
544        v.push("store");
545        #[cfg(feature = "redis")]
546        v.push("redis");
547        #[cfg(feature = "rate-limit")]
548        v.push("rate-limit");
549        #[cfg(feature = "notifications")]
550        v.push("notifications");
551        #[cfg(feature = "acme")]
552        v.push("acme");
553        #[cfg(feature = "console")]
554        v.push("console");
555        #[cfg(feature = "console-redis")]
556        v.push("console-redis");
557        #[cfg(feature = "console-store")]
558        v.push("console-store");
559        #[cfg(feature = "console-graphql")]
560        v.push("console-graphql");
561        #[cfg(feature = "console-tasks")]
562        v.push("console-tasks");
563        #[cfg(feature = "console-mail")]
564        v.push("console-mail");
565        #[cfg(feature = "console-http-external")]
566        v.push("console-http-external");
567        #[cfg(feature = "console-events")]
568        v.push("console-events");
569        #[cfg(feature = "console-rabbit")]
570        v.push("console-rabbit");
571        #[cfg(feature = "console-grpc")]
572        v.push("console-grpc");
573        #[cfg(feature = "console-session")]
574        v.push("console-session");
575        #[cfg(feature = "graphql")]
576        v.push("graphql");
577        #[cfg(feature = "rabbit")]
578        v.push("rabbit");
579        #[cfg(feature = "grpc")]
580        v.push("grpc");
581        #[cfg(feature = "fs")]
582        v.push("fs");
583        v
584    }
585}