use crate::Interests;
use crate::ReceivedAtNode;
use crate::Unwindable;
use crate::channel_metrics::{
ChannelMetricsHandle, NodeCompletionMetrics, NodeInputItemMetrics, NodeInputMetrics,
NodeInputSizeMetrics, NodeOutputItemMetrics, NodeOutputMetrics, NodeOutputSizeMetrics,
};
use crate::completion_emission_metrics::{
CompletionEmissionMetricsHandle, make_completion_emission_metrics,
};
use crate::context::PipelineContext;
use crate::control::{
ControlSenders, Controllable, NodeControlMsg, PipelineCompletionMsgReceiver,
PipelineCompletionMsgSender, RuntimeCtrlMsgReceiver, RuntimeCtrlMsgSender,
};
use crate::entity_context::{
NodeTaskContext, NodeTelemetryGuard, NodeTelemetryHandle, instrument_with_node_context,
};
use crate::error::{Error, TypedError};
use crate::flow_metrics::{
FlowDroppedItemsMetrics, FlowDurationMetricSet, FlowInputItemsMetrics, FlowInputMessageMetrics,
FlowInputSizeMetrics, FlowOutputItemsMetrics, FlowOutputMessageMetrics, FlowOutputSizeMetrics,
build_flow_metric_state,
};
use crate::memory_limiter::MemoryPressureChanged;
use crate::node::{Node, NodeDefs, NodeId, NodeType, NodeWithPDataReceiver, NodeWithPDataSender};
use crate::pipeline_ctrl::{
NodeMetricHandles, PipelineCompletionMsgDispatcher, RuntimeCtrlMsgManager,
snapshot_node_metrics_with_handles,
};
use crate::processor::FlowMetricHook;
use crate::runtime_services::PipelineRuntimeServices;
use crate::terminal_state::{TerminalMetricsDeadline, TerminalState};
use crate::{exporter::ExporterWrapper, processor::ProcessorWrapper, receiver::ReceiverWrapper};
use otel_arrow_dfe_config::DeployedPipelineKey;
use otel_arrow_dfe_config::pipeline::PipelineConfig;
use otel_arrow_dfe_config::policy::TelemetryPolicy;
use otel_arrow_dfe_telemetry::event::ObservedEventReporter;
use otel_arrow_dfe_telemetry::metrics::{MeasurementMetricSet, MetricSetSnapshot};
use otel_arrow_dfe_telemetry::reporter::{MetricsReporter, ReportOutcome};
use std::cell::RefCell;
use std::collections::{HashMap, HashSet};
use std::fmt::Debug;
use std::rc::Rc;
use std::time::Duration;
use tokio::runtime::Builder;
use tokio::sync::watch;
use tokio::task::LocalSet;
const EXTENSION_MONITOR_TICK_INTERVAL: Duration = Duration::from_secs(1);
const EXTENSION_MONITOR_COLLECT_TELEMETRY_INTERVAL: Duration = Duration::from_secs(10);
fn make_output_metrics(
telemetry_handle: &Option<NodeTelemetryHandle>,
pipeline_context: &PipelineContext,
) -> Vec<MeasurementMetricSet<NodeOutputMetrics>> {
telemetry_handle
.as_ref()
.map(|h| {
let mut keys = h.output_channel_keys();
keys.sort_by(|a, b| a.0.cmp(&b.0));
keys.iter()
.map(|(_, key)| {
NodeOutputMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(*key),
)
})
.collect()
})
.unwrap_or_default()
}
fn make_output_completion_metrics(
telemetry_handle: &Option<NodeTelemetryHandle>,
pipeline_context: &PipelineContext,
) -> Vec<MeasurementMetricSet<NodeCompletionMetrics>> {
telemetry_handle
.as_ref()
.map(|h| {
let mut keys = h.output_channel_keys();
keys.sort_by(|a, b| a.0.cmp(&b.0));
keys.iter()
.map(|(_, key)| {
NodeCompletionMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(*key),
)
})
.collect()
})
.unwrap_or_default()
}
fn make_output_item_metrics(
telemetry_handle: &Option<NodeTelemetryHandle>,
pipeline_context: &PipelineContext,
) -> Vec<MeasurementMetricSet<NodeOutputItemMetrics>> {
telemetry_handle
.as_ref()
.map(|h| {
let mut keys = h.output_channel_keys();
keys.sort_by(|a, b| a.0.cmp(&b.0));
keys.iter()
.map(|(_, key)| {
NodeOutputItemMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(*key),
)
})
.collect()
})
.unwrap_or_default()
}
fn make_output_size_metrics(
telemetry_handle: &Option<NodeTelemetryHandle>,
pipeline_context: &PipelineContext,
) -> Vec<MeasurementMetricSet<NodeOutputSizeMetrics>> {
telemetry_handle
.as_ref()
.map(|h| {
let mut keys = h.output_channel_keys();
keys.sort_by(|a, b| a.0.cmp(&b.0));
keys.iter()
.map(|(_, key)| {
NodeOutputSizeMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(*key),
)
})
.collect()
})
.unwrap_or_default()
}
fn make_node_metric_handles(
telemetry_handle: &Option<NodeTelemetryHandle>,
pipeline_context: &PipelineContext,
has_input: bool,
has_outputs: bool,
node_interests: Interests,
completion_emission: Option<CompletionEmissionMetricsHandle>,
) -> NodeMetricHandles {
let input_metrics_enabled = has_input && node_interests.contains(Interests::NODE_INPUT_METRICS);
let output_metrics_enabled =
has_outputs && node_interests.contains(Interests::NODE_OUTPUT_METRICS);
let completion_duration_enabled = node_interests.contains(Interests::NODE_COMPLETION_DURATION);
let item_counts_enabled = node_interests.contains(Interests::NODE_ITEM_COUNTS);
let size_enabled = node_interests.contains(Interests::NODE_SIZE);
let input = if input_metrics_enabled {
telemetry_handle
.as_ref()
.and_then(|h| h.input_channel_key())
.map(|key| {
NodeInputMetrics::register(&pipeline_context.metric_set_registrar_for_entity(key))
})
} else {
None
};
let input_completion = if completion_duration_enabled && has_input {
telemetry_handle
.as_ref()
.and_then(|h| h.input_channel_key())
.map(|key| {
NodeCompletionMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(key),
)
})
} else {
None
};
let input_size = if size_enabled && has_input {
telemetry_handle
.as_ref()
.and_then(|h| h.input_channel_key())
.map(|key| {
NodeInputSizeMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(key),
)
})
} else {
None
};
let input_items = if item_counts_enabled && has_input {
telemetry_handle
.as_ref()
.and_then(|h| h.input_channel_key())
.map(|key| {
NodeInputItemMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(key),
)
})
} else {
None
};
let outputs = if output_metrics_enabled {
make_output_metrics(telemetry_handle, pipeline_context)
} else {
Vec::new()
};
let output_completion = if completion_duration_enabled && has_outputs && !has_input {
make_output_completion_metrics(telemetry_handle, pipeline_context)
} else {
Vec::new()
};
let output_items = if item_counts_enabled && has_outputs {
make_output_item_metrics(telemetry_handle, pipeline_context)
} else {
Vec::new()
};
let output_size = if size_enabled && has_outputs {
make_output_size_metrics(telemetry_handle, pipeline_context)
} else {
Vec::new()
};
NodeMetricHandles {
registry: pipeline_context.metrics_registry(),
input,
input_completion,
input_items,
input_size,
outputs,
output_completion,
output_items,
output_size,
completion_emission,
}
}
pub struct RuntimePipeline<PData: Debug> {
config: PipelineConfig,
receivers: Vec<ReceiverWrapper<PData>>,
processors: Vec<ProcessorWrapper<PData>>,
exporters: Vec<ExporterWrapper<PData>>,
extensions: Vec<(
crate::extension::ExtensionWrapper,
otel_arrow_dfe_telemetry::registry::EntityKey,
)>,
nodes: NodeDefs<PData, PipeNode>,
channel_metrics: Vec<ChannelMetricsHandle>,
admission_metrics: Vec<crate::admission::metrics::AdmissionMetricsHandle>,
telemetry_policy: TelemetryPolicy,
}
async fn flush_metrics_reporter(
metrics_reporter: &MetricsReporter,
phase: &'static str,
terminal_metrics_deadline: &TerminalMetricsDeadline,
) {
if let Err(err) = metrics_reporter
.flush_until(terminal_metrics_deadline.get())
.await
{
otel_arrow_dfe_telemetry::otel_warn!(
"metrics.collection.flush.fail",
phase,
error = err.to_string()
);
}
}
pub(crate) async fn report_terminal_metrics(
metrics_reporter: &MetricsReporter,
terminal_state: TerminalState,
terminal_metrics_deadline: &TerminalMetricsDeadline,
) {
let deadline = terminal_state.deadline();
terminal_metrics_deadline.record(deadline);
report_metric_snapshots(
metrics_reporter,
terminal_state.into_metrics(),
"terminal",
terminal_metrics_deadline.get(),
)
.await;
}
async fn report_metric_snapshots(
metrics_reporter: &MetricsReporter,
snapshots: impl IntoIterator<Item = MetricSetSnapshot>,
phase: &'static str,
deadline: std::time::Instant,
) {
for snapshot in snapshots {
match metrics_reporter
.report_snapshot_reliably_until(snapshot, deadline)
.await
{
Ok(ReportOutcome::Sent) => {}
Ok(ReportOutcome::Deferred) => {
otel_arrow_dfe_telemetry::otel_warn!(
"metrics.terminal.reporting.deferred",
phase,
message = "Terminal metric snapshot was deferred because the standalone reporter channel is full"
);
}
Err(err) => {
otel_arrow_dfe_telemetry::otel_warn!(
"metrics.terminal.reporting.fail",
phase,
error = err.to_string()
);
}
}
}
if let Err(err) = metrics_reporter.flush_until(deadline).await {
otel_arrow_dfe_telemetry::otel_warn!(
"metrics.collection.flush.fail",
phase,
error = err.to_string()
);
}
}
fn connection_edges<'a>(
connections: impl Iterator<Item = &'a otel_arrow_dfe_config::pipeline::PipelineConnection>,
node_name_to_index: &HashMap<String, usize>,
) -> Vec<(usize, usize)> {
let mut edges = Vec::new();
for conn in connections {
let from_indices: Vec<usize> = conn
.from_nodes()
.into_iter()
.filter_map(|name| node_name_to_index.get(name.as_ref()).copied())
.collect();
let to_indices: Vec<usize> = conn
.to_nodes()
.into_iter()
.filter_map(|name| node_name_to_index.get(name.as_ref()).copied())
.collect();
for &src in &from_indices {
for &dst in &to_indices {
edges.push((src, dst));
}
}
}
edges
}
pub(crate) struct PipeNode {
index: usize, }
impl PipeNode {
pub(crate) const fn new(index: usize) -> Self {
Self { index }
}
}
impl<PData: 'static + Debug + Clone> RuntimePipeline<PData> {
#[must_use]
pub(crate) fn new(
config: PipelineConfig,
receivers: Vec<ReceiverWrapper<PData>>,
processors: Vec<ProcessorWrapper<PData>>,
exporters: Vec<ExporterWrapper<PData>>,
extensions: Vec<(
crate::extension::ExtensionWrapper,
otel_arrow_dfe_telemetry::registry::EntityKey,
)>,
nodes: NodeDefs<PData, PipeNode>,
telemetry_policy: TelemetryPolicy,
) -> Self {
Self {
config,
receivers,
processors,
exporters,
extensions,
nodes,
channel_metrics: Default::default(),
admission_metrics: Default::default(),
telemetry_policy,
}
}
pub(crate) fn set_channel_metrics(&mut self, channel_metrics: Vec<ChannelMetricsHandle>) {
self.channel_metrics = channel_metrics;
}
pub(crate) fn set_admission_metrics(
&mut self,
admission_metrics: Vec<crate::admission::metrics::AdmissionMetricsHandle>,
) {
self.admission_metrics = admission_metrics;
}
#[must_use]
pub const fn node_count(&self) -> usize {
self.receivers.len() + self.processors.len() + self.exporters.len()
}
#[must_use]
pub const fn config(&self) -> &PipelineConfig {
&self.config
}
}
impl<PData: 'static + Debug + Clone + ReceivedAtNode + Unwindable + FlowMetricHook>
RuntimePipeline<PData>
{
pub fn run_forever(
self,
pipeline_key: DeployedPipelineKey,
pipeline_context: PipelineContext,
event_reporter: ObservedEventReporter,
metrics_reporter: MetricsReporter,
control_plane_metrics_flush_interval: Duration,
memory_pressure_rx: watch::Receiver<MemoryPressureChanged>,
runtime_ctrl_msg_tx: RuntimeCtrlMsgSender<PData>,
runtime_ctrl_msg_rx: RuntimeCtrlMsgReceiver<PData>,
pipeline_completion_msg_tx: PipelineCompletionMsgSender<PData>,
pipeline_completion_msg_rx: PipelineCompletionMsgReceiver<PData>,
) -> Result<Vec<()>, Error> {
use futures::stream::{FuturesUnordered, StreamExt};
let RuntimePipeline {
config: pipeline_config,
receivers,
processors,
exporters,
extensions,
nodes: _nodes,
channel_metrics,
admission_metrics,
telemetry_policy,
} = self;
let metric_level = telemetry_policy.runtime_metrics;
let rt = Builder::new_current_thread()
.enable_all()
.build()
.expect("Failed to create runtime");
let local_tasks = LocalSet::new();
let runtime_services = PipelineRuntimeServices::new(Default::default())?;
let mut futures = FuturesUnordered::new();
let ext_ctx = pipeline_context.extension_context();
let ext_monitor = {
let _enter = rt.enter();
if telemetry_policy.pipeline_metrics {
crate::extension_monitor::ExtensionMetricsMonitor::new(
ext_ctx.clone(),
EXTENSION_MONITOR_TICK_INTERVAL,
EXTENSION_MONITOR_COLLECT_TELEMETRY_INTERVAL,
)
} else {
crate::extension_monitor::ExtensionMetricsMonitor::disabled(ext_ctx.clone())
}
};
let terminal_metrics_deadline = TerminalMetricsDeadline::default();
let (forced_shutdown_trigger, forced_shutdown_signal) =
crate::forced_shutdown::ForcedShutdownTrigger::pair();
let mut extension_lifecycle = crate::extension_lifecycle::ExtensionLifecycle::spawn(
extensions,
&local_tasks,
metrics_reporter.clone(),
terminal_metrics_deadline.clone(),
&ext_ctx,
ext_monitor,
);
if let Err(barrier_err) =
rt.block_on(local_tasks.run_until(extension_lifecycle.wait_all_spawned()))
{
extension_lifecycle.initiate_shutdown(Some("spawn barrier failed"));
rt.block_on(local_tasks.run_until(extension_lifecycle.drain_until_deadline()));
return Err(barrier_err);
}
if let Err(readiness_err) =
rt.block_on(local_tasks.run_until(extension_lifecycle.wait_all_ready()))
{
extension_lifecycle.initiate_shutdown(Some("extension readiness gate failed"));
rt.block_on(local_tasks.run_until(extension_lifecycle.drain_until_deadline()));
return Err(readiness_err);
}
let mut control_senders = ControlSenders::default();
let mut node_metric_entries: Vec<(usize, NodeMetricHandles)> = Vec::new();
let mut node_telemetry_guards: Vec<NodeTelemetryGuard> = Vec::new();
let node_name_to_index: HashMap<String, usize> = _nodes
.iter()
.map(|(nid, _)| (nid.name.to_string(), nid.index))
.collect();
let processor_indices: HashSet<usize> = processors
.iter()
.map(|p| match p {
ProcessorWrapper::Local { node_id, .. }
| ProcessorWrapper::Shared { node_id, .. } => node_id.index,
})
.collect();
let pipeline_connections =
connection_edges(pipeline_config.connection_iter(), &node_name_to_index);
let mut flow_metric_state = build_flow_metric_state(
&telemetry_policy,
&node_name_to_index,
&processor_indices,
&pipeline_context,
&pipeline_connections,
)?;
for exporter in exporters {
let mut exporter = exporter;
let node_id = exporter.node_id();
let node_config = pipeline_config
.nodes()
.get(node_id.name.as_ref())
.expect("runtime exporter has pipeline configuration");
let node_interests = Interests::for_node(metric_level, node_config);
control_senders.register(
node_id.clone(),
NodeType::Exporter,
exporter.control_sender(),
);
let telemetry_guard = exporter.take_telemetry_guard();
let node_entity_key = telemetry_guard.as_ref().map(|t| t.entity_key());
let telemetry_handle = telemetry_guard.as_ref().map(|t| t.handle());
node_telemetry_guards.extend(telemetry_guard);
let completion_emission_metrics =
make_completion_emission_metrics(&telemetry_handle, metric_level);
node_metric_entries.push((
node_id.index,
make_node_metric_handles(
&telemetry_handle,
&pipeline_context,
true,
false,
node_interests,
completion_emission_metrics.clone(),
),
));
let runtime_ctrl_msg_tx = runtime_ctrl_msg_tx.clone();
let pipeline_completion_msg_tx = pipeline_completion_msg_tx.clone();
let effect_metrics_reporter = metrics_reporter.clone();
let final_metrics_reporter = metrics_reporter.clone();
let exporter_terminal_metrics_deadline = terminal_metrics_deadline.clone();
let exporter_runtime_services = runtime_services.clone();
let fut = async move {
match exporter
.start_with_completion_metrics(
runtime_ctrl_msg_tx,
pipeline_completion_msg_tx,
effect_metrics_reporter,
node_interests,
completion_emission_metrics,
exporter_runtime_services,
)
.await
{
Ok(terminal_state) => {
report_terminal_metrics(
&final_metrics_reporter,
terminal_state,
&exporter_terminal_metrics_deadline,
)
.await;
Ok(())
}
Err(err) => {
flush_metrics_reporter(
&final_metrics_reporter,
"terminal_error",
&exporter_terminal_metrics_deadline,
)
.await;
Err(err)
}
}
};
if let Some(handle) = telemetry_handle {
let input_key = handle.input_channel_key();
let output_keys = handle.output_channel_keys();
let node_ctx =
NodeTaskContext::new(node_entity_key, Some(handle), input_key, output_keys);
futures.push(local_tasks.spawn_local(instrument_with_node_context(node_ctx, fut)));
} else if let Some(key) = node_entity_key {
let node_ctx = NodeTaskContext::new(Some(key), None, None, Vec::new());
futures.push(local_tasks.spawn_local(instrument_with_node_context(node_ctx, fut)));
} else {
futures.push(local_tasks.spawn_local(fut));
}
}
for processor in processors {
let mut processor = processor;
let node_id = processor.node_id();
let node_config = pipeline_config
.nodes()
.get(node_id.name.as_ref())
.expect("runtime processor has pipeline configuration");
let node_interests = Interests::for_node(metric_level, node_config);
control_senders.register(
node_id.clone(),
NodeType::Processor,
processor.control_sender(),
);
let telemetry_guard = processor.take_telemetry_guard();
let node_entity_key = telemetry_guard.as_ref().map(|t| t.entity_key());
let telemetry_handle = telemetry_guard.as_ref().map(|t| t.handle());
node_telemetry_guards.extend(telemetry_guard);
let completion_emission_metrics =
make_completion_emission_metrics(&telemetry_handle, metric_level);
node_metric_entries.push((
node_id.index,
make_node_metric_handles(
&telemetry_handle,
&pipeline_context,
true,
true,
node_interests,
completion_emission_metrics.clone(),
),
));
let runtime_ctrl_msg_tx = runtime_ctrl_msg_tx.clone();
let pipeline_completion_msg_tx = pipeline_completion_msg_tx.clone();
let metrics_reporter = metrics_reporter.clone();
let final_metrics_reporter = metrics_reporter.clone();
let processor_terminal_metrics_deadline = terminal_metrics_deadline.clone();
let processor_forced_shutdown_signal = forced_shutdown_signal.clone();
let processor_runtime_services = runtime_services.clone();
let flow_active = flow_metric_state.is_active();
let flow_needs_timing = flow_metric_state.needs_timing();
let flow_is_start = flow_metric_state.start_nodes.contains_key(&node_id.index);
let flow_is_end = flow_metric_state.end_nodes.contains_key(&node_id.index);
let flow_input_message_metric: Option<MeasurementMetricSet<FlowInputMessageMetrics>> =
flow_metric_state
.start_nodes
.get(&node_id.index)
.and_then(|&id| flow_metric_state.input_message_metrics[id].take());
let flow_input_items_metric: Option<MeasurementMetricSet<FlowInputItemsMetrics>> =
flow_metric_state
.start_nodes
.get(&node_id.index)
.and_then(|&id| flow_metric_state.input_items_metrics[id].take());
let flow_input_size_metric: Option<MeasurementMetricSet<FlowInputSizeMetrics>> =
flow_metric_state
.start_nodes
.get(&node_id.index)
.and_then(|&id| flow_metric_state.input_size_metrics[id].take());
let flow_duration_metric: Option<FlowDurationMetricSet> = flow_metric_state
.end_nodes
.get(&node_id.index)
.and_then(|&id| flow_metric_state.duration_metrics[id].take());
let flow_output_items_metric: Option<MeasurementMetricSet<FlowOutputItemsMetrics>> =
flow_metric_state
.end_nodes
.get(&node_id.index)
.and_then(|&id| flow_metric_state.output_items_metrics[id].take());
let flow_output_message_metric: Option<MeasurementMetricSet<FlowOutputMessageMetrics>> =
flow_metric_state
.end_nodes
.get(&node_id.index)
.and_then(|&id| flow_metric_state.output_message_metrics[id].take());
let flow_output_size_metric: Option<MeasurementMetricSet<FlowOutputSizeMetrics>> =
flow_metric_state
.end_nodes
.get(&node_id.index)
.and_then(|&id| flow_metric_state.output_size_metrics[id].take());
let mut flow_dropped_items_metric: Option<
MeasurementMetricSet<FlowDroppedItemsMetrics>,
> = None;
if processor.runtime_requirements().makes_drop_decisions
&& let Some(candidate) = flow_metric_state.decision_candidates.get(&node_id.index)
{
let mut attrs = candidate.attrs.clone();
attrs.decision = std::borrow::Cow::Owned(node_id.name.to_string());
let entity_key = pipeline_context.metrics_registry().register_entity(attrs);
flow_dropped_items_metric = Some(FlowDroppedItemsMetrics::register(
&pipeline_context.metric_set_registrar_for_entity(entity_key),
));
}
let fut = async move {
let result = processor
.start_with_completion_metrics(
runtime_ctrl_msg_tx,
pipeline_completion_msg_tx,
metrics_reporter,
node_interests,
completion_emission_metrics,
flow_is_start,
flow_is_end,
flow_input_message_metric,
flow_input_items_metric,
flow_input_size_metric,
flow_duration_metric,
flow_output_items_metric,
flow_output_message_metric,
flow_output_size_metric,
flow_dropped_items_metric,
flow_active,
flow_needs_timing,
processor_terminal_metrics_deadline.clone(),
processor_forced_shutdown_signal,
processor_runtime_services,
)
.await;
flush_metrics_reporter(
&final_metrics_reporter,
"terminal",
&processor_terminal_metrics_deadline,
)
.await;
result
};
if let Some(handle) = telemetry_handle {
let input_key = handle.input_channel_key();
let output_keys = handle.output_channel_keys();
let node_ctx =
NodeTaskContext::new(node_entity_key, Some(handle), input_key, output_keys);
futures.push(local_tasks.spawn_local(instrument_with_node_context(node_ctx, fut)));
} else if let Some(key) = node_entity_key {
let node_ctx = NodeTaskContext::new(Some(key), None, None, Vec::new());
futures.push(local_tasks.spawn_local(instrument_with_node_context(node_ctx, fut)));
} else {
futures.push(local_tasks.spawn_local(fut));
}
}
for receiver in receivers {
let mut receiver = receiver;
let node_id = receiver.node_id();
let node_config = pipeline_config
.nodes()
.get(node_id.name.as_ref())
.expect("runtime receiver has pipeline configuration");
let node_interests = Interests::for_node(metric_level, node_config);
control_senders.register(
node_id.clone(),
NodeType::Receiver,
receiver.control_sender(),
);
let telemetry_guard = receiver.take_telemetry_guard();
let node_entity_key = telemetry_guard.as_ref().map(|t| t.entity_key());
let telemetry_handle = telemetry_guard.as_ref().map(|t| t.handle());
node_telemetry_guards.extend(telemetry_guard);
node_metric_entries.push((
node_id.index,
make_node_metric_handles(
&telemetry_handle,
&pipeline_context,
false,
true,
node_interests,
None,
),
));
let runtime_ctrl_msg_tx = runtime_ctrl_msg_tx.clone();
let pipeline_completion_msg_tx = pipeline_completion_msg_tx.clone();
let effect_metrics_reporter = metrics_reporter.clone();
let final_metrics_reporter = metrics_reporter.clone();
let receiver_terminal_metrics_deadline = terminal_metrics_deadline.clone();
let receiver_runtime_services = runtime_services.clone();
let fut = async move {
match receiver
.start(
runtime_ctrl_msg_tx,
pipeline_completion_msg_tx,
effect_metrics_reporter,
node_interests,
receiver_runtime_services,
)
.await
{
Ok(terminal_state) => {
report_terminal_metrics(
&final_metrics_reporter,
terminal_state,
&receiver_terminal_metrics_deadline,
)
.await;
Ok(())
}
Err(err) => {
flush_metrics_reporter(
&final_metrics_reporter,
"terminal_error",
&receiver_terminal_metrics_deadline,
)
.await;
Err(err)
}
}
};
if let Some(handle) = telemetry_handle {
let input_key = handle.input_channel_key();
let output_keys = handle.output_channel_keys();
let node_ctx =
NodeTaskContext::new(node_entity_key, Some(handle), input_key, output_keys);
futures.push(local_tasks.spawn_local(instrument_with_node_context(node_ctx, fut)));
} else if let Some(key) = node_entity_key {
let node_ctx = NodeTaskContext::new(Some(key), None, None, Vec::new());
futures.push(local_tasks.spawn_local(instrument_with_node_context(node_ctx, fut)));
} else {
futures.push(local_tasks.spawn_local(fut));
}
}
let max_node = node_metric_entries
.iter()
.map(|(id, _)| *id)
.max()
.unwrap_or(0);
let mut node_metric_handles: Vec<Option<NodeMetricHandles>> =
(0..=max_node).map(|_| None).collect();
for (id, handles) in node_metric_entries {
node_metric_handles[id] = Some(handles);
}
let node_metric_handles = Rc::new(RefCell::new(node_metric_handles));
drop(runtime_ctrl_msg_tx);
drop(pipeline_completion_msg_tx);
let return_control_senders = control_senders.clone();
let return_node_metric_handles = node_metric_handles.clone();
let final_node_metric_handles = node_metric_handles.clone();
let final_channel_metrics = channel_metrics.clone();
let final_admission_metrics = admission_metrics.clone();
let final_metrics_reporter = metrics_reporter.clone();
let manager_pipeline_context = pipeline_context.clone();
let manager_metrics_reporter = metrics_reporter.clone();
let manager_telemetry_policy = telemetry_policy.clone();
let manager_memory_pressure_rx = memory_pressure_rx;
let manager_terminal_metrics_deadline = terminal_metrics_deadline.clone();
let dispatcher_pipeline_context = pipeline_context.clone();
let dispatcher_metrics_reporter = metrics_reporter.clone();
let dispatcher_telemetry_policy = telemetry_policy.clone();
let dispatcher_terminal_metrics_deadline = terminal_metrics_deadline.clone();
futures.push(local_tasks.spawn_local(async move {
let manager = RuntimeCtrlMsgManager::new(
pipeline_key,
manager_pipeline_context,
runtime_ctrl_msg_rx,
manager_memory_pressure_rx,
control_senders,
event_reporter,
manager_metrics_reporter,
control_plane_metrics_flush_interval,
manager_telemetry_policy,
channel_metrics,
admission_metrics,
node_metric_handles,
manager_terminal_metrics_deadline,
forced_shutdown_trigger,
);
manager.run().await
}));
futures.push(local_tasks.spawn_local(async move {
let dispatcher = PipelineCompletionMsgDispatcher::new(
dispatcher_pipeline_context,
pipeline_completion_msg_rx,
return_control_senders,
return_node_metric_handles,
dispatcher_metrics_reporter,
control_plane_metrics_flush_interval,
dispatcher_telemetry_policy,
dispatcher_terminal_metrics_deadline,
);
dispatcher.run().await
}));
let result = rt.block_on(async {
local_tasks
.run_until(async {
let loop_result: Result<Vec<_>, Error> = async {
let mut task_results = Vec::new();
loop {
tokio::select! {
biased;
Some(result) = futures.next(), if !futures.is_empty() => {
match result {
Ok(Ok(res)) => task_results.push(res),
Ok(Err(e)) => return Err(e),
Err(e) => return Err(Error::JoinTaskError {
is_canceled: e.is_cancelled(),
is_panic: e.is_panic(),
error: e.to_string(),
}),
}
}
event = extension_lifecycle.next_event() => {
match event {
crate::extension_lifecycle::LifecycleEvent::Completion(result) => {
match result {
Ok(Ok(())) => {}
Ok(Err(e)) => return Err(e),
Err(e) => return Err(Error::JoinTaskError {
is_canceled: e.is_cancelled(),
is_panic: e.is_panic(),
error: e.to_string(),
}),
}
}
crate::extension_lifecycle::LifecycleEvent::MonitorTick(now) => {
let mut reporter = metrics_reporter.clone();
extension_lifecycle.monitor_tick(now, &mut reporter);
}
}
}
else => break,
}
if futures.is_empty() {
extension_lifecycle
.initiate_shutdown(Some("pipeline data-path drained"));
break;
}
}
Ok(task_results)
}
.await;
extension_lifecycle
.initiate_shutdown(Some("pipeline data-path drained"));
extension_lifecycle.drain_until_deadline().await;
let mut final_monitor_reporter = metrics_reporter.clone();
extension_lifecycle
.monitor_tick(std::time::Instant::now(), &mut final_monitor_reporter);
if let Err(err) = extension_lifecycle
.finish_metrics_reporting_until(
&final_monitor_reporter,
terminal_metrics_deadline.get(),
)
.await
{
otel_arrow_dfe_telemetry::otel_warn!(
"extension.lifecycle.metrics.final_reporting.fail",
error = err.to_string()
);
}
let final_snapshots = snapshot_node_metrics_with_handles(
&final_node_metric_handles,
)
.into_iter()
.chain(
final_channel_metrics
.iter()
.flat_map(ChannelMetricsHandle::terminal_snapshots),
)
.chain(
final_admission_metrics
.iter()
.flat_map(|metrics| metrics.terminal_snapshots()),
);
report_metric_snapshots(
&final_metrics_reporter,
final_snapshots,
"pipeline_final",
terminal_metrics_deadline.get(),
)
.await;
let task_results = loop_result?;
Ok(task_results)
})
.await
});
drop(node_telemetry_guards);
result
}
}
impl<PData: 'static + Debug + Clone> RuntimePipeline<PData> {
#[must_use]
pub fn get_node(&self, node_id: usize) -> Option<&dyn Node<PData>> {
let ndef = self.nodes.get(node_id)?;
match ndef.ntype {
NodeType::Receiver => self
.receivers
.get(ndef.inner.index)
.map(|r| r as &dyn Node<PData>),
NodeType::Processor => self
.processors
.get(ndef.inner.index)
.map(|p| p as &dyn Node<PData>),
NodeType::Exporter => self
.exporters
.get(ndef.inner.index)
.map(|e| e as &dyn Node<PData>),
}
}
#[must_use]
pub fn get_mut_node_with_pdata_sender(
&mut self,
node_id: usize,
) -> Option<&mut dyn NodeWithPDataSender<PData>> {
let ndef = self.nodes.get(node_id)?;
match ndef.ntype {
NodeType::Receiver => self
.receivers
.get_mut(ndef.inner.index)
.map(|r| r as &mut dyn NodeWithPDataSender<PData>),
NodeType::Processor => self
.processors
.get_mut(ndef.inner.index)
.map(|p| p as &mut dyn NodeWithPDataSender<PData>),
NodeType::Exporter => None,
}
}
#[must_use]
pub fn get_mut_node_with_pdata_receiver(
&mut self,
node_id: usize,
) -> Option<&mut dyn NodeWithPDataReceiver<PData>> {
let ndef = self.nodes.get(node_id)?;
match ndef.ntype {
NodeType::Receiver => None,
NodeType::Processor => self
.processors
.get_mut(ndef.inner.index)
.map(|p| p as &mut dyn NodeWithPDataReceiver<PData>),
NodeType::Exporter => self
.exporters
.get_mut(ndef.inner.index)
.map(|e| e as &mut dyn NodeWithPDataReceiver<PData>),
}
}
pub async fn send_node_control_message(
&self,
node_id: &NodeId,
ctrl_msg: NodeControlMsg<PData>,
) -> Result<(), TypedError<NodeControlMsg<PData>>> {
match self.nodes.get(node_id.index) {
Some(ndef) => match ndef.ntype {
NodeType::Receiver => {
self.receivers
.get(ndef.inner.index)
.expect("precomputed")
.send_control_msg(ctrl_msg)
.await
}
NodeType::Processor => {
self.processors
.get(ndef.inner.index)
.expect("precomputed")
.send_control_msg(ctrl_msg)
.await
}
NodeType::Exporter => {
self.exporters
.get(ndef.inner.index)
.expect("precomputed")
.send_control_msg(ctrl_msg)
.await
}
}
.map_err(|e| TypedError::NodeControlMsgSendError {
node_id: node_id.index,
error: e,
}),
None => Err(TypedError::Error(Error::InternalError {
message: format!("node {node_id:?}"),
})),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::attributes::{
ChannelImplementation, ChannelKind, ChannelMode, ChannelType, EngineEntityAttributeSet,
};
use crate::channel_metrics::ChannelSenderMetrics;
use crate::entity_context::{NodeTelemetryGuard, NodeTelemetryHandle};
use otel_arrow_dfe_config::SignalType;
use otel_arrow_dfe_config::observed_state::SendPolicy;
use otel_arrow_dfe_config::pipeline::telemetry::TelemetryConfig;
use otel_arrow_dfe_telemetry::common_attributes::{Outcome, SignalOutcomeAttributes};
use otel_arrow_dfe_telemetry::registry::TelemetryRegistryHandle;
use otel_arrow_dfe_telemetry::{InternalTelemetrySystem, LogContext};
#[test]
fn optional_node_metrics_do_not_require_direction_interests() {
let (pipeline_context, registry) = crate::testing::test_pipeline_ctx();
let node_entity_key = pipeline_context.register_node_entity();
let telemetry_handle = NodeTelemetryHandle::new(registry.clone(), node_entity_key);
let input_key = pipeline_context.register_node_channel_entity(
"input-channel".into(),
"input".into(),
ChannelKind::Pdata,
ChannelMode::Local,
ChannelType::Mpsc,
ChannelImplementation::Internal,
);
telemetry_handle.set_input_channel_key(input_key);
let output_key = pipeline_context.register_node_channel_entity(
"output-channel".into(),
"default".into(),
ChannelKind::Pdata,
ChannelMode::Local,
ChannelType::Mpsc,
ChannelImplementation::Internal,
);
telemetry_handle.add_output_channel_key("default".into(), output_key);
let handles = make_node_metric_handles(
&Some(telemetry_handle),
&pipeline_context,
true,
true,
Interests::NODE_COMPLETION_DURATION
| Interests::NODE_ITEM_COUNTS
| Interests::NODE_SIZE,
None,
);
assert!(handles.input.is_none());
assert!(handles.input_completion.is_some());
assert!(handles.input_items.is_some());
assert!(handles.input_size.is_some());
assert!(handles.outputs.is_empty());
assert!(handles.output_completion.is_empty());
assert_eq!(handles.output_items.len(), 1);
assert_eq!(handles.output_size.len(), 1);
assert_eq!(registry.metric_set_count(), 5);
}
#[tokio::test(flavor = "current_thread")]
async fn terminal_snapshot_is_aggregated_before_telemetry_cleanup() {
let registry = TelemetryRegistryHandle::new();
let config = TelemetryConfig::default();
let metrics_system = InternalTelemetrySystem::new(
&config,
config.reporting_interval,
registry.clone(),
None,
SendPolicy::default(),
LogContext::new,
None,
)
.expect("ITS telemetry system should initialize");
let reporter = metrics_system.reporter();
let collector_task = tokio::spawn(metrics_system.collector().run_collection_loop());
let entity_key = registry.register_entity(EngineEntityAttributeSet);
let telemetry_handle = NodeTelemetryHandle::new(registry.clone(), entity_key);
let telemetry_guard = NodeTelemetryGuard::new(telemetry_handle.clone());
let mut metric_set = telemetry_handle
.entity_handle()
.register_measurement_metric_set_for_entity::<ChannelSenderMetrics>(entity_key);
metric_set
.with(SignalOutcomeAttributes {
signal: SignalType::Logs,
outcome: Outcome::Success,
})
.messages
.inc();
report_terminal_metrics(
&reporter,
TerminalState::new(std::time::Instant::now(), metric_set.terminal_snapshots()),
&TerminalMetricsDeadline::default(),
)
.await;
let mut message_count = None;
registry.visit_current_metrics(|_, _, metrics| {
for (field, value) in metrics {
if field.name == "messages" {
message_count = Some(value.to_u64_lossy());
}
}
});
assert_eq!(message_count, Some(1));
drop(telemetry_guard);
assert_eq!(registry.metric_set_count(), 1);
assert_eq!(registry.entity_count(), 1);
let export_batch = registry.drain_metric_export_batch();
let exported = export_batch
.metric_sets
.iter()
.find(|metric_set| metric_set.descriptor.name == "channel.sender")
.expect("terminal metric set should remain exportable after cleanup");
let message_count_index = exported
.descriptor
.metrics
.iter()
.position(|field| field.name == "messages")
.expect("messages descriptor should exist");
assert_eq!(exported.values[message_count_index].to_u64_lossy(), 1);
assert_eq!(registry.metric_set_count(), 0);
assert_eq!(registry.entity_count(), 0);
collector_task.abort();
assert!(
collector_task
.await
.expect_err("collector task should be cancelled")
.is_cancelled()
);
}
}