Skip to content

Commit 8a94de9

Browse files
fix(shared-runtime): prevent stale MicroVM runtime data
1 parent e1bf271 commit 8a94de9

12 files changed

Lines changed: 858 additions & 40 deletions

File tree

libdd-data-pipeline-ffi/src/trace_exporter.rs

Lines changed: 77 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -901,6 +901,42 @@ pub unsafe extern "C" fn ddog_trace_exporter_free(handle: Box<TraceExporter>) {
901901
let _ = catch_panic!(handle.shutdown(None), Ok(()));
902902
}
903903

904+
/// Discards buffered data without sending it, then frees the TraceExporter instance.
905+
///
906+
/// This is intended for runtime identity refreshes. Unlike
907+
/// [`ddog_trace_exporter_free`], it does not execute the normal flush-capable shutdown path.
908+
///
909+
/// Returns `None` when every worker has been discarded and sets `*handle` to null. On an ordinary
910+
/// error, the exporter remains owned by `*handle` so the caller can retry after the fork lifecycle
911+
/// has completed. A caught panic sets `*handle` to null.
912+
///
913+
/// # Arguments
914+
///
915+
/// * handle - A non-null pointer to the TraceExporter handle. The pointee must be non-null.
916+
#[no_mangle]
917+
pub unsafe extern "C" fn ddog_trace_exporter_free_without_flush(
918+
handle: &mut *mut TraceExporter,
919+
) -> Option<Box<ExporterError>> {
920+
catch_panic!(
921+
{
922+
let Some(exporter) = NonNull::new(*handle) else {
923+
return gen_error!(ErrorCode::InvalidArgument);
924+
};
925+
let mut exporter = Box::from_raw(exporter.as_ptr());
926+
*handle = std::ptr::null_mut();
927+
928+
match exporter.shutdown_without_flush() {
929+
Ok(()) => None,
930+
Err(err) => {
931+
*handle = Box::into_raw(exporter);
932+
Some(Box::new(ExporterError::from(err)))
933+
}
934+
}
935+
},
936+
gen_error!(ErrorCode::Panic)
937+
)
938+
}
939+
904940
/// Send traces to the Datadog Agent.
905941
///
906942
/// # Arguments
@@ -944,6 +980,7 @@ mod tests {
944980
use crate::error::ddog_trace_exporter_error_free;
945981
use httpmock::prelude::*;
946982
use httpmock::MockServer;
983+
use libdd_shared_runtime::SharedRuntime;
947984
use libdd_trace_utils::span::v04::SpanSlice;
948985
use std::{borrow::Borrow, mem::MaybeUninit};
949986

@@ -1351,6 +1388,46 @@ mod tests {
13511388
}
13521389
}
13531390

