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
4 changes: 4 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -107,3 +107,7 @@ orchagent/p4orch/tests/*_tr.xml

build-env/.env
build-env/custom-setup.sh

# Perf / flamegraph profiling artifacts #
###################
perf.data*
5 changes: 5 additions & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,11 @@ members = [
]
exclude = []

# Enable debug symbols in the bench profile so profilers (cargo-flamegraph, perf)
# can resolve function names.
[profile.bench]
debug = true

[workspace.package]
version = "0.1.0"
authors = ["SONiC"]
Expand Down
35 changes: 24 additions & 11 deletions crates/countersyncd/benches/otel_actor_perf.rs
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,7 @@ fn build_stats_message(counters: usize, seed: u64) -> SAIStatsMessage {
Arc::new(SAIStats::new(seed, stats))
}

async fn run_stream(prepared: PreparedDataset, endpoint: String) -> (std::time::Duration, usize) {
async fn run_stream(messages: Vec<SAIStatsMessage>, total_counters: usize, endpoint: String) -> (std::time::Duration, usize) {
let (tx, rx) = mpsc::channel(1024);
let (shutdown_tx, _shutdown_rx) = oneshot::channel();

Expand All @@ -97,14 +97,11 @@ async fn run_stream(prepared: PreparedDataset, endpoint: String) -> (std::time::

let handle = tokio::spawn(async move { actor.run().await });

let total_counters = prepared.expected_counters;
let start = std::time::Instant::now();

for tmpl in prepared.templates.iter() {
for msg_idx in 0..tmpl.records {
let msg = build_stats_message(tmpl.spec.counters, msg_idx as u64);
tx.send(msg).await.expect("send stats");
}
// Only sending + actor conversion/encode/send is timed
for msg in messages {
tx.send(msg).await.expect("send stats");
}

drop(tx); // close channel so actor exits after processing
Expand Down Expand Up @@ -154,16 +151,32 @@ fn bench_otel_actor(c: &mut Criterion) {
b.to_async(&rt).iter_batched(
{
let spec = spec.clone();
move || PreparedDataset::new(spec.clone())
move || {
// Build all input messages outside the profiled window so the
// flamegraph reflects only conversion + protobuf + gRPC send.
let prepared = PreparedDataset::new(spec.clone());
let total_counters = prepared.expected_counters;
let mut messages = Vec::new();
for tmpl in prepared.templates.iter() {
for msg_idx in 0..tmpl.records {
messages.push(build_stats_message(
tmpl.spec.counters,
msg_idx as u64,
));
}
}
(messages, total_counters)
}
},
move |prepared| {
move |(messages, total_counters)| {
let endpoint = endpoint.clone();
let exports_counter = exports_counter.clone();
let spec = spec.clone();
async move {
let exports_before = exports_counter.load(Ordering::Relaxed);

let (elapsed, counters) = run_stream(prepared, endpoint.clone()).await;
let (elapsed, counters) =
run_stream(messages, total_counters, endpoint.clone()).await;

let exports_after = exports_counter.load(Ordering::Relaxed);
let exported = exports_after.saturating_sub(exports_before);
Expand All @@ -174,7 +187,7 @@ fn bench_otel_actor(c: &mut Criterion) {
);
}
},
BatchSize::SmallInput,
BatchSize::PerIteration,
)
});
}
Expand Down
128 changes: 49 additions & 79 deletions crates/countersyncd/src/actor/otel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -23,16 +23,15 @@ use opentelemetry_proto::tonic::{
KeyValue as ProtoKeyValue,
},
metrics::v1::{
Gauge as ProtoGauge,
Metric,
ResourceMetrics,
ScopeMetrics,
},
resource::v1::Resource as ProtoResource,
};
use crate::message::{
otel::OtelMetrics,
saistats::SAIStatsMessage,
otel::{sai_stats_to_proto_metrics, DisplaySaiStats},
saistats::{SAIStats, SAIStatsMessage},
};
use crate::utilities::{record_comm_stats, ChannelLabel};

Expand Down Expand Up @@ -89,8 +88,9 @@ pub struct OtelActor {
resource: ProtoResource,
instrumentation_scope: InstrumentationScope,

// Batching
buffer: Vec<OtelMetrics>,
// Batching — buffers the raw SAI messages; each is converted straight to
// its final protobuf form at flush time
buffer: Vec<SAIStatsMessage>,
buffered_counters: usize,
flush_deadline: TokioInstant,

Expand Down Expand Up @@ -227,15 +227,14 @@ impl OtelActor {

let was_empty = self.buffer.is_empty();

// Convert to OTel format using message types and buffer
let otel_metrics = OtelMetrics::from_sai_stats(&stats);
let counters_in_message = stats.stats.len();

// The export path converts the buffered SAI message directly
if log::log_enabled!(log::Level::Debug) {
self.print_otel_metrics(&otel_metrics).await;
self.print_stats_report(&stats);
}

self.buffer.push(otel_metrics);
self.buffer.push(stats);
self.buffered_counters += counters_in_message;

// Start timeout when buffer transitions from empty to non-empty
Expand All @@ -252,41 +251,20 @@ impl OtelActor {
Ok(())
}

async fn print_otel_metrics(&mut self, otel_metrics: &OtelMetrics) {
fn print_stats_report(&mut self, stats: &SAIStats) {
self.console_reports += 1;

debug!(
"[OTel Report #{}] Service: {}, Scope: {} v{}, Total Gauges: {}, Messages Received: {}, Exports: {} (Failures: {})",
"[OTel Report #{}] Service: countersyncd, Scope: countersyncd v1.0, Counters: {}, Messages Received: {}, Exports: {} (Failures: {})",
self.console_reports,
otel_metrics.service_name,
otel_metrics.scope_name,
otel_metrics.scope_version,
otel_metrics.len(),
stats.stats.len(),
self.messages_received,
self.exports_performed,
self.export_failures
);

if !otel_metrics.is_empty() {
debug!("Gauge Metrics:");
for (index, gauge) in otel_metrics.gauges.iter().enumerate() {
let data_point = &gauge.data_points[0];

debug!("[{:3}] Gauge: {}", index + 1, gauge.name);
debug!("Value: {}", data_point.value);
debug!("Unit: {}", gauge.unit);
debug!("Time: {}ns", data_point.time_unix_nano);
debug!("Description: {}", gauge.description);

if !data_point.attributes.is_empty() {
debug!("Attributes:");
for attr in &data_point.attributes {
debug!(" - {}={}", attr.key, attr.value);
}
}

debug!("Raw Gauge: {:#?}", gauge);
}
if !stats.stats.is_empty() {
debug!("SAI counters:\n{}", DisplaySaiStats(stats));
}
}

Expand Down Expand Up @@ -314,23 +292,27 @@ impl OtelActor {
self.client.as_mut()
}

async fn send_request(
&mut self,
request: ExportMetricsServiceRequest,
) -> Result<(), Box<dyn ExportError>> {
async fn send_request(&mut self) -> Result<(), Box<dyn ExportError>> {
for attempt in 1..=MAX_EXPORT_RETRIES {
// Ensure we have a client
let client = match self.get_client() {
Some(c) => c, // Use existing or newly created client
_none => { // Failed to create client
self.client = None;
self.backoff(attempt).await; // Wait before retrying
continue;
}
// Ensure a client can be created before doing request
// construction; when get_client() fails the build below is skipped.
if self.get_client().is_none() {
self.client = None;
self.backoff(attempt).await; // Wait before retrying
Comment on lines +297 to +301
continue;
}

// Client is available, build a fresh request
let request = match self.build_export_request() {
Some(r) => r,
None => return Ok(()),
};

// Re-borrow the client for the actual send
let client = self.client.as_mut().expect("client ensured above");

// Attempt to send the request
match client.export(request.clone()).await {
match client.export(request).await {
Comment thread
Janetxxx marked this conversation as resolved.
Ok(_) => { // Successful export
self.exports_performed += 1;
self.consecutive_failures = 0;
Expand All @@ -349,38 +331,18 @@ impl OtelActor {
Err(Box::new(OtelActorExportError("Max export retries exceeded".to_string())))
}

// Export buffered metrics to OpenTelemetry collector
async fn flush_buffer(&mut self) -> Result<(), Box<dyn ExportError>> {
if self.buffer.is_empty() {
return Ok(());
}

/// Build an export request from the currently buffered metrics.
/// Returns `None` when there is nothing to export (empty metrics).
/// The buffer is left intact so it can be rebuilt on retry.
fn build_export_request(&self) -> Option<ExportMetricsServiceRequest> {
let mut proto_metrics: Vec<Metric> = Vec::new();

for otel_metrics in &self.buffer {
for gauge in &otel_metrics.gauges {
let proto_data_points = gauge.data_points.iter()
.map(|dp| dp.to_proto())
.collect();

let proto_gauge = ProtoGauge {
data_points: proto_data_points,
};

proto_metrics.push(Metric {
name: gauge.name.clone(),
description: gauge.description.clone(),
metadata: vec![],
data: Some(opentelemetry_proto::tonic::metrics::v1::metric::Data::Gauge(proto_gauge)),
..Default::default()
});
}
for stats in &self.buffer {
proto_metrics.extend(sai_stats_to_proto_metrics(stats));
}

if proto_metrics.is_empty() {
self.buffer.clear();
self.buffered_counters = 0;
return Ok(());
return None;
}

let resource_metrics = ResourceMetrics {
Expand All @@ -393,12 +355,20 @@ impl OtelActor {
schema_url: String::new(),
};

let request = ExportMetricsServiceRequest {
Some(ExportMetricsServiceRequest {
resource_metrics: vec![resource_metrics],
};
})
}

// Export buffered metrics to OpenTelemetry collector
async fn flush_buffer(&mut self) -> Result<(), Box<dyn ExportError>> {
if self.buffer.is_empty() {
return Ok(());
}

// Send the export request
let result = self.send_request(request).await;
// The request is built lazily inside send_request from the intact
// buffer (rebuilt per retry)
let result = self.send_request().await;

if let Err(e) = &result {
self.export_failures += 1;
Expand Down
Loading
Loading