use std::borrow::Cow;
use std::cell::{Cell, RefCell};
use std::collections::{HashMap, HashSet, VecDeque};
use std::rc::Rc;
use std::sync::{Arc, Mutex};
use std::time::Instant;
use otel_arrow_dfe_telemetry::common_attributes::SignalAttributes;
use otel_arrow_dfe_telemetry::instrument::{Counter, HistogramDetailed, HistogramNormal, Mmsc};
#[cfg(test)]
use otel_arrow_dfe_telemetry::metrics::MetricSetHandler;
use otel_arrow_dfe_telemetry::metrics::{MeasurementMetricSet, MetricSetSnapshot};
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use otel_arrow_dfe_telemetry_macros::{attribute_set, metric_set};
use crate::attributes::PipelineAttributeSet;
use crate::context::PipelineContext;
use otel_arrow_dfe_config::SignalType;
use otel_arrow_dfe_config::policy::{DistributionTier, FlowMetric, TelemetryPolicy};
#[metric_set(name = "flow.input", measurement_attributes = SignalAttributes)]
#[derive(Debug, Default, Clone)]
pub struct FlowInputMessageMetrics {
#[metric(unit = "{message}")]
pub messages: Counter<u64>,
}
#[metric_set(name = "flow.input", measurement_attributes = SignalAttributes)]
#[derive(Debug, Default, Clone)]
pub struct FlowInputItemsMetrics {
#[metric(unit = "{item}")]
pub items: Counter<u64>,
}
#[metric_set(name = "flow.input", measurement_attributes = SignalAttributes)]
#[derive(Debug, Default, Clone)]
pub struct FlowInputSizeMetrics {
#[metric(unit = "By")]
pub size: Counter<u64>,
}
#[metric_set(name = "flow.compute", measurement_attributes = SignalAttributes)]
#[derive(Debug, Default, Clone)]
pub struct FlowDurationBasicMetrics {
#[metric(unit = "s")]
pub duration: Mmsc,
}
#[metric_set(name = "flow.compute", measurement_attributes = SignalAttributes)]
#[derive(Debug, Default, Clone)]
pub struct FlowDurationNormalMetrics {
#[metric(unit = "s")]
pub duration: HistogramNormal,
}
#[metric_set(name = "flow.compute", measurement_attributes = SignalAttributes)]
#[derive(Debug, Default, Clone)]
pub struct FlowDurationDetailedMetrics {
#[metric(unit = "s")]
pub duration: HistogramDetailed,
}
#[derive(Clone)]
pub(crate) enum FlowDurationMetricSet {
Basic(MeasurementMetricSet<FlowDurationBasicMetrics>),
Normal(MeasurementMetricSet<FlowDurationNormalMetrics>),
Detailed(MeasurementMetricSet<FlowDurationDetailedMetrics>),
}
impl From<MeasurementMetricSet<FlowDurationNormalMetrics>> for FlowDurationMetricSet {
fn from(metrics: MeasurementMetricSet<FlowDurationNormalMetrics>) -> Self {
Self::Normal(metrics)
}
}
impl FlowDurationMetricSet {
pub(crate) fn into_measurement(self) -> FlowDurationMeasurement {
match self {
Self::Basic(metrics) => FlowDurationMeasurement::Basic {
metrics,
accumulator: Box::new(std::array::from_fn(|_| Mmsc::default())),
},
Self::Normal(metrics) => FlowDurationMeasurement::Normal {
metrics,
accumulator: Box::new(std::array::from_fn(|_| HistogramNormal::default())),
},
Self::Detailed(metrics) => FlowDurationMeasurement::Detailed {
metrics,
accumulator: Box::new(std::array::from_fn(|_| HistogramDetailed::default())),
},
}
}
}
#[derive(Clone)]
pub(crate) enum FlowDurationMeasurement {
Basic {
metrics: MeasurementMetricSet<FlowDurationBasicMetrics>,
accumulator: Box<[Mmsc; FLOW_SIGNAL_COUNT]>,
},
Normal {
metrics: MeasurementMetricSet<FlowDurationNormalMetrics>,
accumulator: Box<[HistogramNormal; FLOW_SIGNAL_COUNT]>,
},
Detailed {
metrics: MeasurementMetricSet<FlowDurationDetailedMetrics>,
accumulator: Box<[HistogramDetailed; FLOW_SIGNAL_COUNT]>,
},
}
impl FlowDurationMeasurement {
pub(crate) fn record(&mut self, signal: SignalType, value: f64) {
let index = flow_signal_index(signal);
match self {
Self::Basic { accumulator, .. } => accumulator[index].record(value),
Self::Normal { accumulator, .. } => accumulator[index].record(value),
Self::Detailed { accumulator, .. } => accumulator[index].record(value),
}
}
fn merge_pending(&mut self) {
match self {
Self::Basic {
metrics,
accumulator,
} => {
for (duration, signal) in accumulator.iter_mut().zip(FLOW_SIGNALS) {
if !duration.is_empty() {
metrics
.with(SignalAttributes { signal })
.duration
.merge(std::mem::take(duration));
}
}
}
Self::Normal {
metrics,
accumulator,
} => {
for (duration, signal) in accumulator.iter_mut().zip(FLOW_SIGNALS) {
if !duration.is_empty() {
metrics
.with(SignalAttributes { signal })
.duration
.merge(std::mem::take(duration));
}
}
}
Self::Detailed {
metrics,
accumulator,
} => {
for (duration, signal) in accumulator.iter_mut().zip(FLOW_SIGNALS) {
if !duration.is_empty() {
metrics
.with(SignalAttributes { signal })
.duration
.merge(std::mem::take(duration));
}
}
}
}
}
pub(crate) fn report(&mut self, reporter: &mut MetricsReporter) {
self.merge_pending();
match self {
Self::Basic { metrics, .. } => {
let _ = reporter.report_measurement(metrics);
}
Self::Normal { metrics, .. } => {
let _ = reporter.report_measurement(metrics);
}
Self::Detailed { metrics, .. } => {
let _ = reporter.report_measurement(metrics);
}
}
}
pub(crate) fn terminal_snapshots(&mut self) -> Vec<MetricSetSnapshot> {
self.merge_pending();
match self {
Self::Basic { metrics, .. } => metrics.terminal_snapshots(),
Self::Normal { metrics, .. } => metrics.terminal_snapshots(),
Self::Detailed { metrics, .. } => metrics.terminal_snapshots(),
}
}
}
#[metric_set(name = "flow.output", measurement_attributes = SignalAttributes)]
#[derive(Debug, Default, Clone)]
pub struct FlowOutputMessageMetrics {
#[metric(unit = "{message}")]
pub messages: Counter<u64>,
}
#[metric_set(name = "flow.output", measurement_attributes = SignalAttributes)]
#[derive(Debug, Default, Clone)]
pub struct FlowOutputItemsMetrics {
#[metric(unit = "{item}")]
pub items: Counter<u64>,
}
#[metric_set(name = "flow.output", measurement_attributes = SignalAttributes)]
#[derive(Debug, Default, Clone)]
pub struct FlowOutputSizeMetrics {
#[metric(unit = "By")]
pub size: Counter<u64>,
}
#[metric_set(name = "flow.dropped", measurement_attributes = SignalAttributes)]
#[derive(Debug, Default, Clone)]
pub struct FlowDroppedItemsMetrics {
#[metric(unit = "{item}")]
pub items: Counter<u64>,
}
#[attribute_set(scope, name = "flow.attrs")]
#[derive(Debug, Clone, Default, Hash)]
pub struct FlowAttributeSet {
#[attribute_key = "flow.id"]
pub flow_id: Cow<'static, str>,
#[attribute_key = "flow.node.start"]
pub start_node: Cow<'static, str>,
#[attribute_key = "flow.node.end"]
pub end_node: Cow<'static, str>,
#[attribute_key = "flow.purpose"]
pub purpose: Cow<'static, str>,
#[attribute_key = "flow.node.decision"]
pub decision: Cow<'static, str>,
#[compose]
pub pipeline_attrs: PipelineAttributeSet,
}
pub type FlowMetricId = usize;
pub(crate) const FLOW_SIGNAL_COUNT: usize = 3;
pub(crate) type FlowItemAccumulator = [u64; FLOW_SIGNAL_COUNT];
#[must_use]
pub(crate) const fn flow_signal_index(signal: SignalType) -> usize {
match signal {
SignalType::Traces => 0,
SignalType::Metrics => 1,
SignalType::Logs => 2,
}
}
pub(crate) const FLOW_SIGNALS: [SignalType; FLOW_SIGNAL_COUNT] =
[SignalType::Traces, SignalType::Metrics, SignalType::Logs];
bitflags::bitflags! {
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct FlowMetricInterests: u8 {
const COMPUTE_DURATION = 1 << 0;
const INPUT_MESSAGES = 1 << 1;
const INPUT_ITEMS = 1 << 2;
const INPUT_SIZE = 1 << 3;
const OUTPUT_MESSAGES = 1 << 4;
const OUTPUT_ITEMS = 1 << 5;
const OUTPUT_SIZE = 1 << 6;
const DROPPED_ITEMS = 1 << 7;
}
}
pub(crate) struct PipelineFlowMetricState {
pub input_message_metrics: Vec<Option<MeasurementMetricSet<FlowInputMessageMetrics>>>,
pub input_items_metrics: Vec<Option<MeasurementMetricSet<FlowInputItemsMetrics>>>,
pub input_size_metrics: Vec<Option<MeasurementMetricSet<FlowInputSizeMetrics>>>,
pub duration_metrics: Vec<Option<FlowDurationMetricSet>>,
pub output_items_metrics: Vec<Option<MeasurementMetricSet<FlowOutputItemsMetrics>>>,
pub output_message_metrics: Vec<Option<MeasurementMetricSet<FlowOutputMessageMetrics>>>,
pub output_size_metrics: Vec<Option<MeasurementMetricSet<FlowOutputSizeMetrics>>>,
pub end_nodes: HashMap<usize, usize>,
pub start_nodes: HashMap<usize, usize>,
pub decision_candidates: HashMap<usize, DecisionCandidate>,
}
#[derive(Clone)]
pub(crate) struct DecisionCandidate {
pub attrs: FlowAttributeSet,
}
#[derive(Clone)]
pub(crate) struct InputFlowMetrics<ItemAccumulator> {
pub input_messages: Option<(
MeasurementMetricSet<FlowInputMessageMetrics>,
ItemAccumulator,
)>,
pub input_items: Option<(MeasurementMetricSet<FlowInputItemsMetrics>, ItemAccumulator)>,
pub input_size: Option<(MeasurementMetricSet<FlowInputSizeMetrics>, ItemAccumulator)>,
}
#[derive(Clone)]
pub(crate) struct EndFlowMetrics<DurationMeasurement, ItemAccumulator> {
pub duration: Option<DurationMeasurement>,
pub output_items: Option<(
MeasurementMetricSet<FlowOutputItemsMetrics>,
ItemAccumulator,
)>,
pub output_messages: Option<(
MeasurementMetricSet<FlowOutputMessageMetrics>,
ItemAccumulator,
)>,
pub output_size: Option<(MeasurementMetricSet<FlowOutputSizeMetrics>, ItemAccumulator)>,
}
#[derive(Clone)]
pub(crate) struct DecisionFlowMetrics<ItemAccumulator> {
pub dropped_items: Option<(
MeasurementMetricSet<FlowDroppedItemsMetrics>,
ItemAccumulator,
)>,
}
#[derive(Clone)]
pub(crate) struct FlowMetricState<Marker, DurationMeasurement, ItemAccumulator> {
pub last_send_marker: Marker,
pub is_start: bool,
pub is_end: bool,
pub is_decision: bool,
pub interests: FlowMetricInterests,
pub active: bool,
pub needs_timing: bool,
pub input: InputFlowMetrics<ItemAccumulator>,
pub end: EndFlowMetrics<DurationMeasurement, ItemAccumulator>,
pub decision: DecisionFlowMetrics<ItemAccumulator>,
}
pub(crate) type LocalFlowMetricState = FlowMetricState<
Rc<Cell<Option<Instant>>>,
RefCell<FlowDurationMeasurement>,
Cell<FlowItemAccumulator>,
>;
pub(crate) type SharedFlowMetricState = FlowMetricState<
Arc<Mutex<Option<Instant>>>,
Arc<Mutex<FlowDurationMeasurement>>,
Arc<Mutex<FlowItemAccumulator>>,
>;
impl Default for LocalFlowMetricState {
fn default() -> Self {
Self {
last_send_marker: Rc::new(Cell::new(None)),
is_start: false,
is_end: false,
is_decision: false,
interests: FlowMetricInterests::empty(),
active: false,
needs_timing: false,
input: InputFlowMetrics {
input_messages: None,
input_items: None,
input_size: None,
},
end: EndFlowMetrics {
duration: None,
output_messages: None,
output_items: None,
output_size: None,
},
decision: DecisionFlowMetrics {
dropped_items: None,
},
}
}
}
impl Default for SharedFlowMetricState {
fn default() -> Self {
Self {
last_send_marker: Arc::new(Mutex::new(None)),
is_start: false,
is_end: false,
is_decision: false,
interests: FlowMetricInterests::empty(),
active: false,
needs_timing: false,
input: InputFlowMetrics {
input_messages: None,
input_items: None,
input_size: None,
},
end: EndFlowMetrics {
duration: None,
output_messages: None,
output_items: None,
output_size: None,
},
decision: DecisionFlowMetrics {
dropped_items: None,
},
}
}
}
pub(crate) fn build_flow_metric_state(
telemetry_policy: &TelemetryPolicy,
node_name_to_index: &HashMap<String, usize>,
processor_indices: &HashSet<usize>,
pipeline_context: &PipelineContext,
pipeline_connections: &[(usize, usize)],
) -> Result<PipelineFlowMetricState, crate::error::Error> {
let mut input_message_metrics = Vec::new();
let mut input_items_metrics = Vec::new();
let mut input_size_metrics = Vec::new();
let mut duration_metrics = Vec::new();
let mut output_items_metrics = Vec::new();
let mut output_message_metrics = Vec::new();
let mut output_size_metrics = Vec::new();
let mut end_nodes: HashMap<usize, usize> = HashMap::new();
let mut start_nodes: HashMap<usize, usize> = HashMap::new();
let mut decision_candidates: HashMap<usize, DecisionCandidate> = HashMap::new();
let pipeline_attrs = pipeline_context.pipeline_attribute_set();
let adjacency = build_adjacency(pipeline_connections);
let mut resolved_ranges: Vec<(usize, usize, String)> = Vec::new();
for flow_config in &telemetry_policy.flow_metrics {
let start_idx = node_name_to_index
.get(&flow_config.bounds.start_node)
.copied();
let end_idx = node_name_to_index
.get(&flow_config.bounds.end_node)
.copied();
let (Some(start_idx), Some(end_idx)) = (start_idx, end_idx) else {
return Err(invalid_flow_metric_config(format!(
"flow metric `{}` references unknown node(s): start=`{}`, end=`{}`",
flow_config.id, flow_config.bounds.start_node, flow_config.bounds.end_node
)));
};
if !processor_indices.contains(&start_idx) || !processor_indices.contains(&end_idx) {
return Err(invalid_flow_metric_config(format!(
"flow metric `{}` start/end nodes must be processors: start=`{}`, end=`{}`",
flow_config.id, flow_config.bounds.start_node, flow_config.bounds.end_node
)));
}
if start_nodes.contains_key(&start_idx)
|| end_nodes.contains_key(&start_idx)
|| start_nodes.contains_key(&end_idx)
|| end_nodes.contains_key(&end_idx)
{
return Err(invalid_flow_metric_config(format!(
"flow metric `{}` overlaps with another flow metric (non-overlapping ranges only): start=`{}`, end=`{}`",
flow_config.id, flow_config.bounds.start_node, flow_config.bounds.end_node
)));
}
let attrs = FlowAttributeSet {
flow_id: Cow::Owned(flow_config.id.clone()),
start_node: Cow::Owned(flow_config.bounds.start_node.clone()),
end_node: Cow::Owned(flow_config.bounds.end_node.clone()),
purpose: flow_config
.purpose
.clone()
.map_or(Cow::Borrowed(""), Cow::Owned),
decision: Cow::Borrowed(""),
pipeline_attrs: pipeline_attrs.clone(),
};
let entity_key = pipeline_context
.metrics_registry()
.register_entity(attrs.clone());
let input_message_metric = flow_config.has(FlowMetric::InputMessages).then(|| {
FlowInputMessageMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(entity_key),
)
});
let input_items_metric = flow_config.has(FlowMetric::InputItems).then(|| {
FlowInputItemsMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(entity_key),
)
});
let input_size_metric = flow_config.has(FlowMetric::InputSize).then(|| {
FlowInputSizeMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(entity_key),
)
});
let duration_metric = if flow_config.has(FlowMetric::ComputeDuration) {
let registrar = pipeline_context.metric_set_registrar_for_entity(entity_key);
Some(match flow_config.duration_distribution {
DistributionTier::Basic => {
FlowDurationMetricSet::Basic(FlowDurationBasicMetrics::register(®istrar))
}
DistributionTier::Normal => {
FlowDurationMetricSet::Normal(FlowDurationNormalMetrics::register(®istrar))
}
DistributionTier::Detailed => FlowDurationMetricSet::Detailed(
FlowDurationDetailedMetrics::register(®istrar),
),
})
} else {
None
};
let output_message_metric = flow_config.has(FlowMetric::OutputMessages).then(|| {
FlowOutputMessageMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(entity_key),
)
});
let output_items_metric = flow_config.has(FlowMetric::OutputItems).then(|| {
FlowOutputItemsMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(entity_key),
)
});
let output_size_metric = flow_config.has(FlowMetric::OutputSize).then(|| {
FlowOutputSizeMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(entity_key),
)
});
if flow_config.has(FlowMetric::DroppedItems) {
let (mut range_nodes, _) = active_range(start_idx, end_idx, &adjacency);
let _ = range_nodes.insert(start_idx);
let _ = range_nodes.insert(end_idx);
for node_idx in range_nodes {
if processor_indices.contains(&node_idx) {
if let Some(existing) = decision_candidates.get(&node_idx) {
return Err(invalid_flow_metric_config(format!(
"flow metric `{}` shares decision node `{}` with flow metric `{}`: a decision node may belong to at most one flow",
flow_config.id,
node_name_to_index
.iter()
.find(|&(_, &idx)| idx == node_idx)
.map_or("<unknown>", |(name, _)| name.as_str()),
existing.attrs.flow_id,
)));
}
let _ = decision_candidates.insert(
node_idx,
DecisionCandidate {
attrs: attrs.clone(),
},
);
}
}
}
let id = duration_metrics.len();
input_message_metrics.push(input_message_metric);
input_items_metrics.push(input_items_metric);
input_size_metrics.push(input_size_metric);
duration_metrics.push(duration_metric);
output_items_metrics.push(output_items_metric);
output_message_metrics.push(output_message_metric);
output_size_metrics.push(output_size_metric);
let _ = end_nodes.insert(end_idx, id);
let _ = start_nodes.insert(start_idx, id);
resolved_ranges.push((start_idx, end_idx, flow_config.id.clone()));
}
if !resolved_ranges.is_empty() {
validate_metric_ranges(&resolved_ranges, &adjacency)?;
}
Ok(PipelineFlowMetricState {
input_message_metrics,
input_items_metrics,
input_size_metrics,
duration_metrics,
output_items_metrics,
output_message_metrics,
output_size_metrics,
end_nodes,
start_nodes,
decision_candidates,
})
}
fn invalid_flow_metric_config(error: String) -> crate::error::Error {
crate::error::Error::ConfigError(Box::new(
otel_arrow_dfe_config::error::Error::InvalidUserConfig { error },
))
}
fn active_range(
start: usize,
end: usize,
adjacency: &HashMap<usize, Vec<usize>>,
) -> (HashSet<usize>, bool) {
let mut visited = HashSet::new();
let mut queue = VecDeque::new();
let _ = visited.insert(start);
queue.push_back(start);
while let Some(node) = queue.pop_front() {
if node == end {
continue;
}
if let Some(neighbors) = adjacency.get(&node) {
for &next in neighbors {
if visited.insert(next) {
queue.push_back(next);
}
}
}
}
let end_reachable = visited.remove(&end);
(visited, end_reachable)
}
fn validate_metric_ranges(
ranges: &[(usize, usize, String)],
adjacency: &HashMap<usize, Vec<usize>>,
) -> Result<(), crate::error::Error> {
let mut active_ranges: Vec<HashSet<usize>> = Vec::with_capacity(ranges.len());
for &(start, end, ref name) in ranges {
let (set, end_reachable) = active_range(start, end, adjacency);
if !end_reachable {
return Err(invalid_flow_metric_config(format!(
"flow metric `{name}` end node is not reachable from its start node"
)));
}
active_ranges.push(set);
}
for i in 0..ranges.len() {
for j in (i + 1)..ranges.len() {
let (start_j, _, ref name_j) = ranges[j];
if active_ranges[i].contains(&start_j) {
return Err(invalid_flow_metric_config(format!(
"flow metric `{}` interleaves with `{}`: \
start node of `{}` is reachable from start of `{}` \
before reaching its end node",
name_j, ranges[i].2, name_j, ranges[i].2,
)));
}
let (start_i, _, ref name_i) = ranges[i];
if active_ranges[j].contains(&start_i) {
return Err(invalid_flow_metric_config(format!(
"flow metric `{}` interleaves with `{}`: \
start node of `{}` is reachable from start of `{}` \
before reaching its end node",
name_i, name_j, name_i, name_j,
)));
}
}
}
Ok(())
}
fn build_adjacency(edges: &[(usize, usize)]) -> HashMap<usize, Vec<usize>> {
let mut adj: HashMap<usize, Vec<usize>> = HashMap::new();
for &(src, dst) in edges {
adj.entry(src).or_default().push(dst);
}
adj
}
#[inline]
#[must_use]
pub fn nanos_u64(ns: u128) -> u64 {
ns.min(u128::from(u64::MAX)) as u64
}
impl PipelineFlowMetricState {
#[must_use]
#[allow(dead_code)]
pub fn empty() -> Self {
Self {
input_message_metrics: Vec::new(),
input_items_metrics: Vec::new(),
input_size_metrics: Vec::new(),
duration_metrics: Vec::new(),
output_items_metrics: Vec::new(),
output_message_metrics: Vec::new(),
output_size_metrics: Vec::new(),
end_nodes: HashMap::new(),
start_nodes: HashMap::new(),
decision_candidates: HashMap::new(),
}
}
#[must_use]
pub fn is_active(&self) -> bool {
self.input_message_metrics.iter().any(Option::is_some)
|| self.input_items_metrics.iter().any(Option::is_some)
|| self.input_size_metrics.iter().any(Option::is_some)
|| self.duration_metrics.iter().any(Option::is_some)
|| self.output_items_metrics.iter().any(Option::is_some)
|| self.output_message_metrics.iter().any(Option::is_some)
|| self.output_size_metrics.iter().any(Option::is_some)
|| !self.decision_candidates.is_empty()
}
#[must_use]
pub fn needs_timing(&self) -> bool {
self.duration_metrics.iter().any(Option::is_some)
}
}
#[cfg(test)]
impl FlowDurationMeasurement {
pub(crate) fn pending_summary(&self, signal: SignalType) -> (u64, f64, f64, f64) {
let index = flow_signal_index(signal);
match self {
Self::Basic { accumulator, .. } => {
let value = accumulator[index].get();
(value.count, value.sum, value.min, value.max)
}
Self::Normal { accumulator, .. } => accumulator[index].get().summary(),
Self::Detailed { accumulator, .. } => accumulator[index].get().summary(),
}
}
pub(crate) fn reported_summary(&self, signal: SignalType) -> (u64, f64, f64, f64) {
match self {
Self::Basic { metrics, .. } => {
let value = metrics.get(SignalAttributes { signal }).duration.get();
(value.count, value.sum, value.min, value.max)
}
Self::Normal { metrics, .. } => metrics
.get(SignalAttributes { signal })
.duration
.get()
.summary(),
Self::Detailed { metrics, .. } => metrics
.get(SignalAttributes { signal })
.duration
.get()
.summary(),
}
}
pub(crate) fn is_empty(&self, signal: SignalType) -> bool {
self.pending_summary(signal).0 == 0
}
pub(crate) fn accumulator_address(&self) -> *const () {
match self {
Self::Basic { accumulator, .. } => accumulator.as_ptr().cast(),
Self::Normal { accumulator, .. } => accumulator.as_ptr().cast(),
Self::Detailed { accumulator, .. } => accumulator.as_ptr().cast(),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::testing::test_pipeline_ctx;
use otel_arrow_dfe_telemetry::attributes::{AttributeSetHandler, AttributeValue};
use otel_arrow_dfe_telemetry::descriptor::Instrument;
use otel_arrow_dfe_telemetry::metrics::MetricValue;
fn duration_instrument(metrics: &FlowDurationMetricSet) -> Instrument {
let attrs = SignalAttributes {
signal: SignalType::Logs,
};
match metrics {
FlowDurationMetricSet::Basic(metrics) => {
metrics.get(attrs).descriptor().metrics[0].instrument
}
FlowDurationMetricSet::Normal(metrics) => {
metrics.get(attrs).descriptor().metrics[0].instrument
}
FlowDurationMetricSet::Detailed(metrics) => {
metrics.get(attrs).descriptor().metrics[0].instrument
}
}
}
fn one_flow_metric_state(tier: DistributionTier) -> PipelineFlowMetricState {
let (ctx, _) = test_pipeline_ctx();
let entity_key = ctx
.metrics_registry()
.register_entity(FlowAttributeSet::default());
let registrar = ctx.metric_set_registrar_for_entity(entity_key);
let input_items_metric = FlowInputItemsMetrics::register(®istrar);
let duration_metric = match tier {
DistributionTier::Basic => {
FlowDurationMetricSet::Basic(FlowDurationBasicMetrics::register(®istrar))
}
DistributionTier::Normal => {
FlowDurationMetricSet::Normal(FlowDurationNormalMetrics::register(®istrar))
}
DistributionTier::Detailed => {
FlowDurationMetricSet::Detailed(FlowDurationDetailedMetrics::register(®istrar))
}
};
let output_items_metric = FlowOutputItemsMetrics::register(®istrar);
PipelineFlowMetricState {
input_message_metrics: vec![None],
input_items_metrics: vec![Some(input_items_metric)],
input_size_metrics: vec![None],
duration_metrics: vec![Some(duration_metric)],
output_items_metrics: vec![Some(output_items_metric)],
output_message_metrics: vec![None],
output_size_metrics: vec![None],
end_nodes: HashMap::from([(2, 0)]),
start_nodes: HashMap::from([(0, 0)]),
decision_candidates: HashMap::new(),
}
}
#[test]
fn empty_state_is_inactive() {
let state = PipelineFlowMetricState::empty();
assert!(!state.is_active());
}
#[test]
fn nonempty_state_is_active() {
let state = one_flow_metric_state(DistributionTier::Normal);
assert!(state.is_active());
}
#[test]
fn duration_tiers_preserve_summary_statistics() {
for tier in [
DistributionTier::Basic,
DistributionTier::Normal,
DistributionTier::Detailed,
] {
let mut state = one_flow_metric_state(tier);
let metrics = state.duration_metrics[0].take().unwrap();
let mut duration = metrics.into_measurement();
duration.record(SignalType::Logs, 100.0);
duration.record(SignalType::Logs, 200.0);
let (count, sum, min, max) = duration.pending_summary(SignalType::Logs);
assert_eq!(count, 2, "tier: {tier:?}");
assert!((min - 100.0).abs() < f64::EPSILON, "tier: {tier:?}");
assert!((max - 200.0).abs() < f64::EPSILON, "tier: {tier:?}");
assert!((sum - 300.0).abs() < f64::EPSILON, "tier: {tier:?}");
}
}
#[test]
fn duration_tiers_register_matching_descriptors() {
for (tier, expected) in [
(DistributionTier::Basic, Instrument::Mmsc),
(DistributionTier::Normal, Instrument::ExponentialHistogram),
(DistributionTier::Detailed, Instrument::ExponentialHistogram),
] {
let state = one_flow_metric_state(tier);
let duration = state.duration_metrics[0].as_ref().unwrap();
assert_eq!(duration_instrument(duration), expected, "tier: {tier:?}");
}
}
#[test]
fn duration_tiers_report_matching_distribution_values() {
for (tier, expected_instrument, expected_tier) in [
(DistributionTier::Basic, Instrument::Mmsc, "basic"),
(
DistributionTier::Normal,
Instrument::ExponentialHistogram,
"normal",
),
(
DistributionTier::Detailed,
Instrument::ExponentialHistogram,
"detailed",
),
] {
let mut state = one_flow_metric_state(tier);
let metrics = state.duration_metrics[0].take().unwrap();
let mut measurement = metrics.into_measurement();
measurement.record(SignalType::Logs, 1.25);
measurement.record(SignalType::Logs, 2.75);
let (snapshot_rx, mut reporter) = MetricsReporter::create_new_and_receiver(1);
measurement.report(&mut reporter);
assert!(measurement.is_empty(SignalType::Logs), "tier: {tier:?}");
let snapshot = snapshot_rx.try_recv().expect("duration snapshot");
assert_eq!(
snapshot.descriptor().metrics[0].instrument,
expected_instrument
);
let [MetricValue::Distribution(value)] = snapshot.get_metrics() else {
panic!("expected one distribution value");
};
assert_eq!(value.tier_name(), expected_tier);
assert_eq!(value.summary(), (2, 4.0, 1.25, 2.75));
}
}
#[test]
fn duration_tiers_reuse_accumulator_allocation() {
for tier in [
DistributionTier::Basic,
DistributionTier::Normal,
DistributionTier::Detailed,
] {
let mut state = one_flow_metric_state(tier);
let metrics = state.duration_metrics[0].take().unwrap();
let mut measurement = metrics.into_measurement();
let accumulator_address = measurement.accumulator_address();
let (_snapshot_rx, mut reporter) = MetricsReporter::create_new_and_receiver(2);
measurement.record(SignalType::Logs, 1.0);
measurement.report(&mut reporter);
assert_eq!(
measurement.accumulator_address(),
accumulator_address,
"tier: {tier:?}"
);
measurement.report(&mut reporter);
assert_eq!(
measurement.accumulator_address(),
accumulator_address,
"tier: {tier:?}"
);
}
}
#[test]
fn direct_record_increments_items() {
let mut state = one_flow_metric_state(DistributionTier::Normal);
state.input_items_metrics[0]
.as_mut()
.unwrap()
.with(SignalAttributes {
signal: SignalType::Logs,
})
.items
.add(10);
state.input_items_metrics[0]
.as_mut()
.unwrap()
.with(SignalAttributes {
signal: SignalType::Logs,
})
.items
.add(20);
state.output_items_metrics[0]
.as_mut()
.unwrap()
.with(SignalAttributes {
signal: SignalType::Logs,
})
.items
.add(7);
state.output_items_metrics[0]
.as_mut()
.unwrap()
.with(SignalAttributes {
signal: SignalType::Logs,
})
.items
.add(8);
let input = state.input_items_metrics[0]
.as_mut()
.unwrap()
.get(SignalAttributes {
signal: SignalType::Logs,
})
.items
.get();
assert_eq!(input, 30);
let output = state.output_items_metrics[0]
.as_mut()
.unwrap()
.get(SignalAttributes {
signal: SignalType::Logs,
})
.items
.get();
assert_eq!(output, 15);
}
use otel_arrow_dfe_config::policy::{
FlowBounds, FlowMetric, FlowMetricConfig, TelemetryPolicy,
};
fn policy_with(flow_metrics: Vec<FlowMetricConfig>) -> TelemetryPolicy {
TelemetryPolicy {
flow_metrics,
..TelemetryPolicy::default()
}
}
fn sw(name: &str, start: &str, stop: &str) -> FlowMetricConfig {
FlowMetricConfig {
id: name.to_string(),
bounds: FlowBounds {
start_node: start.to_string(),
end_node: stop.to_string(),
},
metrics: None,
duration_distribution: DistributionTier::Normal,
purpose: None,
}
}
fn assert_invalid_user_config(err: &crate::error::Error, sw_name: &str) {
match err {
crate::error::Error::ConfigError(boxed) => match boxed.as_ref() {
otel_arrow_dfe_config::error::Error::InvalidUserConfig { error } => {
assert!(
error.contains(sw_name),
"expected error to mention `{sw_name}`, got: {error}"
);
}
other => panic!("expected InvalidUserConfig, got: {other:?}"),
},
other => panic!("expected ConfigError(InvalidUserConfig), got: {other:?}"),
}
}
fn test_maps(
all_nodes: &[&str],
non_processors: &[&str],
) -> (HashMap<String, usize>, HashSet<usize>) {
let name_to_index: HashMap<String, usize> = all_nodes
.iter()
.enumerate()
.map(|(i, &n)| (n.to_string(), i))
.collect();
let processor_indices: HashSet<usize> = all_nodes
.iter()
.enumerate()
.filter(|&(_, &n)| !non_processors.contains(&n))
.map(|(i, _)| i)
.collect();
(name_to_index, processor_indices)
}
fn test_edges(
edges: &[(&str, &str)],
name_to_index: &HashMap<String, usize>,
) -> Vec<(usize, usize)> {
edges
.iter()
.filter_map(|&(from, to)| {
let src = *name_to_index.get(from)?;
let dst = *name_to_index.get(to)?;
Some((src, dst))
})
.collect()
}
#[test]
fn valid_flow_metric_is_registered() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c"], &[]);
let edges = test_edges(&[("a", "b"), ("b", "c")], &names);
let policy = policy_with(vec![sw("sw1", "a", "c")]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("valid config should build");
assert_eq!(state.duration_metrics.len(), 1);
assert!(state.start_nodes.contains_key(&0)); assert_eq!(state.end_nodes.get(&2), Some(&0)); }
#[test]
fn duration_only_registers_only_duration_metric_set() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b"], &[]);
let edges = test_edges(&[("a", "b")], &names);
let mut flow = sw("duration_only", "a", "b");
flow.metrics = Some(vec![FlowMetric::ComputeDuration]);
let policy = policy_with(vec![flow]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("duration-only config should build");
assert!(state.duration_metrics[0].is_some());
assert!(state.input_message_metrics[0].is_none());
assert!(state.input_items_metrics[0].is_none());
assert!(state.input_size_metrics[0].is_none());
assert!(state.output_message_metrics[0].is_none());
assert!(state.output_items_metrics[0].is_none());
assert!(state.output_size_metrics[0].is_none());
}
#[test]
fn duration_distribution_selects_matching_metric_set() {
for (tier, expected) in [
(DistributionTier::Basic, Instrument::Mmsc),
(DistributionTier::Normal, Instrument::ExponentialHistogram),
(DistributionTier::Detailed, Instrument::ExponentialHistogram),
] {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b"], &[]);
let edges = test_edges(&[("a", "b")], &names);
let mut flow = sw("duration", "a", "b");
flow.metrics = Some(vec![FlowMetric::ComputeDuration]);
flow.duration_distribution = tier;
let state =
build_flow_metric_state(&policy_with(vec![flow]), &names, &procs, &ctx, &edges)
.expect("supported duration tier should build");
assert_eq!(
duration_instrument(state.duration_metrics[0].as_ref().unwrap()),
expected,
"tier: {tier:?}"
);
}
}
#[test]
fn mixed_duration_types_have_distinct_flow_scopes() {
let (ctx, registry) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d"], &[]);
let edges = test_edges(&[("a", "b"), ("c", "d")], &names);
let mut basic = sw("basic_flow", "a", "b");
basic.metrics = Some(vec![FlowMetric::ComputeDuration]);
basic.duration_distribution = DistributionTier::Basic;
let mut normal = sw("normal_flow", "c", "d");
normal.metrics = Some(vec![FlowMetric::ComputeDuration]);
normal.duration_distribution = DistributionTier::Normal;
let _state = build_flow_metric_state(
&policy_with(vec![basic, normal]),
&names,
&procs,
&ctx,
&edges,
)
.expect("mixed duration tiers should build");
let mut scopes = Vec::new();
registry.visit_metrics_and_reset_with_zeroes(
|descriptor, attrs, _| {
if descriptor.name != "flow.compute" {
return;
}
let flow_id = attrs
.iter_attributes()
.find_map(|(key, value)| {
(key == "flow.id").then(|| match value {
AttributeValue::String(value) => value.clone(),
other => panic!("flow.id must be a string, got {other:?}"),
})
})
.expect("flow.compute scope must include flow.id");
scopes.push((flow_id, descriptor.metrics[0].instrument));
},
true,
);
scopes.sort_by(|left, right| left.0.cmp(&right.0));
scopes.dedup();
assert_eq!(
scopes,
[
("basic_flow".to_string(), Instrument::Mmsc),
("normal_flow".to_string(), Instrument::ExponentialHistogram),
]
);
}
#[test]
fn duration_distribution_does_not_allocate_when_duration_is_disabled() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b"], &[]);
let edges = test_edges(&[("a", "b")], &names);
let mut flow = sw("items", "a", "b");
flow.metrics = Some(vec![FlowMetric::InputItems]);
flow.duration_distribution = DistributionTier::Detailed;
let state = build_flow_metric_state(&policy_with(vec![flow]), &names, &procs, &ctx, &edges)
.expect("count-only flow should build");
assert!(state.duration_metrics[0].is_none());
}
#[test]
fn dropped_registers_decision_candidates_for_range_processors() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c"], &[]);
let edges = test_edges(&[("a", "b"), ("b", "c")], &names);
let mut flow = sw("decisions", "a", "c");
flow.metrics = Some(vec![FlowMetric::DroppedItems]);
let policy = policy_with(vec![flow]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("dropped config should build");
assert_eq!(state.decision_candidates.len(), 3);
for idx in [0usize, 1, 2] {
let candidate = state
.decision_candidates
.get(&idx)
.expect("each range node should be a candidate");
assert!(candidate.attrs.decision.is_empty());
}
assert!(state.is_active());
}
#[test]
fn no_dropped_means_no_decision_candidates() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b"], &[]);
let edges = test_edges(&[("a", "b")], &names);
let mut flow = sw("duration_only", "a", "b");
flow.metrics = Some(vec![FlowMetric::ComputeDuration]);
let policy = policy_with(vec![flow]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("config should build");
assert!(state.decision_candidates.is_empty());
}
#[test]
fn single_node_flow_is_supported() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["solo"], &[]);
let mut flow = sw("single", "solo", "solo");
flow.metrics = Some(vec![FlowMetric::DroppedItems]);
let policy = policy_with(vec![flow]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &[])
.expect("single-node flow should build");
assert_eq!(state.decision_candidates.len(), 1);
assert!(state.decision_candidates.contains_key(&0));
}
#[test]
fn shared_interior_decision_node_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "c", "m", "b", "d"], &[]);
let edges = test_edges(&[("a", "m"), ("c", "m"), ("m", "b"), ("m", "d")], &names);
let mut flow1 = sw("flow1", "a", "b");
flow1.metrics = Some(vec![FlowMetric::DroppedItems]);
let mut flow2 = sw("flow2", "c", "d");
flow2.metrics = Some(vec![FlowMetric::DroppedItems]);
let policy = policy_with(vec![flow1, flow2]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.err()
.expect("shared interior decision node should be rejected");
assert_invalid_user_config(&err, "flow2");
}
#[test]
fn unknown_node_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b"], &[]);
let policy = policy_with(vec![sw("sw1", "a", "missing")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &[])
.err()
.expect("expected Err");
assert_invalid_user_config(&err, "sw1");
}
#[test]
fn non_processor_start_node_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["recv", "proc1", "proc2"], &["recv"]);
let policy = policy_with(vec![sw("sw1", "recv", "proc2")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &[])
.err()
.expect("expected Err");
assert_invalid_user_config(&err, "sw1");
}
#[test]
fn non_processor_end_node_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["proc1", "proc2", "exp"], &["exp"]);
let policy = policy_with(vec![sw("sw1", "proc1", "exp")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &[])
.err()
.expect("expected Err");
assert_invalid_user_config(&err, "sw1");
}
#[test]
fn shared_start_node_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d"], &[]);
let policy = policy_with(vec![sw("sw1", "a", "b"), sw("sw2", "a", "d")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &[])
.err()
.expect("expected Err");
assert_invalid_user_config(&err, "sw2");
}
#[test]
fn shared_end_node_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d"], &[]);
let policy = policy_with(vec![sw("sw1", "a", "d"), sw("sw2", "c", "d")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &[])
.err()
.expect("expected Err");
assert_invalid_user_config(&err, "sw2");
}
#[test]
fn stop_of_one_is_start_of_another_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c"], &[]);
let policy = policy_with(vec![sw("sw1", "a", "b"), sw("sw2", "b", "c")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &[])
.err()
.expect("expected Err");
assert_invalid_user_config(&err, "sw2");
}
#[test]
fn disjoint_flow_metrics_are_both_registered() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d"], &[]);
let policy = policy_with(vec![sw("sw1", "a", "b"), sw("sw2", "c", "d")]);
let edges = test_edges(&[("a", "b"), ("b", "c"), ("c", "d")], &names);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("disjoint config should build");
assert_eq!(state.duration_metrics.len(), 2);
assert!(state.start_nodes.contains_key(&0)); assert!(state.start_nodes.contains_key(&2)); assert_eq!(state.end_nodes.get(&1), Some(&0)); assert_eq!(state.end_nodes.get(&3), Some(&1)); }
#[test]
fn interleaved_distinct_endpoints_linear_path_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d"], &[]);
let edges = test_edges(&[("a", "b"), ("b", "c"), ("c", "d")], &names);
let policy = policy_with(vec![sw("sw1", "a", "c"), sw("sw2", "b", "d")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.err()
.expect("expected Err for interleaved ranges");
assert_invalid_user_config(&err, "sw2");
}
#[test]
fn interleaved_reverse_order_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d"], &[]);
let edges = test_edges(&[("a", "b"), ("b", "c"), ("c", "d")], &names);
let policy = policy_with(vec![sw("sw1", "b", "d"), sw("sw2", "a", "c")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.err()
.expect("expected Err for interleaved ranges");
assert_invalid_user_config(&err, "sw2");
}
#[test]
fn disjoint_on_linear_path_is_accepted() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d", "e", "f"], &[]);
let edges = test_edges(
&[("a", "b"), ("b", "c"), ("c", "d"), ("d", "e"), ("e", "f")],
&names,
);
let policy = policy_with(vec![sw("sw1", "a", "c"), sw("sw2", "d", "f")]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("disjoint ranges should build");
assert_eq!(state.duration_metrics.len(), 2);
}
#[test]
fn interleaved_on_branching_path_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d", "e"], &[]);
let edges = test_edges(
&[("a", "b"), ("a", "c"), ("b", "d"), ("c", "d"), ("d", "e")],
&names,
);
let policy = policy_with(vec![sw("sw1", "a", "d"), sw("sw2", "b", "e")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.err()
.expect("expected Err for interleaved ranges on diamond");
assert_invalid_user_config(&err, "sw2");
}
#[test]
fn disjoint_on_separate_branches_is_accepted() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d", "e"], &[]);
let edges = test_edges(&[("a", "b"), ("b", "c"), ("a", "d"), ("d", "e")], &names);
let policy = policy_with(vec![sw("sw1", "b", "c"), sw("sw2", "d", "e")]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("separate-branch ranges should build");
assert_eq!(state.duration_metrics.len(), 2);
}
#[test]
fn nested_ranges_are_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d"], &[]);
let edges = test_edges(&[("a", "b"), ("b", "c"), ("c", "d")], &names);
let policy = policy_with(vec![sw("sw1", "a", "d"), sw("sw2", "b", "c")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.err()
.expect("expected Err for nested ranges");
assert_invalid_user_config(&err, "sw2");
}
#[test]
fn single_flow_metric_accepted() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c"], &[]);
let edges = test_edges(&[("a", "b"), ("b", "c")], &names);
let policy = policy_with(vec![sw("sw1", "a", "c")]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("single flow metric should build");
assert_eq!(state.duration_metrics.len(), 1);
}
#[test]
fn active_range_linear() {
let adj = build_adjacency(&[(0, 1), (1, 2), (2, 3)]);
let (range, _end_reachable) = active_range(0, 2, &adj);
assert!(range.contains(&0));
assert!(range.contains(&1));
assert!(!range.contains(&2));
assert!(!range.contains(&3));
}
#[test]
fn active_range_diamond() {
let adj = build_adjacency(&[(0, 1), (0, 2), (1, 3), (2, 3)]);
let (range, _end_reachable) = active_range(0, 3, &adj);
assert!(range.contains(&0));
assert!(range.contains(&1));
assert!(range.contains(&2));
assert!(!range.contains(&3));
}
#[test]
fn active_range_end_unreachable() {
let adj = build_adjacency(&[(0, 1), (4, 5)]);
let (range, _end_reachable) = active_range(0, 5, &adj);
assert!(range.contains(&0));
assert!(range.contains(&1));
assert!(!range.contains(&5));
}
#[test]
fn unreachable_end_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d"], &[]);
let edges = test_edges(&[("a", "b"), ("b", "c")], &names);
let policy = policy_with(vec![sw("sw1", "a", "d")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.err()
.expect("expected Err for unreachable end node");
assert_invalid_user_config(&err, "sw1");
}
#[test]
fn unreachable_end_on_separate_branch_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d"], &[]);
let edges = test_edges(&[("a", "b"), ("a", "c"), ("c", "d")], &names);
let policy = policy_with(vec![sw("sw1", "b", "d")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.err()
.expect("expected Err for unreachable end node on separate branch");
assert_invalid_user_config(&err, "sw1");
}
#[test]
fn reverse_direction_is_rejected() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c"], &[]);
let edges = test_edges(&[("a", "b"), ("b", "c")], &names);
let policy = policy_with(vec![sw("sw1", "c", "a")]);
let err = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.err()
.expect("expected Err for reverse-direction flow metric");
assert_invalid_user_config(&err, "sw1");
}
#[test]
fn adjacent_nodes_are_accepted() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b"], &[]);
let edges = test_edges(&[("a", "b")], &names);
let policy = policy_with(vec![sw("sw1", "a", "b")]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("adjacent-node flow metric should build");
assert_eq!(state.duration_metrics.len(), 1);
}
#[test]
fn multi_hop_reachable_end_is_accepted() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d", "e"], &[]);
let edges = test_edges(&[("a", "b"), ("b", "c"), ("c", "d"), ("d", "e")], &names);
let policy = policy_with(vec![sw("sw1", "a", "e")]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("multi-hop reachable flow metric should build");
assert_eq!(state.duration_metrics.len(), 1);
}
#[test]
fn diamond_end_reachable_via_either_branch_is_accepted() {
let (ctx, _) = test_pipeline_ctx();
let (names, procs) = test_maps(&["a", "b", "c", "d"], &[]);
let edges = test_edges(&[("a", "b"), ("a", "c"), ("b", "d"), ("c", "d")], &names);
let policy = policy_with(vec![sw("sw1", "a", "d")]);
let state = build_flow_metric_state(&policy, &names, &procs, &ctx, &edges)
.expect("diamond-reachable flow metric should build");
assert_eq!(state.duration_metrics.len(), 1);
}
#[test]
fn flow_attribute_set_exposes_purpose_scope_attribute() {
let descriptor = FlowAttributeSet::default().descriptor();
assert!(
descriptor.fields.iter().any(|f| f.key == "flow.purpose"),
"flow.purpose missing from FlowAttributeSet descriptor"
);
let attrs = FlowAttributeSet {
flow_id: "sampling".into(),
start_node: "log_sampler".into(),
end_node: "log_sampler".into(),
purpose: "filter".into(),
..FlowAttributeSet::default()
};
let set: Vec<(&str, AttributeValue)> = attrs
.iter_attributes()
.map(|(key, value)| (key, value.clone()))
.collect();
assert!(
set.contains(&("flow.purpose", AttributeValue::String("filter".to_string()))),
"expected flow.purpose=filter in {set:?}"
);
let unset = FlowAttributeSet::default();
let unset_set: Vec<(&str, AttributeValue)> = unset
.iter_attributes()
.map(|(key, value)| (key, value.clone()))
.collect();
assert!(
unset_set.contains(&("flow.purpose", AttributeValue::String(String::new()))),
"expected empty flow.purpose in {unset_set:?}"
);
}
}