use connectrpc::client::CallOptions;
use crate::wire::DeclaredCall;
#[must_use]
fn traced_options() -> CallOptions {
let mut headers = http::HeaderMap::new();
polyc_runtime::propagation::inject_current_span_into(&mut headers);
options_from_headers(headers)
}
#[must_use]
pub fn bounded_traced_options_from_headers(
declared: &DeclaredCall,
headers: http::HeaderMap,
) -> CallOptions {
options_from_headers(headers).with_timeout(declared.remaining_budget())
}
fn options_from_headers(headers: http::HeaderMap) -> CallOptions {
CallOptions::default().with_headers(
headers
.into_iter()
.filter_map(|(name, value)| name.map(|name| (name, value))),
)
}
#[must_use]
pub fn bounded_traced_options(declared: &DeclaredCall) -> CallOptions {
traced_options().with_timeout(declared.remaining_budget())
}
#[must_use]
pub(crate) fn streaming_traced_options() -> CallOptions {
traced_options()
}
#[must_use]
pub fn adopt_caller_trace(headers: &http::HeaderMap, method: &'static str) -> tracing::Span {
use opentelemetry::trace::TraceContextExt as _;
let context = polyc_runtime::propagation::extract_context_from(headers);
let remote = context.span().span_context().trace_id();
let span = tracing::info_span!(
"state.call",
rpc = method,
caller_trace_id = %remote
);
span.in_scope(|| {
polyc_runtime::propagation::extract_parent_into_current_span(headers);
});
span
}
#[must_use]
pub fn carried_trace_id(headers: &http::HeaderMap) -> String {
use opentelemetry::trace::TraceContextExt as _;
polyc_runtime::propagation::extract_context_from(headers)
.span()
.span_context()
.trace_id()
.to_string()
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs, clippy::unwrap_used)]
use super::*;
use std::time::Duration;
#[test]
fn what_a_client_injects_is_what_a_listener_extracts() {
opentelemetry::global::set_text_map_propagator(
opentelemetry_sdk::propagation::TraceContextPropagator::new(),
);
let mut carried = http::HeaderMap::new();
carried.insert(
"traceparent",
"00-4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7-01"
.parse()
.unwrap(),
);
assert_eq!(
carried_trace_id(&carried),
"4bf92f3577b34da6a3ce929d0e0e4736",
"the listener reads the trace the caller named"
);
let span = adopt_caller_trace(&carried, "CommitJournalBatch");
span.in_scope(|| {
assert_eq!(
carried_trace_id(&carried),
"4bf92f3577b34da6a3ce929d0e0e4736"
);
});
}
#[test]
fn a_call_without_a_trace_context_still_runs() {
let bare = http::HeaderMap::new();
assert_eq!(
carried_trace_id(&bare),
"00000000000000000000000000000000",
"no context is the invalid trace id, not an error"
);
let _span = adopt_caller_trace(&bare, "GetJournalHead");
}
#[test]
fn the_client_half_produces_call_options() {
let declared = DeclaredCall::bounded(crate::state_audience(), Duration::from_millis(37));
let options = bounded_traced_options(&declared);
assert_eq!(
options.timeout(),
Some(Duration::from_millis(37)),
"the Connect timeout is the exact semantic budget"
);
}
}