use std::{
pin::Pin,
sync::{mpsc, Arc, Condvar, Mutex},
thread,
time::Duration,
};
use arc_swap::ArcSwap;
use libdd_capabilities_impl::NativeCapabilities;
use libdd_data_pipeline::trace_buffer::{
BufferSize, Export, ResponseHandler, TraceBuffer, TraceBufferConfig, TraceBufferError,
TraceChunk,
};
use libdd_data_pipeline::trace_exporter::{
agent_response::AgentResponse, error::TraceExporterError, TelemetryConfig, TraceExporter,
TraceExporterBuilder, TraceExporterOutputFormat,
};
use libdd_shared_runtime::{BasicRuntime, BlockingRuntime, SharedRuntime, SharedRuntimeError};
use opentelemetry_sdk::{trace::SpanData, Resource};
use crate::{
configuration::Config, core::telemetry_session, ddtrace_transform, mappings::CachedConfig,
};
pub use libdd_data_pipeline::trace_buffer::TraceBufferError as DatadogExporterError;
pub type QueueMetricsFetcher = libdd_data_pipeline::trace_buffer::QueueMetricsFetcher<BufferedSpan>;
#[repr(transparent)]
#[derive(Debug)]
pub struct BufferedSpan(SpanData);
const FIXED_SPAN_OVERHEAD: usize = 96;
const ATTRIBUTE_ENTRY_OVERHEAD: usize = 24;
const EVENT_OVERHEAD: usize = 16;
const LINK_OVERHEAD: usize = 32;
fn attr_value_bytes(v: &opentelemetry::Value) -> usize {
use opentelemetry::{Array, Value};
match v {
Value::Bool(_) => 1,
Value::I64(_) | Value::F64(_) => 8,
Value::String(s) => s.as_str().len(),
Value::Array(Array::Bool(a)) => a.len(),
Value::Array(Array::I64(a)) => a.len() * 8,
Value::Array(Array::F64(a)) => a.len() * 8,
Value::Array(Array::String(a)) => a.iter().map(|s| s.as_str().len()).sum(),
_ => 0,
}
}
fn wrap_span_vec(spans: Vec<SpanData>) -> Vec<BufferedSpan> {
unsafe {
let mut spans = std::mem::ManuallyDrop::new(spans);
Vec::from_raw_parts(
spans.as_mut_ptr() as *mut BufferedSpan,
spans.len(),
spans.capacity(),
)
}
}
impl BufferSize for BufferedSpan {
fn byte_size(&self) -> usize {
let s = &self.0;
let mut size: usize = FIXED_SPAN_OVERHEAD;
size += s.name.len();
for kv in &s.attributes {
size += ATTRIBUTE_ENTRY_OVERHEAD + kv.key.as_str().len() + attr_value_bytes(&kv.value);
}
for event in s.events.iter() {
size += EVENT_OVERHEAD + event.name.len();
for kv in &event.attributes {
size +=
ATTRIBUTE_ENTRY_OVERHEAD + kv.key.as_str().len() + attr_value_bytes(&kv.value);
}
}
for link in s.links.iter() {
size += LINK_OVERHEAD;
for kv in &link.attributes {
size +=
ATTRIBUTE_ENTRY_OVERHEAD + kv.key.as_str().len() + attr_value_bytes(&kv.value);
}
}
size
}
}
#[derive(Debug, Default)]
struct PendingState {
pending: usize,
total_exported: u64,
}
type PendingSpans = Arc<(Mutex<PendingState>, Condvar)>;
pub struct DatadogExporter {
trace_buffer: Option<TraceBuffer<BufferedSpan>>,
shared_runtime: Option<Arc<BasicRuntime>>,
otel_resource: Arc<ArcSwap<Resource>>,
shutdown_rx: Mutex<Option<mpsc::Receiver<Result<(), SharedRuntimeError>>>>,
pending_spans: PendingSpans,
}
impl Drop for DatadogExporter {
fn drop(&mut self) {
let trace_buffer = self.trace_buffer.take();
let shared_runtime = self.shared_runtime.take();
let _ = thread::Builder::new()
.name("datadog-trace-drop".into())
.spawn(move || {
drop(trace_buffer);
drop(shared_runtime);
});
}
}
#[derive(Debug)]
pub enum DatadogExporterInitError {
Runtime(SharedRuntimeError),
TraceExporter(TraceExporterError),
BuildThread(std::io::Error),
}
impl std::fmt::Display for DatadogExporterInitError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Runtime(e) => write!(f, "shared runtime init failed: {e}"),
Self::TraceExporter(e) => write!(f, "trace exporter build failed: {e}"),
Self::BuildThread(e) => write!(f, "failed to spawn builder thread: {e}"),
}
}
}
impl std::error::Error for DatadogExporterInitError {}
type AgentResponseHandler = Box<dyn for<'a> Fn(&'a str) + Send + Sync>;
impl DatadogExporter {
pub fn new(
config: Arc<Config>,
agent_response_handler: Option<AgentResponseHandler>,
) -> Result<Self, DatadogExporterInitError> {
let response_handler = build_response_handler(agent_response_handler);
let (tx, rx) = mpsc::sync_channel(1);
thread::Builder::new()
.name("datadog-trace-init".into())
.spawn(move || {
let _ = tx.send(build_on_dedicated_thread(config, response_handler));
})
.map_err(DatadogExporterInitError::BuildThread)?;
rx.recv().map_err(|_| {
DatadogExporterInitError::Runtime(SharedRuntimeError::RuntimeUnavailable)
})?
}
fn trace_buffer(&self) -> &TraceBuffer<BufferedSpan> {
self.trace_buffer
.as_ref()
.expect("trace_buffer accessed after DatadogExporter::drop")
}
pub(crate) fn shared_runtime(&self) -> &Arc<BasicRuntime> {
self.shared_runtime
.as_ref()
.expect("shared_runtime accessed after DatadogExporter::drop")
}
pub fn queue_metrics(&self) -> QueueMetricsFetcher {
self.trace_buffer().queue_metrics()
}
pub fn send_chunk(&self, span_data: Vec<SpanData>) -> Result<(), TraceBufferError> {
if span_data.is_empty() {
return Ok(());
}
let n = span_data.len();
let buffered = wrap_span_vec(span_data);
increment_pending(&self.pending_spans, n);
match self.trace_buffer().send_chunk(buffered) {
Ok(()) => Ok(()),
Err(e) => {
if matches!(
e,
TraceBufferError::BatchFull(_) | TraceBufferError::AlreadyClosed
) {
cancel_pending(&self.pending_spans, n);
}
Err(e)
}
}
}
pub fn flush_and_drain(&self, timeout: Duration) -> Result<(), TraceBufferError> {
self.trace_buffer().force_flush()?;
self.wait_for_drain(timeout)
}
fn wait_for_drain(&self, timeout: Duration) -> Result<(), TraceBufferError> {
wait_for_barrier(&self.pending_spans, timeout)
}
pub fn trigger_shutdown(&self) {
let _ = self.trace_buffer().force_flush();
let mut slot = match self.shutdown_rx.lock() {
Ok(slot) => slot,
Err(_) => return,
};
if slot.is_some() {
return;
}
let (tx, rx) = mpsc::sync_channel(1);
let rt = Arc::clone(self.shared_runtime());
if thread::Builder::new()
.name("datadog-trace-shutdown".into())
.spawn(move || {
let _ = tx.send(shutdown_basic_runtime(&rt));
})
.is_ok()
{
*slot = Some(rx);
}
}
pub fn wait_for_shutdown(&self, timeout: Duration) -> Result<(), TraceBufferError> {
let rx = self
.shutdown_rx
.lock()
.map_err(|_| TraceBufferError::MutexPoisoned)?
.take()
.ok_or(TraceBufferError::AlreadyClosed)?;
let runtime_result = match rx.recv_timeout(timeout) {
Ok(res) => res,
Err(mpsc::RecvTimeoutError::Timeout) => return Err(TraceBufferError::TimedOut(timeout)),
Err(mpsc::RecvTimeoutError::Disconnected) => {
return Err(TraceBufferError::TraceExporter(TraceExporterError::Internal(
libdd_data_pipeline::trace_exporter::error::InternalErrorKind::InvalidWorkerState(
"shutdown thread terminated unexpectedly".to_string(),
),
)))
}
};
match runtime_result {
Ok(()) => Ok(()),
Err(SharedRuntimeError::ShutdownTimedOut(d)) => Err(TraceBufferError::TimedOut(d)),
Err(SharedRuntimeError::LockFailed(_)) => Err(TraceBufferError::MutexPoisoned),
Err(SharedRuntimeError::RuntimeUnavailable) => Err(TraceBufferError::AlreadyClosed),
Err(e) => Err(TraceBufferError::TraceExporter(TraceExporterError::Internal(
libdd_data_pipeline::trace_exporter::error::InternalErrorKind::InvalidWorkerState(
e.to_string(),
),
))),
}
}
pub fn set_resource(&self, r: Resource) {
self.otel_resource.store(Arc::new(r));
}
}
impl std::fmt::Debug for DatadogExporter {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("DatadogExporter").finish()
}
}
#[allow(clippy::type_complexity)]
fn build_response_handler(agent_response_handler: Option<AgentResponseHandler>) -> ResponseHandler {
Box::new(move |result| match result {
Ok(AgentResponse::Changed { body }) => {
if let Some(handler) = agent_response_handler.as_ref() {
handler(&body);
}
}
Ok(AgentResponse::Unchanged) => {}
Err(e) => log_trace_exporter_error(&e),
})
}
fn build_shared_runtime() -> Result<Arc<BasicRuntime>, DatadogExporterInitError> {
let tokio_runtime = tokio::runtime::Builder::new_multi_thread()
.worker_threads(2)
.enable_all()
.build()
.map_err(|e| DatadogExporterInitError::Runtime(SharedRuntimeError::RuntimeCreation(e)))?;
Ok(Arc::new(BasicRuntime::from_handle(Arc::new(tokio_runtime))))
}
fn build_on_dedicated_thread(
config: Arc<Config>,
response_handler: ResponseHandler,
) -> Result<DatadogExporter, DatadogExporterInitError> {
let shared_runtime = build_shared_runtime()?;
let trace_exporter = build_trace_exporter(&config, &shared_runtime)
.map_err(DatadogExporterInitError::TraceExporter)?;
let buffer_config = TraceBufferConfig::default()
.synchronous_export(config.trace_writer_synchronous_write())
.synchronous_export_timeout(Some(config.trace_writer_synchronous_timeout()))
.max_flush_interval(config.trace_writer_max_flush_interval());
let otel_resource = Arc::new(ArcSwap::new(Arc::new(Resource::builder_empty().build())));
let pending_spans: PendingSpans =
Arc::new((Mutex::new(PendingState::default()), Condvar::new()));
let export = SpanDataExport {
trace_exporter,
otel_resource: Arc::clone(&otel_resource),
cached_config: CachedConfig::new(&config),
config: Arc::clone(&config),
pending_spans: Arc::clone(&pending_spans),
};
let (trace_buffer, worker) =
TraceBuffer::new(buffer_config, response_handler, Box::new(export));
let _ = shared_runtime
.spawn_worker(worker, false)
.map_err(DatadogExporterInitError::Runtime)?;
Ok(DatadogExporter {
trace_buffer: Some(trace_buffer),
shared_runtime: Some(shared_runtime),
otel_resource,
shutdown_rx: Mutex::new(None),
pending_spans,
})
}
fn increment_pending(pending: &PendingSpans, n: usize) {
let (lock, _) = &**pending;
if let Ok(mut state) = lock.lock() {
state.pending = state.pending.saturating_add(n);
}
}
fn decrement_pending(pending: &PendingSpans, n: usize) {
let (lock, cvar) = &**pending;
if let Ok(mut state) = lock.lock() {
state.pending = state.pending.saturating_sub(n);
state.total_exported = state.total_exported.saturating_add(n as u64);
cvar.notify_all();
}
}
fn cancel_pending(pending: &PendingSpans, n: usize) {
let (lock, _) = &**pending;
if let Ok(mut state) = lock.lock() {
state.pending = state.pending.saturating_sub(n);
}
}
fn wait_for_barrier(pending: &PendingSpans, timeout: Duration) -> Result<(), TraceBufferError> {
let (lock, cvar) = &**pending;
let guard = lock.lock().map_err(|_| TraceBufferError::MutexPoisoned)?;
let barrier = guard.total_exported.saturating_add(guard.pending as u64);
if guard.total_exported >= barrier {
return Ok(());
}
if timeout.is_zero() {
return Err(TraceBufferError::TimedOut(Duration::ZERO));
}
let (_guard, res) = cvar
.wait_timeout_while(guard, timeout, |state| state.total_exported < barrier)
.map_err(|_| TraceBufferError::MutexPoisoned)?;
if res.timed_out() {
return Err(TraceBufferError::TimedOut(timeout));
}
Ok(())
}
fn build_trace_exporter(
config: &Config,
shared_runtime: &Arc<BasicRuntime>,
) -> Result<TraceExporter<NativeCapabilities, BasicRuntime>, TraceExporterError> {
let mut builder = TraceExporterBuilder::<BasicRuntime>::new();
builder
.set_shared_runtime(Arc::clone(shared_runtime))
.set_url(&config.trace_agent_url())
.set_dogstatsd_url(&config.dogstatsd_agent_url())
.set_tracer_version(config.tracer_version())
.set_language(config.language())
.set_language_version(config.language_version())
.set_service(&config.service())
.set_output_format(TraceExporterOutputFormat::V04)
.enable_health_metrics()
.enable_agent_rates_payload_version();
if config.trace_partial_flush_enabled() {
builder.set_client_computed_top_level();
}
if config.trace_stats_computation_enabled() {
builder.enable_stats(Duration::from_secs(10));
}
if config.trace_stats_computation_experimental_client_obfuscation_enabled() {
builder.enable_client_side_stats_obfuscation();
}
if let Some(env) = config.env() {
builder.set_env(env);
}
if let Some(version) = config.version() {
builder.set_app_version(version);
}
if config.telemetry_enabled() {
builder.enable_telemetry(TelemetryConfig {
heartbeat: (config.telemetry_heartbeat_interval() * 1000.0) as u64,
runtime_id: Some(config.runtime_id().to_string()),
debug_enabled: false,
});
builder.set_telemetry_instrumentation_sessions(
telemetry_session::sessions_from_runtime_id(config.runtime_id()),
);
}
builder.build::<NativeCapabilities>()
}
fn shutdown_basic_runtime(shared_runtime: &BasicRuntime) -> Result<(), SharedRuntimeError> {
shared_runtime
.block_on(shared_runtime.shutdown_async())
.map_err(SharedRuntimeError::RuntimeCreation)
}
#[derive(Debug)]
struct SpanDataExport {
trace_exporter: TraceExporter<NativeCapabilities, BasicRuntime>,
otel_resource: Arc<ArcSwap<Resource>>,
cached_config: CachedConfig,
config: Arc<Config>,
pending_spans: PendingSpans,
}
impl Export<BufferedSpan> for SpanDataExport {
fn export_trace_chunks(
&mut self,
trace_chunks: Vec<TraceChunk<BufferedSpan>>,
) -> Pin<
Box<
dyn std::future::Future<Output = Result<AgentResponse, TraceExporterError>> + Send + '_,
>,
> {
let total_spans: usize = trace_chunks.iter().map(|c| c.len()).sum();
Box::pin(async move {
let resource = self.otel_resource.load_full();
let dd_trace_chunks = trace_chunks
.iter()
.map(|chunk| {
ddtrace_transform::otel_trace_chunk_to_dd_trace_chunk(
&self.cached_config,
chunk.iter().map(|b| &b.0),
&resource,
)
})
.collect::<Vec<_>>();
let services = dd_trace_chunks
.iter()
.flatten()
.map(|s| s.service.as_str())
.filter(|s| !s.is_empty() && *s != "otlpresourcenoservicename");
self.config.add_extra_services(services);
let result = self
.trace_exporter
.send_trace_chunks_async(dd_trace_chunks)
.await;
decrement_pending(&self.pending_spans, total_spans);
result
})
}
#[cfg(feature = "test-utils")]
fn wait_ready(
&mut self,
) -> Pin<Box<dyn std::future::Future<Output = anyhow::Result<()>> + Send + '_>> {
Box::pin(async {
self.trace_exporter
.wait_agent_info_ready(Duration::from_secs(5))
.await
})
}
}
#[track_caller]
fn log_trace_exporter_error(e: &TraceExporterError) {
use libdd_data_pipeline::trace_exporter::error::{
AgentErrorKind, InternalErrorKind, ShutdownError,
};
use crate::{dd_debug, dd_error};
match e {
TraceExporterError::Builder(e) => {
dd_error!("DatadogExporter: Export error: Builder error: {}", e);
}
TraceExporterError::Internal(InternalErrorKind::InvalidWorkerState(state)) => {
dd_error!(
"DatadogExporter: Export error: Internal error: Invalid worker state: {}",
state
);
}
TraceExporterError::Deserialization(e) => {
dd_debug!(
"DatadogExporter: Export error: Deserialization error: {}",
e
);
}
TraceExporterError::Io(error) => {
dd_debug!("DatadogExporter: Export error: IO error: {}", error);
}
TraceExporterError::Network(e) => {
dd_debug!("DatadogExporter: Export error: Network error: {}", e);
}
TraceExporterError::Request(e) => {
dd_debug!("DatadogExporter: Export error: Request error: {}", e);
}
TraceExporterError::Serialization(error) => {
dd_debug!(
"DatadogExporter: Export error: Serialization error: {}",
error
);
}
TraceExporterError::Agent(AgentErrorKind::EmptyResponse) => {
dd_debug!("DatadogExporter: Export error: Agent error: empty response");
}
TraceExporterError::Shutdown(ShutdownError::TimedOut(duration)) => {
dd_debug!(
"DatadogExporter: Export error: Shutdown error: timed out after {}ms",
duration.as_millis()
);
}
TraceExporterError::Telemetry(e) => {
dd_debug!(
"DatadogExporter: Export error: Instrumentation telemetry error: {}",
e
);
}
};
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::{Arc, Mutex};
use std::time::Duration;
fn make_pending() -> PendingSpans {
Arc::new((Mutex::new(PendingState::default()), Condvar::new()))
}
#[test]
fn barrier_wait_returns_when_pre_flush_spans_exported() {
let pending = make_pending();
increment_pending(&pending, 5);
let pending2 = Arc::clone(&pending);
let waiter =
std::thread::spawn(move || wait_for_barrier(&pending2, Duration::from_secs(5)));
std::thread::sleep(Duration::from_millis(50));
increment_pending(&pending, 100);
decrement_pending(&pending, 5);
let res = waiter.join().unwrap();
assert!(
res.is_ok(),
"barrier wait should return once pre-flush spans are exported: {res:?}"
);
let (lock, _) = &*pending;
let state = lock.lock().unwrap();
assert_eq!(state.pending, 100);
assert_eq!(state.total_exported, 5);
}
#[test]
fn barrier_wait_times_out_when_spans_never_exported() {
let pending = make_pending();
increment_pending(&pending, 3);
let res = wait_for_barrier(&pending, Duration::from_millis(100));
assert!(matches!(res, Err(TraceBufferError::TimedOut(_))), "{res:?}");
}
#[test]
fn barrier_wait_no_op_when_nothing_pending() {
let pending = make_pending();
let res = wait_for_barrier(&pending, Duration::from_secs(5));
assert!(res.is_ok());
}
}