diff --git a/Cargo.lock b/Cargo.lock index c366415215..acbf61071a 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3419,6 +3419,8 @@ dependencies = [ "libdd-shared-runtime", "serde", "serde_json", + "strum", + "strum_macros", "sys-info", "tokio", "tokio-util", diff --git a/datadog-sidecar-ffi/src/lib.rs b/datadog-sidecar-ffi/src/lib.rs index c9e132ec73..4365344d9e 100644 --- a/datadog-sidecar-ffi/src/lib.rs +++ b/datadog-sidecar-ffi/src/lib.rs @@ -515,6 +515,8 @@ pub unsafe extern "C" fn ddog_sidecar_telemetry_addDependency( let dependency = TelemetryActions::AddDependency(Dependency { name: dependency_name.to_utf8_lossy().into_owned(), version, + hash: None, + metadata: None, }); try_c!(blocking::enqueue_actions( diff --git a/libdd-data-pipeline/src/telemetry/mod.rs b/libdd-data-pipeline/src/telemetry/mod.rs index ec7950c8ed..c371436996 100644 --- a/libdd-data-pipeline/src/telemetry/mod.rs +++ b/libdd-data-pipeline/src/telemetry/mod.rs @@ -190,6 +190,16 @@ impl std } } +impl TelemetryClient { + /// Allow sharing a telemetry worker with data-pipeline + pub fn with_handle(handle: TelemetryWorkerHandle) -> Self { + TelemetryClient { + metrics: Metrics::new(&handle), + worker: handle, + } + } +} + /// Telemetry describing the sending of a trace payload /// It can be produced from a [`SendWithRetryResult`] or from a [`SendDataResult`]. #[derive(PartialEq, Debug, Default)] diff --git a/libdd-data-pipeline/src/trace_exporter/builder.rs b/libdd-data-pipeline/src/trace_exporter/builder.rs index 65337e46d0..7eed15e578 100644 --- a/libdd-data-pipeline/src/trace_exporter/builder.rs +++ b/libdd-data-pipeline/src/trace_exporter/builder.rs @@ -6,7 +6,7 @@ use crate::agentless::config::{AgentlessTraceConfig, DEFAULT_AGENTLESS_TIMEOUT}; use crate::otlp::config::{OtlpProtocol, DEFAULT_OTLP_TIMEOUT}; use crate::otlp::{OtlpMetricsConfig, OtlpResourceInfo, OtlpTraceConfig}; #[cfg(feature = "telemetry")] -use crate::telemetry::TelemetryClientBuilder; +use crate::telemetry::{TelemetryClient, TelemetryClientBuilder}; use crate::trace_exporter::agent_response::AgentResponsePayloadVersion; use crate::trace_exporter::error::BuilderErrorKind; use crate::trace_exporter::log_writer::DEFAULT_LOG_MAX_LINE_SIZE; @@ -25,6 +25,8 @@ use libdd_dogstatsd_client::DogStatsDClient; use libdd_shared_runtime::SharedRuntime; #[cfg(not(target_arch = "wasm32"))] use libdd_shared_runtime::{BlockingRuntime, ForkSafeRuntime}; +#[cfg(all(not(target_arch = "wasm32"), feature = "telemetry"))] +use libdd_telemetry::worker::TelemetryWorkerHandle; use libdd_trace_utils::trace_filter::TraceFilterer; use std::sync::Arc; use std::time::Duration; @@ -48,6 +50,12 @@ fn build_otlp_header_map(headers: Vec<(String, String)>) -> http::HeaderMap { out } +/// Externally-owned telemetry worker to consolidate through; Boxed to avoid depending on C. +#[cfg(all(feature = "telemetry", not(target_arch = "wasm32")))] +type BoxedTelemetryHandle = Box; +#[cfg(all(feature = "telemetry", target_arch = "wasm32"))] +type BoxedTelemetryHandle = Box; + #[allow(missing_docs)] #[derive(Debug)] pub struct TraceExporterBuilder { @@ -80,6 +88,8 @@ pub struct TraceExporterBuilder { client_side_stats_obfuscation_enabled: bool, #[cfg(feature = "telemetry")] telemetry: Option, + #[cfg(feature = "telemetry")] + telemetry_handle: Option, telemetry_instrumentation_sessions: TelemetryInstrumentationSessions, shared_runtime: Option>, health_metrics_enabled: bool, @@ -149,6 +159,8 @@ impl TraceExporterBuilder { client_side_stats_obfuscation_enabled: false, #[cfg(feature = "telemetry")] telemetry: None, + #[cfg(feature = "telemetry")] + telemetry_handle: None, telemetry_instrumentation_sessions: TelemetryInstrumentationSessions::default(), shared_runtime: None, health_metrics_enabled: false, @@ -362,6 +374,17 @@ impl TraceExporterBuilder { self } + #[cfg(feature = "telemetry")] + pub fn set_telemetry_handle< + C: HttpClientCapability + SleepCapability + MaybeSend + Sync + 'static, + >( + &mut self, + handle: TelemetryWorkerHandle, + ) -> &mut Self { + self.telemetry_handle = Some(Box::new(handle)); + self + } + /// Sets optional instrumentation session headers on telemetry requests (`dd-session-id`, etc.). pub fn set_telemetry_instrumentation_sessions( &mut self, @@ -641,57 +664,65 @@ impl TraceExporterBuilder { } #[cfg(feature = "telemetry")] - let (telemetry_client, telemetry_handle) = { - let sessions = self.telemetry_instrumentation_sessions; - // Telemetry talks to the agent; disable it in agentless and log-export modes. - let telemetry = self - .telemetry - .filter(|_| !(agentless_enabled || self.output_to_log)) - .map(|telemetry_config| -> Result<_, TraceExporterError> { - let mut tb = TelemetryClientBuilder::default() - .set_language(&self.language) - .set_language_version(&self.language_version) - .set_service_name(&self.service) - .set_service_version(&self.app_version) - .set_env(&self.env) - .set_tracer_version(&self.tracer_version) - .set_heartbeat(telemetry_config.heartbeat) - .set_url(base_url) - .set_debug_enabled(telemetry_config.debug_enabled); - if let Some(id) = telemetry_config.runtime_id { - tb = tb.set_runtime_id(&id); - } - if let Some(ref id) = sessions.session_id { - tb = tb.set_session_id(id); - } - if let Some(ref id) = sessions.root_session_id { - tb = tb.set_root_session_id(id); - } - if let Some(ref id) = sessions.parent_session_id { - tb = tb.set_parent_session_id(id); - } - tb.build::().map_err(|e| { - TraceExporterError::Builder(BuilderErrorKind::InvalidConfiguration( - e.to_string(), - )) + let (telemetry_client, telemetry_handle) = + // Consolidated path: report health metrics through the externally-owned worker. + // We do not own it (no spawn, no start, no shutdown handle). + if let Some(shared_handle) = self + .telemetry_handle + .and_then(|h| h.downcast::>().ok()) + { + (Some(TelemetryClient::with_handle(*shared_handle)), None) + } else { + let sessions = self.telemetry_instrumentation_sessions; + // Telemetry talks to the agent; disable it in agentless and log-export modes. + let telemetry = self + .telemetry + .filter(|_| !(agentless_enabled || self.output_to_log)) + .map(|telemetry_config| -> Result<_, TraceExporterError> { + let mut tb = TelemetryClientBuilder::default() + .set_language(&self.language) + .set_language_version(&self.language_version) + .set_service_name(&self.service) + .set_service_version(&self.app_version) + .set_env(&self.env) + .set_tracer_version(&self.tracer_version) + .set_heartbeat(telemetry_config.heartbeat) + .set_url(base_url) + .set_debug_enabled(telemetry_config.debug_enabled); + if let Some(id) = telemetry_config.runtime_id { + tb = tb.set_runtime_id(&id); + } + if let Some(ref id) = sessions.session_id { + tb = tb.set_session_id(id); + } + if let Some(ref id) = sessions.root_session_id { + tb = tb.set_root_session_id(id); + } + if let Some(ref id) = sessions.parent_session_id { + tb = tb.set_parent_session_id(id); + } + tb.build::().map_err(|e| { + TraceExporterError::Builder(BuilderErrorKind::InvalidConfiguration( + e.to_string(), + )) + }) }) - }) - .transpose()?; - match telemetry { - Some((client_tel, worker)) => { - let handle = shared_runtime.spawn_worker(worker, false).map_err(|e| { - TraceExporterError::Builder(BuilderErrorKind::InvalidConfiguration( - e.to_string(), - )) - })?; - if let Err(e) = client_tel.start() { - tracing::warn!("Failed to start telemetry: {e}"); + .transpose()?; + match telemetry { + Some((client_tel, worker)) => { + let handle = shared_runtime.spawn_worker(worker, false).map_err(|e| { + TraceExporterError::Builder(BuilderErrorKind::InvalidConfiguration( + e.to_string(), + )) + })?; + if let Err(e) = client_tel.start() { + tracing::warn!("Failed to start telemetry: {e}"); + } + (Some(client_tel), Some(handle)) } - (Some(client_tel), Some(handle)) + None => (None, None), } - None => (None, None), - } - }; + }; // Transport selection: agentless is mutually exclusive with both OTLP and a // user-supplied agent URL; OTLP and the agent URL may coexist. All exclusion diff --git a/libdd-telemetry-ffi/src/worker_handle.rs b/libdd-telemetry-ffi/src/worker_handle.rs index b3eb93be3e..2afa41d7f8 100644 --- a/libdd-telemetry-ffi/src/worker_handle.rs +++ b/libdd-telemetry-ffi/src/worker_handle.rs @@ -25,7 +25,7 @@ pub unsafe extern "C" fn ddog_telemetry_handle_add_dependency( ) -> MaybeError { let name = crate::try_c!(dependency_name.try_to_string()); let version = crate::try_c!(dependency_version.try_to_string_option()); - crate::try_c!(handle.add_dependency(name, version)); + crate::try_c!(handle.add_dependency(name, version, None)); MaybeError::None } diff --git a/libdd-telemetry/Cargo.toml b/libdd-telemetry/Cargo.toml index 803b7002c2..a3d9bfb86c 100644 --- a/libdd-telemetry/Cargo.toml +++ b/libdd-telemetry/Cargo.toml @@ -25,6 +25,8 @@ futures = { version = "0.3", default-features = false } http = "1" serde = { workspace = true, features = ["derive"] } serde_json = { version = "1.0" } +strum = { version = "0.26", default-features = false } +strum_macros = "0.26" tokio = { workspace = true, features = ["sync", "io-util"] } tokio-util = { version = "0.7", features = ["codec"] } tracing.workspace = true diff --git a/libdd-telemetry/examples/tm-ping.rs b/libdd-telemetry/examples/tm-ping.rs index fc316491d0..8f5d2f9fd4 100644 --- a/libdd-telemetry/examples/tm-ping.rs +++ b/libdd-telemetry/examples/tm-ping.rs @@ -24,6 +24,9 @@ fn build_app_started_payload() -> AppStarted { configuration: Vec::new(), dependencies: Vec::new(), integrations: Vec::new(), + install_signature: None, + products: Default::default(), + error: None, } } diff --git a/libdd-telemetry/src/config.rs b/libdd-telemetry/src/config.rs index 3fc12e014f..519beb34d5 100644 --- a/libdd-telemetry/src/config.rs +++ b/libdd-telemetry/src/config.rs @@ -58,6 +58,15 @@ pub struct Config { pub parent_session_id: Option, #[serde(default)] pub root_session_id: Option, + + /// Whether to emit the `app-started`/`app-closing` lifecycle payloads. + /// Forked processes may not be required to emit these. + #[serde(default = "default_true")] + pub emit_app_lifecycle: bool, +} + +fn default_true() -> bool { + true } fn endpoint_with_telemetry_path( @@ -195,6 +204,7 @@ impl Default for Config { session_id: None, parent_session_id: None, root_session_id: None, + emit_app_lifecycle: true, } } } @@ -334,6 +344,7 @@ impl Config { session_id: None, parent_session_id: None, root_session_id: None, + emit_app_lifecycle: true, }; _ = this.set_endpoint(TelemetryEndpoint { diff --git a/libdd-telemetry/src/data/metrics.rs b/libdd-telemetry/src/data/metrics.rs index e131298a1a..71e7271fe2 100644 --- a/libdd-telemetry/src/data/metrics.rs +++ b/libdd-telemetry/src/data/metrics.rs @@ -36,8 +36,21 @@ pub enum SerializedSketch { B64 { sketch_b64: String }, } -#[derive(Serialize, Deserialize, Debug, Clone, Copy)] +#[derive( + Serialize, + Deserialize, + Debug, + Clone, + Copy, + PartialEq, + Eq, + Hash, + strum_macros::Display, + strum_macros::EnumIter, + strum_macros::IntoStaticStr, +)] #[serde(rename_all = "snake_case")] +#[strum(serialize_all = "snake_case")] #[repr(C)] pub enum MetricNamespace { Tracers, @@ -51,13 +64,30 @@ pub enum MetricNamespace { Telemetry, Apm, Sidecar, + Civisibility, + Mlobs, + Ddtraceapi, } -#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq, Hash)] +#[derive( + Serialize, + Deserialize, + Debug, + Clone, + Copy, + PartialEq, + Eq, + Hash, + strum_macros::Display, + strum_macros::EnumIter, + strum_macros::IntoStaticStr, +)] #[serde(rename_all = "snake_case")] +#[strum(serialize_all = "snake_case")] #[repr(C)] pub enum MetricType { Gauge, Count, Distribution, + Rate, } diff --git a/libdd-telemetry/src/data/payload.rs b/libdd-telemetry/src/data/payload.rs index 609969ce38..815d12a995 100644 --- a/libdd-telemetry/src/data/payload.rs +++ b/libdd-telemetry/src/data/payload.rs @@ -11,6 +11,7 @@ pub enum Payload { AppStarted(AppStarted), AppDependenciesLoaded(AppDependenciesLoaded), AppIntegrationsChange(AppIntegrationsChange), + AppProductChange(AppProductChange), AppClientConfigurationChange(AppClientConfigurationChange), AppEndpoints(AppEndpoints), AppHeartbeat(#[serde(skip_serializing)] ()), @@ -29,6 +30,7 @@ impl Payload { AppStarted(_) => "app-started", AppDependenciesLoaded(_) => "app-dependencies-loaded", AppIntegrationsChange(_) => "app-integrations-change", + AppProductChange(_) => "app-product-change", AppClientConfigurationChange(_) => "app-client-configuration-change", AppEndpoints(_) => "app-endpoints", AppHeartbeat(_) => "app-heartbeat", @@ -68,6 +70,9 @@ mod tests { ], dependencies: Vec::new(), integrations: Vec::new(), + install_signature: None, + products: Default::default(), + error: None, }); let serialized = serde_json::to_value(&payload).unwrap(); @@ -106,10 +111,11 @@ mod tests { Dependency { name: "tokio".to_string(), version: Some("1.32.0".to_string()), + ..Default::default() }, Dependency { name: "serde".to_string(), - version: None, + ..Default::default() }, ], }); @@ -484,6 +490,100 @@ mod tests { assert_eq!(serialized, expected); } + #[test] + fn test_app_product_change_serialization() { + let mut products = std::collections::HashMap::new(); + products.insert( + "appsec".to_string(), + ProductState { + enabled: true, + version: Some("1.2.3".to_string()), + error: None, + }, + ); + let payload = Payload::AppProductChange(AppProductChange { products }); + + let serialized = serde_json::to_value(&payload).unwrap(); + + let expected = json!({ + "request_type": "app-product-change", + "payload": { + "products": { + "appsec": { + "enabled": true, + "version": "1.2.3" + } + } + } + }); + + assert_eq!(serialized, expected); + } + + #[test] + fn test_dependency_metadata_serialization() { + // A plain dependency (no metadata) omits both `hash` and `metadata`. + let plain = Payload::AppDependenciesLoaded(AppDependenciesLoaded { + dependencies: vec![Dependency { + name: "requests".to_string(), + version: Some("2.0".to_string()), + ..Default::default() + }], + }); + assert_eq!( + serde_json::to_value(&plain).unwrap(), + json!({ + "request_type": "app-dependencies-loaded", + "payload": { "dependencies": [{ "name": "requests", "version": "2.0" }] } + }) + ); + + // With SCA metadata: per the telemetry schema, `type` names the metadata kind and + // `value` is an opaque stringified-JSON payload that consumers must parse. + let with_sca = Payload::AppDependenciesLoaded(AppDependenciesLoaded { + dependencies: vec![Dependency { + name: "requests".to_string(), + version: Some("2.0".to_string()), + metadata: Some(vec![DependencyMetadata { + r#type: "reachability".to_string(), + value: "{\"id\":\"CVE-2024-1\",\"reached\":true}".to_string(), + }]), + ..Default::default() + }], + }); + assert_eq!( + serde_json::to_value(&with_sca).unwrap(), + json!({ + "request_type": "app-dependencies-loaded", + "payload": { "dependencies": [{ + "name": "requests", + "version": "2.0", + "metadata": [{ + "type": "reachability", + "value": "{\"id\":\"CVE-2024-1\",\"reached\":true}" + }] + }] } + }) + ); + + // SCA enabled with no findings: `metadata` is present-but-empty (distinct from omitted). + let empty_meta = Payload::AppDependenciesLoaded(AppDependenciesLoaded { + dependencies: vec![Dependency { + name: "requests".to_string(), + version: Some("2.0".to_string()), + metadata: Some(vec![]), + ..Default::default() + }], + }); + assert_eq!( + serde_json::to_value(&empty_meta).unwrap(), + json!({ + "request_type": "app-dependencies-loaded", + "payload": { "dependencies": [{ "name": "requests", "version": "2.0", "metadata": [] }] } + }) + ); + } + #[test] fn test_app_extended_heartbeat_serialization() { let payload = Payload::AppExtendedHeartbeat(AppStarted { @@ -496,6 +596,9 @@ mod tests { }], dependencies: Vec::new(), integrations: Vec::new(), + install_signature: None, + products: Default::default(), + error: None, }); let serialized = serde_json::to_value(&payload).unwrap(); diff --git a/libdd-telemetry/src/data/payloads.rs b/libdd-telemetry/src/data/payloads.rs index fb5c97dc93..6d31e3946c 100644 --- a/libdd-telemetry/src/data/payloads.rs +++ b/libdd-telemetry/src/data/payloads.rs @@ -1,6 +1,7 @@ // Copyright 2021-Present Datadog, Inc. https://www.datadoghq.com/ // SPDX-License-Identifier: Apache-2.0 +use std::collections::HashMap; use std::hash::Hasher; use crate::data::metrics; @@ -11,6 +12,36 @@ use serde::{Deserialize, Serialize}; pub struct Dependency { pub name: String, pub version: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub hash: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub metadata: Option>, +} + +/// SCA metadata may be attached to a dependency after it was first reported; keying the store +/// by this instead of the full struct means that update refreshes the existing entry instead of +/// being stored (and re-sent) as a second entry for the same package/version. +#[derive(Debug, Clone, PartialEq, Eq, Hash)] +pub struct DependencyKey { + name: String, + version: Option, + hash: Option, +} + +impl crate::worker::store::Keyed for Dependency { + fn key(&self) -> DependencyKey { + DependencyKey { + name: self.name.clone(), + version: self.version.clone(), + hash: self.hash.clone(), + } + } +} + +#[derive(Serialize, Deserialize, Debug, Hash, PartialEq, Eq, Clone, Default)] +pub struct DependencyMetadata { + pub r#type: String, + pub value: String, } #[derive(Serialize, Deserialize, Debug, Hash, PartialEq, Eq, Clone, Default)] @@ -31,9 +62,22 @@ pub struct Configuration { pub seq_id: Option, } -#[derive(Serialize, Deserialize, Debug, Hash, PartialEq, Eq, Clone)] +#[derive( + Serialize, + Deserialize, + Debug, + Hash, + PartialEq, + Eq, + Clone, + Copy, + strum_macros::Display, + strum_macros::EnumIter, + strum_macros::IntoStaticStr, +)] #[repr(C)] #[serde(rename_all = "snake_case")] +#[strum(serialize_all = "snake_case")] pub enum ConfigurationOrigin { EnvVar, Code, @@ -43,6 +87,14 @@ pub enum ConfigurationOrigin { LocalStableConfig, FleetStableConfig, Calculated, + OtelEnvVar, + Unknown, +} + +#[derive(Serialize, Deserialize, Debug, PartialEq, Eq, Clone)] +pub struct Error { + pub code: Option, + pub message: Option, } #[derive(Serialize, Debug)] @@ -50,6 +102,19 @@ pub struct AppStarted { pub configuration: Vec, pub dependencies: Vec, pub integrations: Vec, + #[serde(skip_serializing_if = "Option::is_none")] + pub install_signature: Option, + #[serde(skip_serializing_if = "HashMap::is_empty")] + pub products: HashMap, + #[serde(skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +#[derive(Serialize, Debug, Clone)] +pub struct InstallSignature { + pub install_id: Option, + pub install_type: Option, + pub install_time: Option, } #[derive(Serialize, Debug)] @@ -57,6 +122,19 @@ pub struct AppDependenciesLoaded { pub dependencies: Vec, } +#[derive(Serialize, Deserialize, Debug, PartialEq, Eq, Clone)] +pub struct ProductState { + pub enabled: bool, + pub version: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub error: Option, +} + +#[derive(Serialize, Debug)] +pub struct AppProductChange { + pub products: HashMap, +} + #[derive(Serialize, Debug)] pub struct AppIntegrationsChange { pub integrations: Vec, @@ -99,8 +177,22 @@ pub struct Log { pub is_crash: bool, } -#[derive(Serialize, Deserialize, Debug, PartialEq, Eq, Hash, Clone)] +#[derive( + Serialize, + Deserialize, + Debug, + PartialEq, + Eq, + Hash, + Clone, + Copy, + strum_macros::Display, + strum_macros::EnumIter, + // IntoStaticStr: required by the pre-9488 `ConvertToPyO3Enum` macro; see metrics.rs. + strum_macros::IntoStaticStr, +)] #[serde(rename_all = "UPPERCASE")] +#[strum(serialize_all = "UPPERCASE")] #[repr(C)] pub enum LogLevel { Error, diff --git a/libdd-telemetry/src/lib.rs b/libdd-telemetry/src/lib.rs index 09b4390f08..6dc9672f98 100644 --- a/libdd-telemetry/src/lib.rs +++ b/libdd-telemetry/src/lib.rs @@ -18,6 +18,9 @@ pub mod info; pub mod metrics; pub mod worker; +pub use libdd_common::tag::{parse_tags, Tag}; +pub use libdd_common::{parse_uri, Endpoint}; + pub fn build_host() -> data::Host { debug!("Building telemetry host information"); let hostname = info::os::real_hostname().unwrap_or_else(|_| String::from("unknown_hostname")); diff --git a/libdd-telemetry/src/metrics.rs b/libdd-telemetry/src/metrics.rs index c300a1ced9..f8b85bdc45 100644 --- a/libdd-telemetry/src/metrics.rs +++ b/libdd-telemetry/src/metrics.rs @@ -121,7 +121,7 @@ impl MetricBuckets { extra_tags, }; match context_key.1 { - metrics::MetricType::Count => self + metrics::MetricType::Count | metrics::MetricType::Rate => self .buckets .entry(bucket_key) .or_insert_with(|| MetricBucket { diff --git a/libdd-telemetry/src/worker/mod.rs b/libdd-telemetry/src/worker/mod.rs index 212c93730a..382260f902 100644 --- a/libdd-telemetry/src/worker/mod.rs +++ b/libdd-telemetry/src/worker/mod.rs @@ -7,7 +7,10 @@ pub mod store; use crate::{ config::Config, - data::{self, Application, Dependency, Endpoint, Host, Integration, Log, Payload, Telemetry}, + data::{ + self, Application, Dependency, Endpoint, Host, Integration, Log, Payload, ProductState, + Telemetry, + }, metrics::{ContextKey, MetricBuckets, MetricContexts}, }; @@ -85,7 +88,7 @@ macro_rules! telemetry_worker_log { $($arg)* ); if $worker.config.telemetry_debug_logging_enabled { - println!(concat!("{}: Telemetry worker DEBUG: ", $fmt_str), time_now(), $($arg)*); + eprintln!(concat!("{}: Telemetry worker DEBUG: ", $fmt_str), time_now(), $($arg)*); } } }; @@ -97,6 +100,7 @@ pub enum TelemetryActions { AddConfig(data::Configuration), AddDependency(Dependency), AddIntegration(Integration), + AddProductChange((String, ProductState)), AddLog((LogIdentifier, Log)), AddEndpoint(Endpoint), Lifecycle(LifecycleAction), @@ -127,15 +131,18 @@ pub struct LogIdentifier { #[derive(Debug)] struct TelemetryWorkerData { started: bool, - dependencies: store::Store, + dependencies: store::Store, configurations: store::Store, integrations: store::Store, endpoints: HashSet, + products: std::collections::HashMap, + products_pending: HashSet, logs: store::QueueHashMap, metric_contexts: MetricContexts, metric_buckets: MetricBuckets, host: Host, app: Application, + install_signature: Option, } /// `C` is the capability bundle. Leaf crates pin it to a concrete type @@ -238,6 +245,8 @@ impl Wor self.data.integrations.clear(); self.data.configurations.clear(); self.data.endpoints.clear(); + self.data.products.clear(); + self.data.products_pending.clear(); } async fn shutdown(&mut self) { @@ -326,8 +335,14 @@ impl Tel async fn recv_next_action(&mut self) -> TelemetryActions { let action = if let Some((deadline, deadline_action)) = self.deadlines.next_deadline() { let deadline_action = *deadline_action; - // If deadline passed, directly return associated action + // If deadline passed, service any already-queued mailbox action first, then + // return the associated action. + // This avoids pathological cases with a very short heartbeat, which would hang a + // synchronous flush()/stop() (whose FlushData/CollectStats never get processed) then. 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); }; @@ -408,6 +423,7 @@ impl Tel AddConfig(_) | AddDependency(_) | AddIntegration(_) + | AddProductChange(_) | AddEndpoint(_) | Lifecycle(ExtendedHeartbeat) => {} Lifecycle(Stop) => { @@ -450,10 +466,12 @@ impl Tel match action { Lifecycle(Start) => { if !self.data.started { - 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), + 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)] @@ -476,6 +494,10 @@ impl Tel } 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); @@ -538,6 +560,21 @@ impl Tel 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), + } + } // Only re-schedule self. Resetting `FlushData` here would replace its // existing deadline with `now + heartbeat_interval`, starving FlushData // when `extended_heartbeat_interval < heartbeat_interval` because each @@ -554,7 +591,9 @@ impl Tel self.data.metric_buckets.flush_aggregates(); let mut app_events = self.build_app_events_batch(); - app_events.push(data::Payload::AppClosing(())); + if self.config.emit_app_lifecycle { + app_events.push(data::Payload::AppClosing(())); + } let observability_events = self.build_observability_batch(); @@ -615,6 +654,22 @@ impl Tel }, )) } + 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 { @@ -714,6 +769,9 @@ impl Tel 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, } } @@ -723,6 +781,7 @@ impl Tel .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) { @@ -736,6 +795,11 @@ impl Tel 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 @@ -1108,15 +1172,39 @@ impl self.wait_for_shutdown() } - pub fn add_dependency(&self, name: String, version: Option) -> anyhow::Result<()> { + pub fn add_dependency( + &self, + name: String, + version: Option, + metadata: Option>, + ) -> 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, + ) -> anyhow::Result<()> { + self.sender.try_send(TelemetryActions::AddProductChange(( + product, + ProductState { + enabled, + version, + error: None, + }, + )))?; + Ok(()) + } + pub fn add_integration( &self, name: String, @@ -1203,7 +1291,7 @@ pub struct TelemetryWorkerBuilder { pub host: Host, pub application: Application, pub runtime_id: Option, - pub dependencies: store::Store, + pub dependencies: store::Store, pub integrations: store::Store, pub configurations: store::Store, pub endpoints: HashSet, @@ -1211,6 +1299,7 @@ pub struct TelemetryWorkerBuilder { pub rust_shared_lib_deps: bool, pub config: Config, pub flavor: TelemetryWorkerFlavor, + pub install_signature: Option, } impl TelemetryWorkerBuilder { @@ -1262,6 +1351,7 @@ impl TelemetryWorkerBuilder { rust_shared_lib_deps: false, config: Config::default(), flavor: TelemetryWorkerFlavor::default(), + install_signature: None, } } @@ -1294,11 +1384,14 @@ impl TelemetryWorkerBuilder { integrations: self.integrations, configurations: self.configurations, endpoints: self.endpoints, + 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, @@ -1701,7 +1794,7 @@ mod tests { // Populate every data field that reset() should clear. worker.data.dependencies.insert(Dependency { name: "dep".to_string(), - version: None, + ..Default::default() }); worker.data.integrations.insert(Integration { name: "integration".to_string(), @@ -1788,7 +1881,7 @@ mod tests { handle .try_send_msg(TelemetryActions::AddDependency(Dependency { name: "dep".to_string(), - version: None, + ..Default::default() })) .unwrap(); let (id, log) = make_log(1, "pre-fork log"); diff --git a/libdd-telemetry/src/worker/store.rs b/libdd-telemetry/src/worker/store.rs index e085c8a9b1..07b2bd53c3 100644 --- a/libdd-telemetry/src/worker/store.rs +++ b/libdd-telemetry/src/worker/store.rs @@ -65,6 +65,14 @@ mod queuehashmap { Some(&self.items[*idx - self.popped].1) } + pub fn get_mut(&mut self, k: &K) -> Option<&mut V> { + let hash = make_hash(&self.hash_builder, k); + let idx = *self + .table + .find(hash, |other| &self.items[other - self.popped].0 == k)?; + Some(&mut self.items[idx - self.popped].1) + } + pub fn get_idx(&self, idx: usize) -> Option<&(K, V)> { self.items.get(idx - self.popped) } @@ -143,6 +151,17 @@ mod queuehashmap { pub use queuehashmap::QueueHashMap; +/// Stable key to use as hash for stored items. +pub trait Keyed { + fn key(&self) -> K; +} + +impl Keyed for T { + fn key(&self) -> T { + self.clone() + } +} + #[derive(Debug, Default)] /// Stores telemetry data item, like dependencies and integrations /// @@ -150,16 +169,17 @@ pub use queuehashmap::QueueHashMap; /// * Tries to keep a list of items that it has seen (within max number of items) /// * Tries to keep a list of items that haven't been sent to datadog yet /// * Deduplicates items, to make sure we don't send the item twice -pub struct Store { +pub struct Store { // unflushed and set contain indices into unflushed: VecDeque, - items: QueueHashMap, + items: QueueHashMap, max_items: usize, } -impl Store +impl Store where - T: PartialEq + Eq + Hash, + T: Keyed, + K: PartialEq + Eq + Hash + Clone, { pub fn new(max_items: usize) -> Self { Self { @@ -170,13 +190,15 @@ where } pub fn insert(&mut self, item: T) { - if self.items.get(&item).is_some() { + let key = item.key(); + if let Some(existing) = self.items.get_mut(&key) { + *existing = item; return; } if self.items.len() == self.max_items { self.items.pop_front(); } - let (idx, _) = self.items.insert(item, ()); + let (idx, _) = self.items.insert(key, item); if self.unflushed.len() == self.max_items { self.unflushed.pop_front(); } @@ -205,7 +227,7 @@ where pub fn unflushed(&self) -> impl Iterator { self.unflushed .iter() - .flat_map(|i| Some(&self.items.get_idx(*i)?.0)) + .flat_map(|i| Some(&self.items.get_idx(*i)?.1)) } pub fn len_unflushed(&self) -> usize { @@ -223,9 +245,10 @@ where } } -impl Extend for Store +impl Extend for Store where - T: PartialEq + Eq + Hash, + T: Keyed, + K: PartialEq + Eq + Hash + Clone, { fn extend>(&mut self, iter: I) { for i in iter {