Skip to main content

browser_commander/traces/
mutation_stream.rs

1//! The continuous half of a trace (issue #93), in Rust (issue #108).
2//!
3//! `js/src/traces/mutation-stream.js` with the same in-page recorder: it is
4//! installed into the document that exists and registered as an init script for
5//! every document that follows, and each checkpoint drains what it queued. This
6//! module holds the page side; the recorder writes what it returns.
7
8use std::cmp::Ordering;
9
10use super::assets::trace_assets;
11use super::jsonfmt::{Json, JsonObject};
12use super::page::{with_deadline, TracePage};
13
14/// In-page records kept before the recorder starts counting drops instead.
15pub const DEFAULT_MAX_QUEUED_MUTATIONS: u64 = 5000;
16
17/// What one drain found.
18#[derive(Debug, Clone, Default, PartialEq)]
19pub(crate) struct Drained {
20    /// Every batch, owner first, ordered by when it was queued.
21    pub batches: Vec<Json>,
22    /// Records the page could not queue.
23    pub over: f64,
24    /// Frames that answered.
25    pub frames: usize,
26}
27
28/// The page side of the mutation stream.
29#[derive(Debug, Clone)]
30pub(crate) struct MutationStream {
31    enabled: bool,
32    recorder_options: Json,
33    capture_timeout_ms: u64,
34}
35
36impl MutationStream {
37    /// `enabled` is the resolved `dom.mutations`.
38    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    /// One result per frame that answered.
64    ///
65    /// Engines here evaluate in the main frame only, so this is always one
66    /// result; child frames still report through their own init script once a
67    /// later engine evaluates in them.
68    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    /// Install the recorder into the documents that already exist; the error
77    /// to drop as `mutation-recorder`, if any.
78    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    /// Register the recorder for every future document.
89    ///
90    /// `Ok(None)` when there is nothing to undo later; the error is dropped as
91    /// `mutation-recorder-init`.
92    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    /// What every frame had queued, under the capture deadline; `Ok(None)`
101    /// when the stream is off. [`collect_batches`] makes one interval of it.
102    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    /// Switch the recorder off in every document that has one.
118    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
130/// `at ?? 0`; the in-page recorder always writes a number there.
131fn queued_at(batch: &Json) -> f64 {
132    batch.get("at").and_then(Json::as_f64).unwrap_or(0.0)
133}
134
135/// Merge every frame's batches into one interval, ordered by queue time.
136///
137/// `owner` goes ahead of each batch, so a batch that names its own frame or
138/// navigation keeps its own.
139pub(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    // A stable sort, as `Array.prototype.sort` is: equal times keep the order
154    // the frames reported them in.
155    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}