1391+
#[cfg_attr(miri, ignore)]
1392+
#[test]
1393+
fn exporter_free_without_flush_test() {
1394+
unsafe {
1395+
let mut config: MaybeUninit<Box<TraceExporterConfig>> = MaybeUninit::uninit();
1396+
ddog_trace_exporter_config_new(NonNull::new_unchecked(&mut config).cast());
1397+
let cfg = config.assume_init();
1398+
1399+
let mut exporter: MaybeUninit<Box<TraceExporter>> = MaybeUninit::uninit();
1400+
let error = ddog_trace_exporter_new(
1401+
NonNull::new_unchecked(&mut exporter).cast(),
1402+
Some(cfg.borrow()),
1403+
);
1404+
assert!(error.is_none());
1405+
1406+
let mut exporter = Box::into_raw(exporter.assume_init());
1407+
assert!(ddog_trace_exporter_free_without_flush(&mut exporter).is_none());
1408+
assert!(exporter.is_null());
1409+
ddog_trace_exporter_config_free(cfg);
1410+
}
1411+
}
1412+
1413+
#[test]
1414+
fn exporter_free_without_flush_preserves_handle_during_fork() {
1415+
unsafe {
1416+
let runtime = Arc::new(ForkSafeRuntime::new().unwrap());
1417+
let mut builder = TraceExporter::builder();
1418+
builder.set_shared_runtime(runtime.clone());
1419+
let mut exporter = Box::into_raw(Box::new(builder.build().unwrap()));
1420+
1421+
runtime.before_fork();
1422+
assert!(ddog_trace_exporter_free_without_flush(&mut exporter).is_some());
1423+
assert!(!exporter.is_null());
1424+
1425+
runtime.after_fork_parent().unwrap();
1426+
assert!(ddog_trace_exporter_free_without_flush(&mut exporter).is_none());
1427+
assert!(exporter.is_null());
1428+
}
1429+
}
1430+
13541431
#[test]
13551432
fn exporter_send_test_arguments_test() {
13561433
unsafe {

libdd-data-pipeline/src/otlp/metrics.rs

Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -348,6 +348,10 @@ impl<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static> Wor
348348
.flush(SystemTime::now(), true);
349349
}
350350

351+
fn discard(&mut self) {
352+
self.reset();
353+
}
354+
351355
async fn shutdown(&mut self) {
352356
// Single attempt: a long backoff could miss the bounded shutdown window.
353357
if let Err(e) = self.send(true, OTLP_SHUTDOWN_MAX_RETRIES).await {
@@ -487,6 +491,55 @@ mod tests {
487491
str_at(p["attributes"].as_array().unwrap(), "status.code") == Some(STATUS_CODE_ERROR)
488492
}
489493

494+
#[test]
495+
fn discard_clears_buffered_stats_without_export() {
496+
use crate::otlp::config::OtlpProtocol;
497+
use http::HeaderMap;
498+
use libdd_capabilities_impl::NativeCapabilities;
499+
use libdd_trace_utils::span::{trace_utils, v04::SpanSlice};
500+
use std::borrow::Cow;
501+
502+
let concentrator = Arc::new(Mutex::new(SpanConcentrator::new(
503+
Duration::from_secs(10),
504+
SystemTime::now(),
505+
vec![],
506+
vec![],
507+
None,
508+
vec![],
509+
#[cfg(feature = "stats-obfuscation")]
510+
None,
511+
)));
512+
let mut spans = vec![SpanSlice {
513+
service: Cow::Borrowed("svc"),
514+
duration: 1,
515+
..Default::default()
516+
}];
517+
trace_utils::compute_top_level_span(&mut spans);
518+
concentrator.lock_or_panic().add_span(&spans[0]);
519+
520+
let mut exporter = OtlpStatsExporter {
521+
flush_interval: Duration::from_secs(10),
522+
concentrator: concentrator.clone(),
523+
config: OtlpMetricsConfig {
524+
endpoint_url: "http://127.0.0.1:1/v1/metrics".to_string(),
525+
headers: HeaderMap::new(),
526+
timeout: Duration::from_millis(1),
527+
protocol: OtlpProtocol::HttpJson,
528+
otel_trace_semantics_enabled: false,
529+
},
530+
resource: OtlpResourceInfo::default(),
531+
test_token: None,
532+
capabilities: NativeCapabilities::new_client(),
533+
};
534+
535+
exporter.discard();
536+
537+
assert!(concentrator
538+
.lock_or_panic()
539+
.flush_with_otlp_exact(SystemTime::now(), true)
540+
.is_empty());
541+
}
542+
490543
#[test]
491544
fn metric_shape_and_resource_attributes() {
492545
assert!(map_stats_to_otlp_metrics(&[], &resource()).is_none());

libdd-data-pipeline/src/trace_buffer/mod.rs

Lines changed: 19 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -958,7 +958,7 @@ mod tests {
958958
use std::sync::Arc;
959959
use std::time::{Duration, Instant};
960960

961-
use libdd_shared_runtime::{BlockingRuntime, ForkSafeRuntime, SharedRuntime};
961+
use libdd_shared_runtime::{BlockingRuntime, ForkSafeRuntime, SharedRuntime, Worker};
962962

963963
use crate::trace_buffer::{BufferSize, Export, TraceBuffer, TraceBufferConfig};
964964
use crate::trace_exporter::agent_response::AgentResponse;
@@ -1550,4 +1550,22 @@ mod tests {
15501550
assert_eq!(sender.queue_metrics().get_metrics().spans_queued, 2);
15511551
rt.shutdown(None).unwrap();
15521552
}
1553+
1554+
#[test]
1555+
fn test_worker_discard_drops_buffered_chunk() {
1556+
let (sender, mut worker) = TraceBuffer::new(
1557+
TraceBufferConfig::default().flush_threshold_bytes(2),
1558+
Box::new(|_| {}),
1559+
Box::new(AssertExporter(
1560+
Box::new(|_| panic!("discard must not export buffered chunks")),
1561+
Arc::new(tokio::sync::Semaphore::new(0)),
1562+
)),
1563+
);
1564+
1565+
sender.send_chunk(vec![()]).unwrap();
1566+
assert_eq!(sender.queue_metrics().get_metrics().spans_queued, 1);
1567+
1568+
worker.discard();
1569+
assert_eq!(sender.queue_metrics().get_metrics().spans_queued, 0);
1570+
}
15531571
}

0 commit comments

Comments
 (0)