1use 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#[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 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
216pub 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
235pub 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
247pub 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}