use std::time::{Duration, SystemTime};
use axum::extract::{MatchedPath, Request};
use axum::http::HeaderMap;
use axum::middleware::Next;
use axum::response::Response;
use metrique::ServiceMetrics;
use metrique::emf::Emf;
use metrique::local::{LocalFormat, OutputStyle};
use metrique::timers::Timer;
use metrique::unit::{Byte, Count, Millisecond};
use metrique::unit_of_work::metrics;
use metrique::writer::sink::AttachHandle;
use metrique::writer::{AttachGlobalEntrySinkExt, FormatExt};
const CLOUDWATCH_NAMESPACE: &str = "dial9_viewer";
const SESSION_ID_HEADER: &str = "x-dial9-session-id";
const UUID_TEXT_LEN: usize = 36;
const UNMATCHED_OPERATION: &str = "unmatched";
#[metrics(emf::dimension_sets = [["operation"], []])]
struct RequestMetrics {
#[metrics(timestamp)]
timestamp: SystemTime,
operation: String,
status_code: String,
session_id: Option<String>,
#[metrics(unit = Count)]
count: u32,
#[metrics(unit = Count)]
fault: u32,
#[metrics(unit = Count)]
error: u32,
#[metrics(unit = Millisecond)]
latency: Timer,
#[metrics(flatten)]
op_detail: Option<OperationMetrics>,
}
#[derive(Clone)]
#[metrics(tag(name = "op_detail"), subfield)]
pub enum OperationMetrics {
Browse(#[metrics(flatten)] BrowseMetrics),
Flamegraph(#[metrics(flatten)] FlamegraphMetrics),
TokioStats(#[metrics(flatten)] TokioStatsMetrics),
}
#[derive(Clone)]
#[metrics(subfield)]
pub struct BrowseMetrics {
objects_returned: usize,
prefixes_fanned_out: usize,
#[metrics(unit = Count)]
truncated: u32,
#[metrics(unit = Count)]
refined: u32,
}
#[derive(Clone)]
#[metrics(subfield)]
pub struct FlamegraphMetrics {
files_matched: u32,
files_folded: u32,
#[metrics(unit = metrique::unit::Percent)]
coverage_pct: f64,
samples: Option<u64>,
}
#[derive(Clone)]
#[metrics(subfield)]
pub struct TokioStatsMetrics {
files_matched: u32,
files_folded: u32,
notable_polls: Option<u64>,
}
#[metrics(emf::dimension_sets = [["operation"], []])]
pub struct ObjectStreamMetrics {
#[metrics(timestamp)]
timestamp: SystemTime,
operation: String,
#[metrics(unit = Count)]
count: u32,
#[metrics(unit = Count)]
pub truncated_mid_stream: u32,
}
impl ObjectStreamMetrics {
pub fn arm(operation: impl Into<String>) -> ObjectStreamMetricsGuard {
ObjectStreamMetrics {
timestamp: SystemTime::now(),
operation: operation.into(),
count: 1,
truncated_mid_stream: 0,
}
.append_on_drop(ServiceMetrics::sink_or_discard())
}
}
#[metrics(emf::dimension_sets = [["operation"], []])]
pub struct SpanStatsStreamMetrics {
#[metrics(timestamp)]
timestamp: SystemTime,
operation: String,
#[metrics(unit = Count)]
count: u32,
#[metrics(unit = Count)]
pub failed: u32,
#[metrics(unit = Millisecond)]
stream_duration: Timer,
#[metrics(unit = Count)]
pub files_matched: u32,
#[metrics(unit = Count)]
pub files_folded: u32,
#[metrics(unit = Count)]
pub files_seeded: u32,
#[metrics(unit = Count)]
pub files_folded_cold: u32,
#[metrics(unit = Millisecond)]
pub download_duration: Duration,
#[metrics(unit = Millisecond)]
pub parse_duration: Duration,
#[metrics(unit = Millisecond)]
pub query_duration: Duration,
#[metrics(unit = Millisecond)]
pub reader_setup_duration: Duration,
#[metrics(unit = Millisecond)]
pub batch_decode_duration: Duration,
#[metrics(unit = Millisecond)]
pub row_materialize_duration: Duration,
#[metrics(unit = Byte)]
pub parquet_bytes: u64,
#[metrics(unit = Count)]
pub record_batches_decoded: u64,
#[metrics(unit = Count)]
pub rows_materialized: u64,
#[metrics(unit = Count)]
pub attribute_entries: u64,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub(crate) struct SpanStatsPhaseDurations {
pub(crate) download: Duration,
pub(crate) parse: Duration,
pub(crate) query: Duration,
pub(crate) reader_setup: Duration,
pub(crate) batch_decode: Duration,
pub(crate) row_materialize: Duration,
pub(crate) parquet_bytes: u64,
pub(crate) record_batches_decoded: u64,
pub(crate) rows_materialized: u64,
pub(crate) attribute_entries: u64,
}
impl std::ops::AddAssign for SpanStatsPhaseDurations {
fn add_assign(&mut self, rhs: Self) {
self.download += rhs.download;
self.parse += rhs.parse;
self.query += rhs.query;
self.reader_setup += rhs.reader_setup;
self.batch_decode += rhs.batch_decode;
self.row_materialize += rhs.row_materialize;
self.parquet_bytes += rhs.parquet_bytes;
self.record_batches_decoded += rhs.record_batches_decoded;
self.rows_materialized += rhs.rows_materialized;
self.attribute_entries += rhs.attribute_entries;
}
}
impl SpanStatsStreamMetrics {
pub fn arm(operation: impl Into<String>) -> SpanStatsStreamMetricsGuard {
SpanStatsStreamMetrics {
timestamp: SystemTime::now(),
operation: operation.into(),
count: 1,
failed: 0,
stream_duration: Timer::start_now(),
files_matched: 0,
files_folded: 0,
files_seeded: 0,
files_folded_cold: 0,
download_duration: Duration::ZERO,
parse_duration: Duration::ZERO,
query_duration: Duration::ZERO,
reader_setup_duration: Duration::ZERO,
batch_decode_duration: Duration::ZERO,
row_materialize_duration: Duration::ZERO,
parquet_bytes: 0,
record_batches_decoded: 0,
rows_materialized: 0,
attribute_entries: 0,
}
.append_on_drop(ServiceMetrics::sink_or_discard())
}
}
#[metrics(emf::dimension_sets = [[]])]
pub struct FoldFileMetrics {
#[metrics(timestamp)]
timestamp: SystemTime,
#[metrics(unit = Count)]
count: u32,
#[metrics(unit = Count)]
failed: u32,
#[metrics(unit = Byte)]
source_bytes: u64,
#[metrics(unit = Byte)]
decompressed_bytes: u64,
#[metrics(unit = Count)]
events_decoded: u64,
#[metrics(unit = Count)]
span_events_decoded: u64,
#[metrics(unit = Millisecond)]
fetch: Duration,
#[metrics(unit = Millisecond)]
gunzip: Duration,
#[metrics(unit = Millisecond)]
wire_decode: Duration,
#[metrics(unit = Millisecond)]
sort_events: Duration,
#[metrics(unit = Millisecond)]
poll_reconstruct: Duration,
#[metrics(unit = Millisecond)]
sample_resolve: Duration,
#[metrics(unit = Millisecond)]
span_resolve: Duration,
#[metrics(unit = Millisecond)]
sample_attribution: Duration,
#[metrics(unit = Millisecond)]
parquet_encode: Duration,
#[metrics(unit = Millisecond)]
write_parts: Duration,
#[metrics(unit = Millisecond)]
total: Duration,
}
#[derive(Default)]
pub struct FoldFileMetricsBuilder {
source_bytes: u64,
decompressed_bytes: u64,
events_decoded: u64,
span_events_decoded: u64,
fetch: Duration,
gunzip: Duration,
wire_decode: Duration,
sort_events: Duration,
poll_reconstruct: Duration,
sample_resolve: Duration,
span_resolve: Duration,
sample_attribution: Duration,
parquet_encode: Duration,
write_parts: Duration,
total: Duration,
failed: bool,
}
impl FoldFileMetricsBuilder {
pub fn new() -> Self {
Self::default()
}
pub fn source_bytes(&mut self, n: u64) -> &mut Self {
self.source_bytes = n;
self
}
pub fn decompressed_bytes(&mut self, n: u64) -> &mut Self {
self.decompressed_bytes = n;
self
}
pub fn fetch(&mut self, d: Duration) -> &mut Self {
self.fetch = d;
self
}
pub fn gunzip(&mut self, d: Duration) -> &mut Self {
self.gunzip = d;
self
}
pub fn decode_phases(&mut self, s: &crate::ingest::decode::DecodeStats) -> &mut Self {
self.wire_decode = s.wire_decode;
self.sort_events = s.sort_events;
self.poll_reconstruct = s.poll_reconstruct;
self.sample_resolve = s.sample_resolve;
self.span_resolve = s.span_resolve;
self.sample_attribution = s.sample_attribution;
self.events_decoded = s.events_decoded;
self.span_events_decoded = s.span_events_decoded;
self
}
pub fn parquet_encode(&mut self, d: Duration) -> &mut Self {
self.parquet_encode = d;
self
}
pub fn write_parts(&mut self, d: Duration) -> &mut Self {
self.write_parts = d;
self
}
pub fn total(&mut self, d: Duration) -> &mut Self {
self.total = d;
self
}
pub fn failed(&mut self, failed: bool) -> &mut Self {
self.failed = failed;
self
}
pub fn emit(self) {
FoldFileMetrics {
timestamp: SystemTime::now(),
count: 1,
failed: self.failed as u32,
source_bytes: self.source_bytes,
decompressed_bytes: self.decompressed_bytes,
events_decoded: self.events_decoded,
span_events_decoded: self.span_events_decoded,
fetch: self.fetch,
gunzip: self.gunzip,
wire_decode: self.wire_decode,
sort_events: self.sort_events,
poll_reconstruct: self.poll_reconstruct,
sample_resolve: self.sample_resolve,
span_resolve: self.span_resolve,
sample_attribution: self.sample_attribution,
parquet_encode: self.parquet_encode,
write_parts: self.write_parts,
total: self.total,
}
.append_on_drop(ServiceMetrics::sink_or_discard());
}
}
impl OperationMetrics {
pub fn browse(
objects_returned: usize,
prefixes_fanned_out: usize,
truncated: bool,
refined: bool,
) -> Self {
Self::Browse(BrowseMetrics {
objects_returned,
prefixes_fanned_out,
truncated: truncated as u32,
refined: refined as u32,
})
}
pub fn flamegraph(files_matched: u32, files_folded: u32, samples: Option<u64>) -> Self {
let coverage_pct = if files_matched == 0 {
0.0
} else {
(files_folded as f64 / files_matched as f64) * 100.0
};
Self::Flamegraph(FlamegraphMetrics {
files_matched,
files_folded,
coverage_pct,
samples,
})
}
pub fn tokio_stats(files_matched: u32, files_folded: u32, notable_polls: Option<u64>) -> Self {
Self::TokioStats(TokioStatsMetrics {
files_matched,
files_folded,
notable_polls,
})
}
}
pub async fn record_request_metrics(req: Request, next: Next) -> Response {
let operation = req
.extensions()
.get::<MatchedPath>()
.map(|p| p.as_str().to_string())
.unwrap_or_else(|| UNMATCHED_OPERATION.to_string());
let session_id = validated_session_id(req.headers());
let mut metrics = RequestMetrics {
timestamp: SystemTime::now(),
operation,
status_code: String::new(),
session_id,
count: 1,
fault: 0,
error: 0,
latency: Timer::start_now(),
op_detail: None,
}
.append_on_drop(ServiceMetrics::sink_or_discard());
let mut response = next.run(req).await;
metrics.latency.stop();
let status = response.status();
metrics.status_code = status.as_u16().to_string();
if status.is_server_error() {
metrics.fault = 1;
} else if status.is_client_error() {
metrics.error = 1;
}
metrics.op_detail = response.extensions_mut().remove::<OperationMetrics>();
response
}
fn validated_session_id(headers: &HeaderMap) -> Option<String> {
let raw = headers.get(SESSION_ID_HEADER)?.as_bytes();
if raw.len() != UUID_TEXT_LEN {
return None;
}
let text = std::str::from_utf8(raw).ok()?;
let parsed = uuid::Uuid::parse_str(text).ok()?;
(parsed.hyphenated().to_string() == text).then(|| text.to_string())
}
pub fn attach_request_metrics(local: bool) -> AttachHandle {
if local {
ServiceMetrics::attach_to_stream(
LocalFormat::new(OutputStyle::Pretty).output_to_makewriter(|| std::io::stdout().lock()),
)
} else {
ServiceMetrics::attach_to_stream(
Emf::builder(CLOUDWATCH_NAMESPACE.to_string(), vec![vec![]])
.build()
.output_to_makewriter(|| std::io::stdout().lock()),
)
}
}
#[cfg(test)]
mod tests {
use super::*;
use axum::Router;
use axum::body::Body;
use axum::http::{Request, StatusCode};
use axum::routing::get;
use metrique::test_util::{TestEntrySink, test_entry_sink};
use tower::ServiceExt;
async fn run_with_session(
route: &str,
uri: &str,
status: StatusCode,
session_id: Option<&str>,
) -> metrique::test_util::TestEntry {
let TestEntrySink { inspector, sink } = test_entry_sink();
let _guard = ServiceMetrics::set_test_sink_on_current_tokio_runtime(sink);
let app: Router = Router::new()
.route(route, get(move || async move { status }))
.layer(axum::middleware::from_fn(record_request_metrics));
let mut request = Request::builder().uri(uri);
if let Some(session_id) = session_id {
request = request.header(SESSION_ID_HEADER, session_id);
}
let resp = app
.oneshot(request.body(Body::empty()).unwrap())
.await
.unwrap();
assert_eq!(resp.status(), status);
let entries = inspector.entries();
assert_eq!(entries.len(), 1, "exactly one metric entry per request");
entries.into_iter().next().unwrap()
}
async fn run(route: &str, uri: &str, status: StatusCode) -> metrique::test_util::TestEntry {
run_with_session(route, uri, status, None).await
}
#[tokio::test]
async fn server_error_is_a_fault() {
let e = run(
"/items/{id}",
"/items/42",
StatusCode::INTERNAL_SERVER_ERROR,
)
.await;
assert_eq!(e.values["operation"], "/items/{id}");
assert_eq!(e.values["status_code"], "500");
assert_eq!(e.metrics["count"].as_u64(), 1);
assert_eq!(e.metrics["fault"].as_u64(), 1);
assert_eq!(e.metrics["error"].as_u64(), 0);
assert_eq!(e.metrics["latency"].num_observations(), 1);
}
#[tokio::test]
async fn client_error_is_an_error_not_a_fault() {
let e = run("/x", "/x", StatusCode::BAD_REQUEST).await;
assert_eq!(e.metrics["fault"].as_u64(), 0);
assert_eq!(e.metrics["error"].as_u64(), 1);
assert_eq!(e.values["status_code"], "400");
}
#[tokio::test]
async fn success_is_neither_fault_nor_error() {
let e = run("/x", "/x", StatusCode::OK).await;
assert_eq!(e.metrics["fault"].as_u64(), 0);
assert_eq!(e.metrics["error"].as_u64(), 0);
assert_eq!(e.metrics["count"].as_u64(), 1);
}
#[tokio::test]
async fn valid_session_id_is_a_plain_log_field() {
let session_id = "123e4567-e89b-42d3-a456-426614174000";
let e = run_with_session("/x", "/x", StatusCode::OK, Some(session_id)).await;
assert_eq!(e.values["session_id"], session_id);
assert!(!e.metrics.contains_key("session_id"));
}
#[tokio::test]
async fn missing_or_invalid_session_id_is_accepted_and_omitted() {
let missing = run("/x", "/x", StatusCode::OK).await;
assert!(!missing.values.contains_key("session_id"));
let invalid = run_with_session(
"/x",
"/x",
StatusCode::OK,
Some("00000000-0000-0000-0000-00000000000z"),
)
.await;
assert!(!invalid.values.contains_key("session_id"));
let oversized = run_with_session(
"/x",
"/x",
StatusCode::OK,
Some("123e4567-e89b-42d3-a456-426614174000-extra"),
)
.await;
assert!(!oversized.values.contains_key("session_id"));
}
#[tokio::test]
async fn no_op_detail_when_handler_attaches_none() {
let e = run("/x", "/x", StatusCode::OK).await;
assert!(!e.values.contains_key("op_detail"));
assert!(!e.metrics.contains_key("truncated"));
}
async fn run_with_op(route: &str, op: OperationMetrics) -> metrique::test_util::TestEntry {
let TestEntrySink { inspector, sink } = test_entry_sink();
let _guard = ServiceMetrics::set_test_sink_on_current_tokio_runtime(sink);
let app: Router = Router::new()
.route(route, get(move || async move { axum::Extension(op) }))
.layer(axum::middleware::from_fn(record_request_metrics));
app.oneshot(Request::builder().uri(route).body(Body::empty()).unwrap())
.await
.unwrap();
let entries = inspector.entries();
assert_eq!(entries.len(), 1, "exactly one metric entry per request");
entries.into_iter().next().unwrap()
}
#[tokio::test]
async fn browse_op_detail_folds_into_entry() {
let e = run_with_op("/api/browse", OperationMetrics::browse(7, 4, true, false)).await;
assert_eq!(e.values["op_detail"], "Browse");
assert_eq!(e.metrics["objects_returned"].as_u64(), 7);
assert_eq!(e.metrics["prefixes_fanned_out"].as_u64(), 4);
assert_eq!(e.metrics["truncated"].as_u64(), 1);
assert_eq!(e.metrics["refined"].as_u64(), 0);
assert_eq!(e.metrics["count"].as_u64(), 1);
}
#[tokio::test]
async fn flamegraph_op_detail_computes_coverage() {
let e = run_with_op(
"/api/flamegraph",
OperationMetrics::flamegraph(200, 50, Some(9000)),
)
.await;
assert_eq!(e.values["op_detail"], "Flamegraph");
assert_eq!(e.metrics["files_matched"].as_u64(), 200);
assert_eq!(e.metrics["files_folded"].as_u64(), 50);
assert_eq!(e.metrics["coverage_pct"].as_f64(), 25.0);
assert_eq!(e.metrics["samples"].as_u64(), 9000);
}
#[tokio::test]
async fn flamegraph_unknown_samples_are_absent_not_zero() {
let e = run_with_op(
"/api/flamegraph",
OperationMetrics::flamegraph(200, 50, None),
)
.await;
assert!(!e.metrics.contains_key("samples"));
assert_eq!(e.metrics["coverage_pct"].as_f64(), 25.0);
}
#[tokio::test]
async fn flamegraph_coverage_is_zero_when_nothing_matched() {
let e = run_with_op("/api/flamegraph", OperationMetrics::flamegraph(0, 0, None)).await;
assert_eq!(e.metrics["coverage_pct"].as_f64(), 0.0);
}
#[test]
fn span_stats_phase_durations_add_each_phase_independently() {
let mut total = SpanStatsPhaseDurations {
download: Duration::from_millis(10),
parse: Duration::from_millis(20),
query: Duration::from_millis(30),
reader_setup: Duration::from_millis(5),
batch_decode: Duration::from_millis(8),
row_materialize: Duration::from_millis(7),
parquet_bytes: 1000,
record_batches_decoded: 2,
rows_materialized: 50,
attribute_entries: 100,
};
total += SpanStatsPhaseDurations {
download: Duration::from_millis(1),
parse: Duration::from_millis(2),
query: Duration::from_millis(3),
reader_setup: Duration::from_millis(1),
batch_decode: Duration::from_millis(1),
row_materialize: Duration::from_millis(0),
parquet_bytes: 500,
record_batches_decoded: 1,
rows_materialized: 25,
attribute_entries: 40,
};
assert_eq!(total.download, Duration::from_millis(11));
assert_eq!(total.parse, Duration::from_millis(22));
assert_eq!(total.query, Duration::from_millis(33));
assert_eq!(total.reader_setup, Duration::from_millis(6));
assert_eq!(total.batch_decode, Duration::from_millis(9));
assert_eq!(total.row_materialize, Duration::from_millis(7));
assert_eq!(total.parquet_bytes, 1500);
assert_eq!(total.record_batches_decoded, 3);
assert_eq!(total.rows_materialized, 75);
assert_eq!(total.attribute_entries, 140);
}
#[tokio::test]
async fn object_stream_metric_is_a_separate_entry() {
let TestEntrySink { inspector, sink } = test_entry_sink();
let _guard = ServiceMetrics::set_test_sink_on_current_tokio_runtime(sink);
let mut guard = ObjectStreamMetrics::arm("/api/object");
guard.truncated_mid_stream = 1;
drop(guard);
let entries = inspector.entries();
assert_eq!(entries.len(), 1);
let e = &entries[0];
assert_eq!(e.values["operation"], "/api/object");
assert_eq!(e.metrics["count"].as_u64(), 1);
assert_eq!(e.metrics["truncated_mid_stream"].as_u64(), 1);
}
#[tokio::test]
async fn object_stream_metric_clean_when_not_truncated() {
let TestEntrySink { inspector, sink } = test_entry_sink();
let _guard = ServiceMetrics::set_test_sink_on_current_tokio_runtime(sink);
drop(ObjectStreamMetrics::arm("/api/object"));
let e = inspector.get(0);
assert_eq!(e.metrics["truncated_mid_stream"].as_u64(), 0);
}
#[tokio::test]
async fn span_stats_stream_metric_carries_duration_and_final_coverage() {
let TestEntrySink { inspector, sink } = test_entry_sink();
let _guard = ServiceMetrics::set_test_sink_on_current_tokio_runtime(sink);
let mut guard = SpanStatsStreamMetrics::arm("/api/span-stats");
guard.files_matched = 1678;
guard.files_seeded = 154;
guard.files_folded_cold = 20;
guard.files_folded = 174;
guard.download_duration = Duration::from_millis(120);
guard.parse_duration = Duration::from_millis(230);
guard.query_duration = Duration::from_millis(340);
guard.reader_setup_duration = Duration::from_millis(50);
guard.batch_decode_duration = Duration::from_millis(80);
guard.row_materialize_duration = Duration::from_millis(100);
guard.parquet_bytes = 12_345_678;
guard.record_batches_decoded = 42;
guard.rows_materialized = 150_000;
guard.attribute_entries = 300_000;
tokio::time::sleep(Duration::from_millis(5)).await;
drop(guard);
let entries = inspector.entries();
assert_eq!(entries.len(), 1);
let e = &entries[0];
assert_eq!(e.values["operation"], "/api/span-stats");
assert_eq!(e.metrics["count"].as_u64(), 1);
assert_eq!(e.metrics["failed"].as_u64(), 0);
assert_eq!(e.metrics["files_matched"].as_u64(), 1678);
assert_eq!(e.metrics["files_folded"].as_u64(), 174);
assert_eq!(e.metrics["files_seeded"].as_u64(), 154);
assert_eq!(e.metrics["files_folded_cold"].as_u64(), 20);
assert_eq!(e.metrics["download_duration"].num_observations(), 1);
assert_eq!(e.metrics["parse_duration"].num_observations(), 1);
assert_eq!(e.metrics["query_duration"].num_observations(), 1);
assert_eq!(e.metrics["reader_setup_duration"].num_observations(), 1);
assert_eq!(e.metrics["batch_decode_duration"].num_observations(), 1);
assert_eq!(e.metrics["row_materialize_duration"].num_observations(), 1);
assert_eq!(e.metrics["parquet_bytes"].as_u64(), 12_345_678);
assert_eq!(e.metrics["record_batches_decoded"].as_u64(), 42);
assert_eq!(e.metrics["rows_materialized"].as_u64(), 150_000);
assert_eq!(e.metrics["attribute_entries"].as_u64(), 300_000);
assert_eq!(e.metrics["stream_duration"].num_observations(), 1);
}
#[tokio::test]
async fn span_stats_stream_metric_flags_failure() {
let TestEntrySink { inspector, sink } = test_entry_sink();
let _guard = ServiceMetrics::set_test_sink_on_current_tokio_runtime(sink);
let mut guard = SpanStatsStreamMetrics::arm("/api/span-stats");
guard.failed = 1;
drop(guard);
let e = inspector.get(0);
assert_eq!(e.metrics["failed"].as_u64(), 1);
}
#[tokio::test]
async fn fold_file_metric_carries_phases_size_and_counts() {
let TestEntrySink { inspector, sink } = test_entry_sink();
let _guard = ServiceMetrics::set_test_sink_on_current_tokio_runtime(sink);
let mut b = FoldFileMetricsBuilder::new();
b.source_bytes(1_257_686)
.decompressed_bytes(2_314_769)
.fetch(Duration::from_millis(120))
.gunzip(Duration::from_millis(5))
.parquet_encode(Duration::from_millis(8))
.write_parts(Duration::from_millis(30));
let decode = crate::ingest::decode::DecodeStats {
events_decoded: 1_300_000,
span_events_decoded: 4_368,
wire_decode: Duration::from_millis(20),
span_resolve: Duration::from_millis(3),
..Default::default()
};
b.decode_phases(&decode).total(Duration::from_millis(200));
b.emit();
let e = inspector.get(0);
assert_eq!(e.metrics["count"].as_u64(), 1);
assert_eq!(e.metrics["failed"].as_u64(), 0);
assert_eq!(e.metrics["source_bytes"].as_u64(), 1_257_686);
assert_eq!(e.metrics["decompressed_bytes"].as_u64(), 2_314_769);
assert_eq!(e.metrics["events_decoded"].as_u64(), 1_300_000);
assert_eq!(e.metrics["span_events_decoded"].as_u64(), 4_368);
assert_eq!(e.metrics["wire_decode"].num_observations(), 1);
assert_eq!(e.metrics["span_resolve"].num_observations(), 1);
assert_eq!(e.metrics["total"].num_observations(), 1);
}
#[tokio::test]
async fn fold_file_metric_marks_failure() {
let TestEntrySink { inspector, sink } = test_entry_sink();
let _guard = ServiceMetrics::set_test_sink_on_current_tokio_runtime(sink);
let mut b = FoldFileMetricsBuilder::new();
b.failed(true).total(Duration::from_millis(50));
b.emit();
let e = inspector.get(0);
assert_eq!(e.metrics["failed"].as_u64(), 1);
}
}