#![doc(html_logo_url = "https://raw.githubusercontent.com/emit-rs/emit/main/asset/logo.svg")]
#![deny(missing_docs)]
use emit::{
well_known::{
KEY_ERR, KEY_EVT_KIND, KEY_LVL, KEY_SPAN_ID, KEY_SPAN_KIND, KEY_SPAN_LINKS, KEY_SPAN_NAME,
KEY_SPAN_PARENT, KEY_TRACE_ID, LVL_DEBUG, LVL_ERROR, LVL_INFO, LVL_WARN,
},
Filter as _, Props as _,
};
use opentelemetry::trace::{Link, TraceState};
use opentelemetry::{
logs::{AnyValue, LogRecord, Logger, LoggerProvider, Severity},
trace::{
SpanContext, SpanId, SpanKind, Status, TraceContextExt, TraceId, Tracer, TracerProvider,
},
Context, ContextGuard, Key, KeyValue, TraceFlags, Value,
};
use emit::runtime::AmbientSlot;
use std::{borrow::Cow, cell::RefCell, fmt, mem, ops::ControlFlow, sync::Arc};
mod internal_metrics;
pub use internal_metrics::*;
pub fn setup<L: Logger + Send + Sync + 'static, T: Tracer + Send + Sync + 'static>(
logger_provider: impl LoggerProvider<Logger = L>,
tracer_provider: impl TracerProvider<Tracer = T>,
) -> Setup<L, T>
where
T::Span: Send + Sync + 'static,
{
let name = "emit";
let metrics = Arc::new(InternalMetrics::default());
Setup {
inner: emit1::setup()
.with_clock(emit::Empty)
.with_rng(emit::Empty)
.emit_to(OpenTelemetryEmitter::new(
metrics.clone(),
logger_provider.logger(name),
))
.emit_when(OpenTelemetryIncomingFilter {})
.map_ctxt(|ctxt| {
OpenTelemetryCtxt::wrap(metrics.clone(), tracer_provider.tracer(name), ctxt)
}),
metrics,
}
}
pub struct Setup<L, T> {
metrics: Arc<InternalMetrics>,
inner: emit1::Setup<
OpenTelemetryEmitter<L>,
OpenTelemetryIncomingFilter,
OpenTelemetryCtxt<emit1::setup::DefaultCtxt, T>,
emit::Empty,
emit::Empty,
>,
}
impl<L: Logger + Send + Sync + 'static, T: Tracer + Send + Sync + 'static> Setup<L, T>
where
T::Span: Send + Sync + 'static,
{
pub fn with_span_name(
self,
writer: impl Fn(
&emit::event::Event<&dyn emit::props::ErasedProps>,
&mut fmt::Formatter,
) -> fmt::Result
+ Send
+ Sync
+ 'static,
) -> Self {
Setup {
metrics: self.metrics,
inner: self.inner.map_emitter(|mut emitter| {
emitter.span_name = Box::new(writer);
emitter
}),
}
}
pub fn with_log_body(
self,
writer: impl Fn(
&emit::event::Event<&dyn emit::props::ErasedProps>,
&mut fmt::Formatter,
) -> fmt::Result
+ Send
+ Sync
+ 'static,
) -> Self {
Setup {
metrics: self.metrics,
inner: self.inner.map_emitter(|mut emitter| {
emitter.log_body = Box::new(writer);
emitter
}),
}
}
pub fn metric_source(&self) -> EmitOpenTelemetryMetrics {
EmitOpenTelemetryMetrics {
metrics: self.metrics.clone(),
}
}
#[cfg(feature = "implicit_rt")]
pub fn init(self) -> emit1::setup::Init<'static, impl emit::Emitter, impl emit::Ctxt> {
self.inner.init()
}
#[cfg(feature = "implicit_rt")]
pub fn try_init(
self,
) -> Option<emit1::setup::Init<'static, impl emit::Emitter, impl emit::Ctxt>> {
self.inner.try_init()
}
pub fn init_slot<'a>(
self,
slot: &'a AmbientSlot,
) -> emit1::setup::Init<'a, impl emit::Emitter, impl emit::Ctxt> {
self.inner.init_slot(slot)
}
pub fn try_init_slot<'a>(
self,
slot: &'a AmbientSlot,
) -> Option<emit1::setup::Init<'a, impl emit::Emitter, impl emit::Ctxt>> {
self.inner.try_init_slot(slot)
}
}
struct OpenTelemetryCtxt<C, T> {
tracer: T,
metrics: Arc<InternalMetrics>,
inner: C,
}
struct OpenTelemetryFrame<F> {
slot: Option<Context>,
active: bool,
options: OpenTelemetryFrameOptions,
inner: F,
}
#[derive(Debug, Clone, Copy, Default)]
struct OpenTelemetryFrameOptions {
end_on_close: bool,
update_default_name_on_close: bool,
}
struct OpenTelemetryProps<P: ?Sized> {
ctxt: emit::span::SpanCtxt,
inner: *const P,
}
impl<C, T> OpenTelemetryCtxt<C, T> {
fn wrap(metrics: Arc<InternalMetrics>, tracer: T, ctxt: C) -> Self {
OpenTelemetryCtxt {
tracer,
inner: ctxt,
metrics,
}
}
}
impl<P: emit::Props + ?Sized> emit::Props for OpenTelemetryProps<P> {
fn for_each<'kv, F: FnMut(emit::Str<'kv>, emit::Value<'kv>) -> ControlFlow<()>>(
&'kv self,
mut for_each: F,
) -> ControlFlow<()> {
self.ctxt.for_each(&mut for_each)?;
unsafe { &*self.inner }.for_each(|k, v| {
if k != KEY_TRACE_ID && k != KEY_SPAN_ID {
for_each(k, v)?;
}
ControlFlow::Continue(())
})
}
}
thread_local! {
static ACTIVE_FRAME_STACK: RefCell<Vec<ActiveFrame>> = RefCell::new(Vec::new());
}
struct ActiveFrame {
_guard: ContextGuard,
options: OpenTelemetryFrameOptions,
}
fn push(guard: ActiveFrame) {
ACTIVE_FRAME_STACK.with(|stack| stack.borrow_mut().push(guard));
}
fn pop() -> Option<ActiveFrame> {
ACTIVE_FRAME_STACK.with(|stack| stack.borrow_mut().pop())
}
fn with_current(f: impl FnOnce(&mut ActiveFrame)) {
ACTIVE_FRAME_STACK.with(|stack| {
if let Some(frame) = stack.borrow_mut().last_mut() {
f(frame);
}
})
}
impl<C: emit::Ctxt, T: Tracer> emit::Ctxt for OpenTelemetryCtxt<C, T>
where
T::Span: Send + Sync + 'static,
{
type Current = OpenTelemetryProps<C::Current>;
type Frame = OpenTelemetryFrame<C::Frame>;
fn open_root<P: emit::Props>(&self, props: P) -> Self::Frame {
let (incoming, props) = incoming_span_ctxt(&self.tracer, props);
let (slot, options) = match incoming {
IncomingSpanCtxt::None => (Some(Context::new()), Default::default()),
IncomingSpanCtxt::Span(context, options) => (Some(context), options),
};
OpenTelemetryFrame {
active: false,
options,
slot,
inner: self.inner.open_root(props),
}
}
fn open_disabled<P: emit::Props>(&self, props: P) -> Self::Frame {
OpenTelemetryFrame {
active: false,
options: Default::default(),
slot: None,
inner: self.inner.open_disabled(props),
}
}
fn open_push<P: emit::Props>(&self, props: P) -> Self::Frame {
let (incoming, props) = incoming_span_ctxt(&self.tracer, props);
let (slot, options) = match incoming {
IncomingSpanCtxt::None => (None, Default::default()),
IncomingSpanCtxt::Span(context, options) => (Some(context), options),
};
OpenTelemetryFrame {
active: false,
options,
slot,
inner: self.inner.open_push(props),
}
}
fn enter(&self, local: &mut Self::Frame) {
self.inner.enter(&mut local.inner);
if let Some(ctxt) = local.slot.take() {
let guard = ctxt.attach();
push(ActiveFrame {
_guard: guard,
options: local.options,
});
local.active = true;
}
}
fn with_current<R, F: FnOnce(&Self::Current) -> R>(&self, with: F) -> R {
let ctxt = Context::current();
let span = ctxt.span();
let trace_id = span.span_context().trace_id().to_bytes();
let span_id = span.span_context().span_id().to_bytes();
self.inner.with_current(|props| {
let props = OpenTelemetryProps {
ctxt: emit::span::SpanCtxt::new(
emit::TraceId::from_bytes(trace_id),
None,
emit::SpanId::from_bytes(span_id),
),
inner: props as *const C::Current,
};
with(&props)
})
}
fn exit(&self, frame: &mut Self::Frame) {
self.inner.exit(&mut frame.inner);
if frame.active {
frame.slot = Some(Context::current());
if let Some(active) = pop() {
frame.options = active.options;
}
frame.active = false;
}
}
fn close(&self, mut frame: Self::Frame) {
if frame.options.end_on_close {
if let Some(ctxt) = frame.slot.take() {
self.metrics.span_unexpected_close.increment();
let span = ctxt.span();
span.end();
}
}
}
}
thread_local! {
static INCOMING_SPAN_PROPS: RefCell<IncomingSpanProps> = RefCell::new(Default::default());
}
#[derive(Default)]
struct IncomingSpanProps {
trace_id: Option<emit::TraceId>,
span_parent: Option<emit::SpanId>,
span_name: Option<emit::Str<'static>>,
span_kind: Option<emit::span::SpanKind>,
}
fn populate_incoming_span_ctxt(props: impl emit::Props) {
let span_kind = props.pull::<emit::span::SpanKind, _>(KEY_SPAN_KIND);
let span_name = props
.pull::<emit::Str, _>(KEY_SPAN_NAME)
.map(|v| v.to_owned());
let trace_id = props.pull::<emit::TraceId, _>(KEY_TRACE_ID);
let span_parent = props.pull::<emit::SpanId, _>(KEY_SPAN_PARENT);
INCOMING_SPAN_PROPS.with(|incoming| {
*incoming.borrow_mut() = IncomingSpanProps {
trace_id,
span_parent,
span_name,
span_kind,
}
});
}
fn incoming_span_ctxt<T: Tracer>(
tracer: &T,
props: impl emit::Props,
) -> (IncomingSpanCtxt, impl emit::Props)
where
T::Span: Send + Sync + 'static,
{
let IncomingSpanProps {
trace_id,
span_parent,
span_name,
span_kind,
} = INCOMING_SPAN_PROPS.with(|incoming| mem::take(&mut *incoming.borrow_mut()));
let ctxt = Context::current();
let trace_id = trace_id.map(otel_trace_id);
let span_parent = span_parent.map(otel_span_id);
if trace_id != Some(ctxt.span().span_context().trace_id())
|| span_parent != Some(ctxt.span().span_context().span_id())
{
return (
IncomingSpanCtxt::None,
ExcludeSpanCtxtProps {
exclude_ctxt_props: true,
exclude_span_parent: true,
inner: props,
},
);
};
let span_kind = span_kind
.map(|kind| match kind {
emit::span::SpanKind::Client => SpanKind::Client,
emit::span::SpanKind::Server => SpanKind::Server,
emit::span::SpanKind::Producer => SpanKind::Producer,
emit::span::SpanKind::Consumer => SpanKind::Consumer,
_ => SpanKind::Internal,
})
.unwrap_or(SpanKind::Internal);
let span_name = span_name.map(|name| name.into_cow());
let update_default_name_on_close = span_name.is_none();
let span = tracer
.span_builder(span_name.unwrap_or(Cow::Borrowed("emit_span")))
.with_kind(span_kind);
let ctxt = ctxt.with_span(span.start(tracer));
(
IncomingSpanCtxt::Span(
ctxt,
OpenTelemetryFrameOptions {
update_default_name_on_close,
end_on_close: true,
},
),
ExcludeSpanCtxtProps {
exclude_ctxt_props: true,
exclude_span_parent: false,
inner: props,
},
)
}
enum IncomingSpanCtxt {
None,
Span(Context, OpenTelemetryFrameOptions),
}
struct ExcludeSpanCtxtProps<P> {
exclude_ctxt_props: bool,
exclude_span_parent: bool,
inner: P,
}
impl<P: emit::Props> emit::Props for ExcludeSpanCtxtProps<P> {
fn for_each<'kv, F: FnMut(emit::Str<'kv>, emit::Value<'kv>) -> ControlFlow<()>>(
&'kv self,
mut for_each: F,
) -> ControlFlow<()> {
self.inner.for_each(|key, value| match key.get() {
KEY_TRACE_ID | KEY_SPAN_ID | KEY_SPAN_NAME | KEY_SPAN_KIND | KEY_SPAN_LINKS
if self.exclude_ctxt_props =>
{
ControlFlow::Continue(())
}
KEY_SPAN_PARENT if self.exclude_span_parent => ControlFlow::Continue(()),
_ => for_each(key, value),
})
}
}
struct OpenTelemetryEmitter<L> {
logger: L,
span_name: Box<MessageFormatter>,
log_body: Box<MessageFormatter>,
metrics: Arc<InternalMetrics>,
}
type MessageFormatter = dyn Fn(&emit::event::Event<&dyn emit::props::ErasedProps>, &mut fmt::Formatter) -> fmt::Result
+ Send
+ Sync;
fn default_span_name() -> Box<MessageFormatter> {
Box::new(|evt, f| {
if let Some(name) = evt.props().get(KEY_SPAN_NAME) {
write!(f, "{}", name)
} else {
write!(f, "{}", evt.msg())
}
})
}
fn default_log_body() -> Box<MessageFormatter> {
Box::new(|evt, f| write!(f, "{}", evt.msg()))
}
struct MessageRenderer<'a, P> {
pub fmt: &'a MessageFormatter,
pub evt: &'a emit::event::Event<'a, P>,
}
impl<'a, P: emit::props::Props> fmt::Display for MessageRenderer<'a, P> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
(self.fmt)(&self.evt.erase(), f)
}
}
impl<L> OpenTelemetryEmitter<L> {
fn new(metrics: Arc<InternalMetrics>, logger: L) -> Self {
OpenTelemetryEmitter {
logger,
span_name: default_span_name(),
log_body: default_log_body(),
metrics,
}
}
}
impl<L: Logger> emit::Emitter for OpenTelemetryEmitter<L> {
fn emit<E: emit::event::ToEvent>(&self, evt: E) {
let evt = evt.to_event();
if emit::kind::is_span_filter().matches(&evt) {
let mut emitted = false;
with_current(|frame| {
if frame.options.end_on_close {
let ctxt = Context::current();
let span = ctxt.span();
let span_id = span.span_context().span_id();
let evt_span_id = evt
.props()
.pull::<emit::SpanId, _>(KEY_SPAN_ID)
.map(otel_span_id);
if Some(span_id) == evt_span_id {
let mut status = Status::Ok;
let mut status_is_descriptive_err = false;
let mut span_links = Vec::new();
let _ = evt.props().dedup().for_each(|k, v| {
if k == KEY_LVL {
if let Some(emit::Level::Error) = v.by_ref().cast::<emit::Level>() {
if !status_is_descriptive_err {
status = Status::error("error");
}
return ControlFlow::Continue(());
}
}
if k == KEY_ERR {
span.add_event(
"exception",
vec![KeyValue::new("exception.message", v.to_string())],
);
status = Status::error(v.to_string());
status_is_descriptive_err = true;
return ControlFlow::Continue(());
}
if k == KEY_SPAN_LINKS {
span_links = otel_span_links(
v,
span.span_context().trace_flags(),
span.span_context().is_remote(),
span.span_context().trace_state(),
);
return ControlFlow::Continue(());
}
if k == KEY_SPAN_NAME {
span.update_name(
v.by_ref()
.cast::<emit::Str>()
.map(|v| v.into_cow())
.unwrap_or_else(|| Cow::Owned(v.to_string())),
);
frame.options.update_default_name_on_close = false;
}
if k == KEY_TRACE_ID
|| k == KEY_SPAN_ID
|| k == KEY_SPAN_PARENT
|| k == KEY_EVT_KIND
|| k == KEY_SPAN_KIND
{
return ControlFlow::Continue(());
}
if let Some(v) = otel_span_value(v) {
span.set_attribute(KeyValue::new(k.to_cow(), v));
}
ControlFlow::Continue(())
});
if frame.options.update_default_name_on_close {
span.update_name(
MessageRenderer {
fmt: &self.span_name,
evt: &evt,
}
.to_string(),
);
}
span.set_status(status);
for link in span_links {
span.add_link(link.span_context, Default::default());
}
if let Some(extent) = evt.extent().and_then(|ex| ex.as_range()) {
span.end_with_timestamp(extent.end.to_system_time());
} else {
span.end();
}
frame.options.end_on_close = false;
emitted = true;
}
}
});
if !emitted {
self.metrics.span_unexpected_emit.increment();
}
return;
}
let mut record = self.logger.create_log_record();
let body = format!(
"{}",
MessageRenderer {
fmt: &self.log_body,
evt: &evt,
}
);
record.set_body(AnyValue::String(body.into()));
let mut trace_id = None;
let mut span_id = None;
let mut attributes = Vec::new();
{
let _ = evt.props().for_each(|k, v| {
if k == KEY_LVL {
match v.by_ref().cast::<emit::Level>() {
Some(emit::Level::Debug) => {
record.set_severity_number(Severity::Debug);
record.set_severity_text(LVL_DEBUG);
}
Some(emit::Level::Info) => {
record.set_severity_number(Severity::Info);
record.set_severity_text(LVL_INFO);
}
Some(emit::Level::Warn) => {
record.set_severity_number(Severity::Warn);
record.set_severity_text(LVL_WARN);
}
Some(emit::Level::Error) => {
record.set_severity_number(Severity::Error);
record.set_severity_text(LVL_ERROR);
}
None => {
record.set_severity_text("unknown");
}
}
return ControlFlow::Continue(());
}
if k == KEY_TRACE_ID {
if let Some(id) = v.by_ref().cast::<emit::TraceId>() {
trace_id = Some(otel_trace_id(id));
return ControlFlow::Continue(());
}
}
if k == KEY_SPAN_ID {
if let Some(id) = v.by_ref().cast::<emit::SpanId>() {
span_id = Some(otel_span_id(id));
return ControlFlow::Continue(());
}
}
if k == KEY_SPAN_PARENT {
if v.by_ref().cast::<emit::SpanId>().is_some() {
return ControlFlow::Continue(());
}
}
if k == KEY_SPAN_NAME || k == KEY_SPAN_KIND || k == KEY_SPAN_LINKS {
return ControlFlow::Continue(());
}
if let Some(v) = otel_log_value(v) {
attributes.push((Key::new(k.to_cow()), v));
}
ControlFlow::Continue(())
});
}
record.add_attributes(attributes);
if let Some(extent) = evt.extent() {
record.set_timestamp(extent.as_point().to_system_time());
}
let context = Context::current();
let span = context.span();
let span = span.span_context();
let trace_id = trace_id.unwrap_or_else(|| span.trace_id());
let span_id = span_id.unwrap_or_else(|| span.span_id());
let trace_flags = Some(span.trace_flags());
record.set_trace_context(trace_id, span_id, trace_flags);
self.logger.emit(record);
}
fn blocking_flush(&self, _: std::time::Duration) -> bool {
false
}
}
struct OpenTelemetryIncomingFilter {}
impl emit::Filter for OpenTelemetryIncomingFilter {
fn matches<E: emit::event::ToEvent>(&self, evt: E) -> bool {
let evt = evt.to_event();
if emit::kind::is_span_filter().matches(&evt) {
if Context::current().span().span_context().is_sampled() {
populate_incoming_span_ctxt(evt.props());
true
} else {
false
}
} else {
true
}
}
}
fn otel_trace_id(trace_id: emit::TraceId) -> TraceId {
TraceId::from_bytes(trace_id.to_bytes())
}
fn otel_span_id(span_id: emit::SpanId) -> SpanId {
SpanId::from_bytes(span_id.to_bytes())
}
fn otel_span_value(v: emit::Value) -> Option<Value> {
match any_value::serialize(&v) {
Ok(Some(av)) => match av {
AnyValue::Int(v) => Some(Value::I64(v)),
AnyValue::Double(v) => Some(Value::F64(v)),
AnyValue::String(v) => Some(Value::String(v)),
AnyValue::Boolean(v) => Some(Value::Bool(v)),
_ => Some(Value::String(v.to_string().into())),
},
Ok(None) => None,
Err(()) => Some(Value::String(v.to_string().into())),
}
}
fn otel_log_value(v: emit::Value) -> Option<AnyValue> {
match any_value::serialize(&v) {
Ok(v) => v,
Err(()) => Some(AnyValue::String(v.to_string().into())),
}
}
fn otel_span_links(
v: emit::Value,
trace_flags: TraceFlags,
is_remote: bool,
trace_state: &TraceState,
) -> Vec<Link> {
use serde::ser::{
Error, Impossible, Serialize, SerializeSeq, SerializeTuple, SerializeTupleStruct,
SerializeTupleVariant, Serializer, StdError,
};
struct Extract<'a> {
links: Vec<Link>,
trace_flags: TraceFlags,
is_remote: bool,
trace_state: &'a TraceState,
}
#[derive(Debug)]
struct ExtractError;
impl fmt::Display for ExtractError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str("unsupported value for span links")
}
}
impl StdError for ExtractError {}
impl Error for ExtractError {
fn custom<T>(_: T) -> Self
where
T: fmt::Display,
{
ExtractError
}
}
impl<'a> Serializer for Extract<'a> {
type Ok = Vec<Link>;
type Error = ExtractError;
type SerializeSeq = Self;
type SerializeTuple = Self;
type SerializeTupleStruct = Self;
type SerializeTupleVariant = Self;
type SerializeMap = Impossible<Self::Ok, Self::Error>;
type SerializeStruct = Impossible<Self::Ok, Self::Error>;
type SerializeStructVariant = Impossible<Self::Ok, Self::Error>;
fn serialize_bool(self, _: bool) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_i8(self, _: i8) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_i16(self, _: i16) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_i32(self, _: i32) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_i64(self, _: i64) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_u8(self, _: u8) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_u16(self, _: u16) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_u32(self, _: u32) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_u64(self, _: u64) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_f32(self, _: f32) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_f64(self, _: f64) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_char(self, _: char) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_str(self, _: &str) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_bytes(self, _: &[u8]) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_none(self) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_some<T>(self, value: &T) -> Result<Self::Ok, Self::Error>
where
T: ?Sized + Serialize,
{
value.serialize(self)
}
fn serialize_unit(self) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_unit_struct(self, _: &'static str) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_unit_variant(
self,
_: &'static str,
_: u32,
_: &'static str,
) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_newtype_struct<T>(
self,
_: &'static str,
value: &T,
) -> Result<Self::Ok, Self::Error>
where
T: ?Sized + Serialize,
{
value.serialize(self)
}
fn serialize_newtype_variant<T>(
self,
_: &'static str,
_: u32,
_: &'static str,
value: &T,
) -> Result<Self::Ok, Self::Error>
where
T: ?Sized + Serialize,
{
value.serialize(self)
}
fn serialize_seq(self, _: Option<usize>) -> Result<Self::SerializeSeq, Self::Error> {
Ok(self)
}
fn serialize_tuple(self, _: usize) -> Result<Self::SerializeTuple, Self::Error> {
Ok(self)
}
fn serialize_tuple_struct(
self,
_: &'static str,
_: usize,
) -> Result<Self::SerializeTupleStruct, Self::Error> {
Ok(self)
}
fn serialize_tuple_variant(
self,
_: &'static str,
_: u32,
_: &'static str,
_: usize,
) -> Result<Self::SerializeTupleVariant, Self::Error> {
Ok(self)
}
fn serialize_map(self, _: Option<usize>) -> Result<Self::SerializeMap, Self::Error> {
Err(ExtractError)
}
fn serialize_struct(
self,
_: &'static str,
_: usize,
) -> Result<Self::SerializeStruct, Self::Error> {
Err(ExtractError)
}
fn serialize_struct_variant(
self,
_: &'static str,
_: u32,
_: &'static str,
_: usize,
) -> Result<Self::SerializeStructVariant, Self::Error> {
Err(ExtractError)
}
}
impl<'a> SerializeSeq for Extract<'a> {
type Ok = Vec<Link>;
type Error = ExtractError;
fn serialize_element<T>(&mut self, value: &T) -> Result<(), Self::Error>
where
T: ?Sized + Serialize,
{
self.links.push(value.serialize(ParseLink {
trace_flags: self.trace_flags,
is_remote: self.is_remote,
trace_state: self.trace_state,
})?);
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(self.links)
}
}
impl<'a> SerializeTuple for Extract<'a> {
type Ok = Vec<Link>;
type Error = ExtractError;
fn serialize_element<T>(&mut self, value: &T) -> Result<(), Self::Error>
where
T: ?Sized + Serialize,
{
self.links.push(value.serialize(ParseLink {
trace_flags: self.trace_flags,
is_remote: self.is_remote,
trace_state: self.trace_state,
})?);
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(self.links)
}
}
impl<'a> SerializeTupleStruct for Extract<'a> {
type Ok = Vec<Link>;
type Error = ExtractError;
fn serialize_field<T>(&mut self, value: &T) -> Result<(), Self::Error>
where
T: ?Sized + Serialize,
{
self.links.push(value.serialize(ParseLink {
trace_flags: self.trace_flags,
is_remote: self.is_remote,
trace_state: self.trace_state,
})?);
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(self.links)
}
}
impl<'a> SerializeTupleVariant for Extract<'a> {
type Ok = Vec<Link>;
type Error = ExtractError;
fn serialize_field<T>(&mut self, value: &T) -> Result<(), Self::Error>
where
T: ?Sized + Serialize,
{
self.links.push(value.serialize(ParseLink {
trace_flags: self.trace_flags,
is_remote: self.is_remote,
trace_state: &self.trace_state,
})?);
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(self.links)
}
}
struct ParseLink<'a> {
trace_flags: TraceFlags,
is_remote: bool,
trace_state: &'a TraceState,
}
impl<'a> Serializer for ParseLink<'a> {
type Ok = Link;
type Error = ExtractError;
type SerializeSeq = Impossible<Self::Ok, Self::Error>;
type SerializeTuple = Impossible<Self::Ok, Self::Error>;
type SerializeTupleStruct = Impossible<Self::Ok, Self::Error>;
type SerializeTupleVariant = Impossible<Self::Ok, Self::Error>;
type SerializeMap = Impossible<Self::Ok, Self::Error>;
type SerializeStruct = Impossible<Self::Ok, Self::Error>;
type SerializeStructVariant = Impossible<Self::Ok, Self::Error>;
fn serialize_bool(self, _: bool) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_i8(self, _: i8) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_i16(self, _: i16) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_i32(self, _: i32) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_i64(self, _: i64) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_u8(self, _: u8) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_u16(self, _: u16) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_u32(self, _: u32) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_u64(self, _: u64) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_f32(self, _: f32) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_f64(self, _: f64) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_char(self, _: char) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_str(self, v: &str) -> Result<Self::Ok, Self::Error> {
let link = emit::span::SpanLink::try_from_str(v).map_err(|_| ExtractError)?;
Ok(Link::new(
SpanContext::new(
otel_trace_id(*link.trace_id()),
otel_span_id(*link.span_id()),
self.trace_flags,
self.is_remote,
self.trace_state.clone(),
),
Vec::new(),
0,
))
}
fn serialize_bytes(self, _: &[u8]) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_none(self) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_some<T>(self, value: &T) -> Result<Self::Ok, Self::Error>
where
T: ?Sized + Serialize,
{
value.serialize(self)
}
fn serialize_unit(self) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_unit_struct(self, _: &'static str) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_unit_variant(
self,
_: &'static str,
_: u32,
_: &'static str,
) -> Result<Self::Ok, Self::Error> {
Err(ExtractError)
}
fn serialize_newtype_struct<T>(
self,
_: &'static str,
value: &T,
) -> Result<Self::Ok, Self::Error>
where
T: ?Sized + Serialize,
{
value.serialize(self)
}
fn serialize_newtype_variant<T>(
self,
_: &'static str,
_: u32,
_: &'static str,
value: &T,
) -> Result<Self::Ok, Self::Error>
where
T: ?Sized + Serialize,
{
value.serialize(self)
}
fn serialize_seq(self, _: Option<usize>) -> Result<Self::SerializeSeq, Self::Error> {
Err(ExtractError)
}
fn serialize_tuple(self, _: usize) -> Result<Self::SerializeTuple, Self::Error> {
Err(ExtractError)
}
fn serialize_tuple_struct(
self,
_: &'static str,
_: usize,
) -> Result<Self::SerializeTupleStruct, Self::Error> {
Err(ExtractError)
}
fn serialize_tuple_variant(
self,
_: &'static str,
_: u32,
_: &'static str,
_: usize,
) -> Result<Self::SerializeTupleVariant, Self::Error> {
Err(ExtractError)
}
fn serialize_map(self, _: Option<usize>) -> Result<Self::SerializeMap, Self::Error> {
Err(ExtractError)
}
fn serialize_struct(
self,
_: &'static str,
_: usize,
) -> Result<Self::SerializeStruct, Self::Error> {
Err(ExtractError)
}
fn serialize_struct_variant(
self,
_: &'static str,
_: u32,
_: &'static str,
_: usize,
) -> Result<Self::SerializeStructVariant, Self::Error> {
Err(ExtractError)
}
}
v.serialize(Extract {
links: Vec::new(),
trace_flags,
is_remote,
trace_state,
})
.unwrap_or_default()
}
mod any_value {
use std::{collections::HashMap, fmt};
use opentelemetry::{logs::AnyValue, Key, StringValue};
use serde::ser::{
Error, Serialize, SerializeMap, SerializeSeq, SerializeStruct, SerializeStructVariant,
SerializeTuple, SerializeTupleStruct, SerializeTupleVariant, Serializer, StdError,
};
pub(crate) fn serialize(value: impl Serialize) -> Result<Option<AnyValue>, ()> {
value.serialize(ValueSerializer).map_err(|_| ())
}
struct ValueSerializer;
struct ValueSerializeSeq {
value: Vec<AnyValue>,
}
struct ValueSerializeTuple {
value: Vec<AnyValue>,
}
struct ValueSerializeTupleStruct {
value: Vec<AnyValue>,
}
struct ValueSerializeMap {
key: Option<Key>,
value: HashMap<Key, AnyValue>,
}
struct ValueSerializeStruct {
value: HashMap<Key, AnyValue>,
}
struct ValueSerializeTupleVariant {
variant: &'static str,
value: Vec<AnyValue>,
}
struct ValueSerializeStructVariant {
variant: &'static str,
value: HashMap<Key, AnyValue>,
}
#[derive(Debug)]
struct ValueError(String);
impl fmt::Display for ValueError {
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
fmt::Display::fmt(&self.0, f)
}
}
impl Error for ValueError {
fn custom<T>(msg: T) -> Self
where
T: fmt::Display,
{
ValueError(msg.to_string())
}
}
impl StdError for ValueError {}
impl Serializer for ValueSerializer {
type Ok = Option<AnyValue>;
type Error = ValueError;
type SerializeSeq = ValueSerializeSeq;
type SerializeTuple = ValueSerializeTuple;
type SerializeTupleStruct = ValueSerializeTupleStruct;
type SerializeTupleVariant = ValueSerializeTupleVariant;
type SerializeMap = ValueSerializeMap;
type SerializeStruct = ValueSerializeStruct;
type SerializeStructVariant = ValueSerializeStructVariant;
fn serialize_bool(self, v: bool) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::Boolean(v)))
}
fn serialize_i8(self, v: i8) -> Result<Self::Ok, Self::Error> {
self.serialize_i64(v as i64)
}
fn serialize_i16(self, v: i16) -> Result<Self::Ok, Self::Error> {
self.serialize_i64(v as i64)
}
fn serialize_i32(self, v: i32) -> Result<Self::Ok, Self::Error> {
self.serialize_i64(v as i64)
}
fn serialize_i64(self, v: i64) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::Int(v)))
}
fn serialize_i128(self, v: i128) -> Result<Self::Ok, Self::Error> {
if let Ok(v) = v.try_into() {
self.serialize_i64(v)
} else {
self.collect_str(&v)
}
}
fn serialize_u8(self, v: u8) -> Result<Self::Ok, Self::Error> {
self.serialize_i64(v as i64)
}
fn serialize_u16(self, v: u16) -> Result<Self::Ok, Self::Error> {
self.serialize_i64(v as i64)
}
fn serialize_u32(self, v: u32) -> Result<Self::Ok, Self::Error> {
self.serialize_i64(v as i64)
}
fn serialize_u64(self, v: u64) -> Result<Self::Ok, Self::Error> {
if let Ok(v) = v.try_into() {
self.serialize_i64(v)
} else {
self.collect_str(&v)
}
}
fn serialize_u128(self, v: u128) -> Result<Self::Ok, Self::Error> {
if let Ok(v) = v.try_into() {
self.serialize_i64(v)
} else {
self.collect_str(&v)
}
}
fn serialize_f32(self, v: f32) -> Result<Self::Ok, Self::Error> {
self.serialize_f64(v as f64)
}
fn serialize_f64(self, v: f64) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::Double(v)))
}
fn serialize_char(self, v: char) -> Result<Self::Ok, Self::Error> {
self.collect_str(&v)
}
fn serialize_str(self, v: &str) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::String(StringValue::from(v.to_owned()))))
}
fn serialize_bytes(self, v: &[u8]) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::Bytes(Box::new(v.to_owned()))))
}
fn serialize_none(self) -> Result<Self::Ok, Self::Error> {
Ok(None)
}
fn serialize_some<T: Serialize + ?Sized>(self, value: &T) -> Result<Self::Ok, Self::Error> {
value.serialize(self)
}
fn serialize_unit(self) -> Result<Self::Ok, Self::Error> {
Ok(None)
}
fn serialize_unit_struct(self, name: &'static str) -> Result<Self::Ok, Self::Error> {
name.serialize(self)
}
fn serialize_unit_variant(
self,
_: &'static str,
_: u32,
variant: &'static str,
) -> Result<Self::Ok, Self::Error> {
variant.serialize(self)
}
fn serialize_newtype_struct<T: Serialize + ?Sized>(
self,
_: &'static str,
value: &T,
) -> Result<Self::Ok, Self::Error> {
value.serialize(self)
}
fn serialize_newtype_variant<T: Serialize + ?Sized>(
self,
_: &'static str,
_: u32,
variant: &'static str,
value: &T,
) -> Result<Self::Ok, Self::Error> {
let mut map = self.serialize_map(Some(1))?;
map.serialize_entry(variant, value)?;
map.end()
}
fn serialize_seq(self, _: Option<usize>) -> Result<Self::SerializeSeq, Self::Error> {
Ok(ValueSerializeSeq { value: Vec::new() })
}
fn serialize_tuple(self, _: usize) -> Result<Self::SerializeTuple, Self::Error> {
Ok(ValueSerializeTuple { value: Vec::new() })
}
fn serialize_tuple_struct(
self,
_: &'static str,
_: usize,
) -> Result<Self::SerializeTupleStruct, Self::Error> {
Ok(ValueSerializeTupleStruct { value: Vec::new() })
}
fn serialize_tuple_variant(
self,
_: &'static str,
_: u32,
variant: &'static str,
_: usize,
) -> Result<Self::SerializeTupleVariant, Self::Error> {
Ok(ValueSerializeTupleVariant {
variant,
value: Vec::new(),
})
}
fn serialize_map(self, _: Option<usize>) -> Result<Self::SerializeMap, Self::Error> {
Ok(ValueSerializeMap {
key: None,
value: HashMap::new(),
})
}
fn serialize_struct(
self,
_: &'static str,
_: usize,
) -> Result<Self::SerializeStruct, Self::Error> {
Ok(ValueSerializeStruct {
value: HashMap::new(),
})
}
fn serialize_struct_variant(
self,
_: &'static str,
_: u32,
variant: &'static str,
_: usize,
) -> Result<Self::SerializeStructVariant, Self::Error> {
Ok(ValueSerializeStructVariant {
variant,
value: HashMap::new(),
})
}
}
impl SerializeSeq for ValueSerializeSeq {
type Ok = Option<AnyValue>;
type Error = ValueError;
fn serialize_element<T: Serialize + ?Sized>(
&mut self,
value: &T,
) -> Result<(), Self::Error> {
if let Some(value) = value.serialize(ValueSerializer)? {
self.value.push(value);
}
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::ListAny(Box::new(self.value))))
}
}
impl SerializeTuple for ValueSerializeTuple {
type Ok = Option<AnyValue>;
type Error = ValueError;
fn serialize_element<T: Serialize + ?Sized>(
&mut self,
value: &T,
) -> Result<(), Self::Error> {
if let Some(value) = value.serialize(ValueSerializer)? {
self.value.push(value);
}
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::ListAny(Box::new(self.value))))
}
}
impl SerializeTupleStruct for ValueSerializeTupleStruct {
type Ok = Option<AnyValue>;
type Error = ValueError;
fn serialize_field<T: Serialize + ?Sized>(&mut self, value: &T) -> Result<(), Self::Error> {
if let Some(value) = value.serialize(ValueSerializer)? {
self.value.push(value);
}
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::ListAny(Box::new(self.value))))
}
}
impl SerializeTupleVariant for ValueSerializeTupleVariant {
type Ok = Option<AnyValue>;
type Error = ValueError;
fn serialize_field<T: Serialize + ?Sized>(&mut self, value: &T) -> Result<(), Self::Error> {
if let Some(value) = value.serialize(ValueSerializer)? {
self.value.push(value);
}
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::Map({
let mut variant = HashMap::new();
variant.insert(
Key::from(self.variant),
AnyValue::ListAny(Box::new(self.value)),
);
Box::new(variant)
})))
}
}
impl SerializeMap for ValueSerializeMap {
type Ok = Option<AnyValue>;
type Error = ValueError;
fn serialize_key<T: Serialize + ?Sized>(&mut self, key: &T) -> Result<(), Self::Error> {
let key = match key.serialize(ValueSerializer)? {
Some(AnyValue::String(key)) => Key::from(String::from(key)),
key => Key::from(format!("{:?}", key)),
};
self.key = Some(key);
Ok(())
}
fn serialize_value<T: Serialize + ?Sized>(&mut self, value: &T) -> Result<(), Self::Error> {
let key = self
.key
.take()
.ok_or_else(|| Self::Error::custom("missing key"))?;
if let Some(value) = value.serialize(ValueSerializer)? {
self.value.insert(key, value);
}
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::Map(Box::new(self.value))))
}
}
impl SerializeStruct for ValueSerializeStruct {
type Ok = Option<AnyValue>;
type Error = ValueError;
fn serialize_field<T: Serialize + ?Sized>(
&mut self,
key: &'static str,
value: &T,
) -> Result<(), Self::Error> {
let key = Key::from(key);
if let Some(value) = value.serialize(ValueSerializer)? {
self.value.insert(key, value);
}
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::Map(Box::new(self.value))))
}
}
impl SerializeStructVariant for ValueSerializeStructVariant {
type Ok = Option<AnyValue>;
type Error = ValueError;
fn serialize_field<T: Serialize + ?Sized>(
&mut self,
key: &'static str,
value: &T,
) -> Result<(), Self::Error> {
let key = Key::from(key);
if let Some(value) = value.serialize(ValueSerializer)? {
self.value.insert(key, value);
}
Ok(())
}
fn end(self) -> Result<Self::Ok, Self::Error> {
Ok(Some(AnyValue::Map({
let mut variant = HashMap::new();
variant.insert(Key::from(self.variant), AnyValue::Map(Box::new(self.value)));
Box::new(variant)
})))
}
}
}
#[cfg(test)]
mod tests {
use std::io;
use super::*;
use emit::runtime::AmbientSlot;
use opentelemetry_sdk::trace::Sampler;
use opentelemetry_sdk::{
logs::{in_memory_exporter::InMemoryLogExporter, SdkLoggerProvider},
trace::{in_memory_exporter::InMemorySpanExporter, SdkTracerProvider},
};
fn build(
slot: &AmbientSlot,
) -> (
InMemoryLogExporter,
InMemorySpanExporter,
SdkLoggerProvider,
SdkTracerProvider,
) {
build_with_sampler(slot, Sampler::AlwaysOn)
}
fn build_with_sampler(
slot: &AmbientSlot,
sampler: Sampler,
) -> (
InMemoryLogExporter,
InMemorySpanExporter,
SdkLoggerProvider,
SdkTracerProvider,
) {
let logger_exporter = InMemoryLogExporter::default();
let logger_provider = SdkLoggerProvider::builder()
.with_simple_exporter(logger_exporter.clone())
.build();
let tracer_exporter = InMemorySpanExporter::default();
let tracer_provider = SdkTracerProvider::builder()
.with_simple_exporter(tracer_exporter.clone())
.with_sampler(sampler)
.build();
let _ = setup(logger_provider.clone(), tracer_provider.clone()).init_slot(slot);
(
logger_exporter,
tracer_exporter,
logger_provider,
tracer_provider,
)
}
fn otel_span<T>(tracer_provider: &SdkTracerProvider, in_span: impl FnOnce(Context) -> T) -> T {
use opentelemetry::trace::TracerProvider;
tracer_provider
.tracer("otel_span")
.in_span("otel span", in_span)
}
macro_rules! emit_span {
($slot:ident, $in_span:expr) => {{
#[emit::span(rt: $slot.get(), "emit span")]
fn emit_span<T>(in_span: impl FnOnce(emit::SpanCtxt) -> T) -> T {
in_span(emit::span::SpanCtxt::current($slot.get().ctxt()))
}
emit_span($in_span)
}};
}
fn emit_trace_id(trace_id: TraceId) -> emit::TraceId {
emit::TraceId::from_bytes(trace_id.to_bytes()).unwrap()
}
fn emit_span_id(span_id: SpanId) -> emit::SpanId {
emit::SpanId::from_bytes(span_id.to_bytes()).unwrap()
}
#[test]
fn emit_log() {
let slot = AmbientSlot::new();
let (logs, _, _, _) = build(&slot);
emit::emit!(rt: slot.get(), "test {attr}", attr: "log");
let logs = logs.get_emitted_logs().unwrap();
assert_eq!(1, logs.len());
let Some(AnyValue::String(body)) = logs[0].record.body() else {
panic!("unexpected log body value");
};
assert_eq!("test log", body.as_str());
}
#[test]
fn emit_log_trace_context() {
let slot = AmbientSlot::new();
let (logs, _, _, _) = build(&slot);
emit::emit!(rt: slot.get(), "test log", trace_id: "4bf92f3577b34da6a3ce929d0e0e4736", span_id: "00f067aa0ba902b7");
let logs = logs.get_emitted_logs().unwrap();
assert_eq!(1, logs.len());
assert_eq!(
"4bf92f3577b34da6a3ce929d0e0e4736",
logs[0].record.trace_context().unwrap().trace_id.to_string()
);
assert_eq!(
"00f067aa0ba902b7",
logs[0].record.trace_context().unwrap().span_id.to_string()
);
}
#[test]
fn otel_span_emit_span() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build(&SLOT);
let ctxt = otel_span(&tracer_provider, |_| emit_span!(SLOT, |ctxt| ctxt));
let spans = spans.get_finished_spans().unwrap();
assert_eq!(2, spans.len());
assert_eq!("emit span", spans[0].name);
assert_eq!(
ctxt.trace_id().unwrap().to_bytes(),
spans[0].span_context.trace_id().to_bytes()
);
assert_eq!(
ctxt.span_id().unwrap().to_bytes(),
spans[0].span_context.span_id().to_bytes()
);
assert_eq!(
ctxt.span_parent().unwrap().to_bytes(),
spans[0].parent_span_id.to_bytes()
);
}
#[test]
fn otel_span_unsampled_emit_span() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build_with_sampler(&SLOT, Sampler::AlwaysOff);
otel_span(&tracer_provider, |_| emit_span!(SLOT, |ctxt| ctxt));
let spans = spans.get_finished_spans().unwrap();
assert_eq!(0, spans.len());
}
#[test]
fn emit_span() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, _) = build(&SLOT);
let ctxt = emit_span!(SLOT, |ctxt| ctxt);
let spans = spans.get_finished_spans().unwrap();
assert_eq!(0, spans.len());
assert!(ctxt.trace_id().is_none());
assert!(ctxt.span_id().is_none());
assert!(ctxt.span_parent().is_none());
}
#[test]
fn otel_span_emit_span_direct() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build(&SLOT);
otel_span(&tracer_provider, |cx| {
emit::emit!(rt: SLOT.get(), evt: emit::Span::new(
emit::mdl!(),
emit::Empty,
emit::props! {
span_name: "emit span",
trace_id: emit_trace_id(cx.span().span_context().trace_id()),
span_parent: emit_span_id(cx.span().span_context().span_id()),
span_id: "00f067aa0ba902b7",
}
));
});
let spans = spans.get_finished_spans().unwrap();
assert_eq!(1, spans.len());
}
#[test]
fn otel_span_emit_span_name_default() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build(&SLOT);
otel_span(&tracer_provider, |_| emit_span!(SLOT, |_| {}));
let spans = spans.get_finished_spans().unwrap();
assert_eq!(2, spans.len());
assert_eq!("emit span", spans[0].name);
}
#[test]
fn otel_span_emit_span_name_control_param() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build(&SLOT);
otel_span(&tracer_provider, |_| {
#[emit::span(rt: SLOT.get(), name: "custom", "emit span")]
fn emit_span() {}
emit_span();
});
let spans = spans.get_finished_spans().unwrap();
assert_eq!(2, spans.len());
assert_eq!("custom", spans[0].name);
}
#[test]
fn otel_span_emit_span_name_guard_props() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build(&SLOT);
otel_span(&tracer_provider, |_| {
#[emit::span(rt: SLOT.get(), guard: span, "emit span")]
fn emit_span() {
let _span = span.with_name("custom");
}
emit_span();
});
let spans = spans.get_finished_spans().unwrap();
assert_eq!(2, spans.len());
assert_eq!("custom", spans[0].name);
}
#[test]
fn otel_span_emit_span_kind_ctxt_props() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build(&SLOT);
otel_span(&tracer_provider, |_| {
#[emit::span(rt: SLOT.get(), "emit span", span_kind: "client")]
fn emit_span() {}
emit_span();
});
let spans = spans.get_finished_spans().unwrap();
assert_eq!(2, spans.len());
assert_eq!(SpanKind::Client, spans[0].span_kind);
}
#[test]
fn otel_span_emit_span_kind_evt_props() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build(&SLOT);
otel_span(&tracer_provider, |_| {
#[emit::span(rt: SLOT.get(), evt_props: ("span_kind", "client"), "emit span")]
fn emit_span() {}
emit_span();
});
let spans = spans.get_finished_spans().unwrap();
assert_eq!(2, spans.len());
assert_eq!(SpanKind::Client, spans[0].span_kind);
}
#[test]
fn otel_span_emit_span_links_guard_props() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build(&SLOT);
otel_span(&tracer_provider, |_| {
#[emit::span(rt: SLOT.get(), guard: span, "emit span")]
fn emit_span() {
let _span = span.push_prop(
"span_links",
["4bf92f3577b34da6a3ce929d0e0e4736-00f067aa0ba902b7"],
);
}
emit_span();
});
let spans = spans.get_finished_spans().unwrap();
assert_eq!(2, spans.len());
assert_eq!(1, spans[0].links.links.len());
assert_eq!(
"4bf92f3577b34da6a3ce929d0e0e4736",
spans[0].links.links[0].span_context.trace_id().to_string()
);
assert_eq!(
"00f067aa0ba902b7",
spans[0].links.links[0].span_context.span_id().to_string()
);
}
#[test]
fn otel_span_emit_span_ok_control_param() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build(&SLOT);
otel_span(&tracer_provider, |_| {
#[emit::span(rt: SLOT.get(), ok_lvl: "info", "emit span")]
fn emit_span() -> Result<(), io::Error> {
Ok(())
}
let _ = emit_span();
});
let spans = spans.get_finished_spans().unwrap();
assert_eq!(2, spans.len());
assert_eq!(Status::Ok, spans[0].status);
}
#[test]
fn otel_span_emit_span_err_control_param() {
static SLOT: AmbientSlot = AmbientSlot::new();
let (_, spans, _, tracer_provider) = build(&SLOT);
otel_span(&tracer_provider, |_| {
#[emit::span(rt: SLOT.get(), ok_lvl: "info", "emit span")]
fn emit_span() -> Result<(), io::Error> {
Err(io::Error::new(io::ErrorKind::Other, "something failed!"))
}
let _ = emit_span();
});
let spans = spans.get_finished_spans().unwrap();
assert_eq!(2, spans.len());
assert_eq!(Status::error("something failed!"), spans[0].status);
}
#[test]
fn emit_value_to_otel_attribute() {
use opentelemetry::Key;
use std::collections::HashMap;
#[derive(serde::Serialize)]
struct Struct {
a: i32,
b: i32,
c: i32,
}
#[derive(serde::Serialize)]
struct Newtype(i32);
#[derive(serde::Serialize)]
enum Enum {
Unit,
Newtype(i32),
Struct { a: i32, b: i32, c: i32 },
Tuple(i32, i32, i32),
}
struct Bytes<B>(B);
impl<B: AsRef<[u8]>> serde::Serialize for Bytes<B> {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
serializer.serialize_bytes(self.0.as_ref())
}
}
struct Map {
a: i32,
b: i32,
c: i32,
}
impl serde::Serialize for Map {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
use serde::ser::SerializeMap;
let mut map = serializer.serialize_map(Some(3))?;
map.serialize_entry(&"a", &self.a)?;
map.serialize_entry(&"b", &self.b)?;
map.serialize_entry(&"c", &self.c)?;
map.end()
}
}
let slot = AmbientSlot::new();
let (logs, _, _, _) = build(&slot);
emit::emit!(
rt: slot.get(),
"test log",
str_value: "a string",
u8_value: 1u8,
u16_value: 2u16,
u32_value: 42u32,
u64_value: 2147483660u64,
u128_small_value: 2147483660u128,
u128_big_value: 9223372036854775820u128,
i8_value: 1i8,
i16_value: 2i16,
i32_value: 42i32,
i64_value: 2147483660i64,
i128_small_value: 2147483660i128,
i128_big_value: 9223372036854775820i128,
f64_value: 4.2,
bool_value: true,
#[emit::as_serde]
bytes_value: Bytes([1, 1, 1]),
#[emit::as_serde]
unit_value: (),
#[emit::as_serde]
some_value: Some(42),
#[emit::as_serde]
none_value: None::<i32>,
#[emit::as_serde]
slice_value: [1, 1, 1] as [i32; 3],
#[emit::as_serde]
map_value: Map { a: 1, b: 1, c: 1 },
#[emit::as_serde]
struct_value: Struct { a: 1, b: 1, c: 1 },
#[emit::as_serde]
tuple_value: (1, 1, 1),
#[emit::as_serde]
newtype_value: Newtype(42),
#[emit::as_serde]
unit_variant_value: Enum::Unit,
#[emit::as_serde]
unit_variant_value: Enum::Unit,
#[emit::as_serde]
newtype_variant_value: Enum::Newtype(42),
#[emit::as_serde]
struct_variant_value: Enum::Struct { a: 1, b: 1, c: 1 },
#[emit::as_serde]
tuple_variant_value: Enum::Tuple(1, 1, 1),
);
let logs = logs.get_emitted_logs().unwrap();
let get = |needle: &str| -> Option<AnyValue> {
logs[0].record.attributes_iter().find_map(|(k, v)| {
if k.as_str() == needle {
Some(v.clone())
} else {
None
}
})
};
assert_eq!(
AnyValue::String("a string".into()),
get("str_value").unwrap()
);
assert_eq!(AnyValue::Int(1), get("i8_value").unwrap());
assert_eq!(AnyValue::Int(2), get("i16_value").unwrap());
assert_eq!(AnyValue::Int(42), get("i32_value").unwrap());
assert_eq!(AnyValue::Int(2147483660), get("i64_value").unwrap());
assert_eq!(AnyValue::Int(2147483660), get("i128_small_value").unwrap());
assert_eq!(
AnyValue::String("9223372036854775820".into()),
get("i128_big_value").unwrap()
);
assert_eq!(AnyValue::Double(4.2), get("f64_value").unwrap());
assert_eq!(AnyValue::Boolean(true), get("bool_value").unwrap());
assert_eq!(None, get("unit_value"));
assert_eq!(None, get("none_value"));
assert_eq!(AnyValue::Int(42), get("some_value").unwrap());
assert_eq!(
AnyValue::ListAny(Box::new(vec![
AnyValue::Int(1),
AnyValue::Int(1),
AnyValue::Int(1)
])),
get("slice_value").unwrap()
);
assert_eq!(
AnyValue::Map({
let mut map = HashMap::<Key, AnyValue>::default();
map.insert(Key::from("a"), AnyValue::Int(1));
map.insert(Key::from("b"), AnyValue::Int(1));
map.insert(Key::from("c"), AnyValue::Int(1));
Box::new(map)
}),
get("map_value").unwrap()
);
assert_eq!(
AnyValue::Map({
let mut map = HashMap::<Key, AnyValue>::default();
map.insert(Key::from("a"), AnyValue::Int(1));
map.insert(Key::from("b"), AnyValue::Int(1));
map.insert(Key::from("c"), AnyValue::Int(1));
Box::new(map)
}),
get("struct_value").unwrap()
);
assert_eq!(
AnyValue::ListAny(Box::new(vec![
AnyValue::Int(1),
AnyValue::Int(1),
AnyValue::Int(1)
])),
get("tuple_value").unwrap()
);
assert_eq!(
AnyValue::String("Unit".into()),
get("unit_variant_value").unwrap()
);
assert_eq!(
AnyValue::Map({
let mut map = HashMap::new();
map.insert(Key::from("Newtype"), AnyValue::Int(42));
Box::new(map)
}),
get("newtype_variant_value").unwrap()
);
assert_eq!(
AnyValue::Map({
let mut map = HashMap::new();
map.insert(
Key::from("Struct"),
AnyValue::Map({
let mut map = HashMap::new();
map.insert(Key::from("a"), AnyValue::Int(1));
map.insert(Key::from("b"), AnyValue::Int(1));
map.insert(Key::from("c"), AnyValue::Int(1));
Box::new(map)
}),
);
Box::new(map)
}),
get("struct_variant_value").unwrap()
);
assert_eq!(
AnyValue::Map({
let mut map = HashMap::new();
map.insert(
Key::from("Tuple"),
AnyValue::ListAny(Box::new(vec![
AnyValue::Int(1),
AnyValue::Int(1),
AnyValue::Int(1),
])),
);
Box::new(map)
}),
get("tuple_variant_value").unwrap()
);
}
}