opentelemetry/context.rs
1//! Execution-scoped context propagation.
2//!
3//! The `context` module provides mechanisms for propagating values across API boundaries and between
4//! logically associated execution units. It enables cross-cutting concerns to access their data in-process
5//! using a shared context object.
6//!
7//! # Main Types
8//!
9//! - [`Context`]: An immutable, execution-scoped collection of values.
10//!
11
12use crate::otel_warn;
13#[cfg(feature = "trace")]
14use crate::trace::context::SynchronizedSpan;
15use std::any::{Any, TypeId};
16use std::cell::RefCell;
17use std::collections::HashMap;
18use std::fmt;
19use std::hash::{BuildHasherDefault, Hasher};
20use std::marker::PhantomData;
21use std::sync::Arc;
22#[cfg(feature = "experimental_context_observer")]
23use std::sync::OnceLock;
24
25#[cfg(feature = "futures")]
26mod future_ext;
27
28#[cfg(feature = "futures")]
29pub use future_ext::{FutureExt, WithContext};
30
31#[cfg(feature = "experimental_context_observer")]
32pub use context_observer::*;
33
34thread_local! {
35 static CURRENT_CONTEXT: RefCell<ContextStack> = RefCell::new(ContextStack::default());
36}
37
38/// An execution-scoped collection of values.
39///
40/// A [`Context`] is a propagation mechanism which carries execution-scoped
41/// values across API boundaries and between logically associated execution
42/// units. Cross-cutting concerns access their data in-process using the same
43/// shared context object.
44///
45/// [`Context`]s are immutable, and their write operations result in the creation
46/// of a new context containing the original values and the new specified values.
47///
48/// ## Context state
49///
50/// Concerns can create and retrieve their local state in the current execution
51/// state represented by a context through the [`get`] and [`with_value`]
52/// methods. It is recommended to use application-specific types when storing new
53/// context values to avoid unintentionally overwriting existing state.
54///
55/// ## Managing the current context
56///
57/// Contexts can be associated with the caller's current execution unit on a
58/// given thread via the [`attach`] method, and previous contexts can be restored
59/// by dropping the returned [`ContextGuard`]. Context can be nested, and will
60/// restore their parent outer context when detached on drop. To access the
61/// values of the context, a snapshot can be created via the [`Context::current`]
62/// method.
63///
64/// [`Context::current`]: Context::current()
65/// [`get`]: Context::get()
66/// [`with_value`]: Context::with_value()
67/// [`attach`]: Context::attach()
68///
69/// # Examples
70///
71/// ```
72/// use opentelemetry::Context;
73///
74/// // Application-specific `a` and `b` values
75/// #[derive(Debug, PartialEq)]
76/// struct ValueA(&'static str);
77/// #[derive(Debug, PartialEq)]
78/// struct ValueB(u64);
79///
80/// let _outer_guard = Context::new().with_value(ValueA("a")).attach();
81///
82/// // Only value a has been set
83/// let current = Context::current();
84/// assert_eq!(current.get::<ValueA>(), Some(&ValueA("a")));
85/// assert_eq!(current.get::<ValueB>(), None);
86///
87/// {
88/// let _inner_guard = Context::current_with_value(ValueB(42)).attach();
89/// // Both values are set in inner context
90/// let current = Context::current();
91/// assert_eq!(current.get::<ValueA>(), Some(&ValueA("a")));
92/// assert_eq!(current.get::<ValueB>(), Some(&ValueB(42)));
93/// }
94///
95/// // Resets to only the `a` value when inner guard is dropped
96/// let current = Context::current();
97/// assert_eq!(current.get::<ValueA>(), Some(&ValueA("a")));
98/// assert_eq!(current.get::<ValueB>(), None);
99/// ```
100#[derive(Clone, Default)]
101pub struct Context {
102 #[cfg(feature = "trace")]
103 pub(crate) span: Option<Arc<SynchronizedSpan>>,
104 entries: Option<Arc<EntryMap>>,
105 suppress_telemetry: bool,
106 /// A context's observer view of this [Context].
107 ///
108 /// The main application for observers is to share information from the current context with
109 /// external readers (e.g. a full host eBPF profiler). In that case, the observer will attach
110 /// and detach some data on context events using its own mechanism. The shared data has no
111 /// reason to match our internal representation ([Context]).
112 ///
113 /// [observer_view] makes it possible to store the observer's view of the context alongside a
114 /// [Context] value, such that it can be attached and detached on context events.
115 #[cfg(feature = "experimental_context_observer")]
116 observer_view: Option<Arc<dyn ObserverContextView>>,
117}
118
119type EntryMap = HashMap<TypeId, Arc<dyn Any + Sync + Send>, BuildHasherDefault<IdHasher>>;
120
121impl Context {
122 /// Creates an empty `Context`.
123 ///
124 /// The context is initially created with a capacity of 0, so it will not
125 /// allocate. Use [`with_value`] to create a new context that has entries.
126 ///
127 /// [`with_value`]: Context::with_value()
128 pub fn new() -> Self {
129 Context::default()
130 }
131
132 /// Returns an immutable snapshot of the current thread's context.
133 ///
134 /// # Behavior During Context Drop
135 ///
136 /// When called from within a [`Drop`] implementation that is triggered by
137 /// a [`ContextGuard`] being dropped (e.g., when a [`Span`] is dropped as part
138 /// of context cleanup), this function returns **whatever context happens to be
139 /// current after the guard is popped**, not the context being dropped.
140 ///
141 /// **Important**: The returned context may be completely unrelated to the
142 /// context being dropped, as contexts can be activated in any order and are
143 /// not necessarily hierarchical. Do not rely on any relationship between them.
144 ///
145 /// This behavior is by design and prevents panics that would otherwise occur
146 /// from attempting to borrow the context while it's being mutably borrowed
147 /// for cleanup. See [issue #2871](https://github.com/open-telemetry/opentelemetry-rust/issues/2871)
148 /// for details.
149 ///
150 /// # Examples
151 ///
152 /// ```
153 /// use opentelemetry::Context;
154 ///
155 /// #[derive(Debug, PartialEq)]
156 /// struct ValueA(&'static str);
157 ///
158 /// fn do_work() {
159 /// assert_eq!(Context::current().get(), Some(&ValueA("a")));
160 /// }
161 ///
162 /// let _guard = Context::new().with_value(ValueA("a")).attach();
163 /// do_work()
164 /// ```
165 ///
166 /// [`Span`]: crate::trace::Span
167 pub fn current() -> Self {
168 Self::map_current(|cx| cx.clone())
169 }
170
171 /// Applies a function to the current context returning its value.
172 ///
173 /// This can be used to build higher performing algebraic expressions for
174 /// optionally creating a new context without the overhead of cloning the
175 /// current one and dropping it.
176 ///
177 /// Note: This function will panic if you attempt to attach another context
178 /// while the current one is still borrowed.
179 pub fn map_current<T>(f: impl FnOnce(&Context) -> T) -> T {
180 CURRENT_CONTEXT.with(|cx| cx.borrow().map_current_cx(f))
181 }
182
183 /// Returns a clone of the current thread's context with the given value.
184 ///
185 /// This is a more efficient form of `Context::current().with_value(value)`
186 /// as it avoids the intermediate context clone.
187 ///
188 /// # Examples
189 ///
190 /// ```
191 /// use opentelemetry::Context;
192 ///
193 /// // Given some value types defined in your application
194 /// #[derive(Debug, PartialEq)]
195 /// struct ValueA(&'static str);
196 /// #[derive(Debug, PartialEq)]
197 /// struct ValueB(u64);
198 ///
199 /// // You can create and attach context with the first value set to "a"
200 /// let _guard = Context::new().with_value(ValueA("a")).attach();
201 ///
202 /// // And create another context based on the fist with a new value
203 /// let all_current_and_b = Context::current_with_value(ValueB(42));
204 ///
205 /// // The second context now contains all the current values and the addition
206 /// assert_eq!(all_current_and_b.get::<ValueA>(), Some(&ValueA("a")));
207 /// assert_eq!(all_current_and_b.get::<ValueB>(), Some(&ValueB(42)));
208 /// ```
209 pub fn current_with_value<T: 'static + Send + Sync>(value: T) -> Self {
210 Self::map_current(|cx| cx.with_value(value))
211 }
212
213 /// Returns a reference to the entry for the corresponding value type.
214 ///
215 /// # Examples
216 ///
217 /// ```
218 /// use opentelemetry::Context;
219 ///
220 /// // Given some value types defined in your application
221 /// #[derive(Debug, PartialEq)]
222 /// struct ValueA(&'static str);
223 /// #[derive(Debug, PartialEq)]
224 /// struct MyUser();
225 ///
226 /// let cx = Context::new().with_value(ValueA("a"));
227 ///
228 /// // Values can be queried by type
229 /// assert_eq!(cx.get::<ValueA>(), Some(&ValueA("a")));
230 ///
231 /// // And return none if not yet set
232 /// assert_eq!(cx.get::<MyUser>(), None);
233 /// ```
234 pub fn get<T: 'static>(&self) -> Option<&T> {
235 self.entries
236 .as_ref()?
237 .get(&TypeId::of::<T>())?
238 .downcast_ref()
239 }
240
241 /// Returns a copy of the context with the new value included.
242 ///
243 /// # Examples
244 ///
245 /// ```
246 /// use opentelemetry::Context;
247 ///
248 /// // Given some value types defined in your application
249 /// #[derive(Debug, PartialEq)]
250 /// struct ValueA(&'static str);
251 /// #[derive(Debug, PartialEq)]
252 /// struct ValueB(u64);
253 ///
254 /// // You can create a context with the first value set to "a"
255 /// let cx_with_a = Context::new().with_value(ValueA("a"));
256 ///
257 /// // And create another context based on the fist with a new value
258 /// let cx_with_a_and_b = cx_with_a.with_value(ValueB(42));
259 ///
260 /// // The first context is still available and unmodified
261 /// assert_eq!(cx_with_a.get::<ValueA>(), Some(&ValueA("a")));
262 /// assert_eq!(cx_with_a.get::<ValueB>(), None);
263 ///
264 /// // The second context now contains both values
265 /// assert_eq!(cx_with_a_and_b.get::<ValueA>(), Some(&ValueA("a")));
266 /// assert_eq!(cx_with_a_and_b.get::<ValueB>(), Some(&ValueB(42)));
267 /// ```
268 pub fn with_value<T: 'static + Send + Sync>(&self, value: T) -> Self {
269 let entries = if let Some(current_entries) = &self.entries {
270 let mut inner_entries = (**current_entries).clone();
271 inner_entries.insert(TypeId::of::<T>(), Arc::new(value));
272 Some(Arc::new(inner_entries))
273 } else {
274 let mut entries = EntryMap::default();
275 entries.insert(TypeId::of::<T>(), Arc::new(value));
276 Some(Arc::new(entries))
277 };
278
279 Context {
280 entries,
281 #[cfg(feature = "trace")]
282 span: self.span.clone(),
283 suppress_telemetry: self.suppress_telemetry,
284 #[cfg(feature = "experimental_context_observer")]
285 observer_view: None,
286 }
287 .with_observer_view()
288 }
289
290 /// Replaces the current context on this thread with this context.
291 ///
292 /// Dropping the returned [`ContextGuard`] will reset the current context to the
293 /// previous value.
294 ///
295 ///
296 /// # Examples
297 ///
298 /// ```
299 /// use opentelemetry::Context;
300 ///
301 /// #[derive(Debug, PartialEq)]
302 /// struct ValueA(&'static str);
303 ///
304 /// let my_cx = Context::new().with_value(ValueA("a"));
305 ///
306 /// // Set the current thread context
307 /// let cx_guard = my_cx.attach();
308 /// assert_eq!(Context::current().get::<ValueA>(), Some(&ValueA("a")));
309 ///
310 /// // Drop the guard to restore the previous context
311 /// drop(cx_guard);
312 /// assert_eq!(Context::current().get::<ValueA>(), None);
313 /// ```
314 ///
315 /// Guards do not need to be explicitly dropped:
316 ///
317 /// ```
318 /// use opentelemetry::Context;
319 ///
320 /// #[derive(Debug, PartialEq)]
321 /// struct ValueA(&'static str);
322 ///
323 /// fn my_function() -> String {
324 /// // attach a context the duration of this function.
325 /// let my_cx = Context::new().with_value(ValueA("a"));
326 /// // NOTE: a variable name after the underscore is **required** or rust
327 /// // will drop the guard, restoring the previous context _immediately_.
328 /// let _guard = my_cx.attach();
329 ///
330 /// // anything happening in functions we call can still access my_cx...
331 /// my_other_function();
332 ///
333 /// // returning from the function drops the guard, exiting the span.
334 /// return "Hello world".to_owned();
335 /// }
336 ///
337 /// fn my_other_function() {
338 /// // ...
339 /// }
340 /// ```
341 /// Sub-scopes may be created to limit the duration for which the span is
342 /// entered:
343 ///
344 /// ```
345 /// use opentelemetry::Context;
346 ///
347 /// #[derive(Debug, PartialEq)]
348 /// struct ValueA(&'static str);
349 ///
350 /// let my_cx = Context::new().with_value(ValueA("a"));
351 ///
352 /// {
353 /// let _guard = my_cx.attach();
354 ///
355 /// // the current context can access variables in
356 /// assert_eq!(Context::current().get::<ValueA>(), Some(&ValueA("a")));
357 ///
358 /// // exiting the scope drops the guard, detaching the context.
359 /// }
360 ///
361 /// // this is back in the default empty context
362 /// assert_eq!(Context::current().get::<ValueA>(), None);
363 /// ```
364 pub fn attach(self) -> ContextGuard {
365 let cx_id = CURRENT_CONTEXT.with(|cx| cx.borrow_mut().push(self));
366
367 ContextGuard {
368 cx_pos: cx_id,
369 _marker: PhantomData,
370 }
371 }
372
373 /// Returns whether telemetry is suppressed in this context.
374 #[inline]
375 pub fn is_telemetry_suppressed(&self) -> bool {
376 self.suppress_telemetry
377 }
378
379 /// Returns a new context with telemetry suppression enabled.
380 pub fn with_telemetry_suppressed(&self) -> Self {
381 Context {
382 entries: self.entries.clone(),
383 #[cfg(feature = "trace")]
384 span: self.span.clone(),
385 suppress_telemetry: true,
386 #[cfg(feature = "experimental_context_observer")]
387 observer_view: None,
388 }
389 .with_observer_view()
390 }
391
392 /// Enters a scope where telemetry is suppressed.
393 ///
394 /// This method is specifically designed for OpenTelemetry components (like Exporters,
395 /// Processors etc.) to prevent generating recursive or self-referential
396 /// telemetry data when performing their own operations.
397 ///
398 /// Without suppression, we have a telemetry-induced-telemetry situation
399 /// where, operations like exporting telemetry could generate new telemetry
400 /// about the export process itself, potentially causing:
401 /// - Infinite telemetry feedback loops
402 /// - Excessive resource consumption
403 ///
404 /// This method:
405 /// 1. Takes the current context
406 /// 2. Creates a new context from current, with `suppress_telemetry` set to `true`
407 /// 3. Attaches it to the current thread
408 /// 4. Returns a guard that restores the previous context when dropped
409 ///
410 /// OTel SDK components would check `is_current_telemetry_suppressed()` before
411 /// generating new telemetry, but not end users.
412 ///
413 /// # Examples
414 ///
415 /// ```
416 /// use opentelemetry::Context;
417 ///
418 /// // Example: Inside an exporter's implementation
419 /// fn example_export_function() {
420 /// // Prevent telemetry-generating operations from creating more telemetry
421 /// let _guard = Context::enter_telemetry_suppressed_scope();
422 ///
423 /// // Verify suppression is active
424 /// assert_eq!(Context::is_current_telemetry_suppressed(), true);
425 ///
426 /// // Here you would normally perform operations that might generate telemetry
427 /// // but now they won't because the context has suppression enabled
428 /// }
429 ///
430 /// // Demonstrate the function
431 /// example_export_function();
432 /// ```
433 pub fn enter_telemetry_suppressed_scope() -> ContextGuard {
434 Self::map_current(|cx| cx.with_telemetry_suppressed()).attach()
435 }
436
437 /// Returns whether telemetry is suppressed in the current context.
438 ///
439 /// This method is used by OpenTelemetry components to determine whether they should
440 /// generate new telemetry in the current execution context. It provides a performant
441 /// way to check the suppression state.
442 ///
443 /// End-users generally should not use this method directly, as it is primarily intended for
444 /// OpenTelemetry SDK components.
445 ///
446 ///
447 #[inline]
448 pub fn is_current_telemetry_suppressed() -> bool {
449 Self::map_current(|cx| cx.is_telemetry_suppressed())
450 }
451
452 // Initialize `self.observer_view` for a newly created context using the current context
453 // observer, if set.
454 #[cfg(feature = "experimental_context_observer")]
455 fn with_observer_view(mut self) -> Self {
456 if let Some(observer) = GlobalContextObserver::get() {
457 self.observer_view = observer.make_view(&self);
458 }
459
460 self
461 }
462
463 // No-op when the experimental_context_observer feature is not enabled.
464 #[cfg(not(feature = "experimental_context_observer"))]
465 const fn with_observer_view(self) -> Self {
466 self
467 }
468
469 /// Returns the observer's view of this context, if any.
470 #[cfg(feature = "experimental_context_observer")]
471 #[inline]
472 pub fn observer_view(&self) -> &Option<Arc<dyn ObserverContextView>> {
473 &self.observer_view
474 }
475
476 #[cfg(feature = "trace")]
477 pub(crate) fn current_with_synchronized_span(value: SynchronizedSpan) -> Self {
478 Self::map_current(|cx| {
479 Context {
480 span: Some(Arc::new(value)),
481 entries: cx.entries.clone(),
482 suppress_telemetry: cx.suppress_telemetry,
483 #[cfg(feature = "experimental_context_observer")]
484 observer_view: None,
485 }
486 .with_observer_view()
487 })
488 }
489
490 #[cfg(feature = "trace")]
491 pub(crate) fn with_synchronized_span(&self, value: SynchronizedSpan) -> Self {
492 Context {
493 span: Some(Arc::new(value)),
494 entries: self.entries.clone(),
495 suppress_telemetry: self.suppress_telemetry,
496 #[cfg(feature = "experimental_context_observer")]
497 observer_view: None,
498 }
499 .with_observer_view()
500 }
501}
502
503impl fmt::Debug for Context {
504 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
505 let mut dbg = f.debug_struct("Context");
506
507 #[cfg(feature = "trace")]
508 let mut entries = self.entries.as_ref().map_or(0, |e| e.len());
509 #[cfg(feature = "trace")]
510 {
511 if let Some(span) = &self.span {
512 dbg.field("span", &span.span_context());
513 entries += 1;
514 } else {
515 dbg.field("span", &"None");
516 }
517 }
518 #[cfg(not(feature = "trace"))]
519 let entries = self.entries.as_ref().map_or(0, |e| e.len());
520
521 dbg.field("entries count", &entries)
522 .field("suppress_telemetry", &self.suppress_telemetry)
523 .finish()
524 }
525}
526
527/// A guard that resets the current context to the prior context when dropped.
528#[derive(Debug)]
529pub struct ContextGuard {
530 // The position of the context in the stack. This is used to pop the context.
531 cx_pos: u16,
532 // Ensure this type is !Send as it relies on thread locals
533 _marker: PhantomData<*const ()>,
534}
535
536impl Drop for ContextGuard {
537 fn drop(&mut self) {
538 let id = self.cx_pos;
539 if id > ContextStack::BASE_POS && id < ContextStack::MAX_POS {
540 // Extract only the span to drop outside of borrow_mut to avoid panic
541 // when the span's drop implementation calls Context::current()
542 #[cfg(feature = "trace")]
543 let _to_drop =
544 CURRENT_CONTEXT.with(|context_stack| context_stack.borrow_mut().pop_id(id));
545 #[cfg(not(feature = "trace"))]
546 CURRENT_CONTEXT.with(|context_stack| context_stack.borrow_mut().pop_id(id));
547 // Span (if any) is automatically dropped here, outside of borrow_mut scope
548 }
549 }
550}
551
552/// With TypeIds as keys, there's no need to hash them. They are already hashes
553/// themselves, coming from the compiler. The IdHasher holds the u64 of
554/// the TypeId, and then returns it, instead of doing any bit fiddling.
555#[derive(Clone, Default, Debug)]
556struct IdHasher(u64);
557
558impl Hasher for IdHasher {
559 fn write(&mut self, _: &[u8]) {
560 unreachable!("TypeId calls write_u64");
561 }
562
563 #[inline]
564 fn write_u64(&mut self, id: u64) {
565 self.0 = id;
566 }
567
568 #[inline]
569 fn finish(&self) -> u64 {
570 self.0
571 }
572}
573
574/// A stack for keeping track of the [`Context`] instances that have been attached
575/// to a thread.
576///
577/// The stack allows for popping of contexts by position, which is used to do out
578/// of order dropping of [`ContextGuard`] instances. Only when the top of the
579/// stack is popped, the topmost [`Context`] is actually restored.
580///
581/// The stack relies on the fact that it is thread local and that the
582/// [`ContextGuard`] instances that are constructed using ids from it can't be
583/// moved to other threads. That means that the ids are always valid and that
584/// they are always within the bounds of the stack.
585struct ContextStack {
586 /// This is the current [`Context`] that is active on this thread, and the top
587 /// of the [`ContextStack`]. It is always present, and if the `stack` is empty
588 /// it's an empty [`Context`].
589 ///
590 /// Having this here allows for fast access to the current [`Context`].
591 current_cx: Context,
592 /// A `stack` of the other contexts that have been attached to the thread.
593 stack: Vec<Option<Context>>,
594 /// Ensure this type is !Send as it relies on thread locals
595 _marker: PhantomData<*const ()>,
596}
597
598// Type alias for what pop_id returns - only return the span when trace feature is enabled
599#[cfg(feature = "trace")]
600type PopIdReturn = Option<Arc<SynchronizedSpan>>;
601#[cfg(not(feature = "trace"))]
602type PopIdReturn = ();
603
604impl ContextStack {
605 const BASE_POS: u16 = 0;
606 const MAX_POS: u16 = u16::MAX;
607 const INITIAL_CAPACITY: usize = 8;
608
609 #[inline(always)]
610 fn push(&mut self, cx: Context) -> u16 {
611 // The next id is the length of the `stack`, plus one since we have the
612 // top of the [`ContextStack`] as the `current_cx`.
613 let next_id = self.stack.len() + 1;
614 if next_id < ContextStack::MAX_POS.into() {
615 #[cfg(feature = "experimental_context_observer")]
616 if let Some(observer) = GlobalContextObserver::get() {
617 observer.on_context_enter(&self.current_cx, &cx);
618 }
619
620 let current_cx = std::mem::replace(&mut self.current_cx, cx);
621 self.stack.push(Some(current_cx));
622 next_id as u16
623 } else {
624 // This is an overflow, log it and ignore it.
625 otel_warn!(
626 name: "Context.AttachFailed",
627 message = format!("Too many contexts. Max limit is {}. \
628 Context::current() remains unchanged as this attach failed. \
629 Dropping the returned ContextGuard will have no impact on Context::current().",
630 ContextStack::MAX_POS)
631 );
632 ContextStack::MAX_POS
633 }
634 }
635
636 #[inline(always)]
637 fn pop_id(&mut self, pos: u16) -> PopIdReturn {
638 if pos == ContextStack::BASE_POS || pos == ContextStack::MAX_POS {
639 // The empty context is always at the bottom of the [`ContextStack`]
640 // and cannot be popped, and the overflow position is invalid, so do
641 // nothing.
642 otel_warn!(
643 name: "Context.OutOfOrderDrop",
644 position = pos,
645 message = if pos == ContextStack::BASE_POS {
646 "Attempted to pop the base context which is not allowed"
647 } else {
648 "Attempted to pop the overflow position which is not allowed"
649 }
650 );
651 #[cfg(feature = "trace")]
652 return None;
653 }
654 let len: u16 = self.stack.len() as u16;
655 // Are we at the top of the [`ContextStack`]?
656 if pos == len {
657 // Shrink the stack if possible to clear out any out of order pops.
658 while let Some(None) = self.stack.last() {
659 _ = self.stack.pop();
660 }
661 // Restore the previous context. This will always happen since the
662 // empty context is always at the bottom of the stack if the
663 // [`ContextStack`] is not empty.
664 if let Some(Some(next_cx)) = self.stack.pop() {
665 #[cfg(feature = "experimental_context_observer")]
666 if let Some(observer) = GlobalContextObserver::get() {
667 observer.on_context_exit(&self.current_cx, &next_cx);
668 }
669
670 // Extract and return only the span to avoid cloning the entire Context
671 #[cfg(feature = "trace")]
672 {
673 let old_cx = std::mem::replace(&mut self.current_cx, next_cx);
674 return old_cx.span;
675 }
676 #[cfg(not(feature = "trace"))]
677 {
678 self.current_cx = next_cx;
679 }
680 }
681 #[cfg(feature = "trace")]
682 return None;
683 } else {
684 // This is an out of order pop.
685 if pos >= len {
686 // This is an invalid id, ignore it.
687 otel_warn!(
688 name: "Context.PopOutOfBounds",
689 position = pos,
690 stack_length = len,
691 message = "Attempted to pop beyond the end of the context stack"
692 );
693 #[cfg(feature = "trace")]
694 return None;
695 }
696 // Clear out the entry at the given id and extract its span
697 #[cfg(feature = "trace")]
698 return self.stack[pos as usize].take().and_then(|cx| cx.span);
699 #[cfg(not(feature = "trace"))]
700 {
701 self.stack[pos as usize] = None;
702 }
703 }
704 }
705
706 #[inline(always)]
707 fn map_current_cx<T>(&self, f: impl FnOnce(&Context) -> T) -> T {
708 f(&self.current_cx)
709 }
710}
711
712impl Default for ContextStack {
713 fn default() -> Self {
714 ContextStack {
715 current_cx: Context::default(),
716 stack: Vec::with_capacity(ContextStack::INITIAL_CAPACITY),
717 _marker: PhantomData,
718 }
719 }
720}
721
722#[cfg(feature = "experimental_context_observer")]
723mod context_observer {
724 use super::*;
725
726 /// An observer trait for monitoring context transitions.
727 ///
728 /// Implementors of this trait can observe when a context is entered or exited, allowing for
729 /// custom logic to be executed during context switches.
730 ///
731 /// An observer-specific view can be attached to a [`Context`] via
732 /// [`Context::observer_view`]. A view is usually a different representation of the context
733 /// to be published through an alternative channel, but in essence it is arbitrary data computed
734 /// from the context.
735 ///
736 /// # Refcell already borrowed panic
737 ///
738 /// [crate::context] maintains thread-local, internal state in a `RefCell`. This implementation
739 /// detail leaks when the `experimental_context_observer` feature is enabled:
740 /// [Self::on_context_enter] and [Self::on_context_exit] are called from a point where the cell
741 /// is currently borrowed. If a [ContextObserver] implementation calls back into a
742 /// cell-borrowing function from the Context API, typically `Context::current()`, this will
743 /// result in a `Refcell already borrowed` panic.
744 ///
745 /// As general hygiene, with respect to the [Context] API, it is strongly advised to only use
746 /// simple getters and setters on the `from` and `to` context objects provided to the observer.
747 /// Do not call `Context::current()`, nor try to manually attach or detach context objects from
748 /// an observer.
749 ///
750 /// If you get Refcell-related panics within the [crate::context] module, this is the most
751 /// likely cause.
752 ///
753 /// # Example
754 ///
755 /// ```
756 /// # #[cfg(feature = "experimental_context_observer")]
757 /// # {
758 /// use opentelemetry::Context;
759 /// use opentelemetry::context::{ContextObserver, ObserverContextView};
760 /// use std::any::Any;
761 /// use std::sync::Arc;
762 ///
763 /// // An observer's own view of the context, carrying arbitrary data.
764 /// struct MyView {
765 /// correlation_id: u64,
766 /// }
767 ///
768 /// impl ObserverContextView for MyView {
769 /// fn as_any(&self) -> &dyn Any {
770 /// self
771 /// }
772 /// }
773 ///
774 /// struct Observer;
775 ///
776 /// # fn do_something(view: &MyView) { }
777 ///
778 /// impl ContextObserver for Observer {
779 /// fn on_context_enter(&self, from: &Context, to: &Context) {
780 /// let view = to.observer_view().as_ref().unwrap();
781 /// do_something(view.as_any().downcast_ref::<MyView>().unwrap());
782 /// }
783 ///
784 /// fn on_context_exit(&self, from: &Context, to: &Context) {
785 /// let view = to
786 /// .observer_view()
787 /// .as_ref()
788 /// .unwrap()
789 /// .as_any()
790 /// .downcast_ref::<MyView>()
791 /// .unwrap();
792 /// assert_eq!(view.correlation_id, 42);
793 /// }
794 ///
795 /// fn make_view(&self, ctx: &Context) -> Option<Arc<dyn ObserverContextView>> {
796 /// Some(Arc::new(MyView { correlation_id: 42 }))
797 /// }
798 /// }
799 /// # }
800 /// ```
801 pub trait ContextObserver {
802 /// Called when a context is entered, allowing observers to react to the transition
803 /// from one context (`from`) to another (`to`).
804 fn on_context_enter(&self, from: &Context, to: &Context);
805
806 /// Called when a context is exited, allowing observers to react to the transition
807 /// from one context (`from`) to another (`to`).
808 fn on_context_exit(&self, from: &Context, to: &Context);
809
810 /// Derive a fresh view for a newly created context. The provided `ctx` has an unitialized
811 /// view sets to `None`.
812 ///
813 /// The default implementation returns `None` for observers that don't use the view.
814 fn make_view(&self, _ctx: &Context) -> Option<Arc<dyn ObserverContextView>> {
815 None
816 }
817 }
818
819 static GLOBAL_CONTEXT_OBSERVER: OnceLock<Arc<dyn ContextObserver + Send + Sync>> =
820 OnceLock::new();
821
822 /// A global observer for context transitions.
823 ///
824 /// This struct provides static methods to set and get a global context observer.
825 #[derive(Debug)]
826 pub struct GlobalContextObserver;
827
828 impl GlobalContextObserver {
829 /// Sets the global context observer, logging a warning if it was already set.
830 pub fn set(observer: Arc<dyn ContextObserver + Send + Sync>) {
831 if GLOBAL_CONTEXT_OBSERVER.set(observer).is_err() {
832 otel_warn!(
833 name: "GlobalContextObserver.SetFailed",
834 message = "Global context observer was already set. Ignoring new observer."
835 );
836 }
837 }
838
839 /// Returns the global context observer.
840 pub(super) fn get<'a>() -> Option<&'a Arc<dyn ContextObserver + Send + Sync>> {
841 GLOBAL_CONTEXT_OBSERVER.get()
842 }
843 }
844
845 /// A context's observer view of a [Context]. This is an arbitrary data structure. See
846 /// [ContextObserver]'s documentation for examples.
847 #[cfg(feature = "experimental_context_observer")]
848 pub trait ObserverContextView: Any + Send + Sync {
849 /// Returns this view as a `&dyn Any`, enabling downcasting back to the concrete type.
850 ///
851 /// This method is required because this crate's MSRV is below the Rust version (1.86) that
852 /// stabilized trait upcasting. Without it, callers below 1.86 couldn't upcast `&dyn
853 /// ObserverContextView` to `&dyn Any` to perform the downcast themselves.
854 ///
855 /// This is a work-around for the absence of trait upcasting before Rust 1.86.
856 ///
857 /// Implementors should simply return `self`:
858 ///
859 /// ```ignore
860 /// fn as_any(&self) -> &dyn std::any::Any {
861 /// self
862 /// }
863 /// ```
864 fn as_any(&self) -> &dyn Any;
865 }
866}
867
868#[cfg(test)]
869mod tests {
870 use super::*;
871 use std::time::Duration;
872 use tokio::time::sleep;
873
874 #[derive(Debug, PartialEq)]
875 struct ValueA(u64);
876 #[derive(Debug, PartialEq)]
877 struct ValueB(u64);
878
879 #[test]
880 fn context_immutable() {
881 // start with Current, which should be an empty context
882 let cx = Context::current();
883 assert_eq!(cx.get::<ValueA>(), None);
884 assert_eq!(cx.get::<ValueB>(), None);
885
886 // with_value should return a new context,
887 // leaving the original context unchanged
888 let cx_new = cx.with_value(ValueA(1));
889
890 // cx should be unchanged
891 assert_eq!(cx.get::<ValueA>(), None);
892 assert_eq!(cx.get::<ValueB>(), None);
893
894 // cx_new should contain the new value
895 assert_eq!(cx_new.get::<ValueA>(), Some(&ValueA(1)));
896
897 // cx_new should be unchanged
898 let cx_newer = cx_new.with_value(ValueB(1));
899
900 // Cx and cx_new are unchanged
901 assert_eq!(cx.get::<ValueA>(), None);
902 assert_eq!(cx.get::<ValueB>(), None);
903 assert_eq!(cx_new.get::<ValueA>(), Some(&ValueA(1)));
904 assert_eq!(cx_new.get::<ValueB>(), None);
905
906 // cx_newer should contain both values
907 assert_eq!(cx_newer.get::<ValueA>(), Some(&ValueA(1)));
908 assert_eq!(cx_newer.get::<ValueB>(), Some(&ValueB(1)));
909 }
910
911 #[test]
912 fn nested_contexts() {
913 let _outer_guard = Context::new().with_value(ValueA(1)).attach();
914
915 // Only value `a` is set
916 let current = Context::current();
917 assert_eq!(current.get(), Some(&ValueA(1)));
918 assert_eq!(current.get::<ValueB>(), None);
919
920 {
921 let _inner_guard = Context::current_with_value(ValueB(42)).attach();
922 // Both values are set in inner context
923 let current = Context::current();
924 assert_eq!(current.get(), Some(&ValueA(1)));
925 assert_eq!(current.get(), Some(&ValueB(42)));
926
927 assert!(Context::map_current(|cx| {
928 assert_eq!(cx.get(), Some(&ValueA(1)));
929 assert_eq!(cx.get(), Some(&ValueB(42)));
930 true
931 }));
932 }
933
934 // Resets to only value `a` when inner guard is dropped
935 let current = Context::current();
936 assert_eq!(current.get(), Some(&ValueA(1)));
937 assert_eq!(current.get::<ValueB>(), None);
938
939 assert!(Context::map_current(|cx| {
940 assert_eq!(cx.get(), Some(&ValueA(1)));
941 assert_eq!(cx.get::<ValueB>(), None);
942 true
943 }));
944 }
945
946 #[test]
947 fn overlapping_contexts() {
948 let outer_guard = Context::new().with_value(ValueA(1)).attach();
949
950 // Only value `a` is set
951 let current = Context::current();
952 assert_eq!(current.get(), Some(&ValueA(1)));
953 assert_eq!(current.get::<ValueB>(), None);
954
955 let inner_guard = Context::current_with_value(ValueB(42)).attach();
956 // Both values are set in inner context
957 let current = Context::current();
958 assert_eq!(current.get(), Some(&ValueA(1)));
959 assert_eq!(current.get(), Some(&ValueB(42)));
960
961 assert!(Context::map_current(|cx| {
962 assert_eq!(cx.get(), Some(&ValueA(1)));
963 assert_eq!(cx.get(), Some(&ValueB(42)));
964 true
965 }));
966
967 drop(outer_guard);
968
969 // `inner_guard` is still alive so both `ValueA` and `ValueB` should still be accessible
970 let current = Context::current();
971 assert_eq!(current.get(), Some(&ValueA(1)));
972 assert_eq!(current.get(), Some(&ValueB(42)));
973
974 drop(inner_guard);
975
976 // Both guards are dropped and neither value should be accessible.
977 let current = Context::current();
978 assert_eq!(current.get::<ValueA>(), None);
979 assert_eq!(current.get::<ValueB>(), None);
980 }
981
982 #[test]
983 fn too_many_contexts() {
984 let mut guards: Vec<ContextGuard> = Vec::with_capacity(ContextStack::MAX_POS as usize);
985 let stack_max_pos = ContextStack::MAX_POS as u64;
986 // Fill the stack up until the last position
987 for i in 1..stack_max_pos {
988 let cx_guard = Context::current().with_value(ValueB(i)).attach();
989 assert_eq!(Context::current().get(), Some(&ValueB(i)));
990 assert_eq!(cx_guard.cx_pos, i as u16);
991 guards.push(cx_guard);
992 }
993 // Let's overflow the stack a couple of times
994 for _ in 0..16 {
995 let cx_guard = Context::current().with_value(ValueA(1)).attach();
996 assert_eq!(cx_guard.cx_pos, ContextStack::MAX_POS);
997 assert_eq!(Context::current().get::<ValueA>(), None);
998 assert_eq!(Context::current().get(), Some(&ValueB(stack_max_pos - 1)));
999 guards.push(cx_guard);
1000 }
1001 // Drop the overflow contexts
1002 for _ in 0..16 {
1003 guards.pop();
1004 assert_eq!(Context::current().get::<ValueA>(), None);
1005 assert_eq!(Context::current().get(), Some(&ValueB(stack_max_pos - 1)));
1006 }
1007 // Drop one more so we can add a new one
1008 guards.pop();
1009 assert_eq!(Context::current().get::<ValueA>(), None);
1010 assert_eq!(Context::current().get(), Some(&ValueB(stack_max_pos - 2)));
1011 // Push a new context and see that it works
1012 let cx_guard = Context::current().with_value(ValueA(2)).attach();
1013 assert_eq!(cx_guard.cx_pos, ContextStack::MAX_POS - 1);
1014 assert_eq!(Context::current().get(), Some(&ValueA(2)));
1015 assert_eq!(Context::current().get(), Some(&ValueB(stack_max_pos - 2)));
1016 guards.push(cx_guard);
1017 // Let's overflow the stack a couple of times again
1018 for _ in 0..16 {
1019 let cx_guard = Context::current().with_value(ValueA(1)).attach();
1020 assert_eq!(cx_guard.cx_pos, ContextStack::MAX_POS);
1021 assert_eq!(Context::current().get::<ValueA>(), Some(&ValueA(2)));
1022 assert_eq!(Context::current().get(), Some(&ValueB(stack_max_pos - 2)));
1023 guards.push(cx_guard);
1024 }
1025 }
1026
1027 /// Tests that a new ContextStack is created with the correct initial capacity.
1028 #[test]
1029 fn test_initial_capacity() {
1030 let stack = ContextStack::default();
1031 assert_eq!(stack.stack.capacity(), ContextStack::INITIAL_CAPACITY);
1032 }
1033
1034 /// Tests that map_current_cx correctly accesses the current context.
1035 #[test]
1036 fn test_map_current_cx() {
1037 let mut stack = ContextStack::default();
1038 let test_value = ValueA(42);
1039 stack.current_cx = Context::new().with_value(test_value);
1040
1041 let result = stack.map_current_cx(|cx| {
1042 assert_eq!(cx.get::<ValueA>(), Some(&ValueA(42)));
1043 true
1044 });
1045 assert!(result);
1046 }
1047
1048 /// Tests popping contexts in non-sequential order.
1049 #[test]
1050 fn test_pop_id_out_of_order() {
1051 let mut stack = ContextStack::default();
1052
1053 // Push three contexts
1054 let cx1 = Context::new().with_value(ValueA(1));
1055 let cx2 = Context::new().with_value(ValueA(2));
1056 let cx3 = Context::new().with_value(ValueA(3));
1057
1058 let id1 = stack.push(cx1);
1059 let id2 = stack.push(cx2);
1060 let id3 = stack.push(cx3);
1061
1062 // Pop middle context first - should not affect current context
1063 stack.pop_id(id2);
1064 assert_eq!(stack.current_cx.get::<ValueA>(), Some(&ValueA(3)));
1065 assert_eq!(stack.stack.len(), 3); // Length unchanged for middle pops
1066
1067 // Pop last context - should restore previous valid context
1068 stack.pop_id(id3);
1069 assert_eq!(stack.current_cx.get::<ValueA>(), Some(&ValueA(1)));
1070 assert_eq!(stack.stack.len(), 1);
1071
1072 // Pop first context - should restore to empty state
1073 stack.pop_id(id1);
1074 assert_eq!(stack.current_cx.get::<ValueA>(), None);
1075 assert_eq!(stack.stack.len(), 0);
1076 }
1077
1078 /// Tests edge cases in context stack operations. IRL these should log
1079 /// warnings, and definitely not panic.
1080 #[test]
1081 fn test_pop_id_edge_cases() {
1082 let mut stack = ContextStack::default();
1083
1084 // Test popping BASE_POS - should be no-op
1085 stack.pop_id(ContextStack::BASE_POS);
1086 assert_eq!(stack.stack.len(), 0);
1087
1088 // Test popping MAX_POS - should be no-op
1089 stack.pop_id(ContextStack::MAX_POS);
1090 assert_eq!(stack.stack.len(), 0);
1091
1092 // Test popping invalid position - should be no-op
1093 stack.pop_id(1000);
1094 assert_eq!(stack.stack.len(), 0);
1095
1096 // Test popping from empty stack - should be safe
1097 stack.pop_id(1);
1098 assert_eq!(stack.stack.len(), 0);
1099 }
1100
1101 /// Tests stack behavior when reaching maximum capacity.
1102 /// Once we push beyond this point, we should end up with a context
1103 /// that points _somewhere_, but mutating it should not affect the current
1104 /// active context.
1105 #[test]
1106 fn test_push_overflow() {
1107 let mut stack = ContextStack::default();
1108 let max_pos = ContextStack::MAX_POS as usize;
1109
1110 // Fill stack up to max position
1111 for i in 0..max_pos {
1112 let cx = Context::new().with_value(ValueA(i as u64));
1113 let id = stack.push(cx);
1114 assert_eq!(id, (i + 1) as u16);
1115 }
1116
1117 // Try to push beyond capacity
1118 let cx = Context::new().with_value(ValueA(max_pos as u64));
1119 let id = stack.push(cx);
1120 assert_eq!(id, ContextStack::MAX_POS);
1121
1122 // Verify current context remains unchanged after overflow
1123 assert_eq!(
1124 stack.current_cx.get::<ValueA>(),
1125 Some(&ValueA((max_pos - 2) as u64))
1126 );
1127 }
1128
1129 /// Tests that:
1130 /// 1. Parent context values are properly propagated to async operations
1131 /// 2. Values added during async operations do not affect parent context
1132 #[tokio::test]
1133 async fn test_async_context_propagation() {
1134 // A nested async operation we'll use to test propagation
1135 async fn nested_operation() {
1136 // Verify we can see the parent context's value
1137 assert_eq!(
1138 Context::current().get::<ValueA>(),
1139 Some(&ValueA(42)),
1140 "Parent context value should be available in async operation"
1141 );
1142
1143 // Create new context
1144 let cx_with_both = Context::current()
1145 .with_value(ValueA(43)) // override ValueA
1146 .with_value(ValueB(24)); // Add new ValueB
1147
1148 // Run nested async operation with both values
1149 async {
1150 // Verify both values are available
1151 assert_eq!(
1152 Context::current().get::<ValueA>(),
1153 Some(&ValueA(43)),
1154 "Parent value should still be available after adding new value"
1155 );
1156 assert_eq!(
1157 Context::current().get::<ValueB>(),
1158 Some(&ValueB(24)),
1159 "New value should be available in async operation"
1160 );
1161
1162 // Do some async work to simulate real-world scenario
1163 sleep(Duration::from_millis(10)).await;
1164
1165 // Values should still be available after async work
1166 assert_eq!(
1167 Context::current().get::<ValueA>(),
1168 Some(&ValueA(43)),
1169 "Parent value should persist across await points"
1170 );
1171 assert_eq!(
1172 Context::current().get::<ValueB>(),
1173 Some(&ValueB(24)),
1174 "New value should persist across await points"
1175 );
1176 }
1177 .with_context(cx_with_both)
1178 .await;
1179 }
1180
1181 // Set up initial context with ValueA
1182 let parent_cx = Context::new().with_value(ValueA(42));
1183
1184 // Create and run async operation with the parent context explicitly propagated
1185 nested_operation().with_context(parent_cx.clone()).await;
1186
1187 // After async operation completes:
1188 // 1. Parent context should be unchanged
1189 assert_eq!(
1190 parent_cx.get::<ValueA>(),
1191 Some(&ValueA(42)),
1192 "Parent context should be unchanged"
1193 );
1194 assert_eq!(
1195 parent_cx.get::<ValueB>(),
1196 None,
1197 "Parent context should not see values added in async operation"
1198 );
1199
1200 // 2. Current context should be back to default
1201 assert_eq!(
1202 Context::current().get::<ValueA>(),
1203 None,
1204 "Current context should be back to default"
1205 );
1206 assert_eq!(
1207 Context::current().get::<ValueB>(),
1208 None,
1209 "Current context should not have async operation's values"
1210 );
1211 }
1212
1213 ///
1214 /// Tests that unnatural parent->child relationships in nested async
1215 /// operations behave properly.
1216 ///
1217 #[tokio::test]
1218 async fn test_out_of_order_context_detachment_futures() {
1219 // This function returns a future, but doesn't await it
1220 // It will complete before the future that it creates.
1221 async fn create_a_future() -> impl std::future::Future<Output = ()> {
1222 // Create a future that will do some work, referencing our current
1223 // context, but don't await it.
1224 async {
1225 assert_eq!(Context::current().get::<ValueA>(), Some(&ValueA(42)));
1226
1227 // Longer work
1228 sleep(Duration::from_millis(50)).await;
1229 }
1230 .with_context(Context::current())
1231 }
1232
1233 // Create our base context
1234 let parent_cx = Context::new().with_value(ValueA(42));
1235
1236 // await our nested function, which will create and detach a context
1237 let future = create_a_future().with_context(parent_cx).await;
1238
1239 // Execute the future. The future that created it is long gone, but this shouldn't
1240 // cause issues.
1241 let _a = future.await;
1242
1243 // Nothing terrible (e.g., panics!) should happen, and we should definitely not have any
1244 // values attached to our current context that were set in the nested operations.
1245 assert_eq!(Context::current().get::<ValueA>(), None);
1246 assert_eq!(Context::current().get::<ValueB>(), None);
1247 }
1248
1249 #[test]
1250 fn test_is_telemetry_suppressed() {
1251 // Default context has suppression disabled
1252 let cx = Context::new();
1253 assert!(!cx.is_telemetry_suppressed());
1254
1255 // With suppression enabled
1256 let suppressed = cx.with_telemetry_suppressed();
1257 assert!(suppressed.is_telemetry_suppressed());
1258 }
1259
1260 #[test]
1261 fn test_with_telemetry_suppressed() {
1262 // Start with a normal context
1263 let cx = Context::new();
1264 assert!(!cx.is_telemetry_suppressed());
1265
1266 // Create a suppressed context
1267 let suppressed = cx.with_telemetry_suppressed();
1268
1269 // Original should remain unchanged
1270 assert!(!cx.is_telemetry_suppressed());
1271
1272 // New context should be suppressed
1273 assert!(suppressed.is_telemetry_suppressed());
1274
1275 // Test with values to ensure they're preserved
1276 let cx_with_value = cx.with_value(ValueA(42));
1277 let suppressed_with_value = cx_with_value.with_telemetry_suppressed();
1278
1279 assert!(!cx_with_value.is_telemetry_suppressed());
1280 assert!(suppressed_with_value.is_telemetry_suppressed());
1281 assert_eq!(suppressed_with_value.get::<ValueA>(), Some(&ValueA(42)));
1282 }
1283
1284 #[test]
1285 fn test_enter_telemetry_suppressed_scope() {
1286 // Ensure we start with a clean context
1287 let _reset_guard = Context::new().attach();
1288
1289 // Default context should not be suppressed
1290 assert!(!Context::is_current_telemetry_suppressed());
1291
1292 // Add an entry to the current context
1293 let cx_with_value = Context::current().with_value(ValueA(42));
1294 let _guard_with_value = cx_with_value.attach();
1295
1296 // Verify the entry is present and context is not suppressed
1297 assert_eq!(Context::current().get::<ValueA>(), Some(&ValueA(42)));
1298 assert!(!Context::is_current_telemetry_suppressed());
1299
1300 // Enter a suppressed scope
1301 {
1302 let _guard = Context::enter_telemetry_suppressed_scope();
1303
1304 // Verify suppression is active and the entry is still present
1305 assert!(Context::is_current_telemetry_suppressed());
1306 assert!(Context::current().is_telemetry_suppressed());
1307 assert_eq!(Context::current().get::<ValueA>(), Some(&ValueA(42)));
1308 }
1309
1310 // After guard is dropped, should be back to unsuppressed and entry should still be present
1311 assert!(!Context::is_current_telemetry_suppressed());
1312 assert!(!Context::current().is_telemetry_suppressed());
1313 assert_eq!(Context::current().get::<ValueA>(), Some(&ValueA(42)));
1314 }
1315
1316 #[test]
1317 fn test_nested_suppression_scopes() {
1318 // Ensure we start with a clean context
1319 let _reset_guard = Context::new().attach();
1320
1321 // Default context should not be suppressed
1322 assert!(!Context::is_current_telemetry_suppressed());
1323
1324 // First level suppression
1325 {
1326 let _outer = Context::enter_telemetry_suppressed_scope();
1327 assert!(Context::is_current_telemetry_suppressed());
1328
1329 // Second level. This component is unaware of Suppression,
1330 // and just attaches a new context. Since it is from current,
1331 // it'll already have suppression enabled.
1332 {
1333 let _inner = Context::current().with_value(ValueA(1)).attach();
1334 assert!(Context::is_current_telemetry_suppressed());
1335 assert_eq!(Context::current().get::<ValueA>(), Some(&ValueA(1)));
1336 }
1337
1338 // Another scenario. This component is unaware of Suppression,
1339 // and just attaches a new context, not from Current. Since it is
1340 // not from current it will not have suppression enabled.
1341 {
1342 let _inner = Context::new().with_value(ValueA(1)).attach();
1343 assert!(!Context::is_current_telemetry_suppressed());
1344 assert_eq!(Context::current().get::<ValueA>(), Some(&ValueA(1)));
1345 }
1346
1347 // Still suppressed after inner scope
1348 assert!(Context::is_current_telemetry_suppressed());
1349 }
1350
1351 // Back to unsuppressed
1352 assert!(!Context::is_current_telemetry_suppressed());
1353 }
1354
1355 #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
1356 async fn test_async_suppression() {
1357 async fn nested_operation() {
1358 assert!(Context::is_current_telemetry_suppressed());
1359
1360 let cx_with_additional_value = Context::current().with_value(ValueB(24));
1361
1362 async {
1363 assert_eq!(
1364 Context::current().get::<ValueB>(),
1365 Some(&ValueB(24)),
1366 "Parent value should still be available after adding new value"
1367 );
1368 assert!(Context::is_current_telemetry_suppressed());
1369
1370 // Do some async work to simulate real-world scenario
1371 sleep(Duration::from_millis(10)).await;
1372
1373 // Values should still be available after async work
1374 assert_eq!(
1375 Context::current().get::<ValueB>(),
1376 Some(&ValueB(24)),
1377 "Parent value should still be available after adding new value"
1378 );
1379 assert!(Context::is_current_telemetry_suppressed());
1380 }
1381 .with_context(cx_with_additional_value)
1382 .await;
1383 }
1384
1385 // Set up suppressed context, but don't attach it to current
1386 let suppressed_parent = Context::new().with_telemetry_suppressed();
1387 // Current should not be suppressed as we haven't attached it
1388 assert!(!Context::is_current_telemetry_suppressed());
1389
1390 // Create and run async operation with the suppressed context explicitly propagated
1391 nested_operation()
1392 .with_context(suppressed_parent.clone())
1393 .await;
1394
1395 // After async operation completes:
1396 // Suppression should be active
1397 assert!(suppressed_parent.is_telemetry_suppressed());
1398
1399 // Current should still be not suppressed
1400 assert!(!Context::is_current_telemetry_suppressed());
1401 }
1402}
1403
1404#[cfg(all(test, feature = "experimental_context_observer"))]
1405mod observer_tests {
1406 use super::*;
1407 use std::cell::{Cell, RefCell};
1408 use std::sync::Arc;
1409
1410 #[derive(Debug, PartialEq)]
1411 struct V(u64);
1412
1413 #[derive(Debug, PartialEq, Clone)]
1414 enum Event {
1415 Enter { from: Option<u64>, to: Option<u64> },
1416 Exit { from: Option<u64>, to: Option<u64> },
1417 }
1418
1419 thread_local! {
1420 static EVENTS: RefCell<Vec<Event>> = const { RefCell::new(Vec::new()) };
1421 }
1422
1423 fn value_of(cx: &Context) -> Option<u64> {
1424 cx.get::<V>().map(|v| v.0)
1425 }
1426
1427 // When enabled, on each transition it records the event into a thread-local buffer, and on
1428 // enter it also installs an observer view on the entered context and bumps the thread-local
1429 // `VIEWS_CREATED` counter. See `observer_installs_and_reuses_view`.
1430 //
1431 // The observer is a process-global singleton, so its callbacks fire for context activity on
1432 // *every* thread across the whole test binary once installed. To keep that cheap and isolated,
1433 // all of its work is gated behind the per-thread `ENABLE_TEST_OBSERVER` switch: threads that
1434 // didn't opt in (via `setup()`) get an inert observer.
1435 struct RecordingObserver;
1436
1437 impl ContextObserver for RecordingObserver {
1438 fn on_context_enter(&self, from: &Context, to: &Context) {
1439 if !ENABLE_TEST_OBSERVER.with(Cell::get) {
1440 return;
1441 }
1442
1443 EVENTS.with(|e| {
1444 e.borrow_mut().push(Event::Enter {
1445 from: value_of(from),
1446 to: value_of(to),
1447 })
1448 });
1449 }
1450
1451 fn on_context_exit(&self, from: &Context, to: &Context) {
1452 if !ENABLE_TEST_OBSERVER.with(Cell::get) {
1453 return;
1454 }
1455
1456 EVENTS.with(|e| {
1457 e.borrow_mut().push(Event::Exit {
1458 from: value_of(from),
1459 to: value_of(to),
1460 })
1461 });
1462 }
1463
1464 // On context enter, install the observer's view if this context doesn't have one yet, or reuse
1465 // the existing one otherwise. The view is the context's `V` value (or 0 for a value-less
1466 // context).
1467 fn make_view(&self, cx: &Context) -> Option<Arc<dyn ObserverContextView>> {
1468 let correlation_id = value_of(cx).unwrap_or(0);
1469 VIEWS_CREATED.with(|c| c.set(c.get() + 1));
1470 Some(Arc::new(MyView { correlation_id }))
1471 }
1472 }
1473
1474 thread_local! {
1475 // Master switch for the observer per thread. Enabled for the duration of a test that cares,
1476 // and left false everywhere else, so we don't slow down tests unrelated to the observer.
1477 // See `setup`.
1478 static ENABLE_TEST_OBSERVER: Cell<bool> = const { Cell::new(false) };
1479 // Counts how many distinct views the observer created on the current thread.
1480 static VIEWS_CREATED: Cell<u64> = const { Cell::new(0) };
1481 }
1482
1483 // Recovers the correlation id stored in a context's observer view, if any.
1484 fn view_id(cx: &Context) -> Option<u64> {
1485 cx.observer_view()
1486 .as_ref()
1487 .and_then(|v| v.as_any().downcast_ref::<MyView>())
1488 .map(|v| v.correlation_id)
1489 }
1490
1491 // Holds the base context guard and disables the observer on this thread when dropped. The
1492 // struct's own `Drop` runs before its fields, so the flag is cleared before the base context is
1493 // popped, which is thus not observed (symmetrically to `setup`, which installs the base context
1494 // before registering/enabling the observer).
1495 struct ObserverTestGuard {
1496 _cx: ContextGuard,
1497 }
1498
1499 impl Drop for ObserverTestGuard {
1500 fn drop(&mut self) {
1501 ENABLE_TEST_OBSERVER.with(|e| e.set(false));
1502 }
1503 }
1504
1505 // Enables the observer on this thread, establishes a clean base context, and clears any
1506 // previously recorded events, so the assertions below only see transitions this test
1507 // triggers. Dropping the returned guard disables the observer again.
1508 fn setup_observer() -> ObserverTestGuard {
1509 // Idempotent across tests: only the first call installs the observer, but every test
1510 // uses the same `RecordingObserver`, so a lost race is harmless.
1511 GlobalContextObserver::set(Arc::new(RecordingObserver));
1512 ENABLE_TEST_OBSERVER.with(|e| e.set(true));
1513 let guard = Context::new().attach();
1514 EVENTS.with(|e| e.borrow_mut().clear());
1515 ObserverTestGuard { _cx: guard }
1516 }
1517
1518 #[test]
1519 fn global_observer_records_enter_and_exit() {
1520 let _clean = setup_observer();
1521
1522 {
1523 let _g = Context::current().with_value(V(1)).attach();
1524 // Entering fires immediately: from the clean base (no value) to `V(1)`.
1525 EVENTS.with(|e| {
1526 assert_eq!(
1527 *e.borrow(),
1528 vec![Event::Enter {
1529 from: None,
1530 to: Some(1)
1531 }]
1532 )
1533 });
1534 }
1535
1536 // Dropping the (context) guard restores the base and fires the matching exit event.
1537 EVENTS.with(|e| {
1538 assert_eq!(
1539 *e.borrow(),
1540 vec![
1541 Event::Enter {
1542 from: None,
1543 to: Some(1)
1544 },
1545 Event::Exit {
1546 from: Some(1),
1547 to: None
1548 }
1549 ]
1550 )
1551 });
1552 }
1553
1554 #[test]
1555 fn global_observer_records_nested_transitions() {
1556 let _clean = setup_observer();
1557
1558 let g1 = Context::current().with_value(V(1)).attach();
1559 let g2 = Context::current().with_value(V(2)).attach();
1560 drop(g2);
1561 drop(g1);
1562
1563 EVENTS.with(|e| {
1564 assert_eq!(
1565 *e.borrow(),
1566 vec![
1567 Event::Enter {
1568 from: None,
1569 to: Some(1)
1570 },
1571 Event::Enter {
1572 from: Some(1),
1573 to: Some(2)
1574 },
1575 Event::Exit {
1576 from: Some(2),
1577 to: Some(1)
1578 },
1579 Event::Exit {
1580 from: Some(1),
1581 to: None
1582 },
1583 ]
1584 )
1585 });
1586 }
1587
1588 // An observer view carrying arbitrary data, unrelated to the internal `Context`
1589 // representation.
1590 struct MyView {
1591 correlation_id: u64,
1592 }
1593
1594 impl ObserverContextView for MyView {
1595 fn as_any(&self) -> &dyn Any {
1596 self
1597 }
1598 }
1599
1600 // A realistic use of the observer view: the global observer installs its own view on each
1601 // context the first time that context is entered, and reuses it on re-entry. The
1602 // thread-local counter checks that exactly one view is created per distinct context.
1603 #[test]
1604 fn observer_installs_and_reuses_view() {
1605 let _clean = setup_observer();
1606
1607 // Reset the counter for the work below.
1608 VIEWS_CREATED.with(|c| c.set(0));
1609
1610 // Entering two distinct contexts installs a fresh view for each, derived from its value.
1611 let g1 = Context::current().with_value(V(1)).attach();
1612 assert_eq!(view_id(&Context::current()), Some(1));
1613
1614 let g2 = Context::current().with_value(V(2)).attach();
1615 assert_eq!(view_id(&Context::current()), Some(2));
1616
1617 assert_eq!(VIEWS_CREATED.with(Cell::get), 2);
1618
1619 // Re-entering an already-observed context must reuse its view, not create a new one.
1620 // `Context::current()` clones the current context, carrying along its already-set view.
1621 let observed = Context::current();
1622 let g3 = observed.attach();
1623
1624 assert_eq!(view_id(&Context::current()), Some(2));
1625 assert_eq!(
1626 VIEWS_CREATED.with(Cell::get),
1627 2,
1628 "re-entering an observed context must not create a new view"
1629 );
1630
1631 drop(g3);
1632 drop(g2);
1633 drop(g1);
1634 }
1635}