1use std::borrow::Cow;
56use std::fmt::Write as _;
57use std::future::IntoFuture;
58use std::sync::{Mutex, MutexGuard, OnceLock, PoisonError};
59use std::time::{Duration, Instant, SystemTime};
60
61use opentelemetry::trace::{
62 Span as _, SpanBuilder, SpanKind, Status, TraceContextExt as _, Tracer as _,
63 TracerProvider as _,
64};
65use opentelemetry::{Context, KeyValue};
66use opentelemetry_sdk::trace::{SdkTracer, SdkTracerProvider};
67
68pub const MAX_PHASES: usize = 256;
72
73pub const READY: &str = "ready";
75
76const ROOT_SPAN: &str = "boot";
77const LOG_TARGET: &str = "otel_bootstrap::boot";
78
79#[derive(Debug, Clone, Copy, PartialEq, Eq)]
81pub enum Outcome {
82 Ok,
84 Failed,
86}
87
88impl Outcome {
89 #[must_use]
91 pub fn as_str(self) -> &'static str {
92 match self {
93 Self::Ok => "ok",
94 Self::Failed => "failed",
95 }
96 }
97}
98
99#[derive(Debug, Clone, PartialEq, Eq)]
101pub struct Phase {
102 pub name: Cow<'static, str>,
103 pub outcome: Outcome,
104 pub start: SystemTime,
106 pub end: SystemTime,
108 pub took: Duration,
110 pub at: Duration,
113}
114
115struct Recorded {
116 phase: Phase,
117 exported: bool,
118}
119
120#[derive(Clone)]
121struct Exporter {
122 tracer: SdkTracer,
123 root: Context,
124}
125
126#[derive(Default)]
127struct State {
128 service: Option<String>,
129 phases: Vec<Recorded>,
130 dropped: u64,
131 ready: Option<(SystemTime, Duration)>,
132 exporter: Option<Exporter>,
133 root_ended: bool,
134}
135
136pub struct Timeline {
138 origin: Instant,
139 origin_wall: SystemTime,
140 state: Mutex<State>,
141 #[cfg(test)]
142 echoed: Mutex<Vec<String>>,
143}
144
145impl Default for Timeline {
146 fn default() -> Self {
147 Self::new()
148 }
149}
150
151impl Timeline {
152 #[must_use]
155 pub fn new() -> Self {
156 Self::starting_at(Instant::now())
157 }
158
159 fn starting_at(origin: Instant) -> Self {
160 let now = Instant::now();
161 let origin_wall = SystemTime::now()
162 .checked_sub(now.saturating_duration_since(origin))
163 .unwrap_or_else(SystemTime::now);
164 Self {
165 origin,
166 origin_wall,
167 state: Mutex::new(State::default()),
168 #[cfg(test)]
169 echoed: Mutex::new(Vec::new()),
170 }
171 }
172
173 pub fn global() -> &'static Self {
176 static GLOBAL: OnceLock<Timeline> = OnceLock::new();
177 GLOBAL.get_or_init(|| {
178 let now = Instant::now();
179 let origin = process_age()
180 .and_then(|age| now.checked_sub(age))
181 .unwrap_or(now);
182 Self::starting_at(origin)
183 })
184 }
185
186 pub fn set_service_name(&self, name: &str) {
190 self.lock().service = Some(name.to_owned());
191 }
192
193 pub fn phase(&self, name: impl Into<Cow<'static, str>>) -> PhaseGuard<'_> {
197 PhaseGuard {
198 timeline: self,
199 name: Some(name.into()),
200 start: Instant::now(),
201 }
202 }
203
204 pub async fn time<F: IntoFuture>(
207 &self,
208 name: impl Into<Cow<'static, str>>,
209 work: F,
210 ) -> F::Output {
211 let phase = self.phase(name);
212 let output = work.into_future().await;
213 phase.finish();
214 output
215 }
216
217 pub async fn try_time<T, E, F>(
220 &self,
221 name: impl Into<Cow<'static, str>>,
222 work: F,
223 ) -> Result<T, E>
224 where
225 F: IntoFuture<Output = Result<T, E>>,
226 {
227 let phase = self.phase(name);
228 let output = work.into_future().await;
229 match output {
230 Ok(_) => phase.finish(),
231 Err(_) => phase.fail(),
232 }
233 output
234 }
235
236 pub fn ready(&self) {
239 let end = Instant::now();
240 let at = end.saturating_duration_since(self.origin);
241 let wall = self.wall(end);
242 let mut state = self.lock();
243 if state.ready.is_some() {
244 return;
245 }
246 state.ready = Some((wall, at));
247 let end_root = state.take_root_to_end();
248 let service = state.service_label();
249 drop(state);
250
251 self.echo(&line(READY, Outcome::Ok, at, at, &service));
252 if let Some((exporter, count, dropped)) = end_root {
253 end_root_span(&exporter, wall, "ready", count, dropped);
254 tracing::info!(
255 target: LOG_TARGET,
256 took_ms = millis(at),
257 "boot ready"
258 );
259 }
260 }
261
262 pub fn flush(&self, provider: &SdkTracerProvider, service_name: &str) -> usize {
269 self.attach(provider, service_name).unwrap_or(0)
270 }
271
272 pub(crate) fn attach(&self, provider: &SdkTracerProvider, service_name: &str) -> Option<usize> {
275 let mut state = self.lock();
276 if state.exporter.is_some() {
277 return None;
278 }
279 if state.service.is_none() {
280 state.service = Some(service_name.to_owned());
281 }
282 let tracer = provider.tracer(concat!(env!("CARGO_PKG_NAME"), ".boot"));
283 let root_span = tracer.build_with_context(
284 SpanBuilder::from_name(ROOT_SPAN)
285 .with_kind(SpanKind::Internal)
286 .with_start_time(self.origin_wall),
287 &Context::new(),
288 );
289 let exporter = Exporter {
290 tracer,
291 root: Context::new().with_span(root_span),
292 };
293 state.exporter = Some(exporter.clone());
294 let pending: Vec<Phase> = state
295 .phases
296 .iter_mut()
297 .filter(|recorded| !recorded.exported)
298 .map(|recorded| {
299 recorded.exported = true;
300 recorded.phase.clone()
301 })
302 .collect();
303 let ready = state.ready;
304 let end_root = if ready.is_some() {
305 state.take_root_to_end()
306 } else {
307 None
308 };
309 drop(state);
310
311 for phase in &pending {
312 export(&exporter, phase);
313 }
314 if let (Some((wall, at)), Some((exporter, count, dropped))) = (ready, end_root) {
315 end_root_span(&exporter, wall, "ready", count, dropped);
316 tracing::info!(target: LOG_TARGET, took_ms = millis(at), "boot ready");
317 }
318 Some(pending.len())
319 }
320
321 pub(crate) fn close(&self) {
324 let now = self.wall(Instant::now());
325 let end_root = self.lock().take_root_to_end();
326 if let Some((exporter, count, dropped)) = end_root {
327 end_root_span(&exporter, now, "incomplete", count, dropped);
328 }
329 }
330
331 #[must_use]
333 pub fn phases(&self) -> Vec<Phase> {
334 self.lock()
335 .phases
336 .iter()
337 .map(|recorded| recorded.phase.clone())
338 .collect()
339 }
340
341 fn record(&self, name: Cow<'static, str>, start: Instant, end: Instant, outcome: Outcome) {
342 let phase = Phase {
343 name,
344 outcome,
345 start: self.wall(start),
346 end: self.wall(end),
347 took: end.saturating_duration_since(start),
348 at: end.saturating_duration_since(self.origin),
349 };
350 let mut state = self.lock();
351 let exporter = state.exporter.clone();
352 if state.phases.len() < MAX_PHASES {
353 state.phases.push(Recorded {
354 phase: phase.clone(),
355 exported: exporter.is_some(),
356 });
357 } else {
358 state.dropped += 1;
359 }
360 let service = state.service_label();
361 drop(state);
362
363 self.echo(&line(&phase.name, outcome, phase.took, phase.at, &service));
364 if let Some(exporter) = exporter {
365 export(&exporter, &phase);
366 }
367 }
368
369 fn wall(&self, at: Instant) -> SystemTime {
370 self.origin_wall + at.saturating_duration_since(self.origin)
371 }
372
373 fn lock(&self) -> MutexGuard<'_, State> {
374 self.state.lock().unwrap_or_else(PoisonError::into_inner)
375 }
376
377 fn echo(&self, line: &str) {
378 #[cfg(test)]
379 self.echoed.lock().unwrap().push(line.to_owned());
380 eprintln!("{line}");
381 }
382}
383
384impl State {
385 fn service_label(&self) -> String {
386 self.service
387 .clone()
388 .unwrap_or_else(|| default_service_name().to_owned())
389 }
390
391 fn take_root_to_end(&mut self) -> Option<(Exporter, usize, u64)> {
394 if self.root_ended {
395 return None;
396 }
397 let exporter = self.exporter.clone()?;
398 self.root_ended = true;
399 Some((exporter, self.phases.len(), self.dropped))
400 }
401}
402
403#[must_use = "a phase is recorded when its guard is finished or dropped"]
406pub struct PhaseGuard<'t> {
407 timeline: &'t Timeline,
408 name: Option<Cow<'static, str>>,
409 start: Instant,
410}
411
412impl PhaseGuard<'_> {
413 pub fn finish(mut self) {
415 self.end(Outcome::Ok);
416 }
417
418 pub fn fail(mut self) {
420 self.end(Outcome::Failed);
421 }
422
423 fn end(&mut self, outcome: Outcome) {
424 if let Some(name) = self.name.take() {
425 self.timeline
426 .record(name, self.start, Instant::now(), outcome);
427 }
428 }
429}
430
431impl Drop for PhaseGuard<'_> {
432 fn drop(&mut self) {
433 self.end(Outcome::Failed);
434 }
435}
436
437pub fn phase(name: impl Into<Cow<'static, str>>) -> PhaseGuard<'static> {
439 Timeline::global().phase(name)
440}
441
442pub async fn time<F: IntoFuture>(name: impl Into<Cow<'static, str>>, work: F) -> F::Output {
444 Timeline::global().time(name, work).await
445}
446
447pub async fn try_time<T, E, F>(name: impl Into<Cow<'static, str>>, work: F) -> Result<T, E>
449where
450 F: IntoFuture<Output = Result<T, E>>,
451{
452 Timeline::global().try_time(name, work).await
453}
454
455pub fn ready() {
457 Timeline::global().ready();
458}
459
460fn export(exporter: &Exporter, phase: &Phase) {
461 let mut span = exporter.tracer.build_with_context(
462 SpanBuilder::from_name(phase.name.clone())
463 .with_kind(SpanKind::Internal)
464 .with_start_time(phase.start)
465 .with_attributes([
466 KeyValue::new("boot.phase", phase.name.clone()),
467 KeyValue::new("boot.outcome", phase.outcome.as_str()),
468 KeyValue::new("boot.at_ms", millis_i64(phase.at)),
469 ]),
470 &exporter.root,
471 );
472 if phase.outcome == Outcome::Failed {
473 span.set_status(Status::error("boot phase did not complete"));
474 }
475 span.end_with_timestamp(phase.end);
476
477 let (took_ms, at_ms) = (millis(phase.took), millis(phase.at));
478 match phase.outcome {
479 Outcome::Ok => tracing::info!(
480 target: LOG_TARGET,
481 phase = %phase.name,
482 outcome = phase.outcome.as_str(),
483 took_ms,
484 at_ms,
485 "boot phase complete"
486 ),
487 Outcome::Failed => tracing::warn!(
488 target: LOG_TARGET,
489 phase = %phase.name,
490 outcome = phase.outcome.as_str(),
491 took_ms,
492 at_ms,
493 "boot phase failed"
494 ),
495 }
496}
497
498fn end_root_span(
499 exporter: &Exporter,
500 end: SystemTime,
501 outcome: &'static str,
502 count: usize,
503 dropped: u64,
504) {
505 let root = exporter.root.span();
506 root.set_attribute(KeyValue::new("boot.outcome", outcome));
507 root.set_attribute(KeyValue::new(
508 "boot.phase_count",
509 i64::try_from(count).unwrap_or(i64::MAX),
510 ));
511 root.set_attribute(KeyValue::new(
512 "boot.phases_dropped",
513 i64::try_from(dropped).unwrap_or(i64::MAX),
514 ));
515 root.end_with_timestamp(end);
516}
517
518fn line(name: &str, outcome: Outcome, took: Duration, at: Duration, service: &str) -> String {
520 let mut out = String::from("boot phase=");
521 push_value(&mut out, name);
522 let _ = write!(
523 out,
524 " outcome={} took_ms={} at_ms={} service=",
525 outcome.as_str(),
526 millis(took),
527 millis(at)
528 );
529 push_value(&mut out, service);
530 out
531}
532
533fn push_value(out: &mut String, value: &str) {
535 let bare = !value.is_empty()
536 && value
537 .chars()
538 .all(|c| !c.is_whitespace() && !c.is_control() && c != '"' && c != '=' && c != '\\');
539 if bare {
540 out.push_str(value);
541 return;
542 }
543 out.push('"');
544 for c in value.chars() {
545 match c {
546 '"' => out.push_str("\\\""),
547 '\\' => out.push_str("\\\\"),
548 '\n' => out.push_str("\\n"),
549 c if c.is_control() => {
550 let _ = write!(out, "\\u{{{:x}}}", u32::from(c));
551 }
552 c => out.push(c),
553 }
554 }
555 out.push('"');
556}
557
558fn millis(duration: Duration) -> u64 {
559 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
560}
561
562fn millis_i64(duration: Duration) -> i64 {
563 i64::try_from(duration.as_millis()).unwrap_or(i64::MAX)
564}
565
566fn default_service_name() -> &'static str {
567 static NAME: OnceLock<String> = OnceLock::new();
568 NAME.get_or_init(|| {
569 std::env::var("OTEL_SERVICE_NAME")
570 .ok()
571 .filter(|name| !name.is_empty())
572 .or_else(|| {
573 std::env::current_exe().ok().and_then(|exe| {
574 exe.file_stem()
575 .and_then(|stem| stem.to_str())
576 .map(str::to_owned)
577 })
578 })
579 .unwrap_or_else(|| "unknown_service".to_owned())
580 })
581}
582
583#[cfg(target_os = "linux")]
588fn process_age() -> Option<Duration> {
589 let stat = std::fs::read_to_string("/proc/self/stat").ok()?;
590 let uptime = std::fs::read_to_string("/proc/uptime").ok()?;
591 process_age_from(&stat, &uptime)
592}
593
594#[cfg(not(target_os = "linux"))]
595fn process_age() -> Option<Duration> {
596 None
597}
598
599#[cfg_attr(not(target_os = "linux"), allow(dead_code))]
600fn process_age_from(stat: &str, uptime: &str) -> Option<Duration> {
601 const MILLIS_PER_TICK: u64 = 1000 / 100;
602 let fields = stat.rsplit_once(')')?.1;
606 let start_ticks: u64 = fields.split_whitespace().nth(19)?.parse().ok()?;
607 let uptime_secs: f64 = uptime.split_whitespace().next()?.parse().ok()?;
608 let uptime = Duration::try_from_secs_f64(uptime_secs).ok()?;
609 uptime.checked_sub(Duration::from_millis(
610 start_ticks.checked_mul(MILLIS_PER_TICK)?,
611 ))
612}
613
614#[cfg(test)]
615mod tests {
616 use super::*;
617
618 fn capturing() -> Timeline {
619 Timeline::default()
620 }
621
622 fn captured(timeline: &Timeline) -> Vec<String> {
623 timeline.echoed.lock().unwrap().clone()
624 }
625
626 #[test]
627 fn each_completed_phase_is_one_logfmt_line_on_stderr() {
628 let timeline = capturing();
629 timeline.set_service_name("orders");
630 let start = timeline.origin;
631 timeline.record(
632 "config-fetch".into(),
633 start + Duration::from_millis(200),
634 start + Duration::from_millis(1012),
635 Outcome::Ok,
636 );
637
638 assert_eq!(
639 captured(&timeline),
640 ["boot phase=config-fetch outcome=ok took_ms=812 at_ms=1012 service=orders"]
641 );
642 }
643
644 #[test]
645 fn ready_writes_the_whole_boot_once() {
646 let timeline = capturing();
647 timeline.set_service_name("orders");
648 timeline.ready();
649 timeline.ready();
650
651 let lines = captured(&timeline);
652 assert_eq!(lines.len(), 1, "{lines:?}");
653 let ready = &lines[0];
654 assert!(ready.starts_with("boot phase=ready outcome=ok took_ms="));
655 assert!(ready.ends_with(" service=orders"));
656 }
657
658 #[test]
659 fn values_that_would_not_parse_back_are_quoted() {
660 let line = line(
661 "load \"cache\"",
662 Outcome::Failed,
663 Duration::from_millis(5),
664 Duration::from_millis(9),
665 "a=b\\c\nd\u{1}",
666 );
667 assert_eq!(
668 line,
669 r#"boot phase="load \"cache\"" outcome=failed took_ms=5 at_ms=9 service="a=b\\c\nd\u{1}""#
670 );
671 assert_eq!(
672 super::line("", Outcome::Ok, Duration::ZERO, Duration::ZERO, "svc"),
673 r#"boot phase="" outcome=ok took_ms=0 at_ms=0 service=svc"#
674 );
675 }
676
677 #[test]
678 fn an_unfinished_guard_records_a_failure() {
679 let timeline = capturing();
680 drop(timeline.phase("dropped"));
681 timeline.phase("failed").fail();
682 timeline.phase("finished").finish();
683
684 let outcomes: Vec<_> = timeline
685 .phases()
686 .into_iter()
687 .map(|phase| (phase.name.into_owned(), phase.outcome))
688 .collect();
689 assert_eq!(
690 outcomes,
691 [
692 ("dropped".to_owned(), Outcome::Failed),
693 ("failed".to_owned(), Outcome::Failed),
694 ("finished".to_owned(), Outcome::Ok),
695 ]
696 );
697 }
698
699 #[tokio::test]
700 async fn futures_are_timed_and_errors_and_cancellation_fail_the_phase() {
701 let timeline = capturing();
702 timeline
703 .time("sleep", tokio::time::sleep(Duration::from_millis(20)))
704 .await;
705 let ok: Result<u8, &str> = timeline.try_time("ok", async { Ok(1) }).await;
706 let err: Result<u8, &str> = timeline.try_time("err", async { Err("no") }).await;
707 assert_eq!((ok, err), (Ok(1), Err("no")));
708 let cancelled = timeline.time("cancelled", std::future::pending::<()>());
709 let _ = tokio::time::timeout(Duration::from_millis(1), cancelled).await;
710
711 let phases = timeline.phases();
712 let outcomes: Vec<_> = phases.iter().map(|p| (&*p.name, p.outcome)).collect();
713 assert_eq!(
714 outcomes,
715 [
716 ("sleep", Outcome::Ok),
717 ("ok", Outcome::Ok),
718 ("err", Outcome::Failed),
719 ("cancelled", Outcome::Failed),
720 ]
721 );
722 assert!(phases[0].took >= Duration::from_millis(20));
723 assert_eq!(
724 phases[0].end.duration_since(phases[0].start).unwrap(),
725 phases[0].took
726 );
727 assert!(phases[1].at >= phases[0].at);
728 }
729
730 #[test]
731 fn without_telemetry_phases_are_echoed_and_buffered_up_to_the_bound() {
732 let timeline = capturing();
733 for _ in 0..MAX_PHASES + 3 {
734 timeline.phase("step").finish();
735 }
736 assert_eq!(timeline.phases().len(), MAX_PHASES);
737 assert_eq!(captured(&timeline).len(), MAX_PHASES + 3);
738 let state = timeline.lock();
739 assert_eq!(state.dropped, 3);
740 assert!(state.exporter.is_none());
741 assert!(state.phases.iter().all(|recorded| !recorded.exported));
742 }
743
744 #[test]
745 fn close_without_a_provider_does_nothing() {
746 let timeline = capturing();
747 timeline.close();
748 assert!(!timeline.lock().root_ended);
749 }
750
751 #[test]
752 fn the_service_falls_back_to_a_derived_name() {
753 let timeline = capturing();
754 timeline.phase("step").finish();
755 let expected = format!(" service={}", default_service_name());
756 assert!(captured(&timeline)[0].ends_with(&expected));
757 assert!(!default_service_name().is_empty());
758 }
759
760 #[test]
761 fn the_global_timeline_starts_no_later_than_its_first_use() {
762 let before = Instant::now();
763 let global = Timeline::global();
764 assert!(global.origin <= before || process_age().is_none());
767 assert!(std::ptr::eq(global, Timeline::global()));
768 }
769
770 #[test]
771 fn process_age_reads_start_ticks_after_the_command_name() {
772 let stat =
773 "42 (my (odd) cmd) S 1 42 42 0 -1 4194560 100 0 0 0 1 2 0 0 20 0 1 0 1000 1000 10 0";
774 assert_eq!(
775 process_age_from(stat, "15.50 30.00\n"),
776 Some(Duration::from_millis(5500))
777 );
778 assert_eq!(process_age_from(stat, "5.00 1.00\n"), None, "negative age");
779 assert_eq!(process_age_from("garbage", "1.0 1.0"), None);
780 assert_eq!(process_age_from(stat, "-1.0 1.0"), None);
781 assert_eq!(process_age_from(stat, ""), None);
782 }
783}