use std::collections::BTreeMap;
use std::time::SystemTime;
use opentelemetry::KeyValue;
use opentelemetry::metrics::Meter;
use opentelemetry::trace::{Span, SpanBuilder, Tracer};
use subms::{ObservationCtx, SubMsBenchSummary, SubMsStageKind, SubMsStageSummary, SubMsTimer};
pub const HISTOGRAM_NAME: &str = "subms.latency";
pub const HISTOGRAM_UNIT: &str = "s";
pub fn histogram_boundaries(kind: SubMsStageKind) -> &'static [f64] {
const HOT: &[f64] = &[
5e-8, 1e-7, 2e-7, 5e-7, 1e-6, 2e-6, 5e-6, 1e-5, 5e-5, 1e-4, 5e-4, 1e-3,
];
const BATCH: &[f64] = &[1e-6, 1e-5, 1e-4, 1e-3, 1e-2, 1e-1, 1.0];
match kind {
SubMsStageKind::HotPath => HOT,
SubMsStageKind::BatchOp | SubMsStageKind::OneShot => BATCH,
SubMsStageKind::Unspecified => &[],
_ => &[],
}
}
pub fn attributes_from_ctx(ctx: &ObservationCtx<'_>) -> Vec<KeyValue> {
vec![
KeyValue::new("subms.workload", ctx.workload.to_string()),
KeyValue::new("subms.lang", ctx.lang.to_string()),
KeyValue::new("subms.stage", ctx.stage.to_string()),
KeyValue::new("subms.stage.kind", ctx.stage_kind.as_str()),
]
}
pub fn attributes_from_summary(
summary: &SubMsBenchSummary,
stage: &SubMsStageSummary,
kind: SubMsStageKind,
) -> Vec<KeyValue> {
let mut attrs = vec![
KeyValue::new("subms.workload", summary.workload.clone()),
KeyValue::new("subms.lang", summary.lang.clone()),
KeyValue::new("subms.stage", stage.name.clone()),
KeyValue::new("subms.stage.kind", kind.as_str()),
];
push_meta_attrs(&mut attrs, &summary.inputs, &summary.meta);
attrs
}
pub(crate) fn push_meta_attrs(
attrs: &mut Vec<KeyValue>,
inputs: &BTreeMap<String, String>,
meta: &BTreeMap<String, String>,
) {
fn put(
attrs: &mut Vec<KeyValue>,
map: &BTreeMap<String, String>,
src: &str,
dst: &'static str,
) {
if let Some(v) = map.get(src) {
attrs.push(KeyValue::new(dst, v.clone()));
}
}
put(attrs, meta, "subms.recipe.slug", "subms.recipe.slug");
put(
attrs,
meta,
"subms.recipe.category",
"subms.recipe.category",
);
put(
attrs,
meta,
"subms.workload.feature",
"subms.workload.feature",
);
put(attrs, inputs, "entries", "subms.workload.entries");
put(attrs, inputs, "seed", "subms.workload.seed");
put(attrs, meta, "host", "subms.host");
put(attrs, meta, "hardware_tier", "subms.hardware.tier");
put(attrs, meta, "crate_version", "subms.crate.version");
}
pub fn export_summary(summary: &SubMsBenchSummary, meter: &Meter) {
for stage in &summary.stages {
let kind = SubMsStageKind::Unspecified;
let attrs = attributes_from_summary(summary, stage, kind);
let histogram = histogram_for_stage(meter, kind);
record_seconds(&histogram, stage.p50_ns, &attrs);
record_seconds(&histogram, stage.p99_ns, &attrs);
record_seconds(&histogram, stage.p999_ns, &attrs);
record_seconds(&histogram, stage.max_ns, &attrs);
record_seconds(&histogram, stage.mean_ns, &attrs);
if let Some(samples) = &stage.samples_ns {
for ns in samples {
record_seconds(&histogram, *ns, &attrs);
}
}
}
}
fn record_seconds(histogram: &opentelemetry::metrics::Histogram<f64>, ns: u64, attrs: &[KeyValue]) {
histogram.record(ns as f64 / 1e9, attrs);
}
pub(crate) fn histogram_for_stage(
meter: &Meter,
kind: SubMsStageKind,
) -> opentelemetry::metrics::Histogram<f64> {
let bounds = histogram_boundaries(kind);
let mut builder = meter
.f64_histogram(HISTOGRAM_NAME)
.with_unit(HISTOGRAM_UNIT)
.with_description("subms per-stage latency histogram");
if !bounds.is_empty() {
builder = builder.with_boundaries(bounds.to_vec());
}
builder.build()
}
pub fn export_timer<T: Tracer>(timer: &SubMsTimer, tracer: &T)
where
T::Span: Span + Send + Sync + 'static,
{
let total_ns = timer.elapsed_ns();
let now = SystemTime::now();
let parent_start = now
.checked_sub(std::time::Duration::from_nanos(total_ns))
.unwrap_or(now);
let parent_name: std::borrow::Cow<'static, str> = if timer.name().is_empty() {
"subms.timer".into()
} else {
timer.name().to_string().into()
};
let mut parent = SpanBuilder::from_name(parent_name)
.with_start_time(parent_start)
.with_attributes(vec![
KeyValue::new("subms.timer.total_ns", total_ns as i64),
KeyValue::new("subms.timer.checkpoints", timer.checkpoints().len() as i64),
])
.start(tracer);
for cp in timer.checkpoints() {
let child_start = parent_start
.checked_add(std::time::Duration::from_nanos(
cp.since_start_ns.saturating_sub(cp.since_last_ns),
))
.unwrap_or(parent_start);
let child_end = parent_start
.checked_add(std::time::Duration::from_nanos(cp.since_start_ns))
.unwrap_or(parent_start);
let mut child = SpanBuilder::from_name(cp.label.clone())
.with_start_time(child_start)
.with_attributes(vec![
KeyValue::new("subms.timer.since_last_ns", cp.since_last_ns as i64),
KeyValue::new("subms.timer.since_start_ns", cp.since_start_ns as i64),
KeyValue::new("subms.timer.is_stop", cp.is_stop),
])
.start(tracer);
child.end_with_timestamp(child_end);
}
parent.end_with_timestamp(now);
}