1use std::fmt;
9use std::fs::{self, File, OpenOptions};
10use std::io::{self, Write};
11use std::path::{Path, PathBuf};
12use std::sync::{Arc, OnceLock};
13use std::time::{Instant, SystemTime, UNIX_EPOCH};
14
15use super::jsonfmt::{Json, JsonObject};
16use super::links::LinksSink;
17use super::schema::{
18 sequence_name, TraceDropReason, TraceEvent, TraceFiles, TraceOutcome, TRACE_FORMAT,
19 TRACE_SCHEMA_VERSION,
20};
21
22pub const TRACE_FILE_MODE: u32 = 0o600;
24
25pub const TRACE_DIRECTORY_MODE: u32 = 0o700;
27
28pub const DEFAULT_MAX_BUNDLE_BYTES: u64 = 256 * 1024 * 1024;
30
31pub const DEFAULT_MAX_RESOURCE_BYTES: u64 = 32 * 1024 * 1024;
33
34#[derive(Debug, thiserror::Error)]
36pub enum TraceRecordError {
37 #[error(transparent)]
38 Engine(#[from] crate::core::engine::EngineError),
39 #[error("{0}")]
41 Invalid(String),
42 #[error("{0}")]
44 Io(#[from] io::Error),
45 #[error("{0}")]
47 Dropped(String),
48 #[error("this trace has already been stopped")]
50 Stopped,
51}
52
53#[derive(Debug, Clone, Default, PartialEq, Eq)]
55pub struct TraceProblem {
56 pub reason: Option<String>,
58 pub member: Option<String>,
60 pub detail: Option<String>,
62}
63
64#[derive(Debug, Clone, Default, PartialEq, Eq)]
66pub struct TraceLimits {
67 pub max_bundle_bytes: Option<u64>,
69 pub max_resource_bytes: Option<u64>,
71 pub max_event_bytes: Option<u64>,
73 pub max_html_bytes: Option<u64>,
75 pub max_queued_mutations: Option<u64>,
77}
78
79impl TraceLimits {
80 pub fn to_json(&self) -> JsonObject {
82 let mut limits = JsonObject::new();
83 for (name, value) in [
84 ("maxBundleBytes", self.max_bundle_bytes),
85 ("maxResourceBytes", self.max_resource_bytes),
86 ("maxEventBytes", self.max_event_bytes),
87 ("maxHtmlBytes", self.max_html_bytes),
88 ("maxQueuedMutations", self.max_queued_mutations),
89 ] {
90 if let Some(value) = value {
91 limits.insert(name, value);
92 }
93 }
94 limits
95 }
96}
97
98pub type TraceClockFn = Arc<dyn Fn() -> f64 + Send + Sync>;
100
101#[derive(Clone)]
103pub struct TraceClock {
104 pub now: TraceClockFn,
106 pub monotonic: TraceClockFn,
108}
109
110impl Default for TraceClock {
111 fn default() -> Self {
112 static ORIGIN: OnceLock<Instant> = OnceLock::new();
113 ORIGIN.get_or_init(Instant::now);
114 Self {
115 now: Arc::new(|| {
116 SystemTime::now()
117 .duration_since(UNIX_EPOCH)
118 .map(|elapsed| elapsed.as_millis() as f64)
119 .unwrap_or(0.0)
120 }),
121 monotonic: Arc::new(|| {
122 ORIGIN.get_or_init(Instant::now).elapsed().as_secs_f64() * 1000.0
123 }),
124 }
125 }
126}
127
128impl TraceClock {
129 pub fn fixed(now: f64, monotonic: f64) -> Self {
131 Self {
132 now: Arc::new(move || now),
133 monotonic: Arc::new(move || monotonic),
134 }
135 }
136}
137
138impl fmt::Debug for TraceClock {
139 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
140 f.debug_struct("TraceClock").finish_non_exhaustive()
141 }
142}
143
144pub fn iso_timestamp(ms: f64) -> String {
146 let ms = if ms.is_finite() { ms.trunc() as i64 } else { 0 };
147 let days = ms.div_euclid(86_400_000);
148 let in_day = ms.rem_euclid(86_400_000);
149 let z = days + 719_468;
151 let era = z.div_euclid(146_097);
152 let doe = z.rem_euclid(146_097);
153 let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
154 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
155 let mp = (5 * doy + 2) / 153;
156 let day = doy - (153 * mp + 2) / 5 + 1;
157 let month = if mp < 10 { mp + 3 } else { mp - 9 };
158 let year = yoe + era * 400 + i64::from(month <= 2);
159 format!(
160 "{year:04}-{month:02}-{day:02}T{:02}:{:02}:{:02}.{:03}Z",
161 in_day / 3_600_000,
162 in_day / 60_000 % 60,
163 in_day / 1000 % 60,
164 in_day % 1000
165 )
166}
167
168pub fn trace_platform() -> String {
170 format!("{} {}", std::env::consts::OS, std::env::consts::ARCH)
171}
172
173pub fn trace_runtime() -> String {
175 format!("rust {}", env!("CARGO_PKG_VERSION"))
176}
177
178#[derive(Debug, Clone, Default)]
180pub(crate) struct ManifestParts {
181 pub mode: String,
182 pub outcome: String,
183 pub started_at: Option<String>,
184 pub stopped_at: Option<String>,
185 pub commander_version: Option<String>,
186 pub engine: Option<String>,
187 pub events: Vec<Json>,
188 pub dom: JsonObject,
189 pub replay: JsonObject,
190 pub privacy: JsonObject,
191 pub limits: JsonObject,
192 pub counts: JsonObject,
193}
194
195pub(crate) fn create_manifest(parts: ManifestParts) -> JsonObject {
197 let mut replay = JsonObject::new();
198 for name in [
199 "checkpoints",
200 "mutations",
201 "childListPositions",
202 "liveState",
203 "identifiers",
204 ] {
205 replay.insert(name, false);
206 }
207 replay.extend_from(&parts.replay);
208 let mut counts = JsonObject::new()
209 .with("checkpoints", 0)
210 .with("events", 0)
211 .with("mutationBatches", 0);
212 counts.extend_from(&parts.counts);
213 JsonObject::new()
214 .with("schemaVersion", TRACE_SCHEMA_VERSION)
215 .with("format", TRACE_FORMAT)
216 .with("mode", parts.mode)
217 .with("outcome", parts.outcome)
218 .with("startedAt", parts.started_at)
219 .with("stoppedAt", parts.stopped_at)
220 .with("commanderVersion", parts.commander_version)
221 .with("engine", parts.engine)
222 .with("browser", Json::Null)
223 .with("platform", trace_platform())
224 .with("runtime", trace_runtime())
225 .with("events", parts.events)
226 .with("dom", parts.dom)
227 .with("replay", replay)
228 .with("privacy", parts.privacy)
229 .with("limits", parts.limits)
230 .with("counts", counts)
231 .with("dropped", 0)
232}
233
234pub(crate) fn create_dir(path: &Path) -> io::Result<()> {
235 let mut builder = fs::DirBuilder::new();
236 builder.recursive(true);
237 #[cfg(unix)]
238 {
239 use std::os::unix::fs::DirBuilderExt;
240 builder.mode(TRACE_DIRECTORY_MODE);
241 }
242 builder.create(path)
243}
244
245pub(crate) fn open_private(path: &Path, append: bool) -> io::Result<File> {
246 let mut options = OpenOptions::new();
247 options.create(true);
248 if append {
249 options.append(true);
250 } else {
251 options.write(true).truncate(true);
252 }
253 #[cfg(unix)]
254 {
255 use std::os::unix::fs::OpenOptionsExt;
256 options.mode(TRACE_FILE_MODE);
257 }
258 options.open(path)
259}
260
261pub(crate) fn write_private(path: &Path, data: &[u8]) -> io::Result<()> {
262 open_private(path, false)?.write_all(data)
263}
264
265pub(crate) fn resolve_path(path: &Path) -> PathBuf {
267 let joined = if path.is_absolute() {
268 path.to_path_buf()
269 } else {
270 std::env::current_dir()
271 .unwrap_or_else(|_| PathBuf::from("/"))
272 .join(path)
273 };
274 let mut resolved = PathBuf::new();
275 for part in joined.components() {
276 match part {
277 std::path::Component::CurDir => {}
278 std::path::Component::ParentDir => {
279 resolved.pop();
280 }
281 other => resolved.push(other),
282 }
283 }
284 resolved
285}
286
287pub(crate) struct Bundle {
289 pub root: PathBuf,
290 events: Option<File>,
291 written: u64,
292 sequence: u64,
293 dropped: u64,
294 checkpoints: u64,
295 event_count: u64,
296 mutation_batches: u64,
297 pub problems: Vec<TraceProblem>,
298 strict: bool,
299 max_bundle: u64,
300 max_resource: u64,
301 max_event: u64,
302 clock: TraceClock,
303 pub links: Option<LinksSink>,
304}
305
306impl Bundle {
307 pub fn open(
308 output: &Path,
309 limits: &TraceLimits,
310 strict: bool,
311 clock: TraceClock,
312 ) -> Result<Self, TraceRecordError> {
313 if output.as_os_str().is_empty() {
314 return Err(TraceRecordError::Invalid(
315 "trace output must be a path".into(),
316 ));
317 }
318 let root = resolve_path(output);
319 let max_resource = limits
320 .max_resource_bytes
321 .unwrap_or(DEFAULT_MAX_RESOURCE_BYTES);
322 create_dir(&root)?;
323 let events = open_private(&root.join(TraceFiles::EVENTS), true)?;
324 Ok(Self {
325 root,
326 events: Some(events),
327 written: 0,
328 sequence: 0,
329 dropped: 0,
330 checkpoints: 0,
331 event_count: 0,
332 mutation_batches: 0,
333 problems: Vec::new(),
334 strict,
335 max_bundle: limits.max_bundle_bytes.unwrap_or(DEFAULT_MAX_BUNDLE_BYTES),
336 max_resource,
337 max_event: limits.max_event_bytes.unwrap_or(max_resource),
338 clock,
339 links: None,
340 })
341 }
342
343 pub fn now_iso(&self) -> String {
344 iso_timestamp((self.clock.now)())
345 }
346
347 pub fn now_ms(&self) -> f64 {
348 (self.clock.now)()
349 }
350
351 pub fn drop_record(
353 &mut self,
354 reason: &str,
355 member: Option<&str>,
356 detail: Option<&str>,
357 ) -> Result<(), TraceRecordError> {
358 self.dropped += 1;
359 self.problems.push(TraceProblem {
360 reason: Some(reason.to_string()),
361 member: member.map(str::to_string),
362 detail: detail.map(str::to_string),
363 });
364 if self.strict {
365 let detail = detail
366 .filter(|detail| !detail.is_empty())
367 .map(|detail| format!(" ({detail})"))
368 .unwrap_or_default();
369 return Err(TraceRecordError::Dropped(format!(
370 "trace {} dropped: {reason}{detail}",
371 member.unwrap_or("record")
372 )));
373 }
374 let mut record = JsonObject::new()
375 .with("kind", TraceEvent::DROPPED)
376 .with("reason", reason);
377 if let Some(member) = member {
378 record.insert("member", member);
379 }
380 if let Some(detail) = detail {
381 record.insert("detail", detail);
382 }
383 self.append_event(record, false)?;
384 Ok(())
385 }
386
387 pub fn append_event(
389 &mut self,
390 event: JsonObject,
391 retry: bool,
392 ) -> Result<Option<JsonObject>, TraceRecordError> {
393 self.sequence += 1;
394 let mut record = JsonObject::new()
395 .with("sequence", self.sequence)
396 .with("at", self.now_iso())
397 .with("monotonicMs", ((self.clock.monotonic)() + 0.5).floor());
398 record.extend_from(&event);
399 let line = format!("{}\n", Json::Object(record.clone()).to_compact());
400 let bytes = line.len() as u64;
401
402 let fits = self.written + bytes <= self.max_bundle && (!retry || bytes <= self.max_event);
403 if !fits {
404 if retry {
405 let detail = format!("{bytes} bytes");
406 self.drop_record(
407 TraceDropReason::SIZE_LIMIT,
408 Some(TraceFiles::EVENTS),
409 Some(&detail),
410 )?;
411 }
412 return Ok(None);
413 }
414
415 let outcome = match self.events.as_mut() {
416 Some(file) => file.write_all(line.as_bytes()),
417 None => Err(io::Error::other("the timeline is closed")),
418 };
419 match outcome {
420 Ok(()) => {
421 self.written += bytes;
422 self.event_count += 1;
423 if let Some(links) = self.links.as_mut() {
424 links.event(&record);
425 }
426 Ok(Some(record))
427 }
428 Err(error) => {
429 if retry {
430 self.drop_record(
431 TraceDropReason::WRITE_FAILED,
432 Some(TraceFiles::EVENTS),
433 Some(&error.to_string()),
434 )?;
435 }
436 Ok(None)
437 }
438 }
439 }
440
441 pub fn write_member(
443 &mut self,
444 member: &str,
445 data: &[u8],
446 ) -> Result<Option<String>, TraceRecordError> {
447 let length = data.len() as u64;
448 if length > self.max_resource || self.written + length > self.max_bundle {
449 let detail = format!("{length} bytes");
450 self.drop_record(TraceDropReason::SIZE_LIMIT, Some(member), Some(&detail))?;
451 return Ok(None);
452 }
453 let target = self.root.join(member);
454 let outcome = target
455 .parent()
456 .map_or(Ok(()), create_dir)
457 .and_then(|()| write_private(&target, data));
458 match outcome {
459 Ok(()) => {
460 self.written += length;
461 Ok(Some(member.to_string()))
462 }
463 Err(error) => {
464 self.drop_record(
465 TraceDropReason::WRITE_FAILED,
466 Some(member),
467 Some(&error.to_string()),
468 )?;
469 Ok(None)
470 }
471 }
472 }
473
474 pub fn write_checkpoint(
476 &mut self,
477 index: u32,
478 html: Option<&str>,
479 state: Option<&JsonObject>,
480 screenshot: Option<&[u8]>,
481 ) -> Result<JsonObject, TraceRecordError> {
482 let name = sequence_name(index);
483 let dir = TraceFiles::CHECKPOINTS_DIR;
484 let mut members = JsonObject::new();
485 if let Some(html) = html {
486 if let Some(member) =
487 self.write_member(&format!("{dir}/{name}.html"), html.as_bytes())?
488 {
489 members.insert("html", member);
490 }
491 }
492 if let Some(state) = state {
493 let body = format!("{}\n", Json::Object(state.clone()).to_pretty());
494 if let Some(member) =
495 self.write_member(&format!("{dir}/{name}.state.json"), body.as_bytes())?
496 {
497 members.insert("state", member);
498 }
499 }
500 if let Some(shot) = screenshot {
501 if let Some(member) = self.write_member(&format!("{dir}/{name}.png"), shot)? {
502 members.insert("screenshot", member);
503 }
504 }
505 self.checkpoints += 1;
506 Ok(members)
507 }
508
509 pub fn write_mutations(
511 &mut self,
512 index: u32,
513 batches: &[Json],
514 ) -> Result<Option<String>, TraceRecordError> {
515 if batches.is_empty() {
516 return Ok(None);
517 }
518 let member = format!(
519 "{}/{}.ndjson",
520 TraceFiles::MUTATIONS_DIR,
521 sequence_name(index)
522 );
523 let mut body = batches
524 .iter()
525 .map(Json::to_compact)
526 .collect::<Vec<_>>()
527 .join("\n");
528 body.push('\n');
529 let written = self.write_member(&member, body.as_bytes())?;
530 if written.is_some() {
531 self.mutation_batches += batches.len() as u64;
532 }
533 Ok(written)
534 }
535
536 pub fn close(&mut self, mut manifest: JsonObject) -> Result<JsonObject, TraceRecordError> {
538 let counts = JsonObject::new()
539 .with("checkpoints", self.checkpoints)
540 .with("events", self.event_count)
541 .with("mutationBatches", self.mutation_batches);
542 manifest.insert("counts", counts);
543 manifest.insert("dropped", self.dropped);
544 let complete =
545 manifest.get("outcome").and_then(Json::as_str) == Some(TraceOutcome::COMPLETE);
546 if self.dropped > 0 && complete {
547 manifest.insert("outcome", TraceOutcome::PARTIAL);
548 }
549 if let Some(mut file) = self.events.take() {
550 file.flush()?;
551 }
552 let body = format!("{}\n", Json::Object(manifest.clone()).to_pretty());
553 write_private(&self.root.join(TraceFiles::MANIFEST), body.as_bytes())?;
554 Ok(manifest)
555 }
556}
557
558#[cfg(test)]
559mod tests {
560 use super::*;
561
562 #[test]
563 fn timestamps_match_javascript() {
564 assert_eq!(
565 iso_timestamp(1_767_225_600_000.0),
566 "2026-01-01T00:00:00.000Z"
567 );
568 assert_eq!(iso_timestamp(0.0), "1970-01-01T00:00:00.000Z");
569 assert_eq!(iso_timestamp(951_782_400_123.9), "2000-02-29T00:00:00.123Z");
570 assert_eq!(iso_timestamp(-1.0), "1969-12-31T23:59:59.999Z");
571 }
572
573 #[test]
574 fn limits_record_only_what_was_set() {
575 let limits = TraceLimits {
576 max_queued_mutations: Some(100),
577 max_bundle_bytes: Some(1),
578 ..TraceLimits::default()
579 };
580 assert_eq!(
581 Json::Object(limits.to_json()).to_compact(),
582 r#"{"maxBundleBytes":1,"maxQueuedMutations":100}"#
583 );
584 }
585}