1pub mod http_client;
5mod scheduler;
6pub mod store;
7
8use crate::{
9 config::Config,
10 data::{self, Application, Dependency, Endpoint, Host, Integration, Log, Payload, Telemetry},
11 metrics::{ContextKey, MetricBuckets, MetricContexts},
12};
13
14use async_trait::async_trait;
15use libdd_common::{http_common, tag::Tag};
16use libdd_shared_runtime::Worker;
17
18use std::iter::Sum;
19use std::ops::Add;
20use std::{
21 collections::hash_map::DefaultHasher,
22 hash::{Hash, Hasher},
23 ops::ControlFlow,
24 sync::{
25 atomic::{AtomicU64, Ordering},
26 Arc, Condvar, Mutex,
27 },
28 time,
29};
30use std::{collections::HashSet, fmt::Debug, time::Duration};
31
32use crate::metrics::MetricBucketStats;
33use futures::{
34 channel::oneshot,
35 future::{self},
36};
37use http::{header, HeaderValue};
38use serde::{Deserialize, Serialize};
39use tokio::{
40 runtime::{self, Handle},
41 sync::mpsc,
42 task::JoinHandle,
43};
44use tokio_util::sync::CancellationToken;
45use tracing::debug;
46
47const CONTINUE: ControlFlow<()> = ControlFlow::Continue(());
48const BREAK: ControlFlow<()> = ControlFlow::Break(());
49
50fn time_now() -> f64 {
51 #[allow(clippy::unwrap_used)]
52 std::time::SystemTime::UNIX_EPOCH
53 .elapsed()
54 .unwrap_or_default()
55 .as_secs_f64()
56}
57
58macro_rules! telemetry_worker_log {
59 ($worker:expr , ERROR , $fmt_str:tt, $($arg:tt)*) => {
60 {
61 debug!(
62 worker.runtime_id = %$worker.runtime_id,
63 worker.debug_logging = $worker.config.telemetry_debug_logging_enabled,
64 $fmt_str,
65 $($arg)*
66 );
67 if $worker.config.telemetry_debug_logging_enabled {
68 eprintln!(concat!("{}: Telemetry worker ERROR: ", $fmt_str), time_now(), $($arg)*);
69 }
70 }
71 };
72 ($worker:expr , DEBUG , $fmt_str:tt, $($arg:tt)*) => {
73 {
74 debug!(
75 worker.runtime_id = %$worker.runtime_id,
76 worker.debug_logging = $worker.config.telemetry_debug_logging_enabled,
77 $fmt_str,
78 $($arg)*
79 );
80 if $worker.config.telemetry_debug_logging_enabled {
81 println!(concat!("{}: Telemetry worker DEBUG: ", $fmt_str), time_now(), $($arg)*);
82 }
83 }
84 };
85}
86
87#[derive(Debug, Serialize, Deserialize)]
88pub enum TelemetryActions {
89 AddPoint((f64, ContextKey, Vec<Tag>)),
90 AddConfig(data::Configuration),
91 AddDependency(Dependency),
92 AddIntegration(Integration),
93 AddLog((LogIdentifier, Log)),
94 AddEndpoint(Endpoint),
95 Lifecycle(LifecycleAction),
96 #[serde(skip)]
97 CollectStats(oneshot::Sender<TelemetryWorkerStats>),
98}
99
100#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
101pub enum LifecycleAction {
102 Start,
103 Stop,
104 FlushMetricAggr,
105 FlushData,
106 ExtendedHeartbeat,
107}
108
109#[derive(Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
114pub struct LogIdentifier {
115 pub identifier: u64,
117}
118
119#[derive(Debug)]
121struct TelemetryWorkerData {
122 started: bool,
123 dependencies: store::Store<Dependency>,
124 configurations: store::Store<data::Configuration>,
125 integrations: store::Store<data::Integration>,
126 endpoints: HashSet<data::Endpoint>,
127 logs: store::QueueHashMap<LogIdentifier, Log>,
128 metric_contexts: MetricContexts,
129 metric_buckets: MetricBuckets,
130 host: Host,
131 app: Application,
132}
133
134pub struct TelemetryWorker {
135 flavor: TelemetryWorkerFlavor,
136 config: Config,
137 mailbox: mpsc::Receiver<TelemetryActions>,
138 cancellation_token: CancellationToken,
139 seq_id: AtomicU64,
140 runtime_id: String,
141 client: Box<dyn http_client::HttpClient + Sync + Send>,
142 metrics_flush_interval: Duration,
143 deadlines: scheduler::Scheduler<LifecycleAction>,
144 data: TelemetryWorkerData,
145 next_action: Option<TelemetryActions>,
146 stopped: bool,
147}
148impl Debug for TelemetryWorker {
149 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
150 f.debug_struct("TelemetryWorker")
151 .field("flavor", &self.flavor)
152 .field("config", &self.config)
153 .field("mailbox", &self.mailbox)
154 .field("cancellation_token", &self.cancellation_token)
155 .field("seq_id", &self.seq_id)
156 .field("runtime_id", &self.runtime_id)
157 .field("metrics_flush_interval", &self.metrics_flush_interval)
158 .field("deadlines", &self.deadlines)
159 .field("data", &self.data)
160 .finish()
161 }
162}
163
164#[async_trait]
165impl Worker for TelemetryWorker {
166 async fn trigger(&mut self) {
167 if self.next_action.is_some() {
168 return;
170 }
171 if self.stopped {
172 debug!(
176 worker.runtime_id = %self.runtime_id,
177 "Telemetry worker mailbox closed; parking until shutdown"
178 );
179 std::future::pending::<()>().await;
180 }
181 let action = self.recv_next_action().await;
183 self.next_action = Some(action);
184 }
185
186 async fn run(&mut self) {
188 if let Some(action) = self.next_action.take() {
190 debug!(
191 worker.runtime_id = %self.runtime_id,
192 action = ?action,
193 "Received telemetry action"
194 );
195
196 let _action_result = match self.flavor {
199 TelemetryWorkerFlavor::Full => self.dispatch_action(action).await,
200 TelemetryWorkerFlavor::MetricsLogs => {
201 self.dispatch_metrics_logs_action(action).await
202 }
203 };
204 }
205 }
206
207 fn reset(&mut self) {
213 while self.mailbox.try_recv().is_ok() {}
215
216 self.next_action = None;
218
219 self.data.logs = store::QueueHashMap::default();
221 self.data.metric_buckets = MetricBuckets::default();
222 self.data.dependencies.clear();
223 self.data.integrations.clear();
224 self.data.configurations.clear();
225 self.data.endpoints.clear();
226 }
227
228 async fn shutdown(&mut self) {
229 let stop_action = TelemetryActions::Lifecycle(LifecycleAction::Stop);
230 let _action_result = match self.flavor {
231 TelemetryWorkerFlavor::Full => self.dispatch_action(stop_action).await,
232 TelemetryWorkerFlavor::MetricsLogs => {
233 self.dispatch_metrics_logs_action(stop_action).await
234 }
235 };
236 }
237}
238
239#[derive(Debug, Default, Serialize, Deserialize)]
240pub struct TelemetryWorkerStats {
241 pub dependencies_stored: u32,
242 pub dependencies_unflushed: u32,
243 pub configurations_stored: u32,
244 pub configurations_unflushed: u32,
245 pub integrations_stored: u32,
246 pub integrations_unflushed: u32,
247 pub logs: u32,
248 pub metric_contexts: u32,
249 pub metric_buckets: MetricBucketStats,
250}
251
252impl Add for TelemetryWorkerStats {
253 type Output = Self;
254
255 fn add(self, rhs: Self) -> Self::Output {
256 TelemetryWorkerStats {
257 dependencies_stored: self.dependencies_stored + rhs.dependencies_stored,
258 dependencies_unflushed: self.dependencies_unflushed + rhs.dependencies_unflushed,
259 configurations_stored: self.configurations_stored + rhs.configurations_stored,
260 configurations_unflushed: self.configurations_unflushed + rhs.configurations_unflushed,
261 integrations_stored: self.integrations_stored + rhs.integrations_stored,
262 integrations_unflushed: self.integrations_unflushed + rhs.integrations_unflushed,
263 logs: self.logs + rhs.logs,
264 metric_contexts: self.metric_contexts + rhs.metric_contexts,
265 metric_buckets: MetricBucketStats {
266 buckets: self.metric_buckets.buckets + rhs.metric_buckets.buckets,
267 series: self.metric_buckets.series + rhs.metric_buckets.series,
268 series_points: self.metric_buckets.series_points + rhs.metric_buckets.series_points,
269 distributions: self.metric_buckets.distributions + rhs.metric_buckets.distributions,
270 distributions_points: self.metric_buckets.distributions_points
271 + rhs.metric_buckets.distributions_points,
272 },
273 }
274 }
275}
276
277impl Sum for TelemetryWorkerStats {
278 fn sum<I: Iterator<Item = Self>>(iter: I) -> Self {
279 iter.fold(Self::default(), |a, b| a + b)
280 }
281}
282
283mod serialize {
284 use crate::data;
285 use http::HeaderValue;
286 #[allow(clippy::declare_interior_mutable_const)]
287 pub const CONTENT_TYPE_VALUE: HeaderValue = libdd_common::header::APPLICATION_JSON;
288 pub fn serialize(telemetry: &data::Telemetry) -> anyhow::Result<Vec<u8>> {
289 Ok(serde_json::to_vec(telemetry)?)
290 }
291}
292
293impl TelemetryWorker {
294 fn log_err(&self, err: &anyhow::Error) {
295 telemetry_worker_log!(self, ERROR, "{}", err);
296 }
297
298 async fn recv_next_action(&mut self) -> TelemetryActions {
299 let action = if let Some((deadline, deadline_action)) = self.deadlines.next_deadline() {
300 if deadline
302 .checked_duration_since(time::Instant::now())
303 .is_none()
304 {
305 return TelemetryActions::Lifecycle(*deadline_action);
306 };
307
308 match tokio::time::timeout_at(deadline.into(), self.mailbox.recv()).await {
310 Ok(mailbox_action) => mailbox_action,
311 Err(_) => Some(TelemetryActions::Lifecycle(*deadline_action)),
312 }
313 } else {
314 self.mailbox.recv().await
315 };
316
317 action.unwrap_or_else(|| {
319 self.config.restartable = false;
321 self.stopped = true;
322 TelemetryActions::Lifecycle(LifecycleAction::Stop)
323 })
324 }
325
326 async fn dispatch_metrics_logs_action(&mut self, action: TelemetryActions) -> ControlFlow<()> {
327 telemetry_worker_log!(self, DEBUG, "Handling metric action {:?}", action);
328 use LifecycleAction::*;
329 use TelemetryActions::*;
330 match action {
331 Lifecycle(Start) => {
332 if !self.data.started {
333 #[allow(clippy::unwrap_used)]
334 self.deadlines
335 .schedule_event(LifecycleAction::FlushMetricAggr)
336 .unwrap();
337
338 #[allow(clippy::unwrap_used)]
339 self.deadlines
340 .schedule_event(LifecycleAction::FlushData)
341 .unwrap();
342 self.data.started = true;
343 }
344 }
345 AddLog((identifier, log)) => {
346 let (l, new) = self.data.logs.get_mut_or_insert(identifier, log);
347 if !new {
348 l.count += 1;
349 }
350 }
351 AddPoint((point, key, extra_tags)) => {
352 self.data.metric_buckets.add_point(key, point, extra_tags)
353 }
354 Lifecycle(FlushMetricAggr) => {
355 self.data.metric_buckets.flush_aggregates();
356
357 #[allow(clippy::unwrap_used)]
358 self.deadlines
359 .schedule_event(LifecycleAction::FlushMetricAggr)
360 .unwrap();
361 }
362 Lifecycle(FlushData) => {
363 if !(self.data.started || self.config.restartable) {
364 return CONTINUE;
365 }
366
367 #[allow(clippy::unwrap_used)]
368 self.deadlines
369 .schedule_event(LifecycleAction::FlushData)
370 .unwrap();
371
372 let batch = self.build_observability_batch();
373 if !batch.is_empty() {
374 let payload = data::Payload::MessageBatch(batch);
375 match self.send_payload(&payload).await {
376 Ok(()) => self.payload_sent_success(&payload),
377 Err(e) => self.log_err(&e),
378 }
379 }
380 }
381 AddConfig(_)
382 | AddDependency(_)
383 | AddIntegration(_)
384 | AddEndpoint(_)
385 | Lifecycle(ExtendedHeartbeat) => {}
386 Lifecycle(Stop) => {
387 if !self.data.started {
388 return BREAK;
389 }
390 self.data.metric_buckets.flush_aggregates();
391
392 let batch = self.build_observability_batch();
393 if !batch.is_empty() {
394 let payload = data::Payload::MessageBatch(batch);
395 match self.send_payload(&payload).await {
396 Ok(()) => {
397 if self.config.restartable {
398 self.payload_sent_success(&payload)
399 }
400 }
401 Err(e) => self.log_err(&e),
402 }
403 }
404
405 self.data.started = false;
406 if !self.config.restartable {
407 self.deadlines.clear_pending();
408 }
409 return BREAK;
410 }
411 CollectStats(stats_sender) => {
412 stats_sender.send(self.stats()).ok();
413 }
414 };
415 CONTINUE
416 }
417
418 async fn dispatch_action(&mut self, action: TelemetryActions) -> ControlFlow<()> {
419 telemetry_worker_log!(self, DEBUG, "Handling action {:?}", action);
420
421 use LifecycleAction::*;
422 use TelemetryActions::*;
423 match action {
424 Lifecycle(Start) => {
425 if !self.data.started {
426 let app_started = data::Payload::AppStarted(self.build_app_started());
427 match self.send_payload(&app_started).await {
428 Ok(()) => self.payload_sent_success(&app_started),
429 Err(err) => self.log_err(&err),
430 }
431
432 #[allow(clippy::unwrap_used)]
433 self.deadlines
434 .schedule_event(LifecycleAction::FlushMetricAggr)
435 .unwrap();
436
437 #[allow(clippy::unwrap_used)]
438 self.deadlines
440 .schedule_event(LifecycleAction::FlushData)
441 .unwrap();
442
443 #[allow(clippy::unwrap_used)]
444 self.deadlines
445 .schedule_event(LifecycleAction::ExtendedHeartbeat)
446 .unwrap();
447 self.data.started = true;
448 }
449 }
450 AddDependency(dep) => self.data.dependencies.insert(dep),
451 AddIntegration(integration) => self.data.integrations.insert(integration),
452 AddConfig(cfg) => self.data.configurations.insert(cfg),
453 AddEndpoint(endpoint) => {
454 self.data.endpoints.insert(endpoint);
455 }
456 AddLog((identifier, log)) => {
457 let (l, new) = self.data.logs.get_mut_or_insert(identifier, log);
458 if !new {
459 l.count += 1;
460 }
461 }
462 AddPoint((point, key, extra_tags)) => {
463 self.data.metric_buckets.add_point(key, point, extra_tags)
464 }
465 Lifecycle(FlushMetricAggr) => {
466 self.data.metric_buckets.flush_aggregates();
467
468 #[allow(clippy::unwrap_used)]
469 self.deadlines
470 .schedule_event(LifecycleAction::FlushMetricAggr)
471 .unwrap();
472 }
473 Lifecycle(FlushData) => {
474 if !(self.data.started || self.config.restartable) {
475 return CONTINUE;
476 }
477
478 #[allow(clippy::unwrap_used)]
479 self.deadlines
480 .schedule_event(LifecycleAction::FlushData)
481 .unwrap();
482
483 let mut batch = self.build_app_events_batch();
484 let payload = if batch.is_empty() {
485 data::Payload::AppHeartbeat(())
486 } else {
487 batch.push(data::Payload::AppHeartbeat(()));
488 data::Payload::MessageBatch(batch)
489 };
490 match self.send_payload(&payload).await {
491 Ok(()) => self.payload_sent_success(&payload),
492 Err(err) => self.log_err(&err),
493 }
494
495 let batch = self.build_observability_batch();
496 if !batch.is_empty() {
497 let payload = data::Payload::MessageBatch(batch);
498 match self.send_payload(&payload).await {
499 Ok(()) => self.payload_sent_success(&payload),
500 Err(err) => self.log_err(&err),
501 }
502 }
503 }
504 Lifecycle(ExtendedHeartbeat) => {
505 self.data.dependencies.unflush_stored();
506 self.data.integrations.unflush_stored();
507 self.data.configurations.unflush_stored();
508
509 let extended_hb = data::Payload::AppExtendedHeartbeat(self.build_app_started());
510 match self.send_payload(&extended_hb).await {
511 Ok(()) => self.payload_sent_success(&extended_hb),
512 Err(err) => self.log_err(&err),
513 }
514 #[allow(clippy::unwrap_used)]
519 self.deadlines
520 .schedule_event(LifecycleAction::ExtendedHeartbeat)
521 .unwrap();
522 }
523 Lifecycle(Stop) => {
524 if !self.data.started {
525 return BREAK;
526 }
527 self.data.metric_buckets.flush_aggregates();
528
529 let mut app_events = self.build_app_events_batch();
530 app_events.push(data::Payload::AppClosing(()));
531
532 let observability_events = self.build_observability_batch();
533
534 let mut payloads = vec![data::Payload::MessageBatch(app_events)];
535 if !observability_events.is_empty() {
536 payloads.push(data::Payload::MessageBatch(observability_events));
537 }
538
539 let self_arc = Arc::new(tokio::sync::RwLock::new(&mut *self));
540 let futures = payloads.into_iter().map(|payload| {
541 let self_arc = self_arc.clone();
542 async move {
543 let res = {
547 let self_rguard = self_arc.read().await;
548 self_rguard.send_payload(&payload).await
549 };
550 match res {
551 Ok(()) => self_arc.write().await.payload_sent_success(&payload),
552 Err(err) => self_arc.read().await.log_err(&err),
553 }
554 }
555 });
556 future::join_all(futures).await;
557
558 self.data.started = false;
559 if !self.config.restartable {
560 self.deadlines.clear_pending();
561 }
562
563 return BREAK;
564 }
565 CollectStats(stats_sender) => {
566 stats_sender.send(self.stats()).ok();
567 }
568 }
569
570 CONTINUE
571 }
572
573 fn build_app_events_batch(&mut self) -> Vec<Payload> {
575 let mut payloads = Vec::new();
576
577 if self.data.dependencies.flush_not_empty() {
578 payloads.push(data::Payload::AppDependenciesLoaded(
579 data::AppDependenciesLoaded {
580 dependencies: self.data.dependencies.unflushed().cloned().collect(),
581 },
582 ))
583 }
584 if self.data.integrations.flush_not_empty() {
585 payloads.push(data::Payload::AppIntegrationsChange(
586 data::AppIntegrationsChange {
587 integrations: self.data.integrations.unflushed().cloned().collect(),
588 },
589 ))
590 }
591 if self.data.configurations.flush_not_empty() {
592 payloads.push(data::Payload::AppClientConfigurationChange(
593 data::AppClientConfigurationChange {
594 configuration: self.data.configurations.unflushed().cloned().collect(),
595 },
596 ))
597 }
598 if !self.data.endpoints.is_empty() {
599 payloads.push(data::Payload::AppEndpoints(data::AppEndpoints {
600 is_first: true,
601 endpoints: self
602 .data
603 .endpoints
604 .iter()
605 .map(|e| e.to_json_value().unwrap_or_default())
606 .filter(|e| e.is_object())
607 .collect(),
608 }));
609 }
610 payloads
611 }
612
613 fn build_observability_batch(&mut self) -> Vec<Payload> {
615 let mut payloads = Vec::new();
616
617 let logs = self.build_logs();
618 if !logs.logs.is_empty() {
619 payloads.push(data::Payload::Logs(logs));
620 }
621 let metrics = self.build_metrics_series();
622 if !metrics.series.is_empty() {
623 payloads.push(data::Payload::GenerateMetrics(metrics))
624 }
625 let distributions = self.build_metrics_distributions();
626 if !distributions.series.is_empty() {
627 payloads.push(data::Payload::Sketches(distributions))
628 }
629 payloads
630 }
631
632 fn build_metrics_distributions(&mut self) -> data::Distributions {
633 let mut series = Vec::new();
634 let context_guard = self.data.metric_contexts.lock();
635 for (context_key, extra_tags, points) in self.data.metric_buckets.flush_distributions() {
636 let Some(context) = context_guard.read(context_key) else {
637 telemetry_worker_log!(self, ERROR, "Context not found for key {:?}", context_key);
638 continue;
639 };
640 let mut tags = extra_tags;
641 tags.extend(context.tags.iter().cloned());
642 series.push(data::metrics::Distribution {
643 namespace: context.namespace,
644 metric: context.name.clone(),
645 tags,
646 sketch: data::metrics::SerializedSketch::B64 {
647 sketch_b64: base64::Engine::encode(
648 &base64::engine::general_purpose::STANDARD,
649 points.encode_to_vec(),
650 ),
651 },
652 common: context.common,
653 _type: context.metric_type,
654 interval: self.metrics_flush_interval.as_secs(),
655 });
656 }
657 data::Distributions { series }
658 }
659
660 fn build_metrics_series(&mut self) -> data::GenerateMetrics {
661 let mut series = Vec::new();
662 let context_guard = self.data.metric_contexts.lock();
663 for (context_key, extra_tags, points) in self.data.metric_buckets.flush_series() {
664 let Some(context) = context_guard.read(context_key) else {
665 telemetry_worker_log!(self, ERROR, "Context not found for key {:?}", context_key);
666 continue;
667 };
668
669 let mut tags = extra_tags;
670 tags.extend(context.tags.iter().cloned());
671 series.push(data::metrics::Serie {
672 namespace: context.namespace,
673 metric: context.name.clone(),
674 tags,
675 points,
676 common: context.common,
677 _type: context.metric_type,
678 interval: self.metrics_flush_interval.as_secs(),
679 });
680 }
681
682 data::GenerateMetrics { series }
683 }
684
685 fn build_app_started(&mut self) -> data::AppStarted {
686 data::AppStarted {
687 configuration: self.data.configurations.unflushed().cloned().collect(),
688 dependencies: self.data.dependencies.unflushed().cloned().collect(),
689 integrations: self.data.integrations.unflushed().cloned().collect(),
690 }
691 }
692
693 fn app_started_sent_success(&mut self, p: &data::AppStarted) {
694 self.data
695 .configurations
696 .removed_flushed(p.configuration.len());
697 self.data.dependencies.removed_flushed(p.dependencies.len());
698 self.data.integrations.removed_flushed(p.integrations.len());
699 }
700
701 fn payload_sent_success(&mut self, payload: &data::Payload) {
702 use data::Payload::*;
703 match payload {
704 AppStarted(p) => self.app_started_sent_success(p),
705 AppExtendedHeartbeat(p) => self.app_started_sent_success(p),
706 AppDependenciesLoaded(p) => {
707 self.data.dependencies.removed_flushed(p.dependencies.len())
708 }
709 AppIntegrationsChange(p) => {
710 self.data.integrations.removed_flushed(p.integrations.len())
711 }
712 AppClientConfigurationChange(p) => self
713 .data
714 .configurations
715 .removed_flushed(p.configuration.len()),
716 AppEndpoints(_) => self.data.endpoints.clear(),
717 MessageBatch(batch) => {
718 for p in batch {
719 self.payload_sent_success(p);
720 }
721 }
722 Logs(p) => {
723 for _ in &p.logs {
724 self.data.logs.pop_front();
725 }
726 }
727 AppHeartbeat(()) | AppClosing(()) => {}
728 GenerateMetrics(_) | Sketches(_) => {}
729 }
730 }
731
732 fn build_logs(&self) -> data::Logs {
733 let logs = self.data.logs.iter().map(|(_, l)| l.clone()).collect();
735 data::Logs { logs }
736 }
737
738 fn next_seq_id(&self) -> u64 {
739 self.seq_id.fetch_add(1, Ordering::Release)
740 }
741
742 async fn send_payload(&self, payload: &data::Payload) -> anyhow::Result<()> {
743 debug!(
744 worker.runtime_id = %self.runtime_id,
745 payload.type = payload.request_type(),
746 seq_id = self.seq_id.load(Ordering::Acquire),
747 "Sending telemetry payload"
748 );
749 let req = self.build_request(payload)?;
750 let result = self.send_request(req).await;
751 match &result {
752 Ok(resp) => debug!(
753 worker.runtime_id = %self.runtime_id,
754 payload.type = payload.request_type(),
755 response.status = resp.status().as_u16(),
756 "Successfully sent telemetry payload"
757 ),
758 Err(e) => debug!(
759 worker.runtime_id = %self.runtime_id,
760 payload.type = payload.request_type(),
761 error = ?e,
762 "Failed to send telemetry payload"
763 ),
764 }
765 Ok(())
766 }
767
768 fn build_request(&self, payload: &data::Payload) -> anyhow::Result<http_common::HttpRequest> {
769 let seq_id = self.next_seq_id();
770 let tel = Telemetry {
771 api_version: data::ApiVersion::V2,
772 tracer_time: time::SystemTime::UNIX_EPOCH
773 .elapsed()
774 .map_or(0, |d| d.as_secs()),
775 runtime_id: &self.runtime_id,
776 seq_id,
777 host: &self.data.host,
778 origin: None,
779 application: &self.data.app,
780 payload,
781 };
782
783 telemetry_worker_log!(self, DEBUG, "Prepared payload: {:?}", tel);
784
785 let req = http_client::request_builder(&self.config)?
786 .method(http::Method::POST)
787 .header(header::CONTENT_TYPE, serialize::CONTENT_TYPE_VALUE)
788 .header(
789 http_client::header::REQUEST_TYPE,
790 HeaderValue::from_static(payload.request_type()),
791 )
792 .header(
793 http_client::header::API_VERSION,
794 HeaderValue::from_static(data::ApiVersion::V2.to_str()),
795 )
796 .header(
797 http_client::header::LIBRARY_LANGUAGE,
798 tel.application.language_name.clone(),
799 )
800 .header(
801 http_client::header::LIBRARY_VERSION,
802 tel.application.tracer_version.clone(),
803 );
804 let req = http_client::add_instrumentation_session_headers(
805 req,
806 self.config.session_id.as_deref(),
807 self.config.parent_session_id.as_deref(),
808 self.config.root_session_id.as_deref(),
809 );
810
811 let body = http_common::Body::from(serialize::serialize(&tel)?);
812 Ok(req.body(body)?)
813 }
814
815 async fn send_request(
816 &self,
817 req: http_common::HttpRequest,
818 ) -> Result<http_common::HttpResponse, http_common::Error> {
819 let timeout_ms = if let Some(endpoint) = self.config.endpoint.as_ref() {
820 endpoint.timeout_ms
821 } else {
822 libdd_common::Endpoint::DEFAULT_TIMEOUT
823 };
824
825 debug!(
826 worker.runtime_id = %self.runtime_id,
827 http.timeout_ms = timeout_ms,
828 "Sending HTTP request"
829 );
830
831 tokio::select! {
832 _ = self.cancellation_token.cancelled() => {
833 debug!(
834 worker.runtime_id = %self.runtime_id,
835 "Telemetry request cancelled"
836 );
837 Err(http_common::Error::Other(anyhow::anyhow!("Request cancelled")))
838 },
839 _ = tokio::time::sleep(time::Duration::from_millis(timeout_ms)) => {
840 debug!(
841 worker.runtime_id = %self.runtime_id,
842 http.timeout_ms = timeout_ms,
843 "Telemetry request timed out"
844 );
845 Err(http_common::Error::Other(anyhow::anyhow!("Request timed out")))
846 },
847 r = self.client.request(req) => {
848 match r {
849 Ok(resp) => {
850 Ok(resp)
851 }
852 Err(e) => {
853 Err(e)
854 },
855 }
856 }
857 }
858 }
859
860 fn stats(&self) -> TelemetryWorkerStats {
861 TelemetryWorkerStats {
862 dependencies_stored: self.data.dependencies.len_stored() as u32,
863 dependencies_unflushed: self.data.dependencies.len_unflushed() as u32,
864 configurations_stored: self.data.configurations.len_stored() as u32,
865 configurations_unflushed: self.data.configurations.len_unflushed() as u32,
866 integrations_stored: self.data.integrations.len_stored() as u32,
867 integrations_unflushed: self.data.integrations.len_unflushed() as u32,
868 logs: self.data.logs.len() as u32,
869 metric_contexts: self.data.metric_contexts.lock().len() as u32,
870 metric_buckets: self.data.metric_buckets.stats(),
871 }
872 }
873
874 async fn run_loop(mut self) {
877 debug!(
878 worker.flavor = ?self.flavor,
879 worker.runtime_id = %self.runtime_id,
880 "Starting telemetry worker"
881 );
882
883 loop {
884 if self.cancellation_token.is_cancelled() {
885 debug!(
886 worker.runtime_id = %self.runtime_id,
887 "Telemetry worker cancelled, shutting down"
888 );
889 return;
890 }
891
892 let action = self.recv_next_action().await;
893 debug!(
894 worker.runtime_id = %self.runtime_id,
895 action = ?action,
896 "Received telemetry action"
897 );
898
899 let action_result = match self.flavor {
900 TelemetryWorkerFlavor::Full => self.dispatch_action(action).await,
901 TelemetryWorkerFlavor::MetricsLogs => {
902 self.dispatch_metrics_logs_action(action).await
903 }
904 };
905
906 match action_result {
907 ControlFlow::Continue(()) => {}
908 ControlFlow::Break(()) => {
909 debug!(
910 worker.runtime_id = %self.runtime_id,
911 worker.restartable = self.config.restartable,
912 "Telemetry worker received break signal"
913 );
914 if !self.config.restartable {
915 break;
916 }
917 }
918 };
919 }
920
921 debug!(
922 worker.runtime_id = %self.runtime_id,
923 "Telemetry worker stopped"
924 );
925 }
926}
927
928#[derive(Debug)]
929struct InnerTelemetryShutdown {
930 is_shutdown: Mutex<bool>,
931 condvar: Condvar,
932}
933
934impl InnerTelemetryShutdown {
935 fn wait_for_shutdown(&self) {
936 drop(
937 #[allow(clippy::unwrap_used)]
938 self.condvar
939 .wait_while(self.is_shutdown.lock().unwrap(), |is_shutdown| {
940 !*is_shutdown
941 })
942 .unwrap(),
943 )
944 }
945
946 #[allow(clippy::unwrap_used)]
947 fn shutdown_finished(&self) {
948 *self.is_shutdown.lock().unwrap() = true;
949 self.condvar.notify_all();
950 }
951}
952
953#[derive(Clone, Debug)]
954pub struct TelemetryWorkerHandle {
962 sender: mpsc::Sender<TelemetryActions>,
963 shutdown: Arc<InnerTelemetryShutdown>,
964 cancellation_token: CancellationToken,
965 runtime: Option<runtime::Handle>,
968
969 contexts: MetricContexts,
970}
971
972impl TelemetryWorkerHandle {
973 pub fn register_metric_context(
974 &self,
975 name: String,
976 tags: Vec<Tag>,
977 metric_type: data::metrics::MetricType,
978 common: bool,
979 namespace: data::metrics::MetricNamespace,
980 ) -> ContextKey {
981 self.contexts
982 .register_metric_context(name, tags, metric_type, common, namespace)
983 }
984
985 pub fn try_send_msg(&self, msg: TelemetryActions) -> anyhow::Result<()> {
986 Ok(self.sender.try_send(msg)?)
987 }
988
989 pub async fn send_msg(&self, msg: TelemetryActions) -> anyhow::Result<()> {
990 Ok(self.sender.send(msg).await?)
991 }
992
993 pub async fn send_msgs<T>(&self, msgs: T) -> anyhow::Result<()>
994 where
995 T: IntoIterator<Item = TelemetryActions>,
996 {
997 for msg in msgs {
998 self.sender.send(msg).await?;
999 }
1000
1001 Ok(())
1002 }
1003
1004 pub async fn send_msg_timeout(
1005 &self,
1006 msg: TelemetryActions,
1007 timeout: time::Duration,
1008 ) -> anyhow::Result<()> {
1009 Ok(self.sender.send_timeout(msg, timeout).await?)
1010 }
1011
1012 pub fn send_start(&self) -> anyhow::Result<()> {
1013 Ok(self
1014 .sender
1015 .try_send(TelemetryActions::Lifecycle(LifecycleAction::Start))?)
1016 }
1017
1018 pub fn send_stop(&self) -> anyhow::Result<()> {
1019 Ok(self
1020 .sender
1021 .try_send(TelemetryActions::Lifecycle(LifecycleAction::Stop))?)
1022 }
1023
1024 fn cancel_requests_with_deadline(&self, deadline: time::Instant) {
1025 let Some(runtime) = &self.runtime else {
1026 tracing::error!("Cannot schedule cancellation deadline: no runtime handle available");
1027 return;
1028 };
1029 let token = self.cancellation_token.clone();
1030 let f = async move {
1031 tokio::time::sleep_until(deadline.into()).await;
1032 token.cancel()
1033 };
1034 runtime.spawn(f);
1035 }
1036
1037 pub fn wait_for_shutdown_deadline(&self, deadline: time::Instant) {
1038 self.cancel_requests_with_deadline(deadline);
1039 self.wait_for_shutdown()
1040 }
1041
1042 pub fn add_dependency(&self, name: String, version: Option<String>) -> anyhow::Result<()> {
1043 self.sender
1044 .try_send(TelemetryActions::AddDependency(Dependency {
1045 name,
1046 version,
1047 }))?;
1048 Ok(())
1049 }
1050
1051 pub fn add_integration(
1052 &self,
1053 name: String,
1054 enabled: bool,
1055 version: Option<String>,
1056 compatible: Option<bool>,
1057 auto_enabled: Option<bool>,
1058 ) -> anyhow::Result<()> {
1059 self.sender
1060 .try_send(TelemetryActions::AddIntegration(Integration {
1061 name,
1062 version,
1063 compatible,
1064 enabled,
1065 auto_enabled,
1066 }))?;
1067 Ok(())
1068 }
1069
1070 pub fn add_log<T: Hash>(
1071 &self,
1072 identifier: T,
1073 message: String,
1074 level: data::LogLevel,
1075 stack_trace: Option<String>,
1076 ) -> anyhow::Result<()> {
1077 let mut hasher = DefaultHasher::new();
1078 identifier.hash(&mut hasher);
1079 self.sender.try_send(TelemetryActions::AddLog((
1080 LogIdentifier {
1081 identifier: hasher.finish(),
1082 },
1083 data::Log {
1084 message,
1085 level,
1086 stack_trace,
1087 count: 1,
1088 tags: String::new(),
1089 is_sensitive: false,
1090 is_crash: false,
1091 },
1092 )))?;
1093 Ok(())
1094 }
1095
1096 pub fn add_point(
1097 &self,
1098 value: f64,
1099 context: &ContextKey,
1100 extra_tags: Vec<Tag>,
1101 ) -> anyhow::Result<()> {
1102 self.sender
1103 .try_send(TelemetryActions::AddPoint((value, *context, extra_tags)))?;
1104 Ok(())
1105 }
1106
1107 pub fn wait_for_shutdown(&self) {
1108 self.shutdown.wait_for_shutdown();
1109 }
1110
1111 pub fn stats(&self) -> anyhow::Result<oneshot::Receiver<TelemetryWorkerStats>> {
1112 let (sender, receiver) = oneshot::channel();
1113 self.sender
1114 .try_send(TelemetryActions::CollectStats(sender))?;
1115 Ok(receiver)
1116 }
1117}
1118
1119pub const MAX_ITEMS: usize = 5000;
1121
1122#[derive(Debug, Default, Clone, Copy)]
1123pub enum TelemetryWorkerFlavor {
1124 #[default]
1127 Full,
1128 MetricsLogs,
1130}
1131
1132pub struct TelemetryWorkerBuilder {
1133 pub host: Host,
1134 pub application: Application,
1135 pub runtime_id: Option<String>,
1136 pub dependencies: store::Store<data::Dependency>,
1137 pub integrations: store::Store<data::Integration>,
1138 pub configurations: store::Store<data::Configuration>,
1139 pub endpoints: HashSet<data::Endpoint>,
1140 pub native_deps: bool,
1141 pub rust_shared_lib_deps: bool,
1142 pub config: Config,
1143 pub flavor: TelemetryWorkerFlavor,
1144}
1145
1146impl TelemetryWorkerBuilder {
1147 pub fn new_fetch_host(
1149 service_name: String,
1150 language_name: String,
1151 language_version: String,
1152 tracer_version: String,
1153 ) -> Self {
1154 Self {
1155 host: crate::build_host(),
1156 ..Self::new(
1157 String::new(),
1158 service_name,
1159 language_name,
1160 language_version,
1161 tracer_version,
1162 )
1163 }
1164 }
1165
1166 pub fn new(
1168 hostname: String,
1169 service_name: String,
1170 language_name: String,
1171 language_version: String,
1172 tracer_version: String,
1173 ) -> Self {
1174 Self {
1175 host: Host {
1176 hostname,
1177 ..Default::default()
1178 },
1179 application: Application {
1180 service_name,
1181 language_name,
1182 language_version,
1183 tracer_version,
1184 ..Default::default()
1185 },
1186 runtime_id: None,
1187 dependencies: store::Store::new(MAX_ITEMS),
1188 integrations: store::Store::new(MAX_ITEMS),
1189 configurations: store::Store::new(MAX_ITEMS),
1190 endpoints: HashSet::new(),
1191 native_deps: true,
1192 rust_shared_lib_deps: false,
1193 config: Config::default(),
1194 flavor: TelemetryWorkerFlavor::default(),
1195 }
1196 }
1197
1198 pub fn build_worker(
1204 self,
1205 tokio_runtime: Option<Handle>,
1206 ) -> (TelemetryWorkerHandle, TelemetryWorker) {
1207 let (tx, mailbox) = mpsc::channel(5000);
1208 let shutdown = Arc::new(InnerTelemetryShutdown {
1209 is_shutdown: Mutex::new(false),
1210 condvar: Condvar::new(),
1211 });
1212 let contexts = MetricContexts::default();
1213 let token = CancellationToken::new();
1214 let config = self.config;
1215 let telemetry_heartbeat_interval = config.telemetry_heartbeat_interval;
1216 let telemetry_extended_heartbeat_interval = config.telemetry_extended_heartbeat_interval;
1217 let client = http_client::from_config(&config);
1218
1219 let metrics_flush_interval =
1220 telemetry_heartbeat_interval.min(MetricBuckets::METRICS_FLUSH_INTERVAL);
1221
1222 #[allow(clippy::unwrap_used)]
1223 let worker = TelemetryWorker {
1224 flavor: self.flavor,
1225 data: TelemetryWorkerData {
1226 started: false,
1227 dependencies: self.dependencies,
1228 integrations: self.integrations,
1229 configurations: self.configurations,
1230 endpoints: self.endpoints,
1231 logs: store::QueueHashMap::default(),
1232 metric_contexts: contexts.clone(),
1233 metric_buckets: MetricBuckets::default(),
1234 host: self.host,
1235 app: self.application,
1236 },
1237 config,
1238 mailbox,
1239 seq_id: AtomicU64::new(1),
1240 runtime_id: self
1241 .runtime_id
1242 .unwrap_or_else(|| uuid::Uuid::new_v4().to_string()),
1243 client,
1244 metrics_flush_interval,
1245 deadlines: scheduler::Scheduler::new(vec![
1246 (metrics_flush_interval, LifecycleAction::FlushMetricAggr),
1247 (telemetry_heartbeat_interval, LifecycleAction::FlushData),
1248 (
1249 telemetry_extended_heartbeat_interval,
1250 LifecycleAction::ExtendedHeartbeat,
1251 ),
1252 ]),
1253 cancellation_token: token.clone(),
1254 next_action: None,
1255 stopped: false,
1256 };
1257
1258 (
1259 TelemetryWorkerHandle {
1260 sender: tx,
1261 shutdown,
1262 cancellation_token: token,
1263 runtime: tokio_runtime,
1264
1265 contexts,
1266 },
1267 worker,
1268 )
1269 }
1270
1271 pub fn spawn(self) -> (TelemetryWorkerHandle, JoinHandle<()>) {
1274 let tokio_runtime = tokio::runtime::Handle::current();
1275
1276 let (worker_handle, worker) = self.build_worker(Some(tokio_runtime.clone()));
1277
1278 let join_handle = tokio_runtime.spawn(async move { worker.run_loop().await });
1279
1280 (worker_handle, join_handle)
1281 }
1282
1283 pub fn run(self) -> anyhow::Result<TelemetryWorkerHandle> {
1285 let runtime = tokio::runtime::Builder::new_current_thread()
1286 .enable_all()
1287 .build()?;
1288 let (handle, worker) = self.build_worker(Some(runtime.handle().clone()));
1289 let notify_shutdown = handle.shutdown.clone();
1290 std::thread::spawn(move || {
1291 runtime.block_on(worker.run_loop());
1292 runtime.shutdown_background();
1293 notify_shutdown.shutdown_finished();
1294 });
1295
1296 Ok(handle)
1297 }
1298}
1299
1300#[cfg(test)]
1301mod tests {
1302 use crate::config::TelemetryEndpoint;
1303 use crate::data::Payload;
1304 use crate::worker::http_client::header::{
1305 DD_PARENT_SESSION_ID, DD_ROOT_SESSION_ID, DD_SESSION_ID,
1306 };
1307 use crate::worker::{
1308 LifecycleAction, TelemetryActions, TelemetryWorker, TelemetryWorkerBuilder,
1309 TelemetryWorkerFlavor, TelemetryWorkerHandle,
1310 };
1311 use libdd_common::http_common;
1312 use tokio::runtime::Runtime;
1313
1314 fn is_send<T: Send>(_: T) {}
1315 fn is_sync<T: Sync>(_: T) {}
1316
1317 #[test]
1318 fn test_handle_sync_send() {
1319 #[allow(clippy::redundant_closure)]
1320 let _ = |h: TelemetryWorkerHandle| is_send(h);
1321 #[allow(clippy::redundant_closure)]
1322 let _ = |h: TelemetryWorkerHandle| is_sync(h);
1323 }
1324
1325 fn test_worker(
1326 session_id: Option<String>,
1327 root_session_id: Option<String>,
1328 parent_session_id: Option<String>,
1329 ) -> TelemetryWorker {
1330 let mut b = TelemetryWorkerBuilder::new(
1331 "h".into(),
1332 "svc".into(),
1333 "lang".into(),
1334 "1".into(),
1335 "tv".into(),
1336 );
1337 b.config
1338 .set_endpoint(TelemetryEndpoint {
1339 url: Some("http://127.0.0.1:1".to_owned()),
1340 ..Default::default()
1341 })
1342 .unwrap();
1343 b.runtime_id = Some("rid".into());
1344 b.config.session_id = session_id;
1345 b.config.parent_session_id = parent_session_id;
1346 b.config.root_session_id = root_session_id;
1347 let rt = Runtime::new().unwrap();
1348 b.build_worker(Some(rt.handle().clone())).1
1349 }
1350
1351 #[cfg_attr(miri, ignore)] #[test]
1353 fn telemetry_http_includes_dd_session_id() {
1354 let req = test_worker(Some("sess".into()), None, None)
1355 .build_request(&Payload::AppHeartbeat(()))
1356 .unwrap();
1357 assert_eq!(
1358 req.headers().get(DD_SESSION_ID).unwrap().to_str().unwrap(),
1359 "sess"
1360 );
1361 assert!(req.headers().get(DD_ROOT_SESSION_ID).is_none());
1362 assert!(req.headers().get(DD_PARENT_SESSION_ID).is_none());
1363 }
1364
1365 #[cfg_attr(miri, ignore)] #[test]
1367 fn telemetry_http_omits_root_session_id_when_same_as_session_id() {
1368 let req = test_worker(
1369 Some("sess-id".into()),
1370 Some("sess-id".into()),
1371 Some("parent".into()),
1372 )
1373 .build_request(&Payload::AppHeartbeat(()))
1374 .unwrap();
1375 assert_eq!(
1376 req.headers().get(DD_SESSION_ID).unwrap().to_str().unwrap(),
1377 "sess-id"
1378 );
1379 assert!(req.headers().get(DD_ROOT_SESSION_ID).is_none());
1380 assert_eq!(
1381 req.headers()
1382 .get(DD_PARENT_SESSION_ID)
1383 .unwrap()
1384 .to_str()
1385 .unwrap(),
1386 "parent"
1387 );
1388 }
1389
1390 #[cfg_attr(miri, ignore)] #[test]
1392 fn telemetry_http_omits_parent_session_id_when_same_as_session_id() {
1393 let req = test_worker(
1394 Some("sess-id".into()),
1395 Some("root".into()),
1396 Some("sess-id".into()),
1397 )
1398 .build_request(&Payload::AppHeartbeat(()))
1399 .unwrap();
1400 assert_eq!(
1401 req.headers().get(DD_SESSION_ID).unwrap().to_str().unwrap(),
1402 "sess-id"
1403 );
1404 assert_eq!(
1405 req.headers()
1406 .get(DD_ROOT_SESSION_ID)
1407 .unwrap()
1408 .to_str()
1409 .unwrap(),
1410 "root"
1411 );
1412 assert!(req.headers().get(DD_PARENT_SESSION_ID).is_none());
1413 }
1414
1415 #[cfg_attr(miri, ignore)] #[test]
1417 fn telemetry_http_omits_session_family_without_valid_session_id() {
1418 let assert_no_session_headers = |req: &http_common::HttpRequest| {
1419 assert!(req.headers().get(DD_SESSION_ID).is_none());
1420 assert!(req.headers().get(DD_ROOT_SESSION_ID).is_none());
1421 assert!(req.headers().get(DD_PARENT_SESSION_ID).is_none());
1422 };
1423
1424 let req = test_worker(None, Some("root".into()), Some("parent".into()))
1425 .build_request(&Payload::AppHeartbeat(()))
1426 .unwrap();
1427 assert_no_session_headers(&req);
1428
1429 let req = test_worker(
1430 Some(String::new()),
1431 Some("root".into()),
1432 Some("parent".into()),
1433 )
1434 .build_request(&Payload::AppHeartbeat(()))
1435 .unwrap();
1436 assert_no_session_headers(&req);
1437 }
1438
1439 #[cfg_attr(miri, ignore)] #[test]
1441 fn telemetry_http_includes_dd_session_root_and_parent_session_ids() {
1442 let req = test_worker(
1443 Some("sess".into()),
1444 Some("root".into()),
1445 Some("parent".into()),
1446 )
1447 .build_request(&Payload::AppHeartbeat(()))
1448 .unwrap();
1449 assert_eq!(
1450 req.headers().get(DD_SESSION_ID).unwrap().to_str().unwrap(),
1451 "sess"
1452 );
1453 assert_eq!(
1454 req.headers()
1455 .get(DD_ROOT_SESSION_ID)
1456 .unwrap()
1457 .to_str()
1458 .unwrap(),
1459 "root"
1460 );
1461 assert_eq!(
1462 req.headers()
1463 .get(DD_PARENT_SESSION_ID)
1464 .unwrap()
1465 .to_str()
1466 .unwrap(),
1467 "parent"
1468 );
1469 }
1470
1471 fn build_test_worker_with_flavor(flavor: TelemetryWorkerFlavor) -> TelemetryWorker {
1472 let mut b = TelemetryWorkerBuilder::new(
1473 "h".into(),
1474 "svc".into(),
1475 "lang".into(),
1476 "1".into(),
1477 "tv".into(),
1478 );
1479 b.config
1480 .set_endpoint(TelemetryEndpoint {
1481 url: Some("http://127.0.0.1:1".to_owned()),
1482 ..Default::default()
1483 })
1484 .unwrap();
1485 b.runtime_id = Some("rid".into());
1486 b.flavor = flavor;
1487 b.build_worker(Some(tokio::runtime::Handle::current())).1
1488 }
1489
1490 #[tokio::test]
1494 #[cfg_attr(miri, ignore)] async fn full_flavor_start_schedules_every_periodic_action() {
1496 let mut worker = build_test_worker_with_flavor(TelemetryWorkerFlavor::Full);
1497
1498 let _ = worker
1499 .dispatch_action(TelemetryActions::Lifecycle(LifecycleAction::Start))
1500 .await;
1501
1502 let delays: Vec<LifecycleAction> =
1503 worker.deadlines.delays.iter().map(|(_, k)| *k).collect();
1504 let scheduled: Vec<LifecycleAction> =
1505 worker.deadlines.deadlines.iter().map(|(_, k)| *k).collect();
1506
1507 assert!(!delays.is_empty(), "scheduler should have periodic actions");
1508 for ev in &delays {
1509 assert!(
1510 scheduled.contains(ev),
1511 "{ev:?} has a delay but was not scheduled on Start; scheduled={scheduled:?}",
1512 );
1513 }
1514 }
1515
1516 #[tokio::test]
1519 #[cfg_attr(miri, ignore)] async fn metrics_logs_flavor_start_does_not_schedule_extended_heartbeat() {
1521 let mut worker = build_test_worker_with_flavor(TelemetryWorkerFlavor::MetricsLogs);
1522
1523 let _ = worker
1524 .dispatch_metrics_logs_action(TelemetryActions::Lifecycle(LifecycleAction::Start))
1525 .await;
1526
1527 let scheduled: Vec<LifecycleAction> =
1528 worker.deadlines.deadlines.iter().map(|(_, k)| *k).collect();
1529
1530 assert!(scheduled.contains(&LifecycleAction::FlushMetricAggr));
1531 assert!(scheduled.contains(&LifecycleAction::FlushData));
1532 assert!(
1533 !scheduled.contains(&LifecycleAction::ExtendedHeartbeat),
1534 "MetricsLogs should not schedule ExtendedHeartbeat; scheduled={scheduled:?}",
1535 );
1536 }
1537
1538 #[tokio::test]
1543 #[cfg_attr(miri, ignore)] async fn extended_heartbeat_does_not_reset_flush_data() {
1545 let mut worker = build_test_worker_with_flavor(TelemetryWorkerFlavor::Full);
1546
1547 let _ = worker
1548 .dispatch_action(TelemetryActions::Lifecycle(LifecycleAction::Start))
1549 .await;
1550
1551 let flush_data_before = worker
1552 .deadlines
1553 .deadlines
1554 .iter()
1555 .find(|(_, k)| *k == LifecycleAction::FlushData)
1556 .map(|(d, _)| *d)
1557 .expect("FlushData scheduled on Start");
1558
1559 let _ = worker
1560 .dispatch_action(TelemetryActions::Lifecycle(
1561 LifecycleAction::ExtendedHeartbeat,
1562 ))
1563 .await;
1564
1565 let flush_data_after = worker
1566 .deadlines
1567 .deadlines
1568 .iter()
1569 .find(|(_, k)| *k == LifecycleAction::FlushData)
1570 .map(|(d, _)| *d)
1571 .expect("FlushData should still be scheduled after ExtendedHeartbeat fires");
1572
1573 assert_eq!(
1574 flush_data_before, flush_data_after,
1575 "ExtendedHeartbeat must not reset FlushData's deadline",
1576 );
1577 }
1578
1579 mod reset {
1580 use super::super::*;
1581 use crate::data::{
1582 metrics::{MetricNamespace, MetricType},
1583 Configuration, ConfigurationOrigin, Dependency, Endpoint, Integration, Log, LogLevel,
1584 };
1585 use libdd_shared_runtime::Worker;
1586
1587 fn build_test_worker() -> (TelemetryWorkerHandle, TelemetryWorker) {
1588 let builder = TelemetryWorkerBuilder::new(
1589 "hostname".to_string(),
1590 "test-service".to_string(),
1591 "rust".to_string(),
1592 "1.0.0".to_string(),
1593 "1.0.0".to_string(),
1594 );
1595 builder.build_worker(Some(tokio::runtime::Handle::current()))
1597 }
1598
1599 fn make_log(id: u64, message: &str) -> (LogIdentifier, Log) {
1600 (
1601 LogIdentifier { identifier: id },
1602 Log {
1603 message: message.to_string(),
1604 level: LogLevel::Warn,
1605 stack_trace: None,
1606 count: 1,
1607 tags: String::new(),
1608 is_sensitive: false,
1609 is_crash: false,
1610 },
1611 )
1612 }
1613
1614 #[cfg_attr(miri, ignore)] #[tokio::test]
1617 async fn test_reset_clears_buffered_data() {
1618 let (handle, mut worker) = build_test_worker();
1619
1620 worker.data.dependencies.insert(Dependency {
1622 name: "dep".to_string(),
1623 version: None,
1624 });
1625 worker.data.integrations.insert(Integration {
1626 name: "integration".to_string(),
1627 version: None,
1628 enabled: true,
1629 compatible: None,
1630 auto_enabled: None,
1631 });
1632 worker.data.configurations.insert(Configuration {
1633 name: "cfg".to_string(),
1634 value: "true".to_string(),
1635 origin: ConfigurationOrigin::Code,
1636 config_id: None,
1637 seq_id: None,
1638 });
1639 worker.data.endpoints.insert(Endpoint {
1640 operation_name: "GET /health".to_string(),
1641 resource_name: "/health".to_string(),
1642 ..Default::default()
1643 });
1644 let (id, log) = make_log(42, "msg");
1645 worker.data.logs.get_mut_or_insert(id, log);
1646
1647 let key = handle.register_metric_context(
1649 "test.metric".to_string(),
1650 vec![],
1651 MetricType::Count,
1652 false,
1653 MetricNamespace::Tracers,
1654 );
1655 worker.data.metric_buckets.add_point(key, 1.0, vec![]);
1656
1657 worker.reset();
1658
1659 let stats = worker.stats();
1660 assert_eq!(
1661 stats.dependencies_stored, 0,
1662 "dependency dedupe history should be cleared"
1663 );
1664 assert_eq!(
1665 stats.dependencies_unflushed, 0,
1666 "dependency pending queue should be cleared"
1667 );
1668 assert_eq!(
1669 stats.integrations_stored, 0,
1670 "integration dedupe history should be cleared"
1671 );
1672 assert_eq!(
1673 stats.integrations_unflushed, 0,
1674 "integration pending queue should be cleared"
1675 );
1676 assert_eq!(
1677 stats.configurations_stored, 0,
1678 "configuration dedupe history should be cleared"
1679 );
1680 assert_eq!(
1681 stats.configurations_unflushed, 0,
1682 "configuration pending queue should be cleared"
1683 );
1684 assert_eq!(stats.logs, 0, "logs should be cleared");
1685 assert_eq!(
1686 stats.metric_buckets.buckets, 0,
1687 "metric buckets should be cleared"
1688 );
1689 assert_eq!(
1690 stats.metric_buckets.series, 0,
1691 "metric series should be cleared"
1692 );
1693 assert!(
1694 worker.data.endpoints.is_empty(),
1695 "endpoints should be cleared"
1696 );
1697 assert!(worker.next_action.is_none(), "next_action should be None");
1698 }
1699
1700 #[cfg_attr(miri, ignore)] #[tokio::test]
1703 async fn test_reset_drains_mailbox() {
1704 let (handle, mut worker) = build_test_worker();
1705
1706 handle
1708 .try_send_msg(TelemetryActions::AddDependency(Dependency {
1709 name: "dep".to_string(),
1710 version: None,
1711 }))
1712 .unwrap();
1713 let (id, log) = make_log(1, "pre-fork log");
1714 handle
1715 .try_send_msg(TelemetryActions::AddLog((id, log)))
1716 .unwrap();
1717
1718 worker.next_action = Some(TelemetryActions::Lifecycle(LifecycleAction::Start));
1720
1721 worker.reset();
1722
1723 assert!(
1725 worker.mailbox.try_recv().is_err(),
1726 "mailbox should be empty"
1727 );
1728 assert!(worker.next_action.is_none(), "next_action should be None");
1729 let stats = worker.stats();
1731 assert_eq!(
1732 stats.dependencies_stored, 0,
1733 "queued AddDependency must not be applied"
1734 );
1735 assert_eq!(
1736 stats.dependencies_unflushed, 0,
1737 "queued AddDependency must not be pending"
1738 );
1739 assert_eq!(stats.logs, 0, "queued AddLog must be discarded");
1740 }
1741
1742 #[cfg_attr(miri, ignore)] #[tokio::test]
1745 async fn test_worker_accepts_new_data_after_reset() {
1746 let (handle, mut worker) = build_test_worker();
1747 worker.flavor = TelemetryWorkerFlavor::MetricsLogs;
1748
1749 let (id, log) = make_log(99, "pre-fork");
1751 worker.data.logs.get_mut_or_insert(id, log);
1752
1753 worker.reset();
1754
1755 let (id2, log2) = make_log(1, "post-fork");
1757 handle
1758 .try_send_msg(TelemetryActions::AddLog((id2, log2)))
1759 .unwrap();
1760
1761 worker.trigger().await;
1763 worker.run().await;
1764
1765 let stats = worker.stats();
1766 assert_eq!(stats.logs, 1, "only post-fork log should be present");
1768 }
1769
1770 #[cfg_attr(miri, ignore)] #[tokio::test]
1773 async fn test_reset_preserves_started_and_deadlines() {
1774 let (_handle, mut worker) = build_test_worker();
1775
1776 worker.data.started = true;
1777 worker
1778 .deadlines
1779 .schedule_event(LifecycleAction::FlushMetricAggr)
1780 .unwrap();
1781 worker
1782 .deadlines
1783 .schedule_event(LifecycleAction::FlushData)
1784 .unwrap();
1785
1786 let deadlines_before = worker.deadlines.deadlines.clone();
1787
1788 worker.reset();
1789
1790 assert!(worker.data.started, "started flag should be preserved");
1791 assert_eq!(
1792 worker.deadlines.deadlines.len(),
1793 deadlines_before.len(),
1794 "scheduled deadlines should be preserved"
1795 );
1796 for ((_, actual), (_, expected)) in worker
1797 .deadlines
1798 .deadlines
1799 .iter()
1800 .zip(deadlines_before.iter())
1801 {
1802 assert_eq!(
1803 actual, expected,
1804 "deadline kinds should be preserved across reset"
1805 );
1806 }
1807 }
1808 }
1809
1810 #[cfg_attr(miri, ignore)]
1811 #[test]
1812 fn test_channel_close_flushes_and_parks_via_shared_runtime() {
1813 use httpmock::prelude::*;
1814 use libdd_shared_runtime::{BlockingRuntime, ForkSafeRuntime, SharedRuntime};
1815 use std::time::Duration;
1816
1817 const TELEMETRY_PATH: &str = "/telemetry/proxy/api/v2/apmtelemetry";
1818
1819 let server = MockServer::start();
1820 let mock = server.mock(|when, then| {
1821 when.method(POST).path(TELEMETRY_PATH);
1822 then.status(202).body("");
1823 });
1824
1825 let mut builder = TelemetryWorkerBuilder::new(
1826 "host".into(),
1827 "svc".into(),
1828 "lang".into(),
1829 "1".into(),
1830 "tv".into(),
1831 );
1832 builder
1833 .config
1834 .set_endpoint(TelemetryEndpoint {
1835 url: Some(server.url("/")),
1836 ..Default::default()
1837 })
1838 .unwrap();
1839 builder.runtime_id = Some("rid".into());
1840
1841 let shared_runtime = ForkSafeRuntime::new().expect("ForkSafeRuntime::new");
1842 let runtime_handle = shared_runtime
1843 .block_on(async { tokio::runtime::Handle::current() })
1844 .expect("runtime handle");
1845 let (telemetry_handle, worker) = builder.build_worker(Some(runtime_handle));
1846
1847 let _worker_handle = shared_runtime
1848 .spawn_worker(worker, false)
1849 .expect("spawn_worker");
1850
1851 telemetry_handle.send_start().expect("send_start");
1853
1854 for _ in 0..50 {
1857 if mock.calls() >= 1 {
1858 break;
1859 }
1860 std::thread::sleep(Duration::from_millis(20));
1861 }
1862 assert!(
1863 mock.calls() >= 1,
1864 "worker should POST at least once after Start"
1865 );
1866
1867 let hits_before_close = mock.calls();
1869 drop(telemetry_handle);
1870
1871 for _ in 0..50 {
1873 if mock.calls() > hits_before_close {
1874 break;
1875 }
1876 std::thread::sleep(Duration::from_millis(20));
1877 }
1878 assert!(
1879 mock.calls() > hits_before_close,
1880 "worker should flush a final batch after the channel is closed"
1881 );
1882
1883 let stable_hits = mock.calls();
1886 std::thread::sleep(Duration::from_millis(300));
1887 assert_eq!(
1888 mock.calls(),
1889 stable_hits,
1890 "worker must stop POSTing after parking; observed {} extra hits",
1891 mock.calls().saturating_sub(stable_hits),
1892 );
1893 }
1894}