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