Skip to main content

browser_commander/traces/
recorder.rs

1//! The trace recorder (issue #87), native in Rust (issue #108).
2//!
3//! A caller starts a recorder over a page, names checkpoints and stops it, and
4//! gets the same bundle `js/src/traces/recorder.js` writes: the same records in
5//! the same order with the same fields, the same checkpoint members, the same
6//! manifest, and the same Links Notation export written as the trace records.
7//!
8//! Page activity (navigations, console output, page errors, failed requests,
9//! dialogs and downloads) arrives on the engine's [`TraceEngineEvent`] stream
10//! and is written by a background task as it happens. Every recorder call first
11//! waits for that task to write what had already arrived, so a record written
12//! by a call always follows the page activity that came before it.
13
14use std::fmt::Display;
15use std::fs;
16use std::future::Future;
17use std::path::{Path, PathBuf};
18use std::sync::{Arc, Mutex, MutexGuard, PoisonError, Weak};
19
20use futures::stream::BoxStream;
21use futures::StreamExt;
22use tokio::sync::{mpsc, oneshot};
23use tokio::task::JoinHandle;
24
25use super::assets::trace_assets;
26use super::bundle::{create_manifest, Bundle, ManifestParts, TraceLimits, TraceRecordError};
27use super::identity::TraceIdentity;
28use super::jsonfmt::{js_number, Json, JsonObject};
29use super::links::{LinksHeader, LinksSink};
30use super::mutation_stream::{collect_batches, MutationStream};
31use super::page::{with_deadline, TracePage};
32use super::redaction::{
33    normalize_privacy_options, redact_object, redact_url, NormalizedPrivacy, REDACTED,
34};
35use super::schema::{
36    TraceCheckpointReason, TraceDropReason, TraceEvent, TraceMode, TraceOutcome,
37    TRACE_EVENT_SOURCES, TRACE_SCHEMA_VERSION,
38};
39use crate::core::engine::TraceEngineEvent;
40
41pub use super::recorder_options::{
42    TraceCheckpointOptions, TraceDomOptions, TraceFailure, TraceInitialCheckpoint, TraceOptions,
43    TraceResult, TraceScreenshots, TraceStopOptions, DEFAULT_CAPTURE_TIMEOUT_MS,
44};
45
46/// Event sources reported by the engine rather than by the recorder's caller.
47const PAGE_SOURCES: [&str; 6] = [
48    "navigation",
49    "console",
50    "pageerror",
51    "requestfailed",
52    "dialog",
53    "download",
54];
55
56/// Settled once, when the trace starts.
57struct Settings {
58    mode: String,
59    engine: Option<String>,
60    dom: JsonObject,
61    mutations: bool,
62    live_state: bool,
63    capture: Json,
64    events: Vec<String>,
65    privacy: NormalizedPrivacy,
66    limits: TraceLimits,
67    screenshots: TraceScreenshots,
68    capture_timeout_ms: u64,
69    commander_version: Option<String>,
70    started_at: String,
71    root: PathBuf,
72}
73
74/// What changes while the trace runs; never held across an `await`.
75struct State {
76    bundle: Bundle,
77    identity: TraceIdentity,
78    stopped: bool,
79    result: Option<TraceResult>,
80    checkpoint_index: u32,
81    checkpoints: Vec<JsonObject>,
82    init_script: Option<String>,
83}
84
85type FlushRequest = oneshot::Sender<()>;
86
87struct Inner {
88    page: Arc<dyn TracePage>,
89    settings: Settings,
90    mutations: MutationStream,
91    state: Mutex<State>,
92    operations: tokio::sync::Mutex<()>,
93    flush: Option<mpsc::UnboundedSender<FlushRequest>>,
94    pump: Mutex<Option<JoinHandle<()>>>,
95}
96
97/// A running trace; cheap to clone, and every clone is the same trace.
98#[derive(Clone)]
99pub struct TraceRecorder {
100    inner: Arc<Inner>,
101}
102
103fn invalid(message: impl Into<String>) -> TraceRecordError {
104    TraceRecordError::Invalid(message.into())
105}
106
107fn normalize_mode(mode: &str) -> Result<String, TraceRecordError> {
108    let values = [
109        TraceMode::OFF,
110        TraceMode::CHECKPOINTS,
111        TraceMode::CONTINUOUS,
112        TraceMode::RETAIN_ON_FAILURE,
113    ];
114    if values.contains(&mode) {
115        Ok(mode.to_string())
116    } else {
117        Err(invalid(format!(
118            "trace mode must be one of {}",
119            values.join(", ")
120        )))
121    }
122}
123
124fn normalize_events(events: Option<&[String]>) -> Result<Vec<String>, TraceRecordError> {
125    let Some(events) = events else {
126        return Ok(TRACE_EVENT_SOURCES
127            .iter()
128            .map(|name| name.to_string())
129            .collect());
130    };
131    let mut unique: Vec<String> = Vec::new();
132    for name in events {
133        if !TRACE_EVENT_SOURCES.contains(&name.as_str()) {
134            return Err(invalid(format!(
135                "unknown trace event source \"{name}\"; expected one of {}",
136                TRACE_EVENT_SOURCES.join(", ")
137            )));
138        }
139        if !unique.contains(name) {
140            unique.push(name.clone());
141        }
142    }
143    Ok(unique)
144}
145
146fn strings(values: &[String]) -> Json {
147    Json::Array(values.iter().map(Json::from).collect())
148}
149
150/// Start recording `page`.
151///
152/// # Errors
153///
154/// Fails when the options are invalid or the bundle or export cannot be
155/// created.
156pub async fn start_trace(
157    page: Arc<dyn TracePage>,
158    options: TraceOptions,
159) -> Result<TraceRecorder, TraceRecordError> {
160    page.require_feature("tracing")?;
161    let mode = normalize_mode(&options.mode)?;
162    let dom = options.dom;
163    let mutations = match dom.mutations {
164        Some(false) => false,
165        Some(true) => true,
166        None => mode == TraceMode::CONTINUOUS,
167    };
168    let dom_json = JsonObject::new()
169        .with("html", dom.html)
170        .with("liveControlState", dom.live_control_state)
171        .with("liveState", dom.live_state)
172        .with("mutations", mutations)
173        .with("openShadowRoots", dom.open_shadow_roots);
174    let events = normalize_events(options.events.as_deref())?;
175    let privacy =
176        normalize_privacy_options(&options.privacy).map_err(|error| invalid(error.to_string()))?;
177    if let Some(links) = &options.links {
178        if links.output.as_os_str().is_empty() {
179            return Err(invalid("trace links require an output path"));
180        }
181    }
182    let engine = options.engine.clone().or_else(|| page.engine());
183    let identity = TraceIdentity::new(page.context_key(), page.page_key(), page.has_context());
184    let base_checkpoint = match options.initial_checkpoint.clone() {
185        Some(TraceInitialCheckpoint::Skip) => None,
186        Some(TraceInitialCheckpoint::Take) => Some("initial".to_string()),
187        Some(TraceInitialCheckpoint::Named(name)) => Some(name),
188        None => (mode == TraceMode::CONTINUOUS).then(|| "initial".to_string()),
189    };
190
191    let mut bundle = Bundle::open(
192        &options.output,
193        &options.limits,
194        options.strict,
195        options.clock,
196    )?;
197    let started_at = bundle.now_iso();
198    if let Some(links) = &options.links {
199        let header = LinksHeader {
200            bundle: bundle
201                .root
202                .file_name()
203                .map(|name| name.to_string_lossy().into_owned())
204                .unwrap_or_default(),
205            schema_version: Some(Json::from(TRACE_SCHEMA_VERSION)),
206            mode: Some(Json::from(mode.as_str())),
207            engine: Some(Json::from(engine.clone())),
208            started_at: Some(Json::from(started_at.as_str())),
209            commander_version: Some(Json::from(options.commander_version.clone())),
210        };
211        bundle.links = Some(LinksSink::open(links, header).map_err(invalid)?);
212    }
213
214    let capture = JsonObject::new()
215        .with("redactSelectors", strings(&privacy.redact_selectors))
216        .with("redactAttributes", strings(&privacy.redact_attributes))
217        .with("redacted", REDACTED)
218        .with("html", dom.html)
219        .with("liveControlState", dom.live_control_state)
220        .with("openShadowRoots", dom.open_shadow_roots)
221        .with("maxHtmlBytes", options.limits.max_html_bytes.unwrap_or(0));
222    let stream = MutationStream::new(
223        mutations,
224        &privacy.redact_selectors,
225        options.limits.max_queued_mutations,
226        dom.live_state,
227        options.capture_timeout_ms,
228    );
229
230    // The observers attach before anything is recorded, as in JavaScript.
231    let page_events = if PAGE_SOURCES
232        .iter()
233        .any(|source| events.iter().any(|e| e == source))
234    {
235        page.events().await
236    } else {
237        None
238    };
239    let (flush, flushes) = match page_events {
240        Some(_) => {
241            let (sender, receiver) = mpsc::unbounded_channel();
242            (Some(sender), Some(receiver))
243        }
244        None => (None, None),
245    };
246
247    let root = bundle.root.clone();
248    let inner = Arc::new(Inner {
249        page,
250        settings: Settings {
251            mode,
252            engine,
253            dom: dom_json,
254            mutations,
255            live_state: dom.live_state,
256            capture: Json::Object(capture),
257            events,
258            privacy,
259            limits: options.limits,
260            screenshots: options.screenshots,
261            capture_timeout_ms: options.capture_timeout_ms,
262            commander_version: options.commander_version,
263            started_at,
264            root,
265        },
266        mutations: stream,
267        state: Mutex::new(State {
268            bundle,
269            identity,
270            stopped: false,
271            result: None,
272            checkpoint_index: 0,
273            checkpoints: Vec::new(),
274            init_script: None,
275        }),
276        operations: tokio::sync::Mutex::new(()),
277        flush,
278        pump: Mutex::new(None),
279    });
280    if let (Some(events), Some(flushes)) = (page_events, flushes) {
281        let task = tokio::spawn(pump(Arc::downgrade(&inner), events, flushes));
282        *inner.pump.lock().unwrap_or_else(PoisonError::into_inner) = Some(task);
283    }
284
285    // Registered before the first record, so a navigation that starts in the
286    // same moment as the trace is still recorded from its first mutation.
287    match inner
288        .mutations
289        .install_persistent(inner.page.as_ref())
290        .await
291    {
292        Ok(identifier) => inner.lock().init_script = identifier,
293        Err(message) => inner.drop_problem(
294            TraceDropReason::CAPTURE_FAILED,
295            "mutation-recorder-init",
296            &message,
297        )?,
298    }
299
300    let start = JsonObject::new()
301        .with("mode", inner.settings.mode.as_str())
302        .with("engine", Json::from(inner.settings.engine.clone()))
303        .with("dom", inner.settings.dom.clone())
304        .with("events", strings(&inner.settings.events))
305        .with("initialCheckpoint", base_checkpoint.is_some());
306    inner.record(TraceEvent::TRACE_START, &start)?;
307    inner.install().await?;
308
309    if let Some(name) = base_checkpoint {
310        inner
311            .checkpoint(&name, "recorder", TraceCheckpointReason::INITIAL)
312            .await?;
313    }
314    Ok(TraceRecorder { inner })
315}
316
317/// Write page activity as it arrives, and answer flushes once everything that
318/// had arrived is written.
319async fn pump(
320    inner: Weak<Inner>,
321    mut events: BoxStream<'static, TraceEngineEvent>,
322    mut flushes: mpsc::UnboundedReceiver<FlushRequest>,
323) {
324    let mut open = true;
325    loop {
326        tokio::select! {
327            biased;
328            event = events.next(), if open => match event {
329                Some(event) => match inner.upgrade() {
330                    Some(inner) => inner.on_page_event(event),
331                    None => return,
332                },
333                None => open = false,
334            },
335            request = flushes.recv() => match request {
336                Some(done) => {
337                    let _ = done.send(());
338                }
339                None => return,
340            },
341        }
342    }
343}
344
345impl Inner {
346    fn lock(&self) -> MutexGuard<'_, State> {
347        self.state.lock().unwrap_or_else(PoisonError::into_inner)
348    }
349
350    fn records(&self, source: &str) -> bool {
351        self.settings.events.iter().any(|name| name == source)
352    }
353
354    /// Wait until every page event that has arrived is written.
355    async fn flush(&self) {
356        if let Some(sender) = &self.flush {
357            let (done, written) = oneshot::channel();
358            if sender.send(done).is_ok() {
359                let _ = written.await;
360            }
361        }
362    }
363
364    fn stop_pump(&self) {
365        let task = self
366            .pump
367            .lock()
368            .unwrap_or_else(PoisonError::into_inner)
369            .take();
370        if let Some(task) = task {
371            task.abort();
372        }
373    }
374
375    fn record_locked(
376        &self,
377        state: &mut State,
378        kind: &str,
379        payload: &JsonObject,
380    ) -> Result<Option<JsonObject>, TraceRecordError> {
381        if state.stopped && kind != TraceEvent::TRACE_STOP {
382            return Ok(None);
383        }
384        // Who this happened to comes first, so a payload that knows better can
385        // say so without moving the field.
386        let mut event = JsonObject::new().with("kind", kind);
387        event.extend_from(&state.identity.owner());
388        event.extend_from(&redact_object(payload, &self.settings.privacy));
389        state.bundle.append_event(event, true)
390    }
391
392    fn record(
393        &self,
394        kind: &str,
395        payload: &JsonObject,
396    ) -> Result<Option<JsonObject>, TraceRecordError> {
397        self.record_locked(&mut self.lock(), kind, payload)
398    }
399
400    fn drop_problem(
401        &self,
402        reason: &str,
403        member: &str,
404        detail: &str,
405    ) -> Result<(), TraceRecordError> {
406        self.lock()
407            .bundle
408            .drop_record(reason, Some(member), Some(detail))
409    }
410
411    fn on_page_event(&self, event: TraceEngineEvent) {
412        let mut state = self.lock();
413        let (source, kind, payload) = match event {
414            TraceEngineEvent::Navigated { main_frame, url } => {
415                if !self.records("navigation") {
416                    return;
417                }
418                // A record's navigation is the one it happened during, so the
419                // counter moves on before the record of the move is written.
420                let navigation = if main_frame {
421                    state.identity.navigated()
422                } else {
423                    state.identity.navigation_id()
424                };
425                let payload = JsonObject::new()
426                    .with("phase", "framenavigated")
427                    .with("navigationId", navigation)
428                    .with("mainFrame", main_frame)
429                    .with("url", url);
430                ("navigation", TraceEvent::NAVIGATION, payload)
431            }
432            TraceEngineEvent::Console { level, text } => (
433                "console",
434                TraceEvent::CONSOLE,
435                JsonObject::new().with("level", level).with("text", text),
436            ),
437            TraceEngineEvent::PageError { message, stack } => (
438                "pageerror",
439                TraceEvent::PAGE_ERROR,
440                JsonObject::new()
441                    .with("message", message)
442                    .with("stack", Json::from(stack)),
443            ),
444            TraceEngineEvent::RequestFailed {
445                url,
446                method,
447                error_text,
448            } => (
449                "requestfailed",
450                TraceEvent::REQUEST_FAILED,
451                JsonObject::new()
452                    .with("url", url)
453                    .with("method", method)
454                    .with("failure", Json::from(error_text)),
455            ),
456            TraceEngineEvent::Dialog {
457                dialog_type,
458                message,
459            } => (
460                "dialog",
461                TraceEvent::DIALOG,
462                JsonObject::new()
463                    .with("type", dialog_type)
464                    .with("message", message),
465            ),
466            TraceEngineEvent::Download {
467                phase,
468                id,
469                suggested_filename,
470                path,
471                checksum,
472                bytes,
473                url,
474                failure,
475            } => (
476                "download",
477                TraceEvent::DOWNLOAD,
478                JsonObject::new()
479                    .with("phase", phase)
480                    .with("id", Json::from(id))
481                    .with("suggestedFilename", Json::from(suggested_filename))
482                    .with("path", Json::from(path))
483                    .with("checksum", Json::from(checksum))
484                    .with("bytes", Json::from(bytes))
485                    .with("url", Json::from(url))
486                    .with("failure", Json::from(failure)),
487            ),
488        };
489        if self.records(source) {
490            // Nobody is waiting on an observer's record; a strict trace's
491            // failure surfaces at the next checkpoint or at stop.
492            let _ = self.record_locked(&mut state, kind, &payload);
493        }
494    }
495
496    async fn install(&self) -> Result<(), TraceRecordError> {
497        match self.mutations.install(self.page.as_ref()).await {
498            Ok(()) => Ok(()),
499            Err(message) => self.drop_problem(
500                TraceDropReason::CAPTURE_FAILED,
501                "mutation-recorder",
502                &message,
503            ),
504        }
505    }
506
507    /// Write what the page queued as the interval that ends at `index`.
508    async fn drain(&self, index: u32) -> Result<(), TraceRecordError> {
509        let frames = match self.mutations.drain(self.page.as_ref()).await {
510            Ok(None) => return Ok(()),
511            Ok(Some(frames)) => frames,
512            Err(message) => {
513                return self.drop_problem(TraceDropReason::CAPTURE_FAILED, "mutations", &message)
514            }
515        };
516        let mut state = self.lock();
517        let drained = collect_batches(&frames, &state.identity.owner());
518        if drained.over != 0.0 && !drained.over.is_nan() {
519            let detail = format!(
520                "{} records over the in-page queue limit",
521                js_number(drained.over)
522            );
523            state.bundle.drop_record(
524                TraceDropReason::SIZE_LIMIT,
525                Some("mutations"),
526                Some(&detail),
527            )?;
528        }
529        if let Some(member) = state.bundle.write_mutations(index, &drained.batches)? {
530            let payload = JsonObject::new()
531                .with("member", member)
532                .with("batches", drained.batches.len())
533                .with("checkpoint", index)
534                .with("frames", drained.frames);
535            self.record_locked(&mut state, TraceEvent::MUTATIONS, &payload)?;
536        }
537        Ok(())
538    }
539
540    async fn screenshot(&self, reason: &str) -> Result<Option<Vec<u8>>, TraceRecordError> {
541        let wanted = match self.settings.screenshots {
542            TraceScreenshots::Off => false,
543            TraceScreenshots::Checkpoints => true,
544            TraceScreenshots::OnlyOnFailure => reason == TraceCheckpointReason::FAILURE,
545        };
546        if !wanted {
547            return Ok(None);
548        }
549        let shot = with_deadline(
550            self.page.screenshot(),
551            self.settings.capture_timeout_ms,
552            "trace screenshot",
553        )
554        .await;
555        match shot {
556            Ok(shot) => Ok(shot),
557            Err(message) => {
558                self.drop_problem(TraceDropReason::CAPTURE_FAILED, "screenshot", &message)?;
559                Ok(None)
560            }
561        }
562    }
563
564    async fn checkpoint(
565        &self,
566        name: &str,
567        actor: &str,
568        reason: &str,
569    ) -> Result<JsonObject, TraceRecordError> {
570        let index = {
571            let mut state = self.lock();
572            if state.stopped {
573                return Err(TraceRecordError::Stopped);
574            }
575            state.checkpoint_index += 1;
576            state.checkpoint_index
577        };
578
579        // Drained first, so the batches belong to the interval that ended here.
580        self.drain(index.saturating_sub(1)).await?;
581
582        let capture = with_deadline(
583            self.page.evaluate_function(
584                &trace_assets().capture.capture_snapshot,
585                &self.settings.capture,
586            ),
587            self.settings.capture_timeout_ms,
588            "trace checkpoint capture",
589        )
590        .await;
591        let captured = match capture {
592            Ok(captured) => Some(captured),
593            Err(message) => {
594                let dropped = if message.to_lowercase().contains("closed") {
595                    TraceDropReason::PAGE_CLOSED
596                } else {
597                    TraceDropReason::CAPTURE_FAILED
598                };
599                self.drop_problem(dropped, &format!("checkpoints/{index}"), &message)?;
600                None
601            }
602        };
603
604        let shot = self.screenshot(reason).await?;
605        let privacy = &self.settings.privacy;
606        let state = captured
607            .as_ref()
608            .and_then(|captured| captured.get("state"))
609            .filter(|state| state.truthy())
610            .and_then(Json::as_object)
611            .map(|state| redact_object(state, privacy));
612        let html = captured
613            .as_ref()
614            .and_then(|captured| captured.get("html"))
615            .and_then(Json::as_str);
616        let named = state.as_ref().map(|state| {
617            let mut named = state.clone();
618            named.insert("name", name);
619            named.insert("actor", actor);
620            named.insert("reason", reason);
621            named
622        });
623
624        // A block, not `drop()`: the guard must be gone before the `await`.
625        let entry = {
626            let mut locked = self.lock();
627            let members =
628                locked
629                    .bundle
630                    .write_checkpoint(index, html, named.as_ref(), shot.as_deref())?;
631            let mut entry = JsonObject::new()
632                .with("index", index)
633                .with("name", name)
634                .with("actor", actor)
635                .with("reason", reason);
636            match state.as_ref().map(|state| state.get("url")) {
637                None => entry.insert("url", Json::Null),
638                Some(Some(Json::String(url))) => entry.insert("url", redact_url(url, privacy)),
639                Some(Some(other)) => entry.insert("url", other.clone()),
640                // `url: undefined`, which JSON leaves out.
641                Some(None) => {}
642            }
643            entry.insert(
644                "truncated",
645                captured
646                    .as_ref()
647                    .and_then(|captured| captured.get("truncated"))
648                    .is_some_and(Json::truthy),
649            );
650            entry.insert("members", members);
651            locked.checkpoints.push(entry.clone());
652            self.record_locked(&mut locked, TraceEvent::CHECKPOINT, &entry)?;
653            entry
654        };
655
656        // The init script covers every document created from here on; this
657        // covers one that was created without it.
658        self.install().await?;
659        Ok(entry)
660    }
661
662    async fn stop(&self, options: TraceStopOptions) -> Result<TraceResult, TraceRecordError> {
663        if let Some(result) = self.lock().result.clone() {
664            return Ok(result);
665        }
666        let index = self.lock().checkpoint_index;
667        self.drain(index).await?;
668        let init_script = {
669            let mut state = self.lock();
670            if let Some(error) = &options.error {
671                let payload = JsonObject::new()
672                    .with("message", error.message.as_str())
673                    .with("stack", Json::from(error.stack.clone()))
674                    .with("fatal", true);
675                self.record_locked(&mut state, TraceEvent::PAGE_ERROR, &payload)?;
676            }
677            let payload = JsonObject::new().with("discarded", options.discard);
678            self.record_locked(&mut state, TraceEvent::TRACE_STOP, &payload)?;
679            state.stopped = true;
680            state.init_script.take()
681        };
682
683        // The documents that exist stop observing; then the observers detach.
684        let _ = self.mutations.stop(self.page.as_ref()).await;
685        if let Some(identifier) = init_script {
686            let _ = self.page.remove_init_script(&identifier).await;
687        }
688        self.stop_pump();
689
690        let settings = &self.settings;
691        let mut state = self.lock();
692        let stopped_at = state.bundle.now_iso();
693        let manifest = state.bundle.close(create_manifest(ManifestParts {
694            mode: settings.mode.clone(),
695            outcome: TraceOutcome::COMPLETE.to_string(),
696            started_at: Some(settings.started_at.clone()),
697            stopped_at: Some(stopped_at),
698            commander_version: settings.commander_version.clone(),
699            engine: settings.engine.clone(),
700            events: settings.events.iter().map(Json::from).collect(),
701            dom: settings.dom.clone(),
702            replay: JsonObject::new()
703                .with("checkpoints", true)
704                .with("mutations", settings.mutations)
705                .with("childListPositions", settings.mutations)
706                .with("liveState", settings.mutations && settings.live_state)
707                .with("identifiers", true),
708            privacy: JsonObject::new()
709                .with(
710                    "redactSelectors",
711                    strings(&settings.privacy.redact_selectors),
712                )
713                .with(
714                    "redactAttributes",
715                    strings(&settings.privacy.redact_attributes),
716                )
717                .with(
718                    "redactQueryParams",
719                    strings(&settings.privacy.redact_query_params),
720                )
721                .with("hasCallback", settings.privacy.redact.is_some()),
722            limits: settings.limits.to_json(),
723            counts: JsonObject::new(),
724        }))?;
725
726        // Closed after the manifest: its last link reports the settled outcome,
727        // and the control diffs are read back out of the finished bundle.
728        let mut sink = state.bundle.links.take();
729        if let Some(sink) = sink.as_mut() {
730            sink.close(&manifest, Some(&settings.root));
731        }
732        let mut problems = state.bundle.problems.clone();
733        if let Some(sink) = &sink {
734            problems.extend(sink.problems.iter().cloned());
735        }
736        let mut result = TraceResult {
737            path: settings.root.clone(),
738            manifest,
739            checkpoints: state.checkpoints.clone(),
740            problems,
741            links: sink.as_ref().map(|sink| sink.path.clone()),
742            discarded: false,
743        };
744        if options.discard {
745            let _ = fs::remove_dir_all(&settings.root);
746            // An export of a bundle that no longer exists points at nothing.
747            if let Some(sink) = sink.as_mut() {
748                sink.discard();
749            }
750            result.discarded = true;
751        }
752        state.result = Some(result.clone());
753        Ok(result)
754    }
755}
756
757impl TraceRecorder {
758    /// The bundle directory.
759    pub fn path(&self) -> &Path {
760        &self.inner.settings.root
761    }
762
763    /// The trace's mode.
764    pub fn mode(&self) -> &str {
765        &self.inner.settings.mode
766    }
767
768    /// The Links Notation export, when one is being written.
769    pub fn links(&self) -> Option<PathBuf> {
770        let state = self.inner.lock();
771        match &state.result {
772            Some(result) => result.links.clone(),
773            None => state.bundle.links.as_ref().map(|sink| sink.path.clone()),
774        }
775    }
776
777    /// Whether [`TraceRecorder::stop`] has run.
778    pub fn stopped(&self) -> bool {
779        self.inner.lock().stopped
780    }
781
782    /// The checkpoints taken so far.
783    pub fn checkpoints(&self) -> Vec<JsonObject> {
784        self.inner.lock().checkpoints.clone()
785    }
786
787    /// Capture a checkpoint taken by the automation.
788    ///
789    /// # Errors
790    ///
791    /// Fails after [`TraceRecorder::stop`], or when a strict trace drops
792    /// something.
793    pub async fn checkpoint(&self, name: &str) -> Result<JsonObject, TraceRecordError> {
794        self.checkpoint_with(name, TraceCheckpointOptions::default())
795            .await
796    }
797
798    /// Capture a checkpoint; `{index, name, actor, reason, url, truncated, members}`.
799    ///
800    /// # Errors
801    ///
802    /// As [`TraceRecorder::checkpoint`].
803    pub async fn checkpoint_with(
804        &self,
805        name: &str,
806        options: TraceCheckpointOptions,
807    ) -> Result<JsonObject, TraceRecordError> {
808        let _operation = self.inner.operations.lock().await;
809        self.inner.flush().await;
810        let actor = options.actor.as_deref().unwrap_or("automation");
811        let reason = options
812            .reason
813            .as_deref()
814            .unwrap_or(TraceCheckpointReason::CHECKPOINT);
815        self.inner.checkpoint(name, actor, reason).await
816    }
817
818    /// Record something a caller cares about on the same timeline; the record
819    /// written, or `None` once the trace has stopped.
820    ///
821    /// # Errors
822    ///
823    /// Fails when a strict trace drops the record.
824    pub async fn event(
825        &self,
826        name: &str,
827        data: JsonObject,
828    ) -> Result<Option<JsonObject>, TraceRecordError> {
829        let _operation = self.inner.operations.lock().await;
830        self.inner.flush().await;
831        let mut payload = JsonObject::new()
832            .with("action", name)
833            .with("actor", "caller");
834        payload.extend_from(&data);
835        self.inner.record(TraceEvent::INTERACTION, &payload)
836    }
837
838    /// Run one interaction and record it, as a traced commander method is.
839    ///
840    /// `target` is the selector or URL acted on; typed text never belongs
841    /// there. The interaction is recorded only when the trace records the
842    /// `interaction` source, and the work's own outcome is returned unchanged.
843    pub async fn traced<T, E, F>(&self, action: &str, target: Option<&str>, work: F) -> Result<T, E>
844    where
845        E: Display,
846        F: Future<Output = Result<T, E>>,
847    {
848        if !self.inner.records("interaction") {
849            return work.await;
850        }
851        let (started, action_id) = {
852            let mut state = self.inner.lock();
853            (state.bundle.now_ms(), state.identity.next_action_id())
854        };
855        let outcome = work.await;
856        self.inner.flush().await;
857        let mut state = self.inner.lock();
858        let duration = state.bundle.now_ms() - started;
859        let mut payload = JsonObject::new()
860            .with("actionId", action_id)
861            .with("action", action)
862            .with("target", Json::from(target))
863            .with("durationMs", duration)
864            .with("ok", outcome.is_ok());
865        if let Err(error) = &outcome {
866            payload.insert("error", error.to_string());
867        }
868        // The interaction's own result matters more than its record.
869        let _ = self
870            .inner
871            .record_locked(&mut state, TraceEvent::INTERACTION, &payload);
872        outcome
873    }
874
875    /// Stop recording and write the manifest; calling it again returns the
876    /// first result.
877    ///
878    /// # Errors
879    ///
880    /// Fails when the manifest cannot be written or a strict trace drops
881    /// something.
882    pub async fn stop(&self) -> Result<TraceResult, TraceRecordError> {
883        self.stop_with(TraceStopOptions::default()).await
884    }
885
886    /// [`TraceRecorder::stop`], discarding the bundle or recording the error
887    /// the run ended with.
888    ///
889    /// # Errors
890    ///
891    /// As [`TraceRecorder::stop`].
892    pub async fn stop_with(
893        &self,
894        options: TraceStopOptions,
895    ) -> Result<TraceResult, TraceRecordError> {
896        let _operation = self.inner.operations.lock().await;
897        self.inner.flush().await;
898        self.inner.stop(options).await
899    }
900}
901
902impl std::fmt::Debug for TraceRecorder {
903    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
904        f.debug_struct("TraceRecorder")
905            .field("path", &self.inner.settings.root)
906            .field("mode", &self.inner.settings.mode)
907            .field("stopped", &self.stopped())
908            .finish()
909    }
910}
911
912impl Drop for Inner {
913    fn drop(&mut self) {
914        self.stop_pump();
915    }
916}