1use 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
46const PAGE_SOURCES: [&str; 6] = [
48 "navigation",
49 "console",
50 "pageerror",
51 "requestfailed",
52 "dialog",
53 "download",
54];
55
56struct 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
74struct 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#[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
150pub 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 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 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
317async 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 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 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 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 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 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 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 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 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 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 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 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 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 pub fn path(&self) -> &Path {
760 &self.inner.settings.root
761 }
762
763 pub fn mode(&self) -> &str {
765 &self.inner.settings.mode
766 }
767
768 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 pub fn stopped(&self) -> bool {
779 self.inner.lock().stopped
780 }
781
782 pub fn checkpoints(&self) -> Vec<JsonObject> {
784 self.inner.lock().checkpoints.clone()
785 }
786
787 pub async fn checkpoint(&self, name: &str) -> Result<JsonObject, TraceRecordError> {
794 self.checkpoint_with(name, TraceCheckpointOptions::default())
795 .await
796 }
797
798 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 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 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 let _ = self
870 .inner
871 .record_locked(&mut state, TraceEvent::INTERACTION, &payload);
872 outcome
873 }
874
875 pub async fn stop(&self) -> Result<TraceResult, TraceRecordError> {
883 self.stop_with(TraceStopOptions::default()).await
884 }
885
886 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}