Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
91 changes: 91 additions & 0 deletions bin/agent-data-plane/src/state/metrics/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -571,6 +571,14 @@ mod tests {
find("aggregator.number_of_flush"),
Some("Number of flushes done by the aggregator")
);
assert_eq!(
find("aggregator.flush_count"),
Some("Number of items handled by the last flush, by flush type")
);
assert_eq!(
find("aggregator.flush_time"),
Some("Duration in nanoseconds of the last flush, by flush type")
);
assert_eq!(find("filterlist.size"), Some("Metric filter list size"));
assert_eq!(
find("filterlist.updates"),
Expand Down Expand Up @@ -693,6 +701,89 @@ mod tests {
assert!(output.contains("# TYPE aggregator__number_of_flush counter"));
}

#[test]
fn render_rar_telemetry_remaps_aggregator_flush_telemetry() {
let metrics = vec![
Event::Metric(Metric::gauge(
Context::from_static_parts(
"adp.aggregate_last_flush_count",
&["component_id:dsd_agg", "data_type:series"],
),
11.0,
)),
Event::Metric(Metric::gauge(
Context::from_static_parts(
"adp.aggregate_last_flush_count",
&["component_id:dsd_agg", "data_type:sketches"],
),
12.0,
)),
Event::Metric(Metric::gauge(
Context::from_static_parts(
"adp.aggregate_last_flush_duration_nanoseconds",
&["component_id:dsd_agg"],
),
1000.0,
)),
Event::Metric(Metric::counter(
Context::from_static_parts("adp.encoder_flushed_events_total", &["component_id:dd_events_encode"]),
5.0,
)),
Event::Metric(Metric::counter(
Context::from_static_parts(
"adp.encoder_flushed_events_total",
&["component_id:dd_service_checks_encode"],
),
6.0,
)),
Event::Metric(Metric::gauge(
Context::from_static_parts("adp.encoder_last_flush_events", &["component_id:dd_events_encode"]),
2.0,
)),
Event::Metric(Metric::gauge(
Context::from_static_parts(
"adp.encoder_last_flush_events",
&["component_id:dd_service_checks_encode"],
),
3.0,
)),
Event::Metric(Metric::gauge(
Context::from_static_parts(
"adp.encoder_last_flush_duration_nanoseconds",
&["component_id:dd_events_encode"],
),
2000.0,
)),
Event::Metric(Metric::gauge(
Context::from_static_parts(
"adp.encoder_last_flush_duration_nanoseconds",
&["component_id:dd_service_checks_encode"],
),
3000.0,
)),
// The logs encoder shares the same encoder telemetry, but has no Core Agent aggregator counterpart.
Event::Metric(Metric::counter(
Context::from_static_parts("adp.encoder_flushed_events_total", &["component_id:dd_logs_encode"]),
99.0,
)),
];

let output = render_with(get_datadog_agent_remappings(), metrics);

assert!(output.contains("aggregator__flush_count{flush_type=\"series\"} 11"));
assert!(output.contains("aggregator__flush_count{flush_type=\"sketches\"} 12"));
assert!(output.contains("aggregator__flush_count{flush_type=\"events\"} 2"));
assert!(output.contains("aggregator__flush_count{flush_type=\"service_checks\"} 3"));
assert!(output.contains("aggregator__flush_time{flush_type=\"main\"} 1000"));
assert!(output.contains("aggregator__flush_time{flush_type=\"event\"} 2000"));
assert!(output.contains("aggregator__flush_time{flush_type=\"service_check\"} 3000"));
assert!(output.contains("aggregator__flush{data_type=\"events\"} 5"));
assert!(output.contains("aggregator__flush{data_type=\"service_checks\"} 6"));
assert!(!output.contains(" 99"));
assert!(output.contains("# TYPE aggregator__flush_count gauge"));
assert!(output.contains("# TYPE aggregator__flush_time gauge"));
}

#[test]
fn rar_telemetry_remaps_supported_dogstatsd_client_byte_telemetry_as_counters() {
let metrics = [
Expand Down
72 changes: 69 additions & 3 deletions bin/agent-data-plane/src/state/metrics/rules/aggregation.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,12 @@
use super::RemapperRule;

const NO_AGG_SPLIT_COMPONENT_TAG: &str = "component_id:dsd_no_agg_split";
const EVENTS_ENCODER_COMPONENT_TAG: &str = "component_id:dd_events_encode";
const SERVICE_CHECKS_ENCODER_COMPONENT_TAG: &str = "component_id:dd_service_checks_encode";

const FLUSH_HELP_TEXT: &str = "Number of metrics/service checks/events flushed";
const FLUSH_COUNT_HELP_TEXT: &str = "Number of items handled by the last flush, by flush type";
const FLUSH_TIME_HELP_TEXT: &str = "Duration in nanoseconds of the last flush, by flush type";

pub fn get_aggregation_remappings() -> Vec<RemapperRule> {
vec![
Expand Down Expand Up @@ -34,14 +40,14 @@ pub fn get_aggregation_remappings() -> Vec<RemapperRule> {
// Events and service checks don't pass through an aggregator in ADP, so their encoders stand in for it.
RemapperRule::by_name_and_tags(
"adp.component_events_received_total",
&["component_id:dd_events_encode"],
&[EVENTS_ENCODER_COMPONENT_TAG],
"aggregator.processed",
)
.with_additional_tags(["data_type:events"])
.with_help_text("Amount of metrics/services_checks/events processed by the aggregator"),
RemapperRule::by_name_and_tags(
"adp.component_events_received_total",
&["component_id:dd_service_checks_encode"],
&[SERVICE_CHECKS_ENCODER_COMPONENT_TAG],
"aggregator.processed",
)
.with_additional_tags(["data_type:service_checks"])
Expand All @@ -65,12 +71,72 @@ pub fn get_aggregation_remappings() -> Vec<RemapperRule> {
"aggregator.flush",
)
.with_original_tags(["data_type"])
.with_help_text("Number of metrics/service checks/events flushed"),
.with_help_text(FLUSH_HELP_TEXT),
RemapperRule::by_name_and_tags(
"adp.aggregate_flushes_total",
&["component_id:dsd_agg"],
"aggregator.number_of_flush",
)
.with_help_text("Number of flushes done by the aggregator"),
RemapperRule::by_name_and_tags(
"adp.aggregate_last_flush_count",
&["component_id:dsd_agg"],
"aggregator.flush_count",
)
.with_remapped_tags([("data_type", "flush_type")])
.with_help_text(FLUSH_COUNT_HELP_TEXT),
// There's no ADP equivalent of the Core Agent's separate `metric_sketch` and `checks_metric_sample` flush
// timings, which include serialization: ADP encodes metrics in a separate component after the aggregator flush.
RemapperRule::by_name_and_tags(
"adp.aggregate_last_flush_duration_nanoseconds",
&["component_id:dsd_agg"],
"aggregator.flush_time",
)
.with_additional_tags(["flush_type:main"])
.with_help_text(FLUSH_TIME_HELP_TEXT),
// As with `aggregator.processed`, the events and service checks encoders stand in for the aggregator's
// event and service check flushes.
RemapperRule::by_name_and_tags(
"adp.encoder_flushed_events_total",
&[EVENTS_ENCODER_COMPONENT_TAG],
"aggregator.flush",
)
.with_additional_tags(["data_type:events"])
.with_help_text(FLUSH_HELP_TEXT),
RemapperRule::by_name_and_tags(
"adp.encoder_flushed_events_total",
&[SERVICE_CHECKS_ENCODER_COMPONENT_TAG],
"aggregator.flush",
)
.with_additional_tags(["data_type:service_checks"])
.with_help_text(FLUSH_HELP_TEXT),
RemapperRule::by_name_and_tags(
"adp.encoder_last_flush_events",
&[EVENTS_ENCODER_COMPONENT_TAG],
"aggregator.flush_count",
)
.with_additional_tags(["flush_type:events"])
.with_help_text(FLUSH_COUNT_HELP_TEXT),
RemapperRule::by_name_and_tags(
"adp.encoder_last_flush_events",
&[SERVICE_CHECKS_ENCODER_COMPONENT_TAG],
"aggregator.flush_count",
)
.with_additional_tags(["flush_type:service_checks"])
.with_help_text(FLUSH_COUNT_HELP_TEXT),
RemapperRule::by_name_and_tags(
"adp.encoder_last_flush_duration_nanoseconds",
&[EVENTS_ENCODER_COMPONENT_TAG],
"aggregator.flush_time",
)
.with_additional_tags(["flush_type:event"])
.with_help_text(FLUSH_TIME_HELP_TEXT),
RemapperRule::by_name_and_tags(
"adp.encoder_last_flush_duration_nanoseconds",
&[SERVICE_CHECKS_ENCODER_COMPONENT_TAG],
"aggregator.flush_time",
)
.with_additional_tags(["flush_type:service_check"])
.with_help_text(FLUSH_TIME_HELP_TEXT),
]
}
38 changes: 32 additions & 6 deletions lib/saluki-components/src/encoders/buffered_incremental/mod.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use std::time::Duration;
use std::time::{Duration, Instant};

use agent_data_plane_config::defaults::DEFAULT_ENCODER_FLUSH_TIMEOUT;
use async_trait::async_trait;
Expand All @@ -9,6 +9,7 @@ use saluki_core::{
components::{encoders::*, BuildContext},
data_model::{event::EventType, payload::PayloadType},
observability::ComponentMetricsExt,
topology::PayloadsDispatcher,
};
use saluki_error::GenericError;
use saluki_metrics::MetricsBuilder;
Expand Down Expand Up @@ -173,6 +174,7 @@ where
health.mark_ready();

let mut pending_flush = false;
let mut pending_events = 0;
let pending_flush_timeout = sleep(flush_timeout);
pin!(pending_flush_timeout);

Expand All @@ -192,19 +194,22 @@ where
// If we're informed that we need to flush, we'll hold on to this event before triggering a flush and then
// retry processing it after flushing.
let event_to_retry = match encoder.process_event(event).await? {
ProcessResult::Continue => continue,
ProcessResult::Continue => {
pending_events += 1;
continue;
}
ProcessResult::FlushRequired(event) => event,
};

// Flush the encoder, waiting any payloads it has generated.
encoder.flush(context.dispatcher()).await?;
flush_encoder(&mut encoder, context.dispatcher(), &telemetry, &mut pending_events).await?;

// Now try to process the event again.
//
// If this fails, then we drop the event because it's a logical bug to not be able to encode an event after
// flushing, and we don't want to get stuck in an infinite loop.
match encoder.process_event(event_to_retry).await? {
ProcessResult::Continue => {},
ProcessResult::Continue => pending_events += 1,
ProcessResult::FlushRequired(_) => {
error!("Failed to process event after flushing.");
telemetry.events_dropped_encoder().increment(1);
Expand All @@ -225,15 +230,36 @@ where

pending_flush = false;

encoder.flush(context.dispatcher()).await?;
flush_encoder(&mut encoder, context.dispatcher(), &telemetry, &mut pending_events).await?;

debug!("All pending payloads flushed.");
}
}
}

// Do a final flush since we may have had a pending payloads before breaking out of the loop.
encoder.flush(context.dispatcher()).await?;
flush_encoder(&mut encoder, context.dispatcher(), &telemetry, &mut pending_events).await?;

Ok(())
}

/// Flushes the encoder, recording flush telemetry for the events processed since the last flush.
///
/// Flushes with no pending events are not recorded, so the last-flush telemetry always describes a flush that sent
/// something.
async fn flush_encoder<E>(
encoder: &mut E, dispatcher: &PayloadsDispatcher, telemetry: &ComponentTelemetry, pending_events: &mut u64,
) -> Result<(), GenericError>
where
E: IncrementalEncoder,
{
let flush_start = Instant::now();
encoder.flush(dispatcher).await?;

if *pending_events > 0 {
telemetry.record_flush(*pending_events, flush_start.elapsed());
*pending_events = 0;
}

Ok(())
}
Original file line number Diff line number Diff line change
@@ -1,10 +1,15 @@
use metrics::Counter;
use std::time::Duration;

use metrics::{Counter, Gauge};
use saluki_metrics::MetricsBuilder;

/// Incremental encoder-specific telemetry.
#[derive(Clone)]
pub struct ComponentTelemetry {
events_dropped_encoder: Counter,
events_flushed: Counter,
last_flush_events: Gauge,
last_flush_duration: Gauge,
}

impl ComponentTelemetry {
Expand All @@ -15,11 +20,42 @@ impl ComponentTelemetry {
"component_events_dropped_total",
["intentional:false", "drop_reason:encoder_failure"],
),
events_flushed: builder.register_counter("encoder_flushed_events_total"),
last_flush_events: builder.register_gauge("encoder_last_flush_events"),
last_flush_duration: builder.register_gauge("encoder_last_flush_duration_nanoseconds"),
}
}

/// Returns a reference to the "events dropped (encoder)" counter.
pub fn events_dropped_encoder(&self) -> &Counter {
&self.events_dropped_encoder
}

/// Records a completed flush of the given number of events, which took the given duration.
pub fn record_flush(&self, events: u64, duration: Duration) {
self.events_flushed.increment(events);
self.last_flush_events.set(events as f64);
self.last_flush_duration.set(duration.as_nanos() as f64);
}
}

#[cfg(test)]
mod tests {
use saluki_metrics::test::TestRecorder;

use super::*;

#[test]
fn record_flush_tracks_totals_and_last_flush() {
let recorder = TestRecorder::default();
let _recorder_guard = metrics::set_default_local_recorder(&recorder);
let telemetry = ComponentTelemetry::from_builder(&MetricsBuilder::default());

telemetry.record_flush(3, Duration::from_micros(5));
telemetry.record_flush(2, Duration::from_micros(7));

assert_eq!(recorder.counter("encoder_flushed_events_total"), Some(5));
assert_eq!(recorder.gauge("encoder_last_flush_events"), Some(2.0));
assert_eq!(recorder.gauge("encoder_last_flush_duration_nanoseconds"), Some(7_000.0));
}
}
Loading
Loading