1use std::collections::BTreeMap;
9use std::collections::HashMap;
10use std::sync::Arc;
11use std::time::{Duration, Instant};
12
13use serde_json::Value;
14use tokio::sync::{Mutex, Notify};
15use tokio::task::JoinHandle;
16
17use crate::sdk::cdp::ws_client::Connection;
18use crate::sdk::fetch::artifacts::{
19 console::{ConsoleEvent, ConsoleLevel, ConsoleLog},
20 network::{NetworkEntry, NetworkEntryState, NetworkLog, NetworkSummary, NetworkTiming},
21};
22use crate::shared::redact;
23
24struct NetworkSlot {
26 entry: NetworkEntry,
27}
28
29fn default_entry() -> NetworkEntry {
30 NetworkEntry {
31 request_id: String::new(),
32 state: NetworkEntryState::Pending,
33 redirect_from_request_id: None,
34 frame_id: None,
35 loader_id: None,
36 resource_type: "Other".into(),
37 request_url: String::new(),
38 method: "GET".into(),
39 initiator: None,
40 status: None,
41 mime_type: None,
42 request_headers: BTreeMap::new(),
43 response_headers: BTreeMap::new(),
44 request_post_data_present: false,
45 request_post_data_size_bytes: None,
46 body_file: None,
47 timing: NetworkTiming {
48 start_monotonic_ms: 0,
49 end_monotonic_ms: None,
50 },
51 failure: None,
52 hints: BTreeMap::new(),
53 }
54}
55
56pub struct NetworkCollector {
60 inner: Arc<Mutex<NetworkInner>>,
61 task: JoinHandle<()>,
62 redact_headers: bool,
63 main_notify: Arc<Notify>,
67}
68
69#[derive(Default)]
70struct NetworkInner {
71 slots: HashMap<String, NetworkSlot>,
72 main_request_id: Option<String>,
73 finished: Vec<String>,
75 ws_frames: HashMap<String, Vec<Value>>,
77 sse_events: HashMap<String, Vec<Value>>,
79 last_activity: Option<Instant>,
82}
83
84#[derive(Debug, Clone)]
85pub struct NetworkQuietSnapshot {
86 pub quiet: bool,
87 pub idle_for: Duration,
88 pub inflight_total: usize,
89 pub pending_by_resource_type: BTreeMap<String, usize>,
90}
91
92impl NetworkCollector {
93 pub fn start(
94 conn: &Connection,
95 redact_headers: bool,
96 capture_ws: bool,
97 capture_sse: bool,
98 ) -> Self {
99 let inner = Arc::new(Mutex::new(NetworkInner::default()));
100 let main_notify = Arc::new(Notify::new());
101 let mut rx = conn.subscribe();
102 let inner_w = inner.clone();
103 let notify_w = main_notify.clone();
104 let task = tokio::spawn(async move {
105 while let Ok(ev) = rx.recv().await {
106 let mut guard = inner_w.lock().await;
107 let touched_main = handle_event(&mut guard, &ev, redact_headers);
108 if capture_ws {
109 handle_ws_event(&mut guard, &ev);
110 }
111 if capture_sse {
112 handle_sse_event(&mut guard, &ev);
113 }
114 drop(guard);
115 if touched_main {
116 notify_w.notify_waiters();
117 }
118 }
119 });
120 Self {
121 inner,
122 task,
123 redact_headers,
124 main_notify,
125 }
126 }
127
128 pub async fn take_ws_frames(&self) -> HashMap<String, Vec<Value>> {
130 let mut guard = self.inner.lock().await;
131 std::mem::take(&mut guard.ws_frames)
132 }
133
134 pub async fn take_sse_events(&self) -> HashMap<String, Vec<Value>> {
136 let mut guard = self.inner.lock().await;
137 std::mem::take(&mut guard.sse_events)
138 }
139
140 pub async fn wait_for_main_status(&self, timeout: Duration) -> Option<NetworkEntry> {
145 let deadline = tokio::time::Instant::now() + timeout;
146 loop {
147 let notified = self.main_notify.notified();
150 tokio::pin!(notified);
151 if let Some(entry) = self.main_entry().await
152 && (entry.status.is_some() || entry.failure.is_some())
153 {
154 return Some(entry);
155 }
156 let remaining = deadline.saturating_duration_since(tokio::time::Instant::now());
157 if remaining.is_zero() {
158 return self.main_entry().await;
159 }
160 if tokio::time::timeout(remaining, notified).await.is_err() {
161 return self.main_entry().await;
162 }
163 }
164 }
165
166 pub async fn take_finished(&self) -> Vec<String> {
170 let mut guard = self.inner.lock().await;
171 std::mem::take(&mut guard.finished)
172 }
173
174 pub async fn entry(&self, request_id: &str) -> Option<NetworkEntry> {
178 let guard = self.inner.lock().await;
179 guard.slots.get(request_id).map(|s| s.entry.clone())
180 }
181
182 pub async fn main_request_id(&self) -> Option<String> {
183 let guard = self.inner.lock().await;
184 guard.main_request_id.clone()
185 }
186
187 pub async fn main_entry(&self) -> Option<NetworkEntry> {
188 let guard = self.inner.lock().await;
189 let request_id = guard.main_request_id.as_ref()?;
190 guard.slots.get(request_id).map(|s| s.entry.clone())
191 }
192
193 pub async fn quiet_snapshot(&self, idle: Duration) -> NetworkQuietSnapshot {
196 let guard = self.inner.lock().await;
197 let pending_by_resource_type =
198 pending_by_resource_type(guard.slots.values().map(|slot| &slot.entry));
199 let inflight_total: usize = pending_by_resource_type.values().copied().sum();
200 let idle_for = guard
201 .last_activity
202 .map(|last| last.elapsed())
203 .unwrap_or_else(|| Duration::from_secs(u64::MAX / 2));
204 NetworkQuietSnapshot {
205 quiet: inflight_total == 0 && idle_for >= idle,
206 idle_for,
207 inflight_total,
208 pending_by_resource_type,
209 }
210 }
211
212 pub async fn set_body_file(&self, request_id: &str, path: std::path::PathBuf) {
214 let mut guard = self.inner.lock().await;
215 if let Some(s) = guard.slots.get_mut(request_id) {
216 s.entry.body_file = Some(path);
217 }
218 }
219
220 pub async fn set_hint(&self, request_id: &str, key: &str, value: Value) {
222 let mut guard = self.inner.lock().await;
223 if let Some(s) = guard.slots.get_mut(request_id) {
224 s.entry.hints.insert(key.to_string(), value);
225 }
226 }
227
228 pub async fn finish(self) -> NetworkLog {
231 self.task.abort();
232 let mut guard = self.inner.lock().await;
233 let mut entries: Vec<NetworkEntry> = guard.slots.drain().map(|(_, s)| s.entry).collect();
234 entries.sort_by(|a, b| a.request_id.cmp(&b.request_id));
235 let summary = summarize_entries(&entries, self.redact_headers);
236 NetworkLog {
237 schema_version: super::network::NETWORK_SCHEMA_VERSION,
238 main_request_id: guard.main_request_id.clone(),
239 entries,
240 summary,
241 }
242 }
243
244 #[allow(dead_code)]
245 pub fn redact_headers(&self) -> bool {
246 self.redact_headers
247 }
248}
249
250fn handle_event(
254 inner: &mut NetworkInner,
255 ev: &crate::sdk::cdp::ws_client::CdpEvent,
256 redact_headers: bool,
257) -> bool {
258 match ev.method.as_str() {
259 "Network.requestWillBeSent" => {
260 inner.last_activity = Some(Instant::now());
261 let req_id = string_field(&ev.params, "requestId");
262 let frame_id = ev
263 .params
264 .get("frameId")
265 .and_then(|v| v.as_str())
266 .map(str::to_string);
267 let loader_id = ev
268 .params
269 .get("loaderId")
270 .and_then(|v| v.as_str())
271 .map(str::to_string);
272 let resource_type = ev
273 .params
274 .get("type")
275 .and_then(|v| v.as_str())
276 .unwrap_or("Other")
277 .to_string();
278 if resource_type == "Document" && inner.main_request_id.is_none() {
279 inner.main_request_id = Some(req_id.clone());
280 }
281 let timestamp = ev
282 .params
283 .get("timestamp")
284 .and_then(|v| v.as_f64())
285 .unwrap_or(0.0);
286 let request = &ev.params["request"];
287 let initiator = ev.params.get("initiator").cloned();
288 let request_url = request
289 .get("url")
290 .and_then(|v| v.as_str())
291 .map(|url| {
292 crate::sdk::fetch::artifacts::network::redact_request_url(url, redact_headers)
293 })
294 .unwrap_or_default();
295 let method = request
296 .get("method")
297 .and_then(|v| v.as_str())
298 .unwrap_or("GET")
299 .to_string();
300 let mut request_headers = headers_map(request.get("headers"));
301 if redact_headers {
302 redact_inplace(&mut request_headers);
303 }
304 if let Some(redirect) = ev.params.get("redirectResponse")
307 && let Some(existing) = inner.slots.get_mut(&req_id)
308 {
309 let status = redirect
310 .get("status")
311 .and_then(|v| v.as_u64())
312 .map(|s| s as u16);
313 existing.entry.status = status;
314 existing.entry.state = NetworkEntryState::Finished;
315 let mut h = headers_map(redirect.get("headers"));
316 if redact_headers {
317 redact_inplace(&mut h);
318 }
319 existing.entry.response_headers = h;
320 existing
321 .entry
322 .hints
323 .insert("redirected".into(), Value::Bool(true));
324 inner.finished.push(req_id.clone());
325 }
326 let mut entry = default_entry();
327 entry.request_id = req_id.clone();
328 if ev.params.get("redirectResponse").is_some() {
329 entry.redirect_from_request_id = Some(req_id.clone());
330 }
331 entry.frame_id = frame_id;
332 entry.loader_id = loader_id;
333 entry.resource_type = resource_type;
334 entry.request_url = request_url;
335 entry.method = method;
336 entry.initiator = initiator;
337 entry.request_headers = request_headers;
338 entry.timing.start_monotonic_ms = (timestamp * 1000.0) as u64;
339 if let Some(post) = request.get("postData").and_then(|v| v.as_str()) {
341 entry.request_post_data_present = true;
342 entry.request_post_data_size_bytes = Some(post.len());
343 if serde_json::from_str::<Value>(post).is_ok() {
344 entry
345 .hints
346 .insert("request_body_json_valid".into(), Value::Bool(true));
347 }
348 if let Ok(v) = serde_json::from_str::<Value>(post) {
349 if let Some(op) = v.get("operationName").and_then(|v| v.as_str()) {
350 entry.hints.insert(
351 "graphql_operation_name".into(),
352 Value::String(op.to_string()),
353 );
354 }
355 if v.get("query").is_some() {
356 entry.hints.insert(
357 "graphql_operation_type".into(),
358 Value::String("request".into()),
359 );
360 }
361 }
362 }
363 inner.slots.insert(req_id.clone(), NetworkSlot { entry });
364 return inner.main_request_id.as_deref() == Some(req_id.as_str());
365 }
366 "Network.responseReceived" => {
367 inner.last_activity = Some(Instant::now());
368 let req_id = string_field(&ev.params, "requestId");
369 let response = &ev.params["response"];
370 let status = response
371 .get("status")
372 .and_then(|v| v.as_u64())
373 .map(|s| s as u16);
374 let mime_type = response
375 .get("mimeType")
376 .and_then(|v| v.as_str())
377 .map(str::to_string);
378 let mut headers = headers_map(response.get("headers"));
379 if redact_headers {
380 redact_inplace(&mut headers);
381 }
382 let protocol = response
383 .get("protocol")
384 .and_then(|v| v.as_str())
385 .map(str::to_string);
386 let remote = response
387 .get("remoteIPAddress")
388 .and_then(|v| v.as_str())
389 .map(str::to_string);
390 if let Some(s) = inner.slots.get_mut(&req_id) {
391 s.entry.status = status;
392 s.entry.state = NetworkEntryState::Responded;
393 if mime_type.is_some() {
394 s.entry.mime_type = mime_type;
395 }
396 s.entry.response_headers = headers;
397 if let Some(p) = protocol {
398 s.entry.hints.insert("protocol".into(), Value::String(p));
399 }
400 if let Some(r) = remote {
401 s.entry
402 .hints
403 .insert("remote_address".into(), Value::String(r));
404 }
405 }
406 return inner.main_request_id.as_deref() == Some(req_id.as_str());
407 }
408 "Network.loadingFinished" => {
409 inner.last_activity = Some(Instant::now());
410 let req_id = string_field(&ev.params, "requestId");
411 let timestamp = ev
412 .params
413 .get("timestamp")
414 .and_then(|v| v.as_f64())
415 .unwrap_or(0.0);
416 let encoded = ev.params.get("encodedDataLength").and_then(|v| v.as_u64());
417 if let Some(s) = inner.slots.get_mut(&req_id) {
418 s.entry.timing.end_monotonic_ms = Some((timestamp * 1000.0) as u64);
419 s.entry.state = NetworkEntryState::Finished;
420 if let Some(n) = encoded {
421 s.entry
422 .hints
423 .insert("encoded_data_length".into(), Value::from(n));
424 }
425 }
426 let is_main = inner.main_request_id.as_deref() == Some(req_id.as_str());
427 inner.finished.push(req_id);
428 return is_main;
429 }
430 "Network.loadingFailed" => {
431 inner.last_activity = Some(Instant::now());
432 let req_id = string_field(&ev.params, "requestId");
433 let err = ev
434 .params
435 .get("errorText")
436 .and_then(|v| v.as_str())
437 .unwrap_or("")
438 .to_string();
439 if let Some(s) = inner.slots.get_mut(&req_id) {
440 s.entry.failure = Some(err);
441 s.entry.state = NetworkEntryState::Failed;
442 }
443 let is_main = inner.main_request_id.as_deref() == Some(req_id.as_str());
444 inner.finished.push(req_id);
445 return is_main;
446 }
447 "Network.requestServedFromCache" => {
448 inner.last_activity = Some(Instant::now());
449 let req_id = string_field(&ev.params, "requestId");
450 if let Some(s) = inner.slots.get_mut(&req_id) {
451 s.entry
452 .hints
453 .insert("served_from_cache".into(), Value::Bool(true));
454 }
455 }
456 _ => {}
457 }
458 false
459}
460
461fn string_field(params: &Value, key: &str) -> String {
462 params
463 .get(key)
464 .and_then(|v| v.as_str())
465 .unwrap_or("")
466 .to_string()
467}
468
469fn headers_map(value: Option<&Value>) -> BTreeMap<String, String> {
470 let mut out = BTreeMap::new();
471 if let Some(Value::Object(map)) = value {
472 for (k, v) in map {
473 if let Some(s) = v.as_str() {
474 out.insert(k.clone(), s.to_string());
475 } else {
476 out.insert(k.clone(), v.to_string());
477 }
478 }
479 }
480 out
481}
482
483fn redact_inplace(map: &mut BTreeMap<String, String>) {
484 for (name, value) in map.iter_mut() {
485 if redact::should_redact(name) {
486 *value = redact::REDACTED_VALUE.to_string();
487 }
488 }
489}
490
491fn summarize_entries(entries: &[NetworkEntry], redacted: bool) -> NetworkSummary {
492 let requests_total = entries.len();
493 let responses_total = entries.iter().filter(|e| e.status.is_some()).count();
494 let finished_total = entries
495 .iter()
496 .filter(|e| e.state == NetworkEntryState::Finished)
497 .count();
498 let failed_total = entries
499 .iter()
500 .filter(|e| e.state == NetworkEntryState::Failed || e.failure.is_some())
501 .count();
502 let pending_by_resource_type = pending_by_resource_type(entries.iter());
503 let inflight_total_at_capture = pending_by_resource_type.values().copied().sum();
504 let captured_body_files = entries.iter().filter(|e| e.body_file.is_some()).count();
505 NetworkSummary {
506 requests_total,
507 responses_total,
508 finished_total,
509 failed_total,
510 incomplete_total: inflight_total_at_capture,
511 inflight_total_at_capture,
512 pending_by_resource_type,
513 captured_body_files,
514 redacted,
515 }
516}
517
518fn pending_by_resource_type<'a, I>(entries: I) -> BTreeMap<String, usize>
519where
520 I: IntoIterator<Item = &'a NetworkEntry>,
521{
522 let mut out = BTreeMap::new();
523 for entry in entries {
524 if matches!(
525 entry.state,
526 NetworkEntryState::Pending | NetworkEntryState::Responded
527 ) {
528 *out.entry(entry.resource_type.clone()).or_insert(0) += 1;
529 }
530 }
531 out
532}
533
534pub struct ConsoleCollector {
537 inner: Arc<Mutex<Vec<ConsoleEvent>>>,
538 task: JoinHandle<()>,
539}
540
541impl ConsoleCollector {
542 pub fn start(conn: &Connection) -> Self {
543 let inner: Arc<Mutex<Vec<ConsoleEvent>>> = Arc::new(Mutex::new(Vec::new()));
544 let mut rx = conn.subscribe();
545 let inner_w = inner.clone();
546 let task = tokio::spawn(async move {
547 while let Ok(ev) = rx.recv().await {
548 if let Some(console_event) = map_console_event(&ev) {
549 inner_w.lock().await.push(console_event);
550 }
551 }
552 });
553 Self { inner, task }
554 }
555
556 pub async fn finish(self) -> ConsoleLog {
557 self.task.abort();
558 let events = std::mem::take(&mut *self.inner.lock().await);
559 ConsoleLog {
560 schema_version: super::console::CONSOLE_SCHEMA_VERSION,
561 events,
562 }
563 }
564}
565
566fn handle_ws_event(inner: &mut NetworkInner, ev: &crate::sdk::cdp::ws_client::CdpEvent) {
567 let now_ms = std::time::SystemTime::now()
568 .duration_since(std::time::UNIX_EPOCH)
569 .map(|d| d.as_millis() as u64)
570 .unwrap_or(0);
571 match ev.method.as_str() {
572 "Network.webSocketFrameSent" => {
573 let req_id = string_field(&ev.params, "requestId");
574 if req_id.is_empty() {
575 return;
576 }
577 let response = ev.params.get("response").cloned().unwrap_or_default();
578 let frame = serde_json::json!({
579 "type": "sent",
580 "opcode": response.get("opcode").and_then(|v| v.as_u64()).unwrap_or(1),
581 "mask": response.get("mask").and_then(|v| v.as_bool()).unwrap_or(false),
582 "payload": response.get("payloadData").and_then(|v| v.as_str()).unwrap_or(""),
583 "timestamp_epoch_ms": now_ms,
584 });
585 inner.ws_frames.entry(req_id).or_default().push(frame);
586 }
587 "Network.webSocketFrameReceived" => {
588 let req_id = string_field(&ev.params, "requestId");
589 if req_id.is_empty() {
590 return;
591 }
592 let response = ev.params.get("response").cloned().unwrap_or_default();
593 let frame = serde_json::json!({
594 "type": "received",
595 "opcode": response.get("opcode").and_then(|v| v.as_u64()).unwrap_or(1),
596 "mask": response.get("mask").and_then(|v| v.as_bool()).unwrap_or(false),
597 "payload": response.get("payloadData").and_then(|v| v.as_str()).unwrap_or(""),
598 "timestamp_epoch_ms": now_ms,
599 });
600 inner.ws_frames.entry(req_id).or_default().push(frame);
601 }
602 "Network.webSocketFrameError" => {
603 let req_id = string_field(&ev.params, "requestId");
604 if req_id.is_empty() {
605 return;
606 }
607 let error = ev
608 .params
609 .get("errorMessage")
610 .and_then(|v| v.as_str())
611 .unwrap_or("unknown");
612 let frame = serde_json::json!({
613 "type": "error",
614 "error": error,
615 "timestamp_epoch_ms": now_ms,
616 });
617 inner.ws_frames.entry(req_id).or_default().push(frame);
618 }
619 _ => {}
620 }
621}
622
623fn handle_sse_event(inner: &mut NetworkInner, ev: &crate::sdk::cdp::ws_client::CdpEvent) {
624 if ev.method != "Network.eventSourceMessageReceived" {
625 return;
626 }
627 let req_id = string_field(&ev.params, "requestId");
628 if req_id.is_empty() {
629 return;
630 }
631 let event = serde_json::json!({
632 "event_name": ev.params.get("eventName").and_then(|v| v.as_str()).unwrap_or(""),
633 "data": ev.params.get("data").and_then(|v| v.as_str()).unwrap_or(""),
634 "event_id": ev.params.get("eventId").and_then(|v| v.as_str()).unwrap_or(""),
635 "timestamp_monotonic_ms": ev.params.get("timestamp")
636 .and_then(|v| v.as_f64())
637 .map(|t| (t * 1000.0) as u64)
638 .unwrap_or(0),
639 });
640 inner.sse_events.entry(req_id).or_default().push(event);
641}
642
643fn map_console_event(ev: &crate::sdk::cdp::ws_client::CdpEvent) -> Option<ConsoleEvent> {
644 match ev.method.as_str() {
645 "Runtime.consoleAPICalled" => {
646 let level = match ev
647 .params
648 .get("type")
649 .and_then(|v| v.as_str())
650 .unwrap_or("log")
651 {
652 "debug" => ConsoleLevel::Debug,
653 "info" => ConsoleLevel::Info,
654 "warning" | "warn" => ConsoleLevel::Warn,
655 "error" => ConsoleLevel::Error,
656 _ => ConsoleLevel::Log,
657 };
658 let timestamp_epoch_ms = ev
659 .params
660 .get("timestamp")
661 .and_then(|v| v.as_f64())
662 .unwrap_or(0.0);
663 let text = format_args_array(ev.params.get("args"));
664 let (url, line_number) = stack_frame_origin(ev.params.get("stackTrace"));
665 Some(ConsoleEvent {
666 level,
667 timestamp_epoch_ms,
668 text,
669 source_url: url.map(|url| redact::redact_url(&url)),
670 line_number,
671 })
672 }
673 "Runtime.exceptionThrown" => {
674 let details = ev.params.get("exceptionDetails")?;
675 let text = details
676 .get("exception")
677 .and_then(|e| e.get("description").and_then(|v| v.as_str()))
678 .map(str::to_string)
679 .or_else(|| {
680 details
681 .get("text")
682 .and_then(|v| v.as_str())
683 .map(str::to_string)
684 })
685 .unwrap_or_else(|| "exception".to_string());
686 let timestamp_epoch_ms = ev
687 .params
688 .get("timestamp")
689 .and_then(|v| v.as_f64())
690 .unwrap_or(0.0);
691 let url = details
692 .get("url")
693 .and_then(|v| v.as_str())
694 .map(str::to_string);
695 let line_number = details
696 .get("lineNumber")
697 .and_then(|v| v.as_u64())
698 .map(|n| n as u32);
699 Some(ConsoleEvent {
700 level: ConsoleLevel::Exception,
701 timestamp_epoch_ms,
702 text,
703 source_url: url.map(|url| redact::redact_url(&url)),
704 line_number,
705 })
706 }
707 _ => None,
708 }
709}
710
711fn format_args_array(args: Option<&Value>) -> String {
712 let Some(arr) = args.and_then(|v| v.as_array()) else {
713 return String::new();
714 };
715 let mut parts = Vec::with_capacity(arr.len());
716 for a in arr {
717 if let Some(s) = a.get("value").and_then(|v| v.as_str()) {
718 parts.push(s.to_string());
719 } else if let Some(s) = a.get("description").and_then(|v| v.as_str()) {
720 parts.push(s.to_string());
721 } else if let Some(v) = a.get("value") {
722 parts.push(v.to_string());
723 }
724 }
725 parts.join(" ")
726}
727
728fn stack_frame_origin(stack: Option<&Value>) -> (Option<String>, Option<u32>) {
729 let Some(frames) = stack
730 .and_then(|s| s.get("callFrames"))
731 .and_then(|v| v.as_array())
732 else {
733 return (None, None);
734 };
735 let Some(first) = frames.first() else {
736 return (None, None);
737 };
738 let url = first
739 .get("url")
740 .and_then(|v| v.as_str())
741 .map(str::to_string);
742 let line = first
743 .get("lineNumber")
744 .and_then(|v| v.as_u64())
745 .map(|n| n as u32);
746 (url, line)
747}