use crate::visitor::{event_to_values, span_to_values, OpenTelemetryVisitor};
use opentelemetry::api::trace::{self, span_context::SpanId, span_context::TraceId};
use opentelemetry::sdk::trace::config::Config;
use opentelemetry::sdk::trace::evicted_hash_map::EvictedHashMap;
use opentelemetry::sdk::EvictedQueue;
use std::collections::HashMap;
use std::sync::{Arc, Mutex};
use tracing_distributed::{Event, Span, Telemetry};
#[derive(Debug)]
pub struct OpenTelemetry {
pub(crate) exporter: Box<dyn opentelemetry::exporter::trace::SpanExporter>,
pub(crate) events: Mutex<HashMap<SpanId, EvictedQueue<trace::event::Event>>>,
pub(crate) config: Config,
}
impl Telemetry for OpenTelemetry {
type Visitor = OpenTelemetryVisitor;
type TraceId = TraceId;
type SpanId = SpanId;
fn mk_visitor(&self) -> Self::Visitor {
OpenTelemetryVisitor(EvictedHashMap::new(self.config.max_attributes_per_span))
}
fn report_span(&self, span: Span<Self::Visitor, Self::SpanId, Self::TraceId>) {
let mut events = self.events.lock().unwrap();
let events = events
.remove(&span.id)
.unwrap_or_else(|| EvictedQueue::new(0));
let data = span_to_values(span, events);
self.exporter.export(vec![Arc::new(data)]); }
fn report_event(&self, event: Event<Self::Visitor, Self::SpanId, Self::TraceId>) {
match event.parent_id {
Some(id) => {
let mut events = self.events.lock().unwrap();
if let Some(q) = events.get_mut(&id) {
q.append_vec(&mut vec![event_to_values(event)]);
} else {
let mut q = EvictedQueue::new(self.config.max_events_per_span);
q.append_vec(&mut vec![event_to_values(event)]);
events.insert(id, q);
}
}
None => {
}
}
}
}