browser_commander/traces/
mutation_stream.rs1use std::cmp::Ordering;
9
10use super::assets::trace_assets;
11use super::jsonfmt::{Json, JsonObject};
12use super::page::{with_deadline, TracePage};
13
14pub const DEFAULT_MAX_QUEUED_MUTATIONS: u64 = 5000;
16
17#[derive(Debug, Clone, Default, PartialEq)]
19pub(crate) struct Drained {
20 pub batches: Vec<Json>,
22 pub over: f64,
24 pub frames: usize,
26}
27
28#[derive(Debug, Clone)]
30pub(crate) struct MutationStream {
31 enabled: bool,
32 recorder_options: Json,
33 capture_timeout_ms: u64,
34}
35
36impl MutationStream {
37 pub fn new(
39 enabled: bool,
40 redact_selectors: &[String],
41 max_queued: Option<u64>,
42 live_state: bool,
43 capture_timeout_ms: u64,
44 ) -> Self {
45 let assets = trace_assets();
46 let selectors = redact_selectors.iter().map(Json::from).collect::<Vec<_>>();
47 let recorder_options = JsonObject::new()
48 .with("globalName", assets.recorder_global.as_str())
49 .with("redactSelectors", selectors)
50 .with("redacted", assets.redacted.as_str())
51 .with(
52 "maxQueued",
53 max_queued.unwrap_or(DEFAULT_MAX_QUEUED_MUTATIONS),
54 )
55 .with("liveState", live_state);
56 Self {
57 enabled,
58 recorder_options: Json::Object(recorder_options),
59 capture_timeout_ms,
60 }
61 }
62
63 async fn evaluate_in_frames(
69 page: &dyn TracePage,
70 source: &str,
71 argument: &Json,
72 ) -> Result<Vec<Json>, String> {
73 Ok(vec![page.evaluate_function(source, argument).await?])
74 }
75
76 pub async fn install(&self, page: &dyn TracePage) -> Result<(), String> {
79 if !self.enabled {
80 return Ok(());
81 }
82 let source = &trace_assets().capture.install_mutation_recorder;
83 Self::evaluate_in_frames(page, source, &self.recorder_options)
84 .await
85 .map(|_| ())
86 }
87
88 pub async fn install_persistent(&self, page: &dyn TracePage) -> Result<Option<String>, String> {
93 if !self.enabled {
94 return Ok(None);
95 }
96 let source = &trace_assets().capture.install_mutation_recorder;
97 page.add_init_script(source, &self.recorder_options).await
98 }
99
100 pub async fn drain(&self, page: &dyn TracePage) -> Result<Option<Vec<Json>>, String> {
103 if !self.enabled {
104 return Ok(None);
105 }
106 let assets = trace_assets();
107 let global = Json::from(assets.recorder_global.as_str());
108 with_deadline(
109 Self::evaluate_in_frames(page, &assets.capture.drain_mutations, &global),
110 self.capture_timeout_ms,
111 "trace mutation drain",
112 )
113 .await
114 .map(Some)
115 }
116
117 pub async fn stop(&self, page: &dyn TracePage) -> Result<(), String> {
119 if !self.enabled {
120 return Ok(());
121 }
122 let assets = trace_assets();
123 let global = Json::from(assets.recorder_global.as_str());
124 Self::evaluate_in_frames(page, &assets.capture.stop_mutation_recorder, &global)
125 .await
126 .map(|_| ())
127 }
128}
129
130fn queued_at(batch: &Json) -> f64 {
132 batch.get("at").and_then(Json::as_f64).unwrap_or(0.0)
133}
134
135pub(crate) fn collect_batches(drained: &[Json], owner: &JsonObject) -> Drained {
140 let mut batches = Vec::new();
141 let mut over = 0.0;
142 for frame in drained {
143 over += frame.get("dropped").and_then(Json::as_f64).unwrap_or(0.0);
144 let frame_batches = frame.get("batches").and_then(Json::as_array);
145 for batch in frame_batches.into_iter().flatten() {
146 let mut record = owner.clone();
147 if let Some(fields) = batch.as_object() {
148 record.extend_from(fields);
149 }
150 batches.push(Json::Object(record));
151 }
152 }
153 batches.sort_by(|left, right| {
156 queued_at(left)
157 .partial_cmp(&queued_at(right))
158 .unwrap_or(Ordering::Equal)
159 });
160 Drained {
161 batches,
162 over,
163 frames: drained.len(),
164 }
165}
166
167#[cfg(test)]
168mod tests {
169 use super::*;
170
171 #[test]
172 fn batches_are_owned_and_ordered() {
173 let frame = Json::parse(
174 r#"{"batches":[{"at":2,"records":[]},{"navigationId":"nav-9","at":1}],"dropped":3}"#,
175 )
176 .unwrap();
177 let owner = JsonObject::new()
178 .with("traceId", "trace-1")
179 .with("navigationId", "nav-1");
180 let drained = collect_batches(&[frame], &owner);
181 assert_eq!(drained.over, 3.0);
182 assert_eq!(drained.frames, 1);
183 assert_eq!(
184 drained
185 .batches
186 .iter()
187 .map(Json::to_compact)
188 .collect::<Vec<_>>(),
189 vec![
190 r#"{"traceId":"trace-1","navigationId":"nav-9","at":1}"#,
191 r#"{"traceId":"trace-1","navigationId":"nav-1","at":2,"records":[]}"#,
192 ]
193 );
194 }
195
196 #[test]
197 fn the_recorder_options_match_javascript() {
198 let stream = MutationStream::new(true, &["[data-private]".to_string()], None, true, 0);
199 assert_eq!(
200 stream.recorder_options.to_compact(),
201 r#"{"globalName":"__browserCommanderTrace__","redactSelectors":["[data-private]"],"redacted":"[redacted]","maxQueued":5000,"liveState":true}"#
202 );
203 }
204}