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