pub mod http_client;
pub mod metric_ring;
mod scheduler;
pub mod store;
use crate::{
config::Config,
data::{
self, Application, Dependency, Endpoint, Host, Integration, Log, Payload, ProductState,
Telemetry,
},
metrics::{ContextKey, MetricBuckets, MetricContexts},
};
use crate::worker::metric_ring::MetricRing;
use async_trait::async_trait;
use bytes::Bytes;
use libdd_capabilities::{HttpClientCapability, HttpError, MaybeSend, SleepCapability};
use libdd_common::tag::Tag;
use libdd_shared_runtime::Worker;
use std::iter::Sum;
use std::marker::PhantomData;
use std::ops::Add;
use std::{
collections::hash_map::DefaultHasher,
hash::{Hash, Hasher},
ops::ControlFlow,
sync::{
atomic::{AtomicU64, Ordering},
Arc,
},
};
use std::{collections::HashSet, fmt::Debug, time::Duration};
use web_time as time;
#[cfg(not(target_arch = "wasm32"))]
use std::sync::{Condvar, Mutex};
use crate::metrics::MetricBucketStats;
use futures::{
channel::oneshot,
future::{self},
};
use http::{header, HeaderValue};
use serde::{Deserialize, Serialize};
use tokio::sync::mpsc;
#[cfg(not(target_arch = "wasm32"))]
use tokio::{runtime, task::JoinHandle};
use tokio_util::sync::CancellationToken;
use tracing::debug;
const CONTINUE: ControlFlow<()> = ControlFlow::Continue(());
const BREAK: ControlFlow<()> = ControlFlow::Break(());
fn time_now() -> f64 {
time::SystemTime::UNIX_EPOCH
.elapsed()
.unwrap_or_default()
.as_secs_f64()
}
macro_rules! telemetry_worker_log {
($worker:expr , ERROR , $fmt_str:tt, $($arg:tt)*) => {
{
debug!(
worker.runtime_id = %$worker.runtime_id,
worker.debug_logging = $worker.config.telemetry_debug_logging_enabled,
$fmt_str,
$($arg)*
);
if $worker.config.telemetry_debug_logging_enabled {
eprintln!(concat!("{}: Telemetry worker ERROR: ", $fmt_str), time_now(), $($arg)*);
}
}
};
($worker:expr , DEBUG , $fmt_str:tt, $($arg:tt)*) => {
{
debug!(
worker.runtime_id = %$worker.runtime_id,
worker.debug_logging = $worker.config.telemetry_debug_logging_enabled,
$fmt_str,
$($arg)*
);
if $worker.config.telemetry_debug_logging_enabled {
eprintln!(concat!("{}: Telemetry worker DEBUG: ", $fmt_str), time_now(), $($arg)*);
}
}
};
}
#[derive(Debug, Serialize, Deserialize)]
pub enum TelemetryActions {
AddPoint((f64, ContextKey, Vec<Tag>)),
AddConfig(data::Configuration),
AddDependency(Dependency),
AddIntegration(Integration),
AddProductChange((String, ProductState)),
AddLog((LogIdentifier, Log)),
AddEndpoint(Endpoint),
Lifecycle(LifecycleAction),
#[serde(skip)]
CollectStats(oneshot::Sender<TelemetryWorkerStats>),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub enum LifecycleAction {
Start,
Stop,
FlushMetricAggr,
FlushData,
ExtendedHeartbeat,
}
#[derive(Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct LogIdentifier {
pub identifier: u64,
}
#[derive(Debug)]
struct TelemetryWorkerData {
started: bool,
dependencies: store::Store<data::Dependency, data::DependencyKey>,
configurations: store::Store<data::Configuration>,
integrations: store::Store<data::Integration>,
endpoints: store::Store<data::Endpoint>,
endpoints_is_first: bool,
products: std::collections::HashMap<String, ProductState>,
products_pending: HashSet<String>,
logs: store::QueueHashMap<LogIdentifier, Log>,
metric_contexts: MetricContexts,
metric_buckets: MetricBuckets,
host: Host,
app: Application,
install_signature: Option<data::InstallSignature>,
}
pub struct TelemetryWorker<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static> {
flavor: TelemetryWorkerFlavor,
config: Config,
mailbox: mpsc::Receiver<TelemetryActions>,
cancellation_token: CancellationToken,
seq_id: AtomicU64,
runtime_id: String,
capabilities: C,
metrics_flush_interval: Duration,
deadlines: scheduler::Scheduler<LifecycleAction>,
data: TelemetryWorkerData,
next_action: Option<TelemetryActions>,
stopped: bool,
metric_ring: Arc<MetricRing>,
}
impl<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static> Debug
for TelemetryWorker<C>
{
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TelemetryWorker")
.field("flavor", &self.flavor)
.field("config", &self.config)
.field("mailbox", &self.mailbox)
.field("cancellation_token", &self.cancellation_token)
.field("seq_id", &self.seq_id)
.field("runtime_id", &self.runtime_id)
.field("metrics_flush_interval", &self.metrics_flush_interval)
.field("deadlines", &self.deadlines)
.field("data", &self.data)
.finish()
}
}
#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
impl<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static> Worker
for TelemetryWorker<C>
{
async fn trigger(&mut self) {
if self.next_action.is_some() {
return;
}
if self.stopped {
debug!(
worker.runtime_id = %self.runtime_id,
"Telemetry worker mailbox closed; parking until shutdown"
);
std::future::pending::<()>().await;
}
let action = self.recv_next_action().await;
self.next_action = Some(action);
}
async fn run(&mut self) {
if let Some(action) = self.next_action.take() {
debug!(
worker.runtime_id = %self.runtime_id,
action = ?action,
"Received telemetry action"
);
let _action_result = match self.flavor {
TelemetryWorkerFlavor::Full => self.dispatch_action(action).await,
TelemetryWorkerFlavor::MetricsLogs => {
self.dispatch_metrics_logs_action(action).await
}
};
}
}
fn reset(&mut self) {
while self.mailbox.try_recv().is_ok() {}
self.next_action = None;
self.data.logs = store::QueueHashMap::default();
self.data.metric_buckets = MetricBuckets::default();
self.metric_ring.drain(|_, _, _| {});
self.data.dependencies.clear();
self.data.integrations.clear();
self.data.configurations.clear();
self.data.endpoints.clear();
self.data.endpoints_is_first = true;
self.data.products.clear();
self.data.products_pending.clear();
}
async fn shutdown(&mut self) {
for _ in 0..self.mailbox.len() {
if let Ok(action) = self.mailbox.try_recv() {
let _ = match self.flavor {
TelemetryWorkerFlavor::Full => self.dispatch_action(action).await,
TelemetryWorkerFlavor::MetricsLogs => {
self.dispatch_metrics_logs_action(action).await
}
};
}
}
let stop_action = TelemetryActions::Lifecycle(LifecycleAction::Stop);
let _action_result = match self.flavor {
TelemetryWorkerFlavor::Full => self.dispatch_action(stop_action).await,
TelemetryWorkerFlavor::MetricsLogs => {
self.dispatch_metrics_logs_action(stop_action).await
}
};
}
}
#[derive(Debug, Default, Serialize, Deserialize)]
pub struct TelemetryWorkerStats {
pub dependencies_stored: u32,
pub dependencies_unflushed: u32,
pub configurations_stored: u32,
pub configurations_unflushed: u32,
pub integrations_stored: u32,
pub integrations_unflushed: u32,
pub logs: u32,
pub metric_contexts: u32,
pub metric_buckets: MetricBucketStats,
}
impl Add for TelemetryWorkerStats {
type Output = Self;
fn add(self, rhs: Self) -> Self::Output {
TelemetryWorkerStats {
dependencies_stored: self.dependencies_stored + rhs.dependencies_stored,
dependencies_unflushed: self.dependencies_unflushed + rhs.dependencies_unflushed,
configurations_stored: self.configurations_stored + rhs.configurations_stored,
configurations_unflushed: self.configurations_unflushed + rhs.configurations_unflushed,
integrations_stored: self.integrations_stored + rhs.integrations_stored,
integrations_unflushed: self.integrations_unflushed + rhs.integrations_unflushed,
logs: self.logs + rhs.logs,
metric_contexts: self.metric_contexts + rhs.metric_contexts,
metric_buckets: MetricBucketStats {
buckets: self.metric_buckets.buckets + rhs.metric_buckets.buckets,
series: self.metric_buckets.series + rhs.metric_buckets.series,
series_points: self.metric_buckets.series_points + rhs.metric_buckets.series_points,
distributions: self.metric_buckets.distributions + rhs.metric_buckets.distributions,
distributions_points: self.metric_buckets.distributions_points
+ rhs.metric_buckets.distributions_points,
},
}
}
}
impl Sum for TelemetryWorkerStats {
fn sum<I: Iterator<Item = Self>>(iter: I) -> Self {
iter.fold(Self::default(), |a, b| a + b)
}
}
mod serialize {
use crate::data;
use http::HeaderValue;
#[allow(clippy::declare_interior_mutable_const)]
pub const CONTENT_TYPE_VALUE: HeaderValue = libdd_common::header::APPLICATION_JSON;
pub fn serialize(telemetry: &data::Telemetry) -> anyhow::Result<Vec<u8>> {
Ok(serde_json::to_vec(telemetry)?)
}
}
impl<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static> TelemetryWorker<C> {
fn log_err(&self, err: &anyhow::Error) {
telemetry_worker_log!(self, ERROR, "{}", err);
}
fn drain_metric_ring(&mut self) {
let ring = self.metric_ring.clone();
let buckets = &mut self.data.metric_buckets;
ring.drain(|value, key, extra_tags| buckets.add_point(key, value, extra_tags));
}
fn flush_metric_aggregates(&mut self) {
self.drain_metric_ring();
self.data.metric_buckets.flush_aggregates();
}
async fn recv_next_action(&mut self) -> TelemetryActions {
loop {
self.drain_metric_ring();
let action = if let Some((deadline, deadline_action)) = self.deadlines.next_deadline() {
let deadline_action = *deadline_action;
let Some(remaining) = deadline.checked_duration_since(time::Instant::now()) else {
if let Ok(mailbox_action) = self.mailbox.try_recv() {
return mailbox_action;
}
return TelemetryActions::Lifecycle(deadline_action);
};
let sleeper = <C as SleepCapability>::new();
let ring = self.metric_ring.clone();
tokio::select! {
biased;
mailbox_action = self.mailbox.recv() => mailbox_action,
_ = sleeper.sleep(remaining) => Some(TelemetryActions::Lifecycle(deadline_action)),
_ = ring.notified() => continue,
}
} else {
let ring = self.metric_ring.clone();
tokio::select! {
biased;
mailbox_action = self.mailbox.recv() => mailbox_action,
_ = ring.notified() => continue,
}
};
return action.unwrap_or_else(|| {
self.config.restartable = false;
self.stopped = true;
TelemetryActions::Lifecycle(LifecycleAction::Stop)
});
}
}
async fn dispatch_metrics_logs_action(&mut self, action: TelemetryActions) -> ControlFlow<()> {
telemetry_worker_log!(self, DEBUG, "Handling metric action {:?}", action);
use LifecycleAction::*;
use TelemetryActions::*;
match action {
Lifecycle(Start) => {
if !self.data.started {
#[allow(clippy::unwrap_used)]
self.deadlines
.schedule_event(LifecycleAction::FlushMetricAggr)
.unwrap();
#[allow(clippy::unwrap_used)]
self.deadlines
.schedule_event(LifecycleAction::FlushData)
.unwrap();
self.data.started = true;
}
}
AddLog((identifier, log)) => {
let (l, new) = self.data.logs.get_mut_or_insert(identifier, log);
if !new {
l.count += 1;
}
}
AddPoint((point, key, extra_tags)) => {
self.data.metric_buckets.add_point(key, point, extra_tags)
}
Lifecycle(FlushMetricAggr) => {
self.flush_metric_aggregates();
#[allow(clippy::unwrap_used)]
self.deadlines
.schedule_event(LifecycleAction::FlushMetricAggr)
.unwrap();
}
Lifecycle(FlushData) => {
if !(self.data.started || self.config.restartable) {
return CONTINUE;
}
#[allow(clippy::unwrap_used)]
self.deadlines
.schedule_event(LifecycleAction::FlushData)
.unwrap();
let batch = self.build_observability_batch();
if !batch.is_empty() {
let payload = data::Payload::MessageBatch(batch);
match self.send_payload(&payload).await {
Ok(()) => self.payload_sent_success(&payload),
Err(e) => self.log_err(&e),
}
}
}
AddConfig(_)
| AddDependency(_)
| AddIntegration(_)
| AddProductChange(_)
| AddEndpoint(_)
| Lifecycle(ExtendedHeartbeat) => {}
Lifecycle(Stop) => {
if !self.data.started {
return BREAK;
}
self.flush_metric_aggregates();
let batch = self.build_observability_batch();
if !batch.is_empty() {
let payload = data::Payload::MessageBatch(batch);
match self.send_payload(&payload).await {
Ok(()) => {
if self.config.restartable {
self.payload_sent_success(&payload)
}
}
Err(e) => self.log_err(&e),
}
}
self.data.started = false;
if !self.config.restartable {
self.deadlines.clear_pending();
}
return BREAK;
}
CollectStats(stats_sender) => {
stats_sender.send(self.stats()).ok();
}
};
CONTINUE
}
async fn dispatch_action(&mut self, action: TelemetryActions) -> ControlFlow<()> {
telemetry_worker_log!(self, DEBUG, "Handling action {:?}", action);
use LifecycleAction::*;
use TelemetryActions::*;
match action {
Lifecycle(Start) => {
if !self.data.started {
if self.config.emit_app_lifecycle {
let app_started = data::Payload::AppStarted(self.build_app_started());
match self.send_payload(&app_started).await {
Ok(()) => self.payload_sent_success(&app_started),
Err(err) => self.log_err(&err),
}
}
#[allow(clippy::unwrap_used)]
self.deadlines
.schedule_event(LifecycleAction::FlushMetricAggr)
.unwrap();
#[allow(clippy::unwrap_used)]
self.deadlines
.schedule_event(LifecycleAction::FlushData)
.unwrap();
#[allow(clippy::unwrap_used)]
self.deadlines
.schedule_event(LifecycleAction::ExtendedHeartbeat)
.unwrap();
self.data.started = true;
}
}
AddDependency(dep) => self.data.dependencies.insert(dep),
AddIntegration(integration) => self.data.integrations.insert(integration),
AddProductChange((name, state)) => {
self.data.products.insert(name.clone(), state);
self.data.products_pending.insert(name);
}
AddConfig(cfg) => self.data.configurations.insert(cfg),
AddEndpoint(endpoint) => {
self.data.endpoints.insert(endpoint);
}
AddLog((identifier, log)) => {
let (l, new) = self.data.logs.get_mut_or_insert(identifier, log);
if !new {
l.count += 1;
}
}
AddPoint((point, key, extra_tags)) => {
self.data.metric_buckets.add_point(key, point, extra_tags)
}
Lifecycle(FlushMetricAggr) => {
self.flush_metric_aggregates();
#[allow(clippy::unwrap_used)]
self.deadlines
.schedule_event(LifecycleAction::FlushMetricAggr)
.unwrap();
}
Lifecycle(FlushData) => {
if !(self.data.started || self.config.restartable) {
return CONTINUE;
}
#[allow(clippy::unwrap_used)]
self.deadlines
.schedule_event(LifecycleAction::FlushData)
.unwrap();
let mut batch = self.build_app_events_batch();
let payload = if batch.is_empty() {
data::Payload::AppHeartbeat(())
} else {
batch.push(data::Payload::AppHeartbeat(()));
data::Payload::MessageBatch(batch)
};
match self.send_payload(&payload).await {
Ok(()) => self.payload_sent_success(&payload),
Err(err) => self.log_err(&err),
}
let batch = self.build_observability_batch();
if !batch.is_empty() {
let payload = data::Payload::MessageBatch(batch);
match self.send_payload(&payload).await {
Ok(()) => self.payload_sent_success(&payload),
Err(err) => self.log_err(&err),
}
}
}
Lifecycle(ExtendedHeartbeat) => {
let delta = self.build_app_events_batch();
if !delta.is_empty() {
let payload = data::Payload::MessageBatch(delta);
match self.send_payload(&payload).await {
Ok(()) => self.payload_sent_success(&payload),
Err(err) => self.log_err(&err),
}
}
self.data.dependencies.unflush_stored();
self.data.integrations.unflush_stored();
self.data.configurations.unflush_stored();
let extended_hb =
data::Payload::AppExtendedHeartbeat(self.build_extended_heartbeat());
match self.send_payload(&extended_hb).await {
Ok(()) => self.payload_sent_success(&extended_hb),
Err(err) => self.log_err(&err),
}
if !self.data.products.is_empty() {
let products = self
.data
.products
.iter()
.map(|(name, state)| (name.clone(), state.clone()))
.collect();
let product_change =
data::Payload::AppProductChange(data::AppProductChange { products });
match self.send_payload(&product_change).await {
Ok(()) => self.payload_sent_success(&product_change),
Err(err) => self.log_err(&err),
}
}
#[allow(clippy::unwrap_used)]
self.deadlines
.schedule_event(LifecycleAction::ExtendedHeartbeat)
.unwrap();
}
Lifecycle(Stop) => {
if !self.data.started {
return BREAK;
}
self.flush_metric_aggregates();
let mut app_events = self.build_app_events_batch();
if self.config.emit_app_lifecycle {
app_events.push(data::Payload::AppClosing(()));
}
let observability_events = self.build_observability_batch();
let mut payloads = vec![data::Payload::MessageBatch(app_events)];
if !observability_events.is_empty() {
payloads.push(data::Payload::MessageBatch(observability_events));
}
let self_arc = Arc::new(tokio::sync::RwLock::new(&mut *self));
let futures = payloads.into_iter().map(|payload| {
let self_arc = self_arc.clone();
async move {
let res = {
let self_rguard = self_arc.read().await;
self_rguard.send_payload(&payload).await
};
match res {
Ok(()) => self_arc.write().await.payload_sent_success(&payload),
Err(err) => self_arc.read().await.log_err(&err),
}
}
});
future::join_all(futures).await;
self.data.started = false;
if !self.config.restartable {
self.deadlines.clear_pending();
}
return BREAK;
}
CollectStats(stats_sender) => {
stats_sender.send(self.stats()).ok();
}
}
CONTINUE
}
fn build_app_events_batch(&mut self) -> Vec<Payload> {
let mut payloads = Vec::new();
if self.data.dependencies.flush_not_empty() {
payloads.push(data::Payload::AppDependenciesLoaded(
data::AppDependenciesLoaded {
dependencies: self.data.dependencies.unflushed().cloned().collect(),
},
))
}
if self.data.integrations.flush_not_empty() {
payloads.push(data::Payload::AppIntegrationsChange(
data::AppIntegrationsChange {
integrations: self.data.integrations.unflushed().cloned().collect(),
},
))
}
if !self.data.products_pending.is_empty() {
let products = self
.data
.products_pending
.iter()
.filter_map(|name| {
self.data
.products
.get(name)
.map(|state| (name.clone(), state.clone()))
})
.collect();
payloads.push(data::Payload::AppProductChange(data::AppProductChange {
products,
}))
}
if self.data.configurations.flush_not_empty() {
payloads.push(data::Payload::AppClientConfigurationChange(
data::AppClientConfigurationChange {
configuration: self.data.configurations.unflushed().cloned().collect(),
},
))
}
if self.data.endpoints.flush_not_empty() {
payloads.push(data::Payload::AppEndpoints(data::AppEndpoints {
is_first: self.data.endpoints_is_first,
endpoints: self
.data
.endpoints
.unflushed()
.take(self.config.endpoints_message_limit as usize)
.map(|e| e.to_json_value().unwrap_or_default())
.filter(|e| e.is_object())
.collect(),
}));
}
payloads
}
fn build_observability_batch(&mut self) -> Vec<Payload> {
let mut payloads = Vec::new();
let logs = self.build_logs();
if !logs.logs.is_empty() {
payloads.push(data::Payload::Logs(logs));
}
let metrics = self.build_metrics_series();
if !metrics.series.is_empty() {
payloads.push(data::Payload::GenerateMetrics(metrics))
}
let distributions = self.build_metrics_distributions();
if !distributions.series.is_empty() {
payloads.push(data::Payload::Sketches(distributions))
}
payloads
}
fn build_metrics_distributions(&mut self) -> data::Distributions {
let mut series = Vec::new();
let context_guard = self.data.metric_contexts.lock();
for (context_key, extra_tags, points) in self.data.metric_buckets.flush_distributions() {
let Some(context) = context_guard.read(context_key) else {
telemetry_worker_log!(self, ERROR, "Context not found for key {:?}", context_key);
continue;
};
let mut tags = extra_tags;
tags.extend(context.tags.iter().cloned());
series.push(data::metrics::Distribution {
namespace: context.namespace,
metric: context.name.clone(),
tags,
sketch: data::metrics::SerializedSketch::B64 {
sketch_b64: base64::Engine::encode(
&base64::engine::general_purpose::STANDARD,
points.encode_to_vec(),
),
},
common: context.common,
_type: context.metric_type,
interval: self.metrics_flush_interval.as_secs(),
});
}
data::Distributions { series }
}
fn build_metrics_series(&mut self) -> data::GenerateMetrics {
let mut series = Vec::new();
let context_guard = self.data.metric_contexts.lock();
for (context_key, extra_tags, points) in self.data.metric_buckets.flush_series() {
let Some(context) = context_guard.read(context_key) else {
telemetry_worker_log!(self, ERROR, "Context not found for key {:?}", context_key);
continue;
};
let mut tags = extra_tags;
tags.extend(context.tags.iter().cloned());
series.push(data::metrics::Serie {
namespace: context.namespace,
metric: context.name.clone(),
tags,
points,
common: context.common,
_type: context.metric_type,
interval: self.metrics_flush_interval.as_secs(),
});
}
data::GenerateMetrics { series }
}
fn build_app_started(&mut self) -> data::AppStarted {
data::AppStarted {
configuration: self.data.configurations.unflushed().cloned().collect(),
dependencies: Vec::new(),
integrations: Vec::new(),
install_signature: self.data.install_signature.clone(),
products: self.data.products.clone(),
error: None,
}
}
fn build_extended_heartbeat(&mut self) -> data::AppStarted {
data::AppStarted {
configuration: self.data.configurations.unflushed().cloned().collect(),
dependencies: self.data.dependencies.unflushed().cloned().collect(),
integrations: self.data.integrations.unflushed().cloned().collect(),
install_signature: self.data.install_signature.clone(),
products: self.data.products.clone(),
error: None,
}
}
fn app_started_sent_success(&mut self, p: &data::AppStarted) {
self.data
.configurations
.removed_flushed(p.configuration.len());
self.data.dependencies.removed_flushed(p.dependencies.len());
self.data.integrations.removed_flushed(p.integrations.len());
self.data.products_pending.clear();
}
fn payload_sent_success(&mut self, payload: &data::Payload) {
use data::Payload::*;
match payload {
AppStarted(p) => self.app_started_sent_success(p),
AppExtendedHeartbeat(p) => self.app_started_sent_success(p),
AppDependenciesLoaded(p) => {
self.data.dependencies.removed_flushed(p.dependencies.len())
}
AppIntegrationsChange(p) => {
self.data.integrations.removed_flushed(p.integrations.len())
}
AppProductChange(p) => {
for name in p.products.keys() {
self.data.products_pending.remove(name);
}
}
AppClientConfigurationChange(p) => self
.data
.configurations
.removed_flushed(p.configuration.len()),
AppEndpoints(p) => {
self.data.endpoints.removed_flushed(p.endpoints.len());
self.data.endpoints_is_first = false;
}
MessageBatch(batch) => {
for p in batch {
self.payload_sent_success(p);
}
}
Logs(p) => {
for _ in &p.logs {
self.data.logs.pop_front();
}
}
AppHeartbeat(()) | AppClosing(()) => {}
GenerateMetrics(_) | Sketches(_) => {}
}
}
fn build_logs(&self) -> data::Logs {
let logs = self.data.logs.iter().map(|(_, l)| l.clone()).collect();
data::Logs { logs }
}
fn next_seq_id(&self) -> u64 {
self.seq_id.fetch_add(1, Ordering::Release)
}
async fn send_payload(&self, payload: &data::Payload) -> anyhow::Result<()> {
debug!(
worker.runtime_id = %self.runtime_id,
payload.type = payload.request_type(),
seq_id = self.seq_id.load(Ordering::Acquire),
"Sending telemetry payload"
);
let req = self.build_request(payload)?;
let result = self.send_request(req).await;
match &result {
Ok(resp) => debug!(
worker.runtime_id = %self.runtime_id,
payload.type = payload.request_type(),
response.status = resp.status().as_u16(),
"Successfully sent telemetry payload"
),
Err(e) => debug!(
worker.runtime_id = %self.runtime_id,
payload.type = payload.request_type(),
error = ?e,
"Failed to send telemetry payload"
),
}
Ok(())
}
fn build_request(&self, payload: &data::Payload) -> anyhow::Result<http::Request<Bytes>> {
let seq_id = self.next_seq_id();
let tel = Telemetry {
api_version: data::ApiVersion::V2,
tracer_time: time::SystemTime::UNIX_EPOCH
.elapsed()
.map_or(0, |d| d.as_secs()),
runtime_id: &self.runtime_id,
seq_id,
host: &self.data.host,
origin: None,
application: &self.data.app,
payload,
};
telemetry_worker_log!(self, DEBUG, "Prepared payload: {:?}", tel);
let req = http_client::request_builder(&self.config)?
.method(http::Method::POST)
.header(header::CONTENT_TYPE, serialize::CONTENT_TYPE_VALUE)
.header(
http_client::header::REQUEST_TYPE,
HeaderValue::from_static(payload.request_type()),
)
.header(
http_client::header::API_VERSION,
HeaderValue::from_static(data::ApiVersion::V2.to_str()),
)
.header(
http_client::header::LIBRARY_LANGUAGE,
tel.application.language_name.clone(),
)
.header(
http_client::header::LIBRARY_VERSION,
tel.application.tracer_version.clone(),
);
let req = http_client::add_instrumentation_session_headers(
req,
self.config.session_id.as_deref(),
self.config.parent_session_id.as_deref(),
self.config.root_session_id.as_deref(),
);
let body = Bytes::from(serialize::serialize(&tel)?);
Ok(req.body(body)?)
}
async fn send_request(
&self,
req: http::Request<Bytes>,
) -> Result<http::Response<Bytes>, HttpError> {
let timeout_ms = if let Some(endpoint) = self.config.endpoint.as_ref() {
endpoint.timeout_ms
} else {
libdd_common::Endpoint::DEFAULT_TIMEOUT
};
let timeout = time::Duration::from_millis(timeout_ms);
debug!(
worker.runtime_id = %self.runtime_id,
http.timeout_ms = timeout_ms,
"Sending HTTP request"
);
let sleeper = <C as SleepCapability>::new();
tokio::select! {
_ = self.cancellation_token.cancelled() => {
debug!(
worker.runtime_id = %self.runtime_id,
"Telemetry request cancelled"
);
Err(HttpError::Other(anyhow::anyhow!("Request cancelled")))
},
_ = sleeper.sleep(timeout) => {
debug!(
worker.runtime_id = %self.runtime_id,
http.timeout_ms = timeout_ms,
"Telemetry request timed out"
);
Err(HttpError::Other(anyhow::anyhow!("Request timed out")))
},
r = self.capabilities.request(req) => r,
}
}
fn stats(&self) -> TelemetryWorkerStats {
TelemetryWorkerStats {
dependencies_stored: self.data.dependencies.len_stored() as u32,
dependencies_unflushed: self.data.dependencies.len_unflushed() as u32,
configurations_stored: self.data.configurations.len_stored() as u32,
configurations_unflushed: self.data.configurations.len_unflushed() as u32,
integrations_stored: self.data.integrations.len_stored() as u32,
integrations_unflushed: self.data.integrations.len_unflushed() as u32,
logs: self.data.logs.len() as u32,
metric_contexts: self.data.metric_contexts.lock().len() as u32,
metric_buckets: self.data.metric_buckets.stats(),
}
}
async fn run_loop(mut self) {
debug!(
worker.flavor = ?self.flavor,
worker.runtime_id = %self.runtime_id,
"Starting telemetry worker"
);
loop {
if self.cancellation_token.is_cancelled() {
debug!(
worker.runtime_id = %self.runtime_id,
"Telemetry worker cancelled, shutting down"
);
return;
}
let action = self.recv_next_action().await;
debug!(
worker.runtime_id = %self.runtime_id,
action = ?action,
"Received telemetry action"
);
let action_result = match self.flavor {
TelemetryWorkerFlavor::Full => self.dispatch_action(action).await,
TelemetryWorkerFlavor::MetricsLogs => {
self.dispatch_metrics_logs_action(action).await
}
};
match action_result {
ControlFlow::Continue(()) => {}
ControlFlow::Break(()) => {
debug!(
worker.runtime_id = %self.runtime_id,
worker.restartable = self.config.restartable,
"Telemetry worker received break signal"
);
if !self.config.restartable {
break;
}
}
};
}
debug!(
worker.runtime_id = %self.runtime_id,
"Telemetry worker stopped"
);
}
}
#[cfg(not(target_arch = "wasm32"))]
#[derive(Debug)]
struct InnerTelemetryShutdown {
is_shutdown: Mutex<bool>,
condvar: Condvar,
}
#[cfg(not(target_arch = "wasm32"))]
impl InnerTelemetryShutdown {
fn wait_for_shutdown(&self) {
drop(
#[allow(clippy::unwrap_used)]
self.condvar
.wait_while(self.is_shutdown.lock().unwrap(), |is_shutdown| {
!*is_shutdown
})
.unwrap(),
)
}
#[allow(clippy::unwrap_used)]
fn shutdown_finished(&self) {
*self.is_shutdown.lock().unwrap() = true;
self.condvar.notify_all();
}
}
pub struct TelemetryWorkerHandle<
C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static,
> {
sender: mpsc::Sender<TelemetryActions>,
#[cfg(not(target_arch = "wasm32"))]
shutdown: Arc<InnerTelemetryShutdown>,
cancellation_token: CancellationToken,
#[cfg(not(target_arch = "wasm32"))]
runtime: Option<runtime::Handle>,
contexts: MetricContexts,
metric_ring: Arc<MetricRing>,
_phantom: PhantomData<fn() -> C>,
}
impl<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static> Clone
for TelemetryWorkerHandle<C>
{
fn clone(&self) -> Self {
Self {
sender: self.sender.clone(),
#[cfg(not(target_arch = "wasm32"))]
shutdown: self.shutdown.clone(),
cancellation_token: self.cancellation_token.clone(),
#[cfg(not(target_arch = "wasm32"))]
runtime: self.runtime.clone(),
contexts: self.contexts.clone(),
metric_ring: self.metric_ring.clone(),
_phantom: PhantomData,
}
}
}
impl<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static> Debug
for TelemetryWorkerHandle<C>
{
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TelemetryWorkerHandle")
.field("sender", &self.sender)
.field("cancellation_token", &self.cancellation_token)
.finish()
}
}
#[cfg(not(target_arch = "wasm32"))]
fn schedule_deferred_cancel<F>(runtime: Option<&runtime::Handle>, future: F)
where
F: core::future::Future<Output = ()> + Send + 'static,
{
let Some(rt) = runtime else {
tracing::error!("Cannot schedule cancellation deadline: no runtime handle available");
return;
};
rt.spawn(future);
}
impl<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static>
TelemetryWorkerHandle<C>
{
pub fn register_metric_context(
&self,
name: String,
tags: Vec<Tag>,
metric_type: data::metrics::MetricType,
common: bool,
namespace: data::metrics::MetricNamespace,
) -> ContextKey {
self.contexts
.register_metric_context(name, tags, metric_type, common, namespace)
}
pub fn try_send_msg(&self, msg: TelemetryActions) -> anyhow::Result<()> {
Ok(self.sender.try_send(msg)?)
}
pub async fn send_msg(&self, msg: TelemetryActions) -> anyhow::Result<()> {
Ok(self.sender.send(msg).await?)
}
pub async fn send_msgs<T>(&self, msgs: T) -> anyhow::Result<()>
where
T: IntoIterator<Item = TelemetryActions>,
{
for msg in msgs {
self.sender.send(msg).await?;
}
Ok(())
}
pub async fn send_msg_timeout(
&self,
msg: TelemetryActions,
timeout: time::Duration,
) -> anyhow::Result<()> {
Ok(self.sender.send_timeout(msg, timeout).await?)
}
pub fn send_start(&self) -> anyhow::Result<()> {
Ok(self
.sender
.try_send(TelemetryActions::Lifecycle(LifecycleAction::Start))?)
}
pub fn send_stop(&self) -> anyhow::Result<()> {
Ok(self
.sender
.try_send(TelemetryActions::Lifecycle(LifecycleAction::Stop))?)
}
#[cfg(not(target_arch = "wasm32"))]
pub fn cancel_requests_with_deadline(&self, deadline: time::Instant) {
let token = self.cancellation_token.clone();
let remaining = deadline.saturating_duration_since(time::Instant::now());
let sleeper = <C as SleepCapability>::new();
let future = async move {
sleeper.sleep(remaining).await;
token.cancel();
};
schedule_deferred_cancel(self.runtime.as_ref(), future);
}
#[cfg(not(target_arch = "wasm32"))]
pub fn wait_for_shutdown_deadline(&self, deadline: time::Instant) {
self.cancel_requests_with_deadline(deadline);
self.wait_for_shutdown()
}
pub fn add_dependency(
&self,
name: String,
version: Option<String>,
metadata: Option<Vec<data::DependencyMetadata>>,
) -> anyhow::Result<()> {
self.sender
.try_send(TelemetryActions::AddDependency(Dependency {
name,
version,
hash: None,
metadata,
}))?;
Ok(())
}
pub fn add_product_change(
&self,
product: String,
enabled: bool,
version: Option<String>,
) -> anyhow::Result<()> {
self.sender.try_send(TelemetryActions::AddProductChange((
product,
ProductState {
enabled,
version,
error: None,
},
)))?;
Ok(())
}
pub fn add_integration(
&self,
name: String,
enabled: bool,
version: Option<String>,
compatible: Option<bool>,
auto_enabled: Option<bool>,
error: Option<String>,
) -> anyhow::Result<()> {
self.sender
.try_send(TelemetryActions::AddIntegration(Integration {
name,
version,
compatible,
enabled,
auto_enabled,
error,
}))?;
Ok(())
}
pub fn add_log<T: Hash>(
&self,
identifier: T,
message: String,
level: data::LogLevel,
stack_trace: Option<String>,
) -> anyhow::Result<()> {
let mut hasher = DefaultHasher::new();
identifier.hash(&mut hasher);
self.sender.try_send(TelemetryActions::AddLog((
LogIdentifier {
identifier: hasher.finish(),
},
data::Log {
message,
level,
stack_trace,
count: 1,
tags: String::new(),
is_sensitive: false,
is_crash: false,
},
)))?;
Ok(())
}
pub fn add_point(
&self,
value: f64,
context: &ContextKey,
extra_tags: Vec<Tag>,
) -> anyhow::Result<()> {
self.metric_ring.push(value, *context, extra_tags);
Ok(())
}
#[cfg(not(target_arch = "wasm32"))]
pub fn wait_for_shutdown(&self) {
self.shutdown.wait_for_shutdown();
}
pub fn stats(&self) -> anyhow::Result<oneshot::Receiver<TelemetryWorkerStats>> {
let (sender, receiver) = oneshot::channel();
self.sender
.try_send(TelemetryActions::CollectStats(sender))?;
Ok(receiver)
}
}
pub const MAX_ITEMS: usize = 5000;
#[derive(Debug, Default, Clone, Copy)]
pub enum TelemetryWorkerFlavor {
#[default]
Full,
MetricsLogs,
}
pub struct TelemetryWorkerBuilder {
pub host: Host,
pub application: Application,
pub runtime_id: Option<String>,
pub dependencies: store::Store<data::Dependency, data::DependencyKey>,
pub integrations: store::Store<data::Integration>,
pub configurations: store::Store<data::Configuration>,
pub endpoints: store::Store<data::Endpoint>,
pub native_deps: bool,
pub rust_shared_lib_deps: bool,
pub config: Config,
pub flavor: TelemetryWorkerFlavor,
pub install_signature: Option<data::InstallSignature>,
}
impl TelemetryWorkerBuilder {
pub fn new_fetch_host(
service_name: String,
language_name: String,
language_version: String,
tracer_version: String,
) -> Self {
Self {
host: crate::build_host(),
..Self::new(
String::new(),
service_name,
language_name,
language_version,
tracer_version,
)
}
}
pub fn new(
hostname: String,
service_name: String,
language_name: String,
language_version: String,
tracer_version: String,
) -> Self {
Self {
host: Host {
hostname,
..Default::default()
},
application: Application {
service_name,
language_name,
language_version,
tracer_version,
..Default::default()
},
runtime_id: None,
dependencies: store::Store::new(MAX_ITEMS),
integrations: store::Store::new(MAX_ITEMS),
configurations: store::Store::new(MAX_ITEMS),
endpoints: store::Store::new(10000),
native_deps: true,
rust_shared_lib_deps: false,
config: Config::default(),
flavor: TelemetryWorkerFlavor::default(),
install_signature: None,
}
}
pub fn build_worker<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static>(
self,
#[cfg(not(target_arch = "wasm32"))] tokio_runtime: Option<runtime::Handle>,
) -> (TelemetryWorkerHandle<C>, TelemetryWorker<C>) {
let (tx, mailbox) = mpsc::channel(5000);
#[cfg(not(target_arch = "wasm32"))]
let shutdown = Arc::new(InnerTelemetryShutdown {
is_shutdown: Mutex::new(false),
condvar: Condvar::new(),
});
let contexts = MetricContexts::default();
let metric_ring = Arc::new(MetricRing::new());
let token = CancellationToken::new();
let config = self.config;
let telemetry_heartbeat_interval = config.telemetry_heartbeat_interval;
let telemetry_extended_heartbeat_interval = config.telemetry_extended_heartbeat_interval;
let capabilities = C::new_without_connection_pooling();
let metrics_flush_interval =
telemetry_heartbeat_interval.min(MetricBuckets::METRICS_FLUSH_INTERVAL);
let worker = TelemetryWorker {
flavor: self.flavor,
data: TelemetryWorkerData {
started: false,
dependencies: self.dependencies,
integrations: self.integrations,
configurations: self.configurations,
endpoints: self.endpoints,
endpoints_is_first: true,
products: std::collections::HashMap::new(),
products_pending: HashSet::new(),
logs: store::QueueHashMap::default(),
metric_contexts: contexts.clone(),
metric_buckets: MetricBuckets::default(),
host: self.host,
app: self.application,
install_signature: self.install_signature,
},
config,
mailbox,
seq_id: AtomicU64::new(1),
runtime_id: self
.runtime_id
.unwrap_or_else(|| uuid::Uuid::new_v4().to_string()),
capabilities,
metrics_flush_interval,
deadlines: scheduler::Scheduler::new(vec![
(metrics_flush_interval, LifecycleAction::FlushMetricAggr),
(telemetry_heartbeat_interval, LifecycleAction::FlushData),
(
telemetry_extended_heartbeat_interval,
LifecycleAction::ExtendedHeartbeat,
),
]),
cancellation_token: token.clone(),
next_action: None,
stopped: false,
metric_ring: metric_ring.clone(),
};
(
TelemetryWorkerHandle {
sender: tx,
#[cfg(not(target_arch = "wasm32"))]
shutdown,
cancellation_token: token,
#[cfg(not(target_arch = "wasm32"))]
runtime: tokio_runtime,
contexts,
metric_ring,
_phantom: PhantomData,
},
worker,
)
}
#[cfg(not(target_arch = "wasm32"))]
pub fn spawn<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static>(
self,
) -> (TelemetryWorkerHandle<C>, JoinHandle<()>) {
let tokio_runtime = tokio::runtime::Handle::current();
let (worker_handle, worker) = self.build_worker::<C>(Some(tokio_runtime.clone()));
let join_handle = tokio_runtime.spawn(async move { worker.run_loop().await });
(worker_handle, join_handle)
}
#[cfg(not(target_arch = "wasm32"))]
pub fn run<C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static>(
self,
) -> anyhow::Result<TelemetryWorkerHandle<C>> {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
let (handle, worker) = self.build_worker::<C>(Some(runtime.handle().clone()));
let notify_shutdown = handle.shutdown.clone();
std::thread::spawn(move || {
runtime.block_on(worker.run_loop());
runtime.shutdown_background();
notify_shutdown.shutdown_finished();
});
Ok(handle)
}
}
#[cfg(test)]
mod tests {
use crate::config::TelemetryEndpoint;
use crate::data::Payload;
use crate::worker::http_client::header::{
DD_PARENT_SESSION_ID, DD_ROOT_SESSION_ID, DD_SESSION_ID,
};
use crate::worker::{
LifecycleAction, TelemetryActions, TelemetryWorker, TelemetryWorkerBuilder,
TelemetryWorkerFlavor, TelemetryWorkerHandle,
};
use libdd_capabilities_impl::NativeCapabilities;
use tokio::runtime::Runtime;
fn is_send<T: Send>(_: T) {}
fn is_sync<T: Sync>(_: T) {}
#[test]
fn test_handle_sync_send() {
#[allow(clippy::redundant_closure)]
let _ = |h: TelemetryWorkerHandle<NativeCapabilities>| is_send(h);
#[allow(clippy::redundant_closure)]
let _ = |h: TelemetryWorkerHandle<NativeCapabilities>| is_sync(h);
}
fn test_worker(
session_id: Option<String>,
root_session_id: Option<String>,
parent_session_id: Option<String>,
) -> TelemetryWorker<NativeCapabilities> {
let mut b = TelemetryWorkerBuilder::new(
"h".into(),
"svc".into(),
"lang".into(),
"1".into(),
"tv".into(),
);
b.config
.set_endpoint(TelemetryEndpoint {
url: Some("http://127.0.0.1:1".to_owned()),
..Default::default()
})
.unwrap();
b.runtime_id = Some("rid".into());
b.config.session_id = session_id;
b.config.parent_session_id = parent_session_id;
b.config.root_session_id = root_session_id;
let rt = Runtime::new().unwrap();
b.build_worker::<NativeCapabilities>(Some(rt.handle().clone()))
.1
}
#[cfg_attr(miri, ignore)] #[test]
fn telemetry_http_includes_dd_session_id() {
let req = test_worker(Some("sess".into()), None, None)
.build_request(&Payload::AppHeartbeat(()))
.unwrap();
assert_eq!(
req.headers().get(DD_SESSION_ID).unwrap().to_str().unwrap(),
"sess"
);
assert!(req.headers().get(DD_ROOT_SESSION_ID).is_none());
assert!(req.headers().get(DD_PARENT_SESSION_ID).is_none());
}
#[cfg_attr(miri, ignore)] #[test]
fn telemetry_http_omits_root_session_id_when_same_as_session_id() {
let req = test_worker(
Some("sess-id".into()),
Some("sess-id".into()),
Some("parent".into()),
)
.build_request(&Payload::AppHeartbeat(()))
.unwrap();
assert_eq!(
req.headers().get(DD_SESSION_ID).unwrap().to_str().unwrap(),
"sess-id"
);
assert!(req.headers().get(DD_ROOT_SESSION_ID).is_none());
assert_eq!(
req.headers()
.get(DD_PARENT_SESSION_ID)
.unwrap()
.to_str()
.unwrap(),
"parent"
);
}
#[cfg_attr(miri, ignore)] #[test]
fn telemetry_http_omits_parent_session_id_when_same_as_session_id() {
let req = test_worker(
Some("sess-id".into()),
Some("root".into()),
Some("sess-id".into()),
)
.build_request(&Payload::AppHeartbeat(()))
.unwrap();
assert_eq!(
req.headers().get(DD_SESSION_ID).unwrap().to_str().unwrap(),
"sess-id"
);
assert_eq!(
req.headers()
.get(DD_ROOT_SESSION_ID)
.unwrap()
.to_str()
.unwrap(),
"root"
);
assert!(req.headers().get(DD_PARENT_SESSION_ID).is_none());
}
#[cfg_attr(miri, ignore)] #[test]
fn telemetry_http_omits_session_family_without_valid_session_id() {
let assert_no_session_headers = |req: &http::Request<bytes::Bytes>| {
assert!(req.headers().get(DD_SESSION_ID).is_none());
assert!(req.headers().get(DD_ROOT_SESSION_ID).is_none());
assert!(req.headers().get(DD_PARENT_SESSION_ID).is_none());
};
let req = test_worker(None, Some("root".into()), Some("parent".into()))
.build_request(&Payload::AppHeartbeat(()))
.unwrap();
assert_no_session_headers(&req);
let req = test_worker(
Some(String::new()),
Some("root".into()),
Some("parent".into()),
)
.build_request(&Payload::AppHeartbeat(()))
.unwrap();
assert_no_session_headers(&req);
}
#[cfg_attr(miri, ignore)] #[test]
fn telemetry_http_includes_dd_session_root_and_parent_session_ids() {
let req = test_worker(
Some("sess".into()),
Some("root".into()),
Some("parent".into()),
)
.build_request(&Payload::AppHeartbeat(()))
.unwrap();
assert_eq!(
req.headers().get(DD_SESSION_ID).unwrap().to_str().unwrap(),
"sess"
);
assert_eq!(
req.headers()
.get(DD_ROOT_SESSION_ID)
.unwrap()
.to_str()
.unwrap(),
"root"
);
assert_eq!(
req.headers()
.get(DD_PARENT_SESSION_ID)
.unwrap()
.to_str()
.unwrap(),
"parent"
);
}
fn build_test_worker_with_flavor(
flavor: TelemetryWorkerFlavor,
) -> TelemetryWorker<NativeCapabilities> {
let mut b = TelemetryWorkerBuilder::new(
"h".into(),
"svc".into(),
"lang".into(),
"1".into(),
"tv".into(),
);
b.config
.set_endpoint(TelemetryEndpoint {
url: Some("http://127.0.0.1:1".to_owned()),
..Default::default()
})
.unwrap();
b.runtime_id = Some("rid".into());
b.flavor = flavor;
b.build_worker::<NativeCapabilities>(Some(tokio::runtime::Handle::current()))
.1
}
#[tokio::test]
#[cfg_attr(miri, ignore)] async fn endpoints_message_limit_chunks_payloads_and_flags_only_the_first() {
let mut worker = build_test_worker_with_flavor(TelemetryWorkerFlavor::Full);
worker.config.endpoints_message_limit = 2;
for i in 0..5 {
worker.data.endpoints.insert(crate::data::Endpoint {
operation_name: "http.request".to_string(),
resource_name: format!("GET /r{i}"),
..Default::default()
});
}
let mut chunks = Vec::new();
while worker.data.endpoints.flush_not_empty() {
let payloads = worker.build_app_events_batch();
let endpoints = payloads
.iter()
.find_map(|p| match p {
crate::data::Payload::AppEndpoints(e) => Some(e),
_ => None,
})
.expect("an app-endpoints payload while endpoints are queued");
chunks.push((endpoints.is_first, endpoints.endpoints.len()));
let sent = crate::data::Payload::AppEndpoints(crate::data::AppEndpoints {
is_first: endpoints.is_first,
endpoints: endpoints.endpoints.clone(),
});
worker.payload_sent_success(&sent);
}
assert_eq!(
chunks,
vec![(true, 2), (false, 2), (false, 1)],
"5 endpoints at a limit of 2 should be 2+2+1 with is_first only on the first payload"
);
}
#[tokio::test]
#[cfg_attr(miri, ignore)] async fn full_flavor_start_schedules_every_periodic_action() {
let mut worker = build_test_worker_with_flavor(TelemetryWorkerFlavor::Full);
let _ = worker
.dispatch_action(TelemetryActions::Lifecycle(LifecycleAction::Start))
.await;
let delays: Vec<LifecycleAction> =
worker.deadlines.delays.iter().map(|(_, k)| *k).collect();
let scheduled: Vec<LifecycleAction> =
worker.deadlines.deadlines.iter().map(|(_, k)| *k).collect();
assert!(!delays.is_empty(), "scheduler should have periodic actions");
for ev in &delays {
assert!(
scheduled.contains(ev),
"{ev:?} has a delay but was not scheduled on Start; scheduled={scheduled:?}",
);
}
}
#[tokio::test]
#[cfg_attr(miri, ignore)] async fn metrics_logs_flavor_start_does_not_schedule_extended_heartbeat() {
let mut worker = build_test_worker_with_flavor(TelemetryWorkerFlavor::MetricsLogs);
let _ = worker
.dispatch_metrics_logs_action(TelemetryActions::Lifecycle(LifecycleAction::Start))
.await;
let scheduled: Vec<LifecycleAction> =
worker.deadlines.deadlines.iter().map(|(_, k)| *k).collect();
assert!(scheduled.contains(&LifecycleAction::FlushMetricAggr));
assert!(scheduled.contains(&LifecycleAction::FlushData));
assert!(
!scheduled.contains(&LifecycleAction::ExtendedHeartbeat),
"MetricsLogs should not schedule ExtendedHeartbeat; scheduled={scheduled:?}",
);
}
#[tokio::test]
#[cfg_attr(miri, ignore)] async fn extended_heartbeat_does_not_reset_flush_data() {
let mut worker = build_test_worker_with_flavor(TelemetryWorkerFlavor::Full);
let _ = worker
.dispatch_action(TelemetryActions::Lifecycle(LifecycleAction::Start))
.await;
let flush_data_before = worker
.deadlines
.deadlines
.iter()
.find(|(_, k)| *k == LifecycleAction::FlushData)
.map(|(d, _)| *d)
.expect("FlushData scheduled on Start");
let _ = worker
.dispatch_action(TelemetryActions::Lifecycle(
LifecycleAction::ExtendedHeartbeat,
))
.await;
let flush_data_after = worker
.deadlines
.deadlines
.iter()
.find(|(_, k)| *k == LifecycleAction::FlushData)
.map(|(d, _)| *d)
.expect("FlushData should still be scheduled after ExtendedHeartbeat fires");
assert_eq!(
flush_data_before, flush_data_after,
"ExtendedHeartbeat must not reset FlushData's deadline",
);
}
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn app_started_omits_dependencies_and_integrations() {
let mut worker = build_test_worker_with_flavor(TelemetryWorkerFlavor::Full);
let _ = worker
.dispatch_action(TelemetryActions::AddDependency(crate::data::Dependency {
name: "monolog/monolog".into(),
version: Some("3.5.0".into()),
..Default::default()
}))
.await;
let _ = worker
.dispatch_action(TelemetryActions::AddIntegration(crate::data::Integration {
name: "curl".into(),
enabled: true,
..Default::default()
}))
.await;
let app_started = worker.build_app_started();
assert!(
app_started.dependencies.is_empty(),
"app-started must not carry dependencies; the intake rejects the whole payload",
);
assert!(
app_started.integrations.is_empty(),
"app-started must not carry integrations; the intake rejects the whole payload",
);
let batch = worker.build_app_events_batch();
assert!(
batch.iter().any(|p| matches!(
p,
crate::data::Payload::AppDependenciesLoaded(d) if !d.dependencies.is_empty()
)),
"dependencies registered before Start must still be reported via \
app-dependencies-loaded, got {batch:?}",
);
assert!(
batch.iter().any(|p| matches!(
p,
crate::data::Payload::AppIntegrationsChange(i) if !i.integrations.is_empty()
)),
"integrations registered before Start must still be reported via \
app-integrations-change, got {batch:?}",
);
}
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn extended_heartbeat_carries_dependencies_and_integrations() {
let mut worker = build_test_worker_with_flavor(TelemetryWorkerFlavor::Full);
let _ = worker
.dispatch_action(TelemetryActions::AddDependency(crate::data::Dependency {
name: "monolog/monolog".into(),
version: Some("3.5.0".into()),
..Default::default()
}))
.await;
let _ = worker
.dispatch_action(TelemetryActions::AddIntegration(crate::data::Integration {
name: "curl".into(),
enabled: true,
..Default::default()
}))
.await;
let hb = worker.build_extended_heartbeat();
assert_eq!(1, hb.dependencies.len(), "{hb:?}");
assert_eq!(1, hb.integrations.len(), "{hb:?}");
}
mod reset {
use super::super::*;
use crate::data::{
metrics::{MetricNamespace, MetricType},
Configuration, ConfigurationOrigin, Dependency, Endpoint, Integration, Log, LogLevel,
};
use libdd_capabilities_impl::NativeCapabilities;
use libdd_shared_runtime::Worker;
fn build_test_worker() -> (
TelemetryWorkerHandle<NativeCapabilities>,
TelemetryWorker<NativeCapabilities>,
) {
let builder = TelemetryWorkerBuilder::new(
"hostname".to_string(),
"test-service".to_string(),
"rust".to_string(),
"1.0.0".to_string(),
"1.0.0".to_string(),
);
builder.build_worker::<NativeCapabilities>(Some(tokio::runtime::Handle::current()))
}
fn make_log(id: u64, message: &str) -> (LogIdentifier, Log) {
(
LogIdentifier { identifier: id },
Log {
message: message.to_string(),
level: LogLevel::Warn,
stack_trace: None,
count: 1,
tags: String::new(),
is_sensitive: false,
is_crash: false,
},
)
}
#[cfg_attr(miri, ignore)] #[tokio::test]
async fn test_reset_clears_buffered_data() {
let (handle, mut worker) = build_test_worker();
worker.data.dependencies.insert(Dependency {
name: "dep".to_string(),
..Default::default()
});
worker.data.integrations.insert(Integration {
name: "integration".to_string(),
version: None,
enabled: true,
compatible: None,
auto_enabled: None,
..Default::default()
});
worker.data.configurations.insert(Configuration {
name: "cfg".to_string(),
value: Some("true".to_string()),
origin: ConfigurationOrigin::Code,
config_id: None,
seq_id: None,
});
worker.data.endpoints.insert(Endpoint {
operation_name: "GET /health".to_string(),
resource_name: "/health".to_string(),
..Default::default()
});
let (id, log) = make_log(42, "msg");
worker.data.logs.get_mut_or_insert(id, log);
let key = handle.register_metric_context(
"test.metric".to_string(),
vec![],
MetricType::Count,
false,
MetricNamespace::Tracers,
);
worker.data.metric_buckets.add_point(key, 1.0, vec![]);
worker.reset();
let stats = worker.stats();
assert_eq!(
stats.dependencies_stored, 0,
"dependency dedupe history should be cleared"
);
assert_eq!(
stats.dependencies_unflushed, 0,
"dependency pending queue should be cleared"
);
assert_eq!(
stats.integrations_stored, 0,
"integration dedupe history should be cleared"
);
assert_eq!(
stats.integrations_unflushed, 0,
"integration pending queue should be cleared"
);
assert_eq!(
stats.configurations_stored, 0,
"configuration dedupe history should be cleared"
);
assert_eq!(
stats.configurations_unflushed, 0,
"configuration pending queue should be cleared"
);
assert_eq!(stats.logs, 0, "logs should be cleared");
assert_eq!(
stats.metric_buckets.buckets, 0,
"metric buckets should be cleared"
);
assert_eq!(
stats.metric_buckets.series, 0,
"metric series should be cleared"
);
assert_eq!(
worker.data.endpoints.len_stored(),
0,
"endpoints should be cleared"
);
assert!(
worker.data.endpoints_is_first,
"the child's first app-endpoints payload is a first one again"
);
assert!(worker.next_action.is_none(), "next_action should be None");
}
#[cfg_attr(miri, ignore)] #[tokio::test]
async fn test_reset_drains_mailbox() {
let (handle, mut worker) = build_test_worker();
handle
.try_send_msg(TelemetryActions::AddDependency(Dependency {
name: "dep".to_string(),
..Default::default()
}))
.unwrap();
let (id, log) = make_log(1, "pre-fork log");
handle
.try_send_msg(TelemetryActions::AddLog((id, log)))
.unwrap();
worker.next_action = Some(TelemetryActions::Lifecycle(LifecycleAction::Start));
worker.reset();
assert!(
worker.mailbox.try_recv().is_err(),
"mailbox should be empty"
);
assert!(worker.next_action.is_none(), "next_action should be None");
let stats = worker.stats();
assert_eq!(
stats.dependencies_stored, 0,
"queued AddDependency must not be applied"
);
assert_eq!(
stats.dependencies_unflushed, 0,
"queued AddDependency must not be pending"
);
assert_eq!(stats.logs, 0, "queued AddLog must be discarded");
}
#[cfg_attr(miri, ignore)] #[tokio::test]
async fn test_worker_accepts_new_data_after_reset() {
let (handle, mut worker) = build_test_worker();
worker.flavor = TelemetryWorkerFlavor::MetricsLogs;
let (id, log) = make_log(99, "pre-fork");
worker.data.logs.get_mut_or_insert(id, log);
worker.reset();
let (id2, log2) = make_log(1, "post-fork");
handle
.try_send_msg(TelemetryActions::AddLog((id2, log2)))
.unwrap();
worker.trigger().await;
worker.run().await;
let stats = worker.stats();
assert_eq!(stats.logs, 1, "only post-fork log should be present");
}
#[cfg_attr(miri, ignore)] #[tokio::test]
async fn test_reset_preserves_started_and_deadlines() {
let (_handle, mut worker) = build_test_worker();
worker.data.started = true;
worker
.deadlines
.schedule_event(LifecycleAction::FlushMetricAggr)
.unwrap();
worker
.deadlines
.schedule_event(LifecycleAction::FlushData)
.unwrap();
let deadlines_before = worker.deadlines.deadlines.clone();
worker.reset();
assert!(worker.data.started, "started flag should be preserved");
assert_eq!(
worker.deadlines.deadlines.len(),
deadlines_before.len(),
"scheduled deadlines should be preserved"
);
for ((_, actual), (_, expected)) in worker
.deadlines
.deadlines
.iter()
.zip(deadlines_before.iter())
{
assert_eq!(
actual, expected,
"deadline kinds should be preserved across reset"
);
}
}
}
#[cfg_attr(miri, ignore)]
#[test]
fn test_channel_close_flushes_and_parks_via_shared_runtime() {
use httpmock::prelude::*;
use libdd_shared_runtime::{BlockingRuntime, ForkSafeRuntime, SharedRuntime};
use std::time::Duration;
const TELEMETRY_PATH: &str = "/telemetry/proxy/api/v2/apmtelemetry";
let server = MockServer::start();
let mock = server.mock(|when, then| {
when.method(POST).path(TELEMETRY_PATH);
then.status(202).body("");
});
let mut builder = TelemetryWorkerBuilder::new(
"host".into(),
"svc".into(),
"lang".into(),
"1".into(),
"tv".into(),
);
builder
.config
.set_endpoint(TelemetryEndpoint {
url: Some(server.url("/")),
..Default::default()
})
.unwrap();
builder.runtime_id = Some("rid".into());
let shared_runtime = ForkSafeRuntime::new().expect("ForkSafeRuntime::new");
let runtime_handle = shared_runtime
.block_on(async { tokio::runtime::Handle::current() })
.expect("runtime handle");
let (telemetry_handle, worker) =
builder.build_worker::<NativeCapabilities>(Some(runtime_handle));
let _worker_handle = shared_runtime
.spawn_worker(worker, false)
.expect("spawn_worker");
telemetry_handle.send_start().expect("send_start");
for _ in 0..50 {
if mock.calls() >= 1 {
break;
}
std::thread::sleep(Duration::from_millis(20));
}
assert!(
mock.calls() >= 1,
"worker should POST at least once after Start"
);
let hits_before_close = mock.calls();
drop(telemetry_handle);
for _ in 0..50 {
if mock.calls() > hits_before_close {
break;
}
std::thread::sleep(Duration::from_millis(20));
}
assert!(
mock.calls() > hits_before_close,
"worker should flush a final batch after the channel is closed"
);
let stable_hits = mock.calls();
std::thread::sleep(Duration::from_millis(300));
assert_eq!(
mock.calls(),
stable_hits,
"worker must stop POSTing after parking; observed {} extra hits",
mock.calls().saturating_sub(stable_hits),
);
}
#[cfg_attr(miri, ignore)]
#[tokio::test]
async fn shutdown_drains_pending_actions_before_stop() {
use crate::data::metrics::{MetricNamespace, MetricType};
use httpmock::prelude::*;
const TELEMETRY_PATH: &str = "/telemetry/proxy/api/v2/apmtelemetry";
const METRIC_NAME: &str = "regression.drain_before_stop";
let server = MockServer::start();
let metric_mock = server.mock(|when, then| {
when.method(POST)
.path(TELEMETRY_PATH)
.body_includes(format!(r#""metric":"{METRIC_NAME}""#));
then.status(202).body("");
});
let _lifecycle = server.mock(|when, then| {
when.method(POST).path(TELEMETRY_PATH);
then.status(202).body("");
});
let mut builder = TelemetryWorkerBuilder::new(
"host".into(),
"svc".into(),
"lang".into(),
"1".into(),
"tv".into(),
);
builder
.config
.set_endpoint(TelemetryEndpoint {
url: Some(server.url("/")),
..Default::default()
})
.unwrap();
builder.runtime_id = Some("rid".into());
let (handle, mut worker) =
builder.build_worker::<NativeCapabilities>(Some(tokio::runtime::Handle::current()));
let context = handle.register_metric_context(
METRIC_NAME.into(),
Vec::new(),
MetricType::Count,
false,
MetricNamespace::Tracers,
);
let _ = worker
.dispatch_action(TelemetryActions::Lifecycle(LifecycleAction::Start))
.await;
handle
.add_point(1.0, &context, Vec::new())
.expect("add_point");
<TelemetryWorker<_> as libdd_shared_runtime::Worker>::shutdown(&mut worker).await;
metric_mock.assert_calls(1);
}
}