1use crate::{CallContext, DbScopeDiagnosticIdentity};
4use serde::{Serialize, Serializer, ser::SerializeStruct};
5use std::sync::{
6 Arc, Mutex, OnceLock,
7 atomic::{AtomicU64, Ordering},
8};
9
10mod unified;
11pub use unified::{UnifiedContext, SelectedContextValue};
12
13static NEXT_ROOT: AtomicU64 = AtomicU64::new(1);
14
15#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
17#[serde(tag = "state", content = "value", rename_all = "snake_case")]
18pub enum ContextFact<T> {
19 Present(T),
20 NotApplicable,
21 NotEstablished,
22 Unavailable,
23}
24
25#[derive(Clone, Copy, Eq, PartialEq)]
27pub struct ContextIdentity {
28 bytes: [u8; 256],
29 len: u16,
30}
31impl ContextIdentity {
32 pub fn checked(value: &str) -> Result<Self, ContextConflict> {
33 if value.is_empty() || value.len() > 256 || value.chars().any(char::is_control) {
34 return Err(ContextConflict::InvalidIdentity);
35 }
36 let mut out = Self {
37 bytes: [0; 256],
38 len: value.len() as u16,
39 };
40 out.bytes[..value.len()].copy_from_slice(value.as_bytes());
41 Ok(out)
42 }
43 fn as_str(&self) -> &str {
44 std::str::from_utf8(&self.bytes[..usize::from(self.len)]).expect("validated UTF-8")
45 }
46}
47impl Serialize for ContextIdentity {
48 fn serialize<S: Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
49 s.serialize_str(self.as_str())
50 }
51}
52
53#[derive(Clone, Copy, Eq, PartialEq, Serialize)]
55#[serde(transparent)]
56pub struct ContextLabel(ContextIdentity);
57impl ContextLabel {
58 pub fn checked(value: &str) -> Result<Self, ContextConflict> {
59 let value = ContextIdentity::checked(value)?;
60 if value.as_str().contains("://")
61 || !value
62 .as_str()
63 .chars()
64 .all(|c| c.is_alphanumeric() || "_.:/{}*-".contains(c))
65 {
66 return Err(ContextConflict::UnsafeMetadata);
67 }
68 Ok(Self(value))
69 }
70}
71
72#[derive(Clone, Copy, Debug, Eq, PartialEq)]
73pub enum ContextConflict {
74 InvalidIdentity,
75 UnsafeMetadata,
76 IdentityGroup,
77 ForeignRoot,
78 ChildRelation,
79 CounterExhausted,
80}
81
82#[derive(Eq, PartialEq)]
85pub struct RequestIdentityGroup {
86 application: ContextLabel,
87 module: ContextLabel,
88 service: ContextLabel,
89 operation: ContextLabel,
90 trace: ContextIdentity,
91 rpc: ContextFact<ContextIdentity>,
92 span: u64,
93 request: ContextIdentity,
94 route: ContextLabel,
95 attempt: u32,
96 zone: ContextFact<ContextLabel>,
97}
98impl RequestIdentityGroup {
99 pub fn from_validated(
100 call: &CallContext,
101 request: &str,
102 route: &str,
103 attempt: u32,
104 zone: ContextFact<ContextLabel>,
105 ) -> Result<Self, ContextConflict> {
106 if attempt == 0 {
107 return Err(ContextConflict::InvalidIdentity);
108 }
109 Ok(Self {
110 application: ContextLabel::checked(call.application().as_str())?,
111 module: ContextLabel::checked(call.module().as_str())?,
112 service: ContextLabel::checked(call.service().as_str())?,
113 operation: ContextLabel::checked(call.operation().as_str())?,
114 trace: ContextIdentity::checked(call.trace_correlation_id().as_str())?,
115 rpc: match call.rpc_correlation_id() {
116 Some(id) => ContextFact::Present(ContextIdentity::checked(id.as_str())?),
117 None => ContextFact::Unavailable,
118 },
119 span: call.span_id().as_u64(),
120 request: ContextIdentity::checked(request)?,
121 route: ContextLabel::checked(route)?,
122 attempt,
123 zone,
124 })
125 }
126}
127
128struct RequestRoot {
129 local: u64,
130 next_revision: AtomicU64,
131 application: ContextLabel,
132 initial: ContextFact<()>,
133 identity: OnceLock<RequestIdentityGroup>,
134 ingress: OnceLock<Mutex<Option<IngressSnapshot>>>,
135}
136impl RequestRoot {
137 fn ingress_snapshot(&self) -> Option<IngressSnapshot> {
138 self.ingress.get().and_then(|slot|
139 slot.lock().unwrap_or_else(|p| p.into_inner()).clone())
140 }
141}
142impl Drop for RequestRoot {
143 fn drop(&mut self) {
144 observe(self.local, "root_drop");
145 }
146}
147
148pub struct RequestRootPublisher {
155 root: Arc<RequestRoot>,
156}
157
158pub struct RequestRootRef {
160 root: Arc<RequestRoot>,
161}
162impl Clone for RequestRootRef {
163 fn clone(&self) -> Self {
164 observe(self.root.local, "root_share");
165 Self {
166 root: Arc::clone(&self.root),
167 }
168 }
169}
170impl Drop for RequestRootRef {
171 fn drop(&mut self) {
172 observe(self.root.local, "root_release");
173 }
174}
175
176impl RequestRootPublisher {
177 pub fn create(
180 application: ContextLabel,
181 initial: ContextFact<()>,
182 ) -> Result<Self, ContextConflict> {
183 if matches!(initial, ContextFact::Present(())) {
184 return Err(ContextConflict::InvalidIdentity);
185 }
186 let local = NEXT_ROOT
187 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |n| n.checked_add(1))
188 .map_err(|_| ContextConflict::CounterExhausted)?;
189 let root = Arc::new(RequestRoot {
190 local,
191 next_revision: AtomicU64::new(1),
192 application,
193 initial,
194 identity: OnceLock::new(),
195 ingress: OnceLock::new(),
196 });
197 observe(local, "root_create");
198 Ok(Self { root })
199 }
200 pub fn publish_ingress(
202 &mut self,
203 source: impl Into<crate::context_json::ContextSourceRef>,
204 call_id: &str,
205 deadline_unix_ms: i64,
206 ) -> Result<(), ContextConflict> {
207 if self.root.identity.get().is_none() { return Err(ContextConflict::IdentityGroup); }
208 let ingress = IngressSnapshot { source: source.into(), call_id: ContextIdentity::checked(call_id)?, deadline_unix_ms };
209 self.root.ingress.set(Mutex::new(Some(ingress))).map_err(|_| ContextConflict::IdentityGroup)
210 }
211 #[doc(hidden)]
215 pub fn retire_ingress(&mut self) {
216 if let Some(slot) = self.root.ingress.get() {
217 drop(slot.lock().unwrap_or_else(|p| p.into_inner()).take());
218 }
219 }
220 pub fn reference(&self) -> RequestRootRef {
221 observe(self.root.local, "root_share");
222 RequestRootRef {
223 root: Arc::clone(&self.root),
224 }
225 }
226 #[allow(clippy::result_large_err)] pub fn publish(
228 &mut self,
229 group: RequestIdentityGroup,
230 ) -> Result<(), (ContextConflict, RequestIdentityGroup)> {
231 if group.application != self.root.application {
232 return Err((ContextConflict::IdentityGroup, group));
233 }
234 if let Some(old) = self.root.identity.get() {
235 return if old == &group {
236 Ok(())
237 } else {
238 Err((ContextConflict::IdentityGroup, group))
239 };
240 }
241 self.root
242 .identity
243 .set(group)
244 .map_err(|group| (ContextConflict::IdentityGroup, group))
245 }
246}
247
248#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
250#[serde(rename_all = "snake_case")]
251pub enum RequestViewPhase {
252 SocketAccepted,
253 Reading,
254 Admitted,
255 Dispatch,
256 Handler,
257 Database,
258 Outbound,
259 Response,
260 Finalizing,
261 Finished,
262}
263
264#[derive(Clone, Copy, Eq, PartialEq, Serialize)]
266pub struct RegisteredContextOperation(&'static str);
267impl RegisteredContextOperation {
268 pub fn checked(value: &'static str) -> Result<Self, ContextConflict> {
269 ContextLabel::checked(value)?;
270 Ok(Self(value))
271 }
272}
273
274#[derive(Clone, Copy)]
275pub struct RequestLocalFacts {
276 task: ContextFact<u64>,
277 scope: ContextFact<DbScopeDiagnosticIdentity>,
278 db_operation: ContextFact<RegisteredContextOperation>,
279 phase: RequestViewPhase,
280}
281impl RequestLocalFacts {
282 pub fn new(phase: RequestViewPhase) -> Self {
283 Self {
284 task: ContextFact::NotEstablished,
285 scope: ContextFact::NotEstablished,
286 db_operation: ContextFact::NotApplicable,
287 phase,
288 }
289 }
290 pub fn with_task(mut self, task: ContextFact<u64>) -> Self {
292 self.task = task;
293 self
294 }
295 pub fn without_scope(mut self, state: ContextFact<()>) -> Self {
296 self.scope = absent(state);
297 self
298 }
299 pub fn with_db_operation(mut self, operation: RegisteredContextOperation) -> Self {
300 self.db_operation = ContextFact::Present(operation);
301 self
302 }
303}
304
305struct ChildCall {
308 #[cfg(test)]
309 observation: LocalObservation,
310 application: ContextLabel,
311 module: ContextLabel,
312 service: ContextLabel,
313 operation: ContextLabel,
314 rpc: ContextIdentity,
315 span: u64,
316 route: ContextLabel,
317 attempt: u32,
318}
319#[derive(Clone)]
320struct IngressSnapshot {
321 source: crate::context_json::ContextSourceRef,
322 call_id: ContextIdentity,
323 deadline_unix_ms: i64,
324}
325struct LocalView {
326 #[cfg(test)]
327 observation: LocalObservation,
328 root: Arc<RequestRoot>,
329 published: bool,
330 facts: RequestLocalFacts,
331 child: Option<Arc<ChildCall>>,
332 ingress: Option<IngressSnapshot>,
333 business: Option<crate::context_json::ContextSourceRef>,
334 revision: u64,
335}
336impl Drop for LocalView {
337 fn drop(&mut self) {
338 observe(self.root.local, "view_drop");
339 #[cfg(test)]
340 self.observation.record("destroy");
341 }
342}
343impl Drop for ChildCall {
344 fn drop(&mut self) {
345 #[cfg(test)]
346 self.observation.record("destroy");
347 }
348}
349
350pub struct RequestExecutionView {
358 inner: Arc<LocalView>,
359}
360
361#[derive(Clone, Copy, Serialize)]
364pub struct RequestContextSnapshot {
365 schema_version: u8,
366 local_request: u64,
367 publication: u8,
368 application: ContextFact<ContextLabel>,
369 call_application: ContextFact<ContextLabel>,
370 module: ContextFact<ContextLabel>,
371 service: ContextFact<ContextLabel>,
372 operation: ContextFact<ContextLabel>,
373 trace_id: ContextFact<ContextIdentity>,
374 request: ContextFact<ContextIdentity>,
375 span_id: ContextFact<SpanProjection>,
376 route: ContextFact<ContextLabel>,
377 attempt: ContextFact<u32>,
378 rpc_id: ContextFact<ContextIdentity>,
379 zone: ContextFact<ContextLabel>,
380 db_operation: ContextFact<RegisteredContextOperation>,
381 scope: ContextFact<DbScopeDiagnosticIdentity>,
382 task: ContextFact<u64>,
383 lifecycle: ContextFact<RequestViewPhase>,
384 target: ContextFact<ContextLabel>,
385}
386impl RequestContextSnapshot {
387 #[doc(hidden)]
388 pub fn local_request_id(&self) -> u64 { self.local_request }
389}
390pub struct OutboundSourceContext<'a> {
393 view: &'a RequestExecutionView,
394 business_unit: &'a str,
395}
396impl Clone for RequestExecutionView {
397 fn clone(&self) -> Self {
398 #[cfg(test)]
399 self.inner.observation.record("share");
400 Self {
401 inner: Arc::clone(&self.inner),
402 }
403 }
404}
405impl Drop for RequestExecutionView {
406 fn drop(&mut self) {
407 #[cfg(test)]
408 self.inner.observation.record("release");
409 }
410}
411impl RequestRootRef {
412 pub fn view(&self, facts: RequestLocalFacts) -> RequestExecutionView {
413 RequestExecutionView::allocate(
414 Arc::clone(&self.root),
415 self.root.identity.get().is_some(),
416 facts,
417 None,
418 self.root.ingress_snapshot(),
419 None,
420 )
421 }
422 pub fn same_request(&self, view: &RequestExecutionView) -> bool {
423 Arc::ptr_eq(&self.root, &view.inner.root)
424 }
425}
426impl RequestExecutionView {
427 pub fn diagnostic_snapshot(&self) -> RequestContextSnapshot {
428 let identity = self.identity();
429 let child = self.inner.child.as_deref();
430 let initial = self.inner.root.initial;
431 macro_rules! root {
432 ($field:ident) => { identity.map_or_else(|| absent(initial), |g| ContextFact::Present(g.$field)) };
433 }
434 macro_rules! call {
435 ($field:ident) => { child.map(|c| c.$field).or_else(|| identity.map(|g| g.$field))
436 .map_or_else(|| absent(initial), ContextFact::Present) };
437 }
438 RequestContextSnapshot {
439 schema_version: 2,
440 local_request: self.inner.root.local,
441 publication: u8::from(self.inner.published),
442 application: ContextFact::Present(self.inner.root.application),
443 call_application: ContextFact::Present(child.map_or(self.inner.root.application, |c| c.application)),
444 module: call!(module), service: call!(service), operation: call!(operation),
445 trace_id: root!(trace), request: root!(request),
446 span_id: child.map(|c| c.span).or_else(|| identity.map(|g| g.span))
447 .map_or_else(|| absent(initial), |span| ContextFact::Present(SpanProjection(span))),
448 route: call!(route), attempt: call!(attempt),
449 rpc_id: child.map(|c| ContextFact::Present(c.rpc)).or_else(|| identity.map(|g| g.rpc))
450 .unwrap_or_else(|| absent(initial)),
451 zone: identity.map_or_else(|| absent(initial), |g| g.zone),
452 db_operation: self.inner.facts.db_operation, scope: self.inner.facts.scope,
453 task: self.inner.facts.task, lifecycle: ContextFact::Present(self.inner.facts.phase),
454 target: child.map_or(ContextFact::NotApplicable, |c| ContextFact::Present(c.route)),
455 }
456 }
457 pub fn with_business(&self, source: impl Into<crate::context_json::ContextSourceRef>) -> Self {
460 let source = source.into();
461 Self::allocate(Arc::clone(&self.inner.root), self.inner.published, self.inner.facts,
462 self.inner.child.clone(), self.inner.ingress.clone(), Some(source))
463 }
464 pub fn unified_context(&self) -> UnifiedContext<'_> { UnifiedContext(self) }
465 pub fn outbound_source_context<'a>(&'a self, business_unit: &'a str) -> OutboundSourceContext<'a> {
466 OutboundSourceContext { view: self, business_unit }
467 }
468 fn allocate(
469 root: Arc<RequestRoot>,
470 published: bool,
471 facts: RequestLocalFacts,
472 child: Option<Arc<ChildCall>>,
473 ingress: Option<IngressSnapshot>,
474 business: Option<crate::context_json::ContextSourceRef>,
475 ) -> Self {
476 let revision = root.next_revision.fetch_update(Ordering::Relaxed, Ordering::Relaxed,
477 |next| next.checked_add(1)).expect("request context revision counter exhausted");
478 observe(root.local, "view_create");
479 Self {
480 inner: Arc::new(LocalView {
481 #[cfg(test)]
482 observation: LocalObservation::create(root.local, "view"),
483 root,
484 published,
485 facts,
486 child,
487 ingress,
488 business,
489 revision,
490 }),
491 }
492 }
493 fn local(&self, facts: RequestLocalFacts) -> Self {
495 Self::allocate(
496 Arc::clone(&self.inner.root),
497 self.inner.published,
498 facts,
499 self.inner.child.clone(),
500 self.inner.ingress.clone(),
501 self.inner.business.clone(),
502 )
503 }
504 pub fn with_phase(&self, phase: RequestViewPhase) -> Self {
505 let mut facts = self.inner.facts;
506 facts.phase = phase;
507 self.local(facts)
508 }
509 pub fn with_db_operation(&self, operation: RegisteredContextOperation) -> Self {
510 self.local(self.inner.facts.with_db_operation(operation))
511 }
512 pub fn in_task(&self, task: u64) -> Self {
515 self.local(self.inner.facts.with_task(ContextFact::Present(task)))
516 }
517 pub fn in_db_scope(
518 &self,
519 scope: &crate::DbScopeDiagnosticContext<Self>,
520 ) -> Result<Self, ContextConflict> {
521 let (original, identity) = scope.diagnostic_context();
522 if !self.same_request(original) {
523 return Err(ContextConflict::ForeignRoot);
524 }
525 let mut facts = self.inner.facts;
526 facts.scope = ContextFact::Present(identity);
527 Ok(self.local(facts))
528 }
529 pub fn observed_db_scope(
530 &self,
531 scope: &crate::DbScopeObservation<Self>,
532 ) -> Result<Self, ContextConflict> {
533 let (original, identity) = scope.diagnostic_context();
534 if !self.same_request(original) {
535 return Err(ContextConflict::ForeignRoot);
536 }
537 let mut facts = self.inner.facts;
538 facts.scope = ContextFact::Present(identity);
539 Ok(self.local(facts))
540 }
541 pub fn same_request(&self, other: &Self) -> bool {
542 Arc::ptr_eq(&self.inner.root, &other.inner.root)
543 }
544 pub fn same_view(&self, other: &Self) -> bool {
545 Arc::ptr_eq(&self.inner, &other.inner)
546 }
547 pub fn refresh(&self, root: &RequestRootRef) -> Result<Self, ContextConflict> {
549 if !root.same_request(self) {
550 return Err(ContextConflict::ForeignRoot);
551 }
552 Ok(Self::allocate(
553 Arc::clone(&self.inner.root),
554 self.inner.root.identity.get().is_some(),
555 self.inner.facts,
556 self.inner.child.clone(),
557 self.inner.root.ingress_snapshot(),
558 self.inner.business.clone(),
559 ))
560 }
561 pub fn child(
562 &self,
563 call: &CallContext,
564 request: &str,
565 route: &str,
566 attempt: u32,
567 ) -> Result<Self, ContextConflict> {
568 let identity = self.identity().ok_or(ContextConflict::ChildRelation)?;
569 if identity.trace.as_str() != call.trace_correlation_id().as_str()
570 || identity.request.as_str() != request
571 {
572 return Err(ContextConflict::ForeignRoot);
573 }
574 let parent_rpc = if let Some(child) = &self.inner.child {
575 &child.rpc
576 } else if let ContextFact::Present(rpc) = &identity.rpc {
577 rpc
578 } else {
579 return Err(ContextConflict::ChildRelation);
580 };
581 let rpc = call
582 .rpc_correlation_id()
583 .ok_or(ContextConflict::ChildRelation)?;
584 let suffix = rpc
585 .as_str()
586 .strip_prefix(parent_rpc.as_str())
587 .and_then(|s| s.strip_prefix('.'))
588 .ok_or(ContextConflict::ChildRelation)?;
589 let span = self
590 .inner
591 .child
592 .as_ref()
593 .map_or(identity.span, |child| child.span);
594 if suffix.is_empty()
595 || !suffix.bytes().all(|b| b.is_ascii_digit())
596 || span == call.span_id().as_u64()
597 || attempt == 0
598 {
599 return Err(ContextConflict::ChildRelation);
600 }
601 let child = Arc::new(ChildCall {
602 application: ContextLabel::checked(call.application().as_str())?,
603 module: ContextLabel::checked(call.module().as_str())?,
604 service: ContextLabel::checked(call.service().as_str())?,
605 operation: ContextLabel::checked(call.operation().as_str())?,
606 rpc: ContextIdentity::checked(rpc.as_str())?,
607 span: call.span_id().as_u64(),
608 route: ContextLabel::checked(route)?,
609 attempt,
610 #[cfg(test)]
611 observation: LocalObservation::create(self.inner.root.local, "child"),
612 });
613 Ok(Self::allocate(
614 Arc::clone(&self.inner.root),
615 self.inner.published,
616 self.inner.facts,
617 Some(child),
618 self.inner.ingress.clone(),
619 self.inner.business.clone(),
620 ))
621 }
622 fn identity(&self) -> Option<&RequestIdentityGroup> {
623 self.inner
624 .published
625 .then(|| self.inner.root.identity.get())
626 .flatten()
627 }
628}
629
630fn absent<T>(state: ContextFact<()>) -> ContextFact<T> {
631 match state {
632 ContextFact::NotApplicable => ContextFact::NotApplicable,
633 ContextFact::NotEstablished => ContextFact::NotEstablished,
634 _ => ContextFact::Unavailable,
635 }
636}
637impl Serialize for RequestExecutionView {
638 fn serialize<S: Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
639 self.serialize_context(s, None)
640 }
641}
642impl Serialize for OutboundSourceContext<'_> {
643 fn serialize<S: Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
644 self.view.serialize_context(s, Some(self.business_unit))
645 }
646}
647impl RequestExecutionView {
648 fn serialize_context<S: Serializer>(&self, s: S, business_unit: Option<&str>) -> Result<S::Ok, S::Error> {
649 let mut out = s.serialize_struct("RequestContext", 21 + usize::from(business_unit.is_some()))?;
650 let identity = self.identity();
651 let child = self.inner.child.as_deref();
652 let initial = self.inner.root.initial;
653 out.serialize_field("schema_version", &2u8)?;
654 out.serialize_field("local_request", &self.inner.root.local)?;
655 out.serialize_field("publication", &u8::from(self.inner.published))?;
656 out.serialize_field(
657 "application",
658 &ContextFact::Present(&self.inner.root.application),
659 )?;
660 out.serialize_field(
661 "call_application",
662 &ContextFact::Present(child.map_or(&self.inner.root.application, |c| &c.application)),
663 )?;
664 macro_rules! root_field {
665 ($name:literal, $field:ident) => {
666 out.serialize_field(
667 $name,
668 &identity.map_or_else(|| absent(initial), |g| ContextFact::Present(&g.$field)),
669 )?;
670 };
671 }
672 macro_rules! call_field {
673 ($name:literal, $field:ident) => {
674 out.serialize_field(
675 $name,
676 &child
677 .map(|c| &c.$field)
678 .or_else(|| identity.map(|g| &g.$field))
679 .map_or_else(|| absent(initial), ContextFact::Present),
680 )?;
681 };
682 }
683 call_field!("module", module);
684 call_field!("service", service);
685 call_field!("operation", operation);
686 root_field!("trace_id", trace);
687 root_field!("request", request);
688 let span = child.map(|c| c.span).or_else(|| identity.map(|g| g.span));
689 out.serialize_field(
690 "span_id",
691 &span.map_or_else(
692 || absent(initial),
693 |s| ContextFact::Present(SpanProjection(s)),
694 ),
695 )?;
696 call_field!("route", route);
697 call_field!("attempt", attempt);
698 let rpc = child
699 .map(|c| ContextFact::Present(c.rpc))
700 .or_else(|| identity.map(|g| g.rpc))
701 .unwrap_or_else(|| absent(initial));
702 out.serialize_field("rpc_id", &rpc)?;
703 out.serialize_field(
704 "zone",
705 &identity.map_or_else(|| absent(initial), |g| g.zone),
706 )?;
707 out.serialize_field("db_operation", &self.inner.facts.db_operation)?;
708 out.serialize_field("scope", &self.inner.facts.scope)?;
709 out.serialize_field("task", &self.inner.facts.task)?;
710 out.serialize_field("lifecycle", &ContextFact::Present(self.inner.facts.phase))?;
711 out.serialize_field(
713 "target",
714 &child.map_or(ContextFact::NotApplicable, |c| {
715 ContextFact::Present(&c.route)
716 }),
717 )?;
718 if let Some(business_unit) = business_unit {
719 out.serialize_field("business_unit", &ContextFact::Present(business_unit))?;
720 }
721 out.end()
722 }
723}
724
725pub fn request_context_layouts() -> [(std::alloc::Layout, std::alloc::Layout); 3] {
729 fn pair<T>() -> (std::alloc::Layout, std::alloc::Layout) {
730 let payload = std::alloc::Layout::new::<T>();
731 let header = std::alloc::Layout::new::<[std::sync::atomic::AtomicUsize; 2]>();
732 (
733 payload,
734 header
735 .extend(payload)
736 .expect("fixed layout")
737 .0
738 .pad_to_align(),
739 )
740 }
741 [
742 pair::<RequestRoot>(),
743 pair::<LocalView>(),
744 pair::<ChildCall>(),
745 ]
746}
747
748#[cfg(not(test))]
749fn observe(_: u64, _: &'static str) {}
750#[cfg(test)]
751fn observe(root: u64, event: &'static str) {
752 EVENTS.lock().unwrap().push((root, event));
753 record_observation(root, root, "root", event);
754}
755#[cfg(test)]
756static EVENTS: std::sync::Mutex<Vec<(u64, &'static str)>> = std::sync::Mutex::new(Vec::new());
757
758#[derive(Clone, Copy)]
759struct SpanProjection(u64);
760impl Serialize for SpanProjection {
761 fn serialize<S: Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
762 let mut bytes = [b'0'; 16];
763 for (i, byte) in bytes.iter_mut().enumerate() {
764 *byte = b"0123456789abcdef"[((self.0 >> ((15 - i) * 4)) & 15) as usize];
765 }
766 serializer.serialize_str(std::str::from_utf8(&bytes).expect("hex"))
767 }
768}
769
770#[cfg(test)]
773struct LocalObservation {
774 id: u64,
775 root: u64,
776 kind: &'static str,
777}
778#[cfg(test)]
779static LOCAL_EVENTS: std::sync::Mutex<Vec<(u64, u64, &'static str, &'static str)>> =
780 std::sync::Mutex::new(Vec::new());
781#[cfg(test)]
782impl LocalObservation {
783 fn create(root: u64, kind: &'static str) -> Self {
784 static NEXT: AtomicU64 = AtomicU64::new(1);
785 let observation = Self {
786 id: NEXT.fetch_add(1, Ordering::Relaxed),
787 root,
788 kind,
789 };
790 observation.record("create");
791 observation
792 }
793 fn record(&self, event: &'static str) {
794 LOCAL_EVENTS
795 .lock()
796 .unwrap()
797 .push((self.id, self.root, self.kind, event));
798 record_observation(self.root, self.id, self.kind, event);
799 }
800}
801
802#[cfg(test)]
803type ObservationEvent = (u64, u64, u64, &'static str, &'static str);
804#[cfg(test)]
805static ORDERED_EVENTS: std::sync::Mutex<Vec<ObservationEvent>> = std::sync::Mutex::new(Vec::new());
806#[cfg(test)]
807fn record_observation(root: u64, object: u64, kind: &'static str, event: &'static str) {
808 static SEQUENCE: AtomicU64 = AtomicU64::new(1);
809 let sequence = SEQUENCE.fetch_add(1, Ordering::Relaxed);
810 ORDERED_EVENTS
811 .lock()
812 .unwrap()
813 .push((sequence, root, object, kind, event));
814}
815
816#[cfg(test)]
817mod tests;
818
819mod field_paths;
820pub const CONTEXT_FIELDS_JSON: &str = include_str!("request_context/context-fields.json");
822pub const CONTEXT_RECORD_SCHEMA_JSON: &str = include_str!("request_context/context-record.schema.json");
824pub fn is_known_context_record_path(path: &str) -> bool {
826 let Ok(pointer) = crate::json_pointer::JsonPointer::parse(path) else { return false; };
827 if path.is_empty() { return true; }
828 field_paths::CONTEXT_RECORD_PATHS.iter().any(|(declared, dynamic)| {
829 let mut requested = pointer.segments();
830 for part in declared[1..].split('/') {
831 match requested.next() {
832 None => return true, Some(segment) if segment.matches(part) => {},
834 Some(_) => return false,
835 }
836 }
837 *dynamic || requested.next().is_none()
838 })
839}