use crate::Interests;
use crate::ReceivedAtNode;
use crate::channel_metrics::ChannelMetricsRegistry;
use crate::channel_mode::{LocalMode, SharedMode, wrap_node_control_channel_metrics};
use crate::completion_emission_metrics::CompletionEmissionMetricsHandle;
use crate::config::ProcessorConfig;
use crate::context::PipelineContext;
use crate::control::{
Controllable, NodeControlMsg, PipelineCompletionMsgSender, RuntimeCtrlMsgSender,
};
use crate::effect_handler::SourceTagging;
use crate::entity_context::NodeTelemetryGuard;
use crate::error::{Error, ProcessorErrorKind};
use crate::flow_metrics::{
FlowDroppedItemsMetrics, FlowDurationMetricSet, FlowInputItemsMetrics, FlowInputMessageMetrics,
FlowInputSizeMetrics, FlowOutputItemsMetrics, FlowOutputMessageMetrics, FlowOutputSizeMetrics,
};
use crate::forced_shutdown::ForcedShutdownSignal;
use crate::local::message::{LocalReceiver, LocalSender};
use crate::local::processor as local;
use crate::message::{Message, ProcessorInbox, Receiver, Sender};
use crate::node::{Node, NodeId, NodeWithPDataReceiver, NodeWithPDataSender};
use crate::node_local_scheduler::NodeLocalSchedulerHandle;
use crate::runtime_services::PipelineRuntimeServices;
use crate::shared::message::{SharedReceiver, SharedSender};
use crate::shared::processor as shared;
use crate::terminal_state::TerminalMetricsDeadline;
use otel_arrow_dfe_channel::error::SendError;
use otel_arrow_dfe_channel::mpsc;
use otel_arrow_dfe_config::node::NodeUserConfig;
use otel_arrow_dfe_config::{PortName, SignalType};
use otel_arrow_dfe_telemetry::metrics::MeasurementMetricSet;
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use std::collections::HashMap;
use std::sync::Arc;
pub trait FlowMetricEffectHandler {
fn is_flow_start(&self) -> bool;
fn is_flow_end(&self) -> bool;
fn flow_metric_interests(&self) -> crate::flow_metrics::FlowMetricInterests;
fn take_elapsed_since_send_marker_ns(&self) -> u64;
fn record_flow_duration(&self, signal: SignalType, total: u64);
fn record_flow_input_items(&self, signal: SignalType, items: u64);
fn record_flow_input_message(&self, signal: SignalType);
fn record_flow_input_size(&self, signal: SignalType, size: u64);
fn record_flow_output_items(&self, signal: SignalType, items: u64);
fn record_flow_output_message(&self, signal: SignalType);
fn record_flow_output_size(&self, signal: SignalType, size: u64);
}
pub trait FlowMetricHook: Sized {
fn before_processor_send<H: FlowMetricEffectHandler>(&mut self, _handler: &H) {}
fn complete_processor_without_output<H: FlowMetricEffectHandler>(&mut self, handler: &H) {
self.before_processor_send(handler);
}
fn after_processor_receive<H: FlowMetricEffectHandler>(&mut self, _handler: &H) {}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub struct LocalWakeupRequirements {
pub live_slots: usize,
}
impl LocalWakeupRequirements {
#[must_use]
pub const fn new(live_slots: usize) -> Self {
Self { live_slots }
}
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct ProcessorRuntimeRequirements {
pub local_wakeups: Option<LocalWakeupRequirements>,
pub makes_drop_decisions: bool,
}
impl ProcessorRuntimeRequirements {
#[must_use]
pub const fn none() -> Self {
Self {
local_wakeups: None,
makes_drop_decisions: false,
}
}
#[must_use]
pub const fn with_local_wakeups(live_slots: usize) -> Self {
Self {
local_wakeups: Some(LocalWakeupRequirements::new(live_slots)),
makes_drop_decisions: false,
}
}
#[must_use]
pub const fn with_drop_decisions(mut self) -> Self {
self.makes_drop_decisions = true;
self
}
}
pub enum ProcessorWrapper<PData> {
Local {
node_id: NodeId,
user_config: Arc<NodeUserConfig>,
runtime_config: ProcessorConfig,
processor: Box<dyn local::Processor<PData>>,
control_sender: LocalSender<NodeControlMsg<PData>>,
control_receiver: LocalReceiver<NodeControlMsg<PData>>,
pdata_senders: HashMap<PortName, Sender<PData>>,
pdata_receiver: Option<Receiver<PData>>,
telemetry: Option<NodeTelemetryGuard>,
source_tag: SourceTagging,
},
Shared {
node_id: NodeId,
user_config: Arc<NodeUserConfig>,
runtime_config: ProcessorConfig,
processor: Box<dyn shared::Processor<PData>>,
control_sender: SharedSender<NodeControlMsg<PData>>,
control_receiver: SharedReceiver<NodeControlMsg<PData>>,
pdata_senders: HashMap<PortName, SharedSender<PData>>,
pdata_receiver: Option<SharedReceiver<PData>>,
telemetry: Option<NodeTelemetryGuard>,
source_tag: SourceTagging,
},
}
#[allow(clippy::large_enum_variant)]
pub enum ProcessorWrapperRuntime<PData> {
Local {
processor: Box<dyn local::Processor<PData>>,
inbox: ProcessorInbox<PData>,
effect_handler: local::EffectHandler<PData>,
},
Shared {
processor: Box<dyn shared::Processor<PData>>,
inbox: ProcessorInbox<PData>,
effect_handler: shared::EffectHandler<PData>,
},
}
impl<PData> ProcessorWrapper<PData> {
pub fn local<P>(
processor: P,
node_id: NodeId,
user_config: Arc<NodeUserConfig>,
config: &ProcessorConfig,
) -> Self
where
P: local::Processor<PData> + 'static,
{
let runtime_config = config.clone();
let (control_sender, control_receiver) =
mpsc::Channel::new(config.control_channel.capacity);
ProcessorWrapper::Local {
node_id,
user_config,
runtime_config,
processor: Box::new(processor),
control_sender: LocalSender::mpsc(control_sender),
control_receiver: LocalReceiver::mpsc(control_receiver),
pdata_senders: HashMap::new(),
pdata_receiver: None,
telemetry: None,
source_tag: SourceTagging::Disabled,
}
}
pub fn shared<P>(
processor: P,
node_id: NodeId,
user_config: Arc<NodeUserConfig>,
config: &ProcessorConfig,
) -> Self
where
P: shared::Processor<PData> + 'static,
{
let runtime_config = config.clone();
let (control_sender, control_receiver) =
tokio::sync::mpsc::channel(config.control_channel.capacity);
ProcessorWrapper::Shared {
node_id,
user_config,
runtime_config,
processor: Box::new(processor),
control_sender: SharedSender::mpsc(control_sender),
control_receiver: SharedReceiver::mpsc(control_receiver),
pdata_senders: HashMap::new(),
pdata_receiver: None,
telemetry: None,
source_tag: SourceTagging::Disabled,
}
}
pub(crate) fn with_node_telemetry_guard(self, guard: NodeTelemetryGuard) -> Self {
match self {
ProcessorWrapper::Local {
node_id,
user_config,
runtime_config,
processor,
control_sender,
control_receiver,
pdata_senders,
pdata_receiver,
source_tag,
..
} => ProcessorWrapper::Local {
node_id,
user_config,
runtime_config,
processor,
control_sender,
control_receiver,
pdata_senders,
pdata_receiver,
telemetry: Some(guard),
source_tag,
},
ProcessorWrapper::Shared {
node_id,
user_config,
runtime_config,
processor,
control_sender,
control_receiver,
pdata_senders,
pdata_receiver,
source_tag,
..
} => ProcessorWrapper::Shared {
node_id,
user_config,
runtime_config,
processor,
control_sender,
control_receiver,
pdata_senders,
pdata_receiver,
telemetry: Some(guard),
source_tag,
},
}
}
pub(crate) const fn take_telemetry_guard(&mut self) -> Option<NodeTelemetryGuard> {
match self {
ProcessorWrapper::Local { telemetry, .. } => telemetry.take(),
ProcessorWrapper::Shared { telemetry, .. } => telemetry.take(),
}
}
pub(crate) fn runtime_requirements(&self) -> ProcessorRuntimeRequirements {
match self {
ProcessorWrapper::Local { processor, .. } => processor.runtime_requirements(),
ProcessorWrapper::Shared { processor, .. } => processor.runtime_requirements(),
}
}
pub(crate) fn with_control_channel_metrics(
self,
pipeline_ctx: &PipelineContext,
channel_metrics: &mut ChannelMetricsRegistry,
channel_metrics_enabled: bool,
) -> Self {
match self {
ProcessorWrapper::Local {
node_id,
runtime_config,
control_sender,
control_receiver,
user_config,
processor,
pdata_senders,
pdata_receiver,
telemetry,
source_tag,
} => {
let (control_sender, control_receiver) =
wrap_node_control_channel_metrics::<LocalMode, NodeControlMsg<PData>>(
node_id.name.as_ref(),
pipeline_ctx,
channel_metrics,
channel_metrics_enabled,
runtime_config.control_channel.capacity as u64,
control_sender,
control_receiver,
);
ProcessorWrapper::Local {
node_id,
user_config,
runtime_config,
processor,
control_sender,
control_receiver,
pdata_senders,
pdata_receiver,
telemetry,
source_tag,
}
}
ProcessorWrapper::Shared {
node_id,
runtime_config,
control_sender,
control_receiver,
user_config,
processor,
pdata_senders,
pdata_receiver,
telemetry,
source_tag,
} => {
let (control_sender, control_receiver) =
wrap_node_control_channel_metrics::<SharedMode, NodeControlMsg<PData>>(
node_id.name.as_ref(),
pipeline_ctx,
channel_metrics,
channel_metrics_enabled,
runtime_config.control_channel.capacity as u64,
control_sender,
control_receiver,
);
ProcessorWrapper::Shared {
node_id,
user_config,
runtime_config,
processor,
control_sender,
control_receiver,
pdata_senders,
pdata_receiver,
telemetry,
source_tag,
}
}
}
}
pub async fn prepare_runtime(
self,
metrics_reporter: MetricsReporter,
node_interests: Interests,
runtime_services: PipelineRuntimeServices,
) -> Result<ProcessorWrapperRuntime<PData>, Error> {
match self {
ProcessorWrapper::Local {
node_id,
runtime_config,
processor,
control_receiver,
pdata_senders,
pdata_receiver,
user_config,
source_tag,
..
} => {
let runtime_requirements = processor.runtime_requirements();
let pdata_receiver = pdata_receiver.ok_or_else(|| Error::ProcessorError {
processor: node_id.clone(),
kind: ProcessorErrorKind::Configuration,
error: "The pdata receiver must be defined at this stage".to_owned(),
source_detail: String::new(),
})?;
validate_local_wakeup_requirements(&node_id, runtime_requirements)?;
let local_scheduler = NodeLocalSchedulerHandle::new(
runtime_config.input_pdata_channel.capacity,
runtime_requirements
.local_wakeups
.map(|requirements| requirements.live_slots)
.unwrap_or(0),
);
let inbox = ProcessorInbox::new_with_local_scheduler(
Receiver::Local(control_receiver),
pdata_receiver,
local_scheduler.clone(),
node_id.index,
node_interests,
);
let default_port = user_config.default_output.clone();
let mut effect_handler = local::EffectHandler::new(
node_id,
pdata_senders,
default_port,
metrics_reporter,
runtime_services.clone(),
);
effect_handler.set_source_tagging(source_tag);
effect_handler.core.set_local_scheduler(local_scheduler);
Ok(ProcessorWrapperRuntime::Local {
processor,
effect_handler,
inbox,
})
}
ProcessorWrapper::Shared {
node_id,
runtime_config,
processor,
control_receiver,
pdata_senders,
pdata_receiver,
user_config,
source_tag,
..
} => {
let runtime_requirements = processor.runtime_requirements();
let pdata_receiver =
Receiver::Shared(pdata_receiver.ok_or_else(|| Error::ProcessorError {
processor: node_id.clone(),
kind: ProcessorErrorKind::Configuration,
error: "The pdata receiver must be defined at this stage".to_owned(),
source_detail: String::new(),
})?);
validate_local_wakeup_requirements(&node_id, runtime_requirements)?;
let local_scheduler = NodeLocalSchedulerHandle::new(
runtime_config.input_pdata_channel.capacity,
runtime_requirements
.local_wakeups
.map(|requirements| requirements.live_slots)
.unwrap_or(0),
);
let inbox = ProcessorInbox::new_with_local_scheduler(
Receiver::Shared(control_receiver),
pdata_receiver,
local_scheduler.clone(),
node_id.index,
node_interests,
);
let default_port = user_config.default_output.clone();
let mut effect_handler = shared::EffectHandler::new(
node_id,
pdata_senders,
default_port,
metrics_reporter,
runtime_services,
);
effect_handler.set_source_tagging(source_tag);
effect_handler.core.set_local_scheduler(local_scheduler);
Ok(ProcessorWrapperRuntime::Shared {
processor,
effect_handler,
inbox,
})
}
}
}
pub(crate) async fn start_with_completion_metrics(
self,
runtime_ctrl_msg_tx: RuntimeCtrlMsgSender<PData>,
pipeline_completion_msg_tx: PipelineCompletionMsgSender<PData>,
metrics_reporter: MetricsReporter,
node_interests: Interests,
completion_emission_metrics: Option<CompletionEmissionMetricsHandle>,
flow_is_start: bool,
flow_is_end: bool,
flow_input_message_metric: Option<MeasurementMetricSet<FlowInputMessageMetrics>>,
flow_input_items_metric: Option<MeasurementMetricSet<FlowInputItemsMetrics>>,
flow_input_size_metric: Option<MeasurementMetricSet<FlowInputSizeMetrics>>,
flow_duration_metric: Option<FlowDurationMetricSet>,
flow_output_items_metric: Option<MeasurementMetricSet<FlowOutputItemsMetrics>>,
flow_output_message_metric: Option<MeasurementMetricSet<FlowOutputMessageMetrics>>,
flow_output_size_metric: Option<MeasurementMetricSet<FlowOutputSizeMetrics>>,
flow_dropped_items_metric: Option<MeasurementMetricSet<FlowDroppedItemsMetrics>>,
flow_metrics_active: bool,
flow_needs_timing: bool,
terminal_metrics_deadline: TerminalMetricsDeadline,
forced_shutdown_signal: ForcedShutdownSignal,
runtime_services: PipelineRuntimeServices,
) -> Result<(), Error>
where
PData: ReceivedAtNode + FlowMetricHook,
{
let runtime = self
.prepare_runtime(metrics_reporter.clone(), node_interests, runtime_services)
.await?;
let mut processing_error: Option<Error> = None;
let run = async {
match runtime {
ProcessorWrapperRuntime::Local {
mut processor,
mut inbox,
mut effect_handler,
} => {
effect_handler
.core
.set_runtime_ctrl_msg_sender(runtime_ctrl_msg_tx);
effect_handler
.core
.set_pipeline_completion_msg_sender(pipeline_completion_msg_tx);
effect_handler.core.set_node_interests(node_interests);
effect_handler
.core
.set_completion_emission_metrics(completion_emission_metrics.clone());
effect_handler.set_flow_roles(
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_metrics_active,
flow_needs_timing,
);
while let Ok(mut msg) = inbox.recv_when(processor.accept_pdata()).await {
if effect_handler.flow_metrics_active() {
match &mut msg {
Message::Control(NodeControlMsg::CollectTelemetry { .. })
if effect_handler.is_flow_start()
|| effect_handler.is_flow_end()
|| effect_handler.is_flow_decision() =>
{
effect_handler.report_flow_metrics();
}
Message::PData(data) => {
data.after_processor_receive(&effect_handler);
effect_handler.begin_process_timing();
}
_ => {}
}
}
if let Err(err) = processor.process(msg, &mut effect_handler).await {
processing_error = Some(err);
break;
}
}
let terminal_metrics_deadline = terminal_metrics_deadline.get();
if (effect_handler.is_flow_start()
|| effect_handler.is_flow_end()
|| effect_handler.is_flow_decision())
&& let Err(error) = effect_handler
.report_flow_metrics_reliably(terminal_metrics_deadline)
.await
{
otel_arrow_dfe_telemetry::otel_warn!(
"processor.flow_metrics.final_reporting.fail",
error = error.to_string()
);
}
let (terminal_metrics_tx, terminal_metrics_rx) = flume::unbounded();
let terminal_metrics_reporter = MetricsReporter::new(terminal_metrics_tx);
let collect_result = processor
.process(
Message::Control(NodeControlMsg::CollectTelemetry {
metrics_reporter: terminal_metrics_reporter,
}),
&mut effect_handler,
)
.await;
while let Ok(snapshot) = terminal_metrics_rx.try_recv() {
if let Err(error) = metrics_reporter
.report_snapshot_reliably_until(snapshot, terminal_metrics_deadline)
.await
{
otel_arrow_dfe_telemetry::otel_warn!(
"processor.metrics.final_reporting.fail",
error = error.to_string()
);
}
}
collect_result?
}
ProcessorWrapperRuntime::Shared {
mut processor,
mut inbox,
mut effect_handler,
} => {
effect_handler
.core
.set_runtime_ctrl_msg_sender(runtime_ctrl_msg_tx);
effect_handler
.core
.set_pipeline_completion_msg_sender(pipeline_completion_msg_tx);
effect_handler.core.set_node_interests(node_interests);
effect_handler
.core
.set_completion_emission_metrics(completion_emission_metrics);
effect_handler.set_flow_roles(
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_metrics_active,
flow_needs_timing,
);
while let Ok(mut msg) = inbox.recv_when(processor.accept_pdata()).await {
if effect_handler.flow_metrics_active() {
match &mut msg {
Message::Control(NodeControlMsg::CollectTelemetry { .. })
if effect_handler.is_flow_start()
|| effect_handler.is_flow_end()
|| effect_handler.is_flow_decision() =>
{
effect_handler.report_flow_metrics();
}
Message::PData(data) => {
data.after_processor_receive(&effect_handler);
effect_handler.begin_process_timing();
}
_ => {}
}
}
if let Err(err) = processor.process(msg, &mut effect_handler).await {
processing_error = Some(err);
break;
}
}
let terminal_metrics_deadline = terminal_metrics_deadline.get();
if (effect_handler.is_flow_start()
|| effect_handler.is_flow_end()
|| effect_handler.is_flow_decision())
&& let Err(error) = effect_handler
.report_flow_metrics_reliably(terminal_metrics_deadline)
.await
{
otel_arrow_dfe_telemetry::otel_warn!(
"processor.flow_metrics.final_reporting.fail",
error = error.to_string()
);
}
let (terminal_metrics_tx, terminal_metrics_rx) = flume::unbounded();
let terminal_metrics_reporter = MetricsReporter::new(terminal_metrics_tx);
let collect_result = processor
.process(
Message::Control(NodeControlMsg::CollectTelemetry {
metrics_reporter: terminal_metrics_reporter,
}),
&mut effect_handler,
)
.await;
while let Ok(snapshot) = terminal_metrics_rx.try_recv() {
if let Err(error) = metrics_reporter
.report_snapshot_reliably_until(snapshot, terminal_metrics_deadline)
.await
{
otel_arrow_dfe_telemetry::otel_warn!(
"processor.metrics.final_reporting.fail",
error = error.to_string()
);
}
}
collect_result?
}
}
Ok(())
};
let result = tokio::select! {
biased;
_ = forced_shutdown_signal.triggered() => Ok(()),
result = run => result,
};
processing_error.map_or(result, Err)
}
pub const fn take_pdata_receiver(&mut self) -> Receiver<PData> {
match self {
ProcessorWrapper::Local { pdata_receiver, .. } => {
pdata_receiver.take().expect("pdata_receiver is None")
}
ProcessorWrapper::Shared { pdata_receiver, .. } => {
Receiver::Shared(pdata_receiver.take().expect("pdata_receiver is None"))
}
}
}
}
#[async_trait::async_trait(?Send)]
impl<PData> Node<PData> for ProcessorWrapper<PData> {
fn is_shared(&self) -> bool {
match self {
ProcessorWrapper::Local { .. } => false,
ProcessorWrapper::Shared { .. } => true,
}
}
fn node_id(&self) -> NodeId {
match self {
ProcessorWrapper::Local { node_id, .. } => node_id.clone(),
ProcessorWrapper::Shared { node_id, .. } => node_id.clone(),
}
}
fn user_config(&self) -> Arc<NodeUserConfig> {
match self {
ProcessorWrapper::Local {
user_config: config,
..
} => config.clone(),
ProcessorWrapper::Shared {
user_config: config,
..
} => config.clone(),
}
}
async fn send_control_msg(
&self,
msg: NodeControlMsg<PData>,
) -> Result<(), SendError<NodeControlMsg<PData>>> {
match self {
ProcessorWrapper::Local { control_sender, .. } => control_sender.send(msg).await,
ProcessorWrapper::Shared { control_sender, .. } => control_sender.send(msg).await,
}
}
}
pub(crate) fn validate_local_wakeup_requirements(
node_id: &NodeId,
requirements: ProcessorRuntimeRequirements,
) -> Result<(), Error> {
let Some(local_wakeups) = requirements.local_wakeups else {
return Ok(());
};
if local_wakeups.live_slots == 0 {
return Err(Error::ProcessorError {
processor: node_id.clone(),
kind: ProcessorErrorKind::Configuration,
error: "processor-local wakeup requirement must declare at least one live slot"
.to_owned(),
source_detail: String::new(),
});
}
Ok(())
}
#[async_trait::async_trait(?Send)]
impl<PData> Controllable<PData> for ProcessorWrapper<PData> {
fn control_sender(&self) -> Sender<NodeControlMsg<PData>> {
match self {
ProcessorWrapper::Local { control_sender, .. } => Sender::Local(control_sender.clone()),
ProcessorWrapper::Shared { control_sender, .. } => {
Sender::Shared(control_sender.clone())
}
}
}
}
impl<PData> NodeWithPDataSender<PData> for ProcessorWrapper<PData> {
fn set_pdata_sender(
&mut self,
node_id: NodeId,
port: PortName,
sender: Sender<PData>,
) -> Result<(), Error> {
match (self, sender) {
(ProcessorWrapper::Local { pdata_senders, .. }, sender) => {
let _ = pdata_senders.insert(port, sender);
Ok(())
}
(ProcessorWrapper::Shared { pdata_senders, .. }, Sender::Shared(sender)) => {
let _ = pdata_senders.insert(port, sender);
Ok(())
}
(ProcessorWrapper::Shared { .. }, _) => Err(Error::ProcessorError {
processor: node_id,
kind: ProcessorErrorKind::Configuration,
error: "Expected a shared sender for PData".to_owned(),
source_detail: String::new(),
}),
}
}
fn set_source_tagging(&mut self, value: SourceTagging) {
match self {
ProcessorWrapper::Local { source_tag, .. } => *source_tag = value,
ProcessorWrapper::Shared { source_tag, .. } => *source_tag = value,
}
}
}
impl<PData> NodeWithPDataReceiver<PData> for ProcessorWrapper<PData> {
fn set_pdata_receiver(
&mut self,
node_id: NodeId,
receiver: Receiver<PData>,
) -> Result<(), Error> {
match (self, receiver) {
(ProcessorWrapper::Local { pdata_receiver, .. }, receiver) => {
*pdata_receiver = Some(receiver);
Ok(())
}
(ProcessorWrapper::Shared { pdata_receiver, .. }, Receiver::Shared(receiver)) => {
*pdata_receiver = Some(receiver);
Ok(())
}
(ProcessorWrapper::Shared { .. }, _) => Err(Error::ProcessorError {
processor: node_id,
kind: ProcessorErrorKind::Configuration,
error: "Expected a shared receiver for PData".to_owned(),
source_detail: String::new(),
}),
}
}
}
#[cfg(test)]
mod tests {
use crate::config::ProcessorConfig;
use crate::control::{
Controllable, NodeControlMsg,
NodeControlMsg::{Config, Shutdown, TimerTick},
pipeline_completion_msg_channel, runtime_ctrl_msg_channel,
};
use crate::error::ProcessorErrorKind;
use crate::flow_metrics::{
FlowAttributeSet, FlowDroppedItemsMetrics, FlowDurationNormalMetrics,
FlowInputItemsMetrics, FlowOutputItemsMetrics,
};
use crate::local::message::{LocalReceiver, LocalSender};
use crate::local::processor as local;
use crate::message::{Message, Receiver, Sender};
use crate::node::{Node, NodeWithPDataReceiver, NodeWithPDataSender};
use crate::processor::{
Error, ProcessorRuntimeRequirements, ProcessorWrapper, validate_local_wakeup_requirements,
};
use crate::shared::message::{SharedReceiver, SharedSender};
use crate::shared::processor as shared;
use crate::testing::processor::TestRuntime;
use crate::testing::processor::{TestContext, ValidateContext};
use crate::testing::{CtrlMsgCounters, TestMsg, test_node};
use async_trait::async_trait;
use otel_arrow_dfe_config::{SignalType, node::NodeUserConfig};
use otel_arrow_dfe_telemetry::common_attributes::SignalAttributes;
use otel_arrow_dfe_telemetry::metrics::{MeasurementMetricSet, MetricValue};
use serde_json::Value;
use std::ops::Add;
use std::pin::Pin;
use std::sync::Arc;
use std::time::{Duration, Instant};
pub struct TestProcessor {
ctrl_msg_counters: CtrlMsgCounters,
}
impl TestProcessor {
pub fn new(ctrl_msg_counters: CtrlMsgCounters) -> Self {
TestProcessor { ctrl_msg_counters }
}
}
#[async_trait(?Send)]
impl local::Processor<TestMsg> for TestProcessor {
async fn process(
&mut self,
msg: Message<TestMsg>,
effect_handler: &mut local::EffectHandler<TestMsg>,
) -> Result<(), Error> {
match msg {
Message::Control(control) => match control {
TimerTick {} => {
self.ctrl_msg_counters.increment_timer_tick();
}
Config { .. } => {
self.ctrl_msg_counters.increment_config();
}
Shutdown { .. } => {
self.ctrl_msg_counters.increment_shutdown();
}
_ => {}
},
Message::PData(data) => {
self.ctrl_msg_counters.increment_message();
effect_handler
.send_message(TestMsg(format!("{} RECEIVED", data.0)))
.await?;
}
}
Ok(())
}
}
#[async_trait]
impl shared::Processor<TestMsg> for TestProcessor {
async fn process(
&mut self,
msg: Message<TestMsg>,
effect_handler: &mut shared::EffectHandler<TestMsg>,
) -> Result<(), Error> {
match msg {
Message::Control(control) => match control {
TimerTick {} => {
self.ctrl_msg_counters.increment_timer_tick();
}
Config { .. } => {
self.ctrl_msg_counters.increment_config();
}
Shutdown { .. } => {
self.ctrl_msg_counters.increment_shutdown();
}
_ => {}
},
Message::PData(data) => {
self.ctrl_msg_counters.increment_message();
effect_handler
.send_message(TestMsg(format!("{} RECEIVED", data.0)))
.await?;
}
}
Ok(())
}
}
fn scenario() -> impl FnOnce(TestContext<TestMsg>) -> Pin<Box<dyn Future<Output = ()>>> {
move |mut ctx| {
Box::pin(async move {
ctx.process(Message::timer_tick_ctrl_msg())
.await
.expect("Processor failed on TimerTick");
assert!(ctx.drain_pdata().await.is_empty());
ctx.process(Message::data_msg(TestMsg("Hello".to_owned())))
.await
.expect("Processor failed on Message");
let msgs = ctx.drain_pdata().await;
assert_eq!(msgs.len(), 1);
assert_eq!(msgs[0], TestMsg("Hello RECEIVED".to_string()));
ctx.process(Message::config_ctrl_msg(Value::Null))
.await
.expect("Processor failed on Config");
assert!(ctx.drain_pdata().await.is_empty());
ctx.process(Message::shutdown_ctrl_msg(
Instant::now().add(Duration::from_millis(200)),
"no reason",
))
.await
.expect("Processor failed on Shutdown");
assert!(ctx.drain_pdata().await.is_empty());
})
}
}
fn validation_procedure() -> impl FnOnce(ValidateContext) -> Pin<Box<dyn Future<Output = ()>>> {
|ctx| {
Box::pin(async move {
ctx.counters().assert(
1, 1, 1, 1, );
})
}
}
#[test]
fn test_processor_local() {
let test_runtime = TestRuntime::new();
let user_config = Arc::new(NodeUserConfig::new_processor_config("test_processor"));
let processor = ProcessorWrapper::local(
TestProcessor::new(test_runtime.counters()),
test_node(test_runtime.config().name.clone()),
user_config,
test_runtime.config(),
);
test_runtime
.set_processor(processor)
.run_test(scenario())
.validate(validation_procedure());
}
#[test]
fn test_processor_shared() {
let test_runtime = TestRuntime::new();
let user_config = Arc::new(NodeUserConfig::new_processor_config("test_processor"));
let processor = ProcessorWrapper::shared(
TestProcessor::new(test_runtime.counters()),
test_node(test_runtime.config().name.clone()),
user_config,
test_runtime.config(),
);
test_runtime
.set_processor(processor)
.run_test(scenario())
.validate(validation_procedure());
}
#[test]
fn validate_local_wakeup_requirements_accepts_processors_without_wakeups() {
assert!(
validate_local_wakeup_requirements(
&test_node("test_processor"),
ProcessorRuntimeRequirements::none(),
)
.is_ok()
);
}
#[test]
fn validate_local_wakeup_requirements_rejects_zero_live_slots() {
let err = validate_local_wakeup_requirements(
&test_node("test_processor"),
ProcessorRuntimeRequirements::with_local_wakeups(0),
)
.expect_err("zero live slots must be rejected");
let Error::ProcessorError { error, .. } = err else {
panic!("expected processor configuration error");
};
assert_eq!(
error,
"processor-local wakeup requirement must declare at least one live slot"
);
}
#[test]
fn validate_local_wakeup_requirements_accepts_positive_live_slots() {
assert!(
validate_local_wakeup_requirements(
&test_node("test_processor"),
ProcessorRuntimeRequirements::with_local_wakeups(6),
)
.is_ok()
);
}
#[derive(Clone, Debug, Default)]
struct FlowMetricTestPData {
flow_compute_ns: u64,
flow_metric_active: bool,
}
impl crate::ReceivedAtNode for FlowMetricTestPData {
fn received_at_node(&mut self, _node_id: usize, _node_interests: crate::Interests) {}
}
impl crate::processor::FlowMetricHook for FlowMetricTestPData {
fn before_processor_send<H: crate::processor::FlowMetricEffectHandler>(
&mut self,
handler: &H,
) {
if !handler.is_flow_start() && !handler.is_flow_end() && !self.flow_metric_active {
return;
}
if handler.is_flow_start() {
self.flow_metric_active = true;
}
self.flow_compute_ns = self
.flow_compute_ns
.saturating_add(handler.take_elapsed_since_send_marker_ns());
if handler.is_flow_end() && self.flow_metric_active && self.flow_compute_ns > 0 {
handler.record_flow_duration(SignalType::Logs, self.flow_compute_ns);
handler.record_flow_output_items(SignalType::Logs, 1);
self.flow_compute_ns = 0;
self.flow_metric_active = false;
}
}
fn after_processor_receive<H: crate::processor::FlowMetricEffectHandler>(
&mut self,
handler: &H,
) {
if handler.is_flow_start() {
handler.record_flow_input_items(SignalType::Logs, 1);
}
}
}
struct SyncOnlyFlowMetricProcessor;
#[async_trait(?Send)]
impl local::Processor<FlowMetricTestPData> for SyncOnlyFlowMetricProcessor {
async fn process(
&mut self,
msg: Message<FlowMetricTestPData>,
effect_handler: &mut local::EffectHandler<FlowMetricTestPData>,
) -> Result<(), Error> {
let Message::PData(data) = msg else {
return Ok(());
};
let mut value = 0u64;
for i in 0..50_000 {
value = value.wrapping_add(std::hint::black_box(i));
}
let _ = std::hint::black_box(value);
tokio::task::yield_now().await;
effect_handler.send_message(data).await?;
Ok(())
}
}
#[test]
fn flow_opt_in_input_items_reports_only_start_metric() {
let (pipeline_ctx, _) = crate::testing::test_pipeline_ctx();
let entity_key = pipeline_ctx
.metrics_registry()
.register_entity(FlowAttributeSet::default());
let incoming_metric = FlowInputItemsMetrics::register(
&pipeline_ctx.metric_set_registrar_for_entity(entity_key),
);
let (metrics_rx, metrics_reporter) =
otel_arrow_dfe_telemetry::reporter::MetricsReporter::create_new_and_receiver(4);
let mut handler = local::EffectHandler::<TestMsg>::new(
test_node("proc"),
std::collections::HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
handler.set_flow_roles(
true,
false,
None,
Some(incoming_metric),
None,
None,
None,
None,
None,
None,
true,
false,
);
handler.record_flow_input_items(SignalType::Logs, 3);
handler.record_flow_duration(SignalType::Logs, 10);
handler.record_flow_output_items(SignalType::Logs, 4);
handler.report_flow_metrics();
let snapshot = metrics_rx
.try_recv()
.expect("incoming metric should report");
let [MetricValue::U64(input_items)] = snapshot.get_metrics() else {
panic!("expected input item metric only");
};
assert_eq!(*input_items, 3);
assert!(metrics_rx.try_recv().is_err());
}
#[test]
fn flow_opt_in_duration_and_output_items_reports_only_end_metrics() {
let (pipeline_ctx, _) = crate::testing::test_pipeline_ctx();
let entity_key = pipeline_ctx
.metrics_registry()
.register_entity(FlowAttributeSet::default());
let registrar = pipeline_ctx.metric_set_registrar_for_entity(entity_key);
let duration_metric = FlowDurationNormalMetrics::register(®istrar);
let output_items_metric = FlowOutputItemsMetrics::register(®istrar);
let (metrics_rx, metrics_reporter) =
otel_arrow_dfe_telemetry::reporter::MetricsReporter::create_new_and_receiver(4);
let mut handler = local::EffectHandler::<TestMsg>::new(
test_node("proc"),
std::collections::HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
handler.set_flow_roles(
false,
true,
None,
None,
None,
Some(duration_metric.into()),
Some(output_items_metric),
None,
None,
None,
true,
true,
);
handler.record_flow_input_items(SignalType::Logs, 3);
handler.record_flow_duration(SignalType::Logs, 10);
handler.record_flow_output_items(SignalType::Logs, 4);
handler.report_flow_metrics();
let duration_snapshot = metrics_rx
.try_recv()
.expect("duration metric should report");
let [MetricValue::Distribution(duration)] = duration_snapshot.get_metrics() else {
panic!("expected duration metric");
};
assert_eq!(duration.summary().0, 1);
let output_items_snapshot = metrics_rx
.try_recv()
.expect("output item metric should report");
let [MetricValue::U64(output_items)] = output_items_snapshot.get_metrics() else {
panic!("expected output item metric");
};
assert_eq!(*output_items, 4);
assert!(metrics_rx.try_recv().is_err());
}
#[test]
fn flow_decision_node_reports_dropped() {
let (pipeline_ctx, _) = crate::testing::test_pipeline_ctx();
let entity_key = pipeline_ctx
.metrics_registry()
.register_entity(FlowAttributeSet::default());
let dropped_metric = FlowDroppedItemsMetrics::register(
&pipeline_ctx.metric_set_registrar_for_entity(entity_key),
);
let (metrics_rx, metrics_reporter) =
otel_arrow_dfe_telemetry::reporter::MetricsReporter::create_new_and_receiver(4);
let mut handler = local::EffectHandler::<TestMsg>::new(
test_node("proc"),
std::collections::HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
handler.set_flow_roles(
false,
false,
None,
None,
None,
None,
None,
None,
None,
Some(dropped_metric),
true,
false,
);
assert!(handler.is_flow_decision());
handler.record_flow_dropped_items(SignalType::Logs, 3);
handler.record_flow_input_items(SignalType::Logs, 99);
handler.record_flow_output_items(SignalType::Logs, 99);
handler.report_flow_metrics();
let dropped_snapshot = metrics_rx.try_recv().expect("dropped metric should report");
let [MetricValue::U64(dropped_items)] = dropped_snapshot.get_metrics() else {
panic!("expected dropped metric");
};
assert_eq!(*dropped_items, 3);
assert!(metrics_rx.try_recv().is_err());
}
#[tokio::test]
async fn flow_metric_auto_measures_process_without_timed() {
let (pipeline_ctx, _) = crate::testing::test_pipeline_ctx();
let attrs = FlowAttributeSet {
flow_id: "auto_measure".into(),
start_node: "auto_measure_processor".into(),
end_node: "auto_measure_processor".into(),
purpose: "".into(),
decision: "".into(),
pipeline_attrs: pipeline_ctx.pipeline_attribute_set(),
};
let entity_key = pipeline_ctx.metrics_registry().register_entity(attrs);
let registrar = pipeline_ctx.metric_set_registrar_for_entity(entity_key);
let start_metric_set = FlowInputItemsMetrics::register(®istrar);
let duration_metric_set = FlowDurationNormalMetrics::register(®istrar);
let outgoing_metric_set = FlowOutputItemsMetrics::register(®istrar);
let config = ProcessorConfig::new("auto_measure_processor");
let node_id = test_node(config.name.clone());
let user_config = Arc::new(NodeUserConfig::new_processor_config(
"auto_measure_processor",
));
let mut processor = ProcessorWrapper::local(
SyncOnlyFlowMetricProcessor,
node_id.clone(),
user_config,
&config,
);
let (input_tx, input_rx) = otel_arrow_dfe_channel::mpsc::Channel::new(1);
processor
.set_pdata_receiver(
node_id.clone(),
Receiver::Local(LocalReceiver::mpsc(input_rx)),
)
.expect("input receiver should be accepted");
let (output_tx, output_rx) = otel_arrow_dfe_channel::mpsc::Channel::new(1);
processor
.set_pdata_sender(
node_id,
"out".into(),
Sender::Local(LocalSender::mpsc(output_tx)),
)
.expect("output sender should be accepted");
let control_sender = processor.control_sender();
let (metrics_rx, metrics_reporter) =
otel_arrow_dfe_telemetry::reporter::MetricsReporter::create_new_and_receiver(8);
let collect_metrics_reporter = metrics_reporter.clone();
let (runtime_ctrl_tx, _runtime_ctrl_rx) = runtime_ctrl_msg_channel(1);
let (completion_tx, _completion_rx) = pipeline_completion_msg_channel(1);
let local_tasks = tokio::task::LocalSet::new();
local_tasks
.run_until(async move {
let processor_task = tokio::task::spawn_local(async move {
let (_, forced_shutdown_signal) =
crate::forced_shutdown::ForcedShutdownTrigger::pair();
processor
.start_with_completion_metrics(
runtime_ctrl_tx,
completion_tx,
metrics_reporter,
crate::Interests::NODE_LOCAL_DURATION,
None,
true,
true,
None,
Some(start_metric_set),
None,
Some(duration_metric_set.into()),
Some(outgoing_metric_set),
None,
None,
None,
true,
true,
crate::terminal_state::TerminalMetricsDeadline::default(),
forced_shutdown_signal,
crate::testing::test_pipeline_runtime_services(),
)
.await
});
input_tx
.send(FlowMetricTestPData::default())
.expect("test input should enqueue");
let _ = output_rx
.recv()
.await
.expect("processor should forward the test message");
control_sender
.send(NodeControlMsg::CollectTelemetry {
metrics_reporter: collect_metrics_reporter,
})
.await
.expect("collect telemetry should enqueue");
let snapshot =
tokio::time::timeout(Duration::from_secs(1), metrics_rx.recv_async())
.await
.expect("flow_metric metric should be reported")
.expect("metrics channel should remain open");
processor_task.abort();
let _ = processor_task.await;
let [MetricValue::U64(input_items)] = snapshot.get_metrics() else {
panic!("expected one flow input-item metric");
};
assert_eq!(*input_items, 1);
let snapshot =
tokio::time::timeout(Duration::from_secs(1), metrics_rx.recv_async())
.await
.expect("flow_metric stop metric should be reported")
.expect("metrics channel should remain open");
let [MetricValue::Distribution(compute_duration)] = snapshot.get_metrics() else {
panic!("expected flow duration histogram");
};
let (count, sum, _, _) = compute_duration.summary();
assert!(
count >= 1,
"flow_metric compute duration should have at least one observation"
);
assert!(
sum > 0.0,
"flow_metric compute duration sum should be non-zero"
);
let snapshot =
tokio::time::timeout(Duration::from_secs(1), metrics_rx.recv_async())
.await
.expect("flow output-item metric should be reported")
.expect("metrics channel should remain open");
let [MetricValue::U64(output_items)] = snapshot.get_metrics() else {
panic!("expected flow output-item metric");
};
assert_eq!(*output_items, 1);
})
.await;
}
struct DelayedProcessor {
started: Option<tokio::sync::oneshot::Sender<()>>,
fail_before_final_metrics: bool,
}
#[async_trait(?Send)]
impl local::Processor<FlowMetricTestPData> for DelayedProcessor {
async fn process(
&mut self,
msg: Message<FlowMetricTestPData>,
effect_handler: &mut local::EffectHandler<FlowMetricTestPData>,
) -> Result<(), Error> {
if let Message::PData(pdata) = msg {
self.started.take().expect("first batch").send(()).unwrap();
if self.fail_before_final_metrics {
return Err(Error::ProcessorError {
processor: test_node("test_processor"),
kind: ProcessorErrorKind::Other,
error: "error before shutdown deadline".to_owned(),
source_detail: String::new(),
});
}
tokio::time::sleep(Duration::from_secs(10)).await;
effect_handler.send_message(pdata).await?;
} else if self.fail_before_final_metrics
&& matches!(
msg,
Message::Control(NodeControlMsg::CollectTelemetry { .. })
)
{
std::future::pending::<()>().await;
}
Ok(())
}
}
#[async_trait]
impl shared::Processor<FlowMetricTestPData> for DelayedProcessor {
async fn process(
&mut self,
msg: Message<FlowMetricTestPData>,
effect_handler: &mut shared::EffectHandler<FlowMetricTestPData>,
) -> Result<(), Error> {
if let Message::PData(pdata) = msg {
self.started.take().expect("first batch").send(()).unwrap();
if self.fail_before_final_metrics {
return Err(Error::ProcessorError {
processor: test_node("test_processor"),
kind: ProcessorErrorKind::Other,
error: "error before shutdown deadline".to_owned(),
source_detail: String::new(),
});
}
tokio::time::sleep(Duration::from_secs(10)).await;
effect_handler.send_message(pdata).await?;
} else if self.fail_before_final_metrics
&& matches!(
msg,
Message::Control(NodeControlMsg::CollectTelemetry { .. })
)
{
std::future::pending::<()>().await;
}
Ok(())
}
}
async fn run_delayed_shutdown_scenario(
shared: bool,
deadline_secs: u64,
fail_before_final_metrics: bool,
) {
let config = ProcessorConfig::new("test_processor");
let node_id = test_node(config.name.clone());
let user_config = Arc::new(NodeUserConfig::new_processor_config("test_processor"));
let (started_tx, started_rx) = tokio::sync::oneshot::channel();
let processor = DelayedProcessor {
started: Some(started_tx),
fail_before_final_metrics,
};
let mut wrapper = if shared {
ProcessorWrapper::shared(processor, node_id.clone(), user_config, &config)
} else {
ProcessorWrapper::local(processor, node_id.clone(), user_config, &config)
};
let (input_tx, input_rx) = tokio::sync::mpsc::channel(4);
let (output_tx, mut output_rx) = tokio::sync::mpsc::channel(4);
wrapper
.set_pdata_receiver(
node_id.clone(),
Receiver::Shared(SharedReceiver::mpsc(input_rx)),
)
.unwrap();
wrapper
.set_pdata_sender(
node_id,
"default".into(),
Sender::Shared(SharedSender::mpsc(output_tx)),
)
.unwrap();
input_tx.send(FlowMetricTestPData::default()).await.unwrap();
drop(input_tx);
let (_metrics_rx, reporter) =
otel_arrow_dfe_telemetry::reporter::MetricsReporter::create_new_and_receiver(16);
let (runtime_tx, _runtime_rx) = runtime_ctrl_msg_channel(4);
let (completion_tx, _completion_rx) = pipeline_completion_msg_channel(4);
let _control_keepalive = wrapper.control_sender();
let deadline = crate::terminal_state::TerminalMetricsDeadline::default();
let (forced_shutdown_trigger, forced_shutdown_signal) =
crate::forced_shutdown::ForcedShutdownTrigger::pair();
let start = tokio::time::Instant::now();
let run = wrapper.start_with_completion_metrics(
runtime_tx,
completion_tx,
reporter,
crate::Interests::empty(),
None,
false,
false,
None,
None,
None,
None,
None,
None,
None,
None,
false,
false,
deadline.clone(),
forced_shutdown_signal,
crate::testing::test_pipeline_runtime_services(),
);
let shutdown = tokio::spawn(async move {
started_rx.await.expect("handler started before shutdown");
tokio::time::sleep_until(start + Duration::from_secs(deadline_secs)).await;
forced_shutdown_trigger.trigger();
});
let result = run.await;
shutdown.abort();
if fail_before_final_metrics {
let Error::ProcessorError { error, .. } = result.expect_err("original error survives")
else {
panic!("expected the original processor error");
};
assert_eq!(error, "error before shutdown deadline");
} else {
result.expect("shutdown must not become a processing error");
}
assert_eq!(start.elapsed(), Duration::from_secs(deadline_secs.min(10)));
if deadline_secs < 10 {
assert!(
output_rx.try_recv().is_err(),
"expired batch must not be forwarded"
);
} else {
let _ = output_rx
.try_recv()
.expect("graceful drain forwards the batch");
}
}
#[tokio::test(start_paused = true)]
async fn local_processor_deadline_cancels_inflight_handler() {
run_delayed_shutdown_scenario(false, 2, false).await;
}
#[tokio::test(start_paused = true)]
async fn shared_processor_deadline_cancels_inflight_handler() {
run_delayed_shutdown_scenario(true, 2, false).await;
}
#[tokio::test(start_paused = true)]
async fn local_processor_deadline_allows_graceful_completion() {
run_delayed_shutdown_scenario(false, 20, false).await;
}
#[tokio::test(start_paused = true)]
async fn shared_processor_deadline_allows_graceful_completion() {
run_delayed_shutdown_scenario(true, 20, false).await;
}
#[tokio::test(start_paused = true)]
async fn local_processor_deadline_preserves_processing_error() {
run_delayed_shutdown_scenario(false, 2, true).await;
}
#[tokio::test(start_paused = true)]
async fn shared_processor_deadline_preserves_processing_error() {
run_delayed_shutdown_scenario(true, 2, true).await;
}
struct ErrorOnPDataProcessor {
node_id: String,
snapshot_metric: MeasurementMetricSet<FlowInputItemsMetrics>,
}
#[async_trait(?Send)]
impl local::Processor<FlowMetricTestPData> for ErrorOnPDataProcessor {
async fn process(
&mut self,
msg: Message<FlowMetricTestPData>,
_effect_handler: &mut local::EffectHandler<FlowMetricTestPData>,
) -> Result<(), Error> {
match msg {
Message::Control(NodeControlMsg::CollectTelemetry {
mut metrics_reporter,
..
}) => {
self.snapshot_metric
.with(SignalAttributes {
signal: SignalType::Logs,
})
.items
.add(7);
metrics_reporter
.report_measurement(&mut self.snapshot_metric)
.expect("test: failed to emit processor-local snapshot");
Ok(())
}
Message::PData(_) => Err(Error::ProcessorError {
processor: test_node(self.node_id.clone()),
kind: ProcessorErrorKind::Other,
error: "deliberate test error".to_owned(),
source_detail: String::new(),
}),
_ => Ok(()),
}
}
}
#[async_trait]
impl shared::Processor<FlowMetricTestPData> for ErrorOnPDataProcessor {
async fn process(
&mut self,
msg: Message<FlowMetricTestPData>,
_effect_handler: &mut shared::EffectHandler<FlowMetricTestPData>,
) -> Result<(), Error> {
match msg {
Message::Control(NodeControlMsg::CollectTelemetry {
mut metrics_reporter,
..
}) => {
self.snapshot_metric
.with(SignalAttributes {
signal: SignalType::Logs,
})
.items
.add(7);
metrics_reporter
.report_measurement(&mut self.snapshot_metric)
.expect("test: failed to emit processor-local snapshot");
Ok(())
}
Message::PData(_) => Err(Error::ProcessorError {
processor: test_node(self.node_id.clone()),
kind: ProcessorErrorKind::Other,
error: "deliberate test error".to_owned(),
source_detail: String::new(),
}),
_ => Ok(()),
}
}
}
async fn run_error_on_pdata_scenario(
processor: ProcessorWrapper<FlowMetricTestPData>,
input_metric: MeasurementMetricSet<FlowInputItemsMetrics>,
) -> (
Error,
flume::Receiver<otel_arrow_dfe_telemetry::metrics::MetricSetSnapshot>,
) {
let config = ProcessorConfig::new("test_processor");
let node_id = test_node(config.name.clone());
let mut p = processor;
let is_shared = p.is_shared();
if !is_shared {
let (tx, rx) = otel_arrow_dfe_channel::mpsc::Channel::new(4);
let (out_tx, _out_rx) = otel_arrow_dfe_channel::mpsc::Channel::new(4);
p.set_pdata_receiver(node_id.clone(), Receiver::Local(LocalReceiver::mpsc(rx)))
.expect("set pdata receiver");
p.set_pdata_sender(
node_id,
"default".into(),
Sender::Local(LocalSender::mpsc(out_tx)),
)
.expect("set pdata sender");
tx.send(FlowMetricTestPData::default())
.expect("pdata should enqueue");
} else {
let (tx, rx) = tokio::sync::mpsc::channel(4);
let (out_tx, _out_rx) = tokio::sync::mpsc::channel(4);
p.set_pdata_receiver(node_id.clone(), Receiver::Shared(SharedReceiver::mpsc(rx)))
.expect("set pdata receiver");
p.set_pdata_sender(
node_id,
"default".into(),
Sender::Shared(SharedSender::mpsc(out_tx)),
)
.expect("set pdata sender");
tx.send(FlowMetricTestPData::default())
.await
.expect("pdata should enqueue");
}
let (metrics_rx, metrics_reporter) =
otel_arrow_dfe_telemetry::reporter::MetricsReporter::create_new_and_receiver(16);
let (runtime_ctrl_tx, _runtime_ctrl_rx) = runtime_ctrl_msg_channel(4);
let (completion_tx, _completion_rx) = pipeline_completion_msg_channel(4);
let _ctrl_keepalive = p.control_sender();
let (_, forced_shutdown_signal) = crate::forced_shutdown::ForcedShutdownTrigger::pair();
let result = p
.start_with_completion_metrics(
runtime_ctrl_tx,
completion_tx,
metrics_reporter,
crate::Interests::empty(),
None,
true, false,
None,
Some(input_metric),
None,
None,
None,
None,
None,
None,
true, false,
crate::terminal_state::TerminalMetricsDeadline::default(),
forced_shutdown_signal,
crate::testing::test_pipeline_runtime_services(),
)
.await;
drop(_ctrl_keepalive);
let err = result.expect_err("run loop must return the processing error");
(err, metrics_rx)
}
#[tokio::test]
async fn local_processor_error_flushes_flow_and_local_metrics() {
let (pipeline_ctx, _) = crate::testing::test_pipeline_ctx();
let entity_key = pipeline_ctx
.metrics_registry()
.register_entity(FlowAttributeSet::default());
let registrar = pipeline_ctx.metric_set_registrar_for_entity(entity_key);
let input_metric = FlowInputItemsMetrics::register(®istrar);
let local_metric = FlowInputItemsMetrics::register(®istrar);
let config = ProcessorConfig::new("test_processor");
let user_config = Arc::new(NodeUserConfig::new_processor_config("test_processor"));
let proc = ErrorOnPDataProcessor {
node_id: config.name.to_string(),
snapshot_metric: local_metric,
};
let wrapper =
ProcessorWrapper::local(proc, test_node(config.name.clone()), user_config, &config);
let (err, metrics_rx) = run_error_on_pdata_scenario(wrapper, input_metric).await;
let Error::ProcessorError { error, .. } = err else {
panic!("expected ProcessorError, got {err:?}");
};
assert_eq!(
error, "deliberate test error",
"original error must be returned"
);
let flow_snapshot = metrics_rx
.try_recv()
.expect("flow input-items snapshot must be delivered after error");
let [MetricValue::U64(input)] = flow_snapshot.get_metrics() else {
panic!(
"expected U64 input-items metric, got {:?}",
flow_snapshot.get_metrics()
);
};
assert_eq!(*input, 1, "one PData message entered before the error");
let local_snapshot = metrics_rx
.try_recv()
.expect("processor-local snapshot from CollectTelemetry must be delivered");
let [MetricValue::U64(local_val)] = local_snapshot.get_metrics() else {
panic!(
"expected U64 local metric, got {:?}",
local_snapshot.get_metrics()
);
};
assert_eq!(
*local_val, 7,
"processor emitted one local snapshot on CollectTelemetry"
);
assert!(
metrics_rx.try_recv().is_err(),
"no extra snapshots expected"
);
}
#[tokio::test]
async fn shared_processor_error_flushes_flow_and_local_metrics() {
let (pipeline_ctx, _) = crate::testing::test_pipeline_ctx();
let entity_key = pipeline_ctx
.metrics_registry()
.register_entity(FlowAttributeSet::default());
let registrar = pipeline_ctx.metric_set_registrar_for_entity(entity_key);
let input_metric = FlowInputItemsMetrics::register(®istrar);
let local_metric = FlowInputItemsMetrics::register(®istrar);
let config = ProcessorConfig::new("test_processor");
let user_config = Arc::new(NodeUserConfig::new_processor_config("test_processor"));
let proc = ErrorOnPDataProcessor {
node_id: config.name.to_string(),
snapshot_metric: local_metric,
};
let wrapper =
ProcessorWrapper::shared(proc, test_node(config.name.clone()), user_config, &config);
let (err, metrics_rx) = run_error_on_pdata_scenario(wrapper, input_metric).await;
let Error::ProcessorError { error, .. } = err else {
panic!("expected ProcessorError, got {err:?}");
};
assert_eq!(
error, "deliberate test error",
"original error must be returned"
);
let flow_snapshot = metrics_rx
.try_recv()
.expect("flow input-items snapshot must be delivered after error");
let [MetricValue::U64(input)] = flow_snapshot.get_metrics() else {
panic!(
"expected U64 input-items metric, got {:?}",
flow_snapshot.get_metrics()
);
};
assert_eq!(*input, 1, "one PData message entered before the error");
let local_snapshot = metrics_rx
.try_recv()
.expect("processor-local snapshot from CollectTelemetry must be delivered");
let [MetricValue::U64(local_val)] = local_snapshot.get_metrics() else {
panic!(
"expected U64 local metric, got {:?}",
local_snapshot.get_metrics()
);
};
assert_eq!(
*local_val, 7,
"processor emitted one local snapshot on CollectTelemetry"
);
assert!(
metrics_rx.try_recv().is_err(),
"no extra snapshots expected"
);
}
#[tokio::test]
async fn local_processor_error_takes_precedence_over_collect_telemetry_error() {
struct ErrorOnPDataAndCollectProcessor;
#[async_trait(?Send)]
impl local::Processor<FlowMetricTestPData> for ErrorOnPDataAndCollectProcessor {
async fn process(
&mut self,
msg: Message<FlowMetricTestPData>,
_effect_handler: &mut local::EffectHandler<FlowMetricTestPData>,
) -> Result<(), Error> {
match msg {
Message::Control(NodeControlMsg::CollectTelemetry { .. }) => {
Err(Error::ProcessorError {
processor: test_node("err_collect_proc"),
kind: ProcessorErrorKind::Other,
error: "secondary collect error".to_owned(),
source_detail: String::new(),
})
}
Message::PData(_) => Err(Error::ProcessorError {
processor: test_node("err_collect_proc"),
kind: ProcessorErrorKind::Other,
error: "original processing error".to_owned(),
source_detail: String::new(),
}),
_ => Ok(()),
}
}
}
let (pipeline_ctx, _) = crate::testing::test_pipeline_ctx();
let entity_key = pipeline_ctx
.metrics_registry()
.register_entity(FlowAttributeSet::default());
let input_metric = FlowInputItemsMetrics::register(
&pipeline_ctx.metric_set_registrar_for_entity(entity_key),
);
let config = ProcessorConfig::new("test_processor");
let node_id = test_node(config.name.clone());
let user_config = Arc::new(NodeUserConfig::new_processor_config("test_processor"));
let (input_tx, input_rx) = otel_arrow_dfe_channel::mpsc::Channel::new(4);
let (out_tx, _out_rx) = otel_arrow_dfe_channel::mpsc::Channel::new(4);
let mut p = ProcessorWrapper::local(
ErrorOnPDataAndCollectProcessor,
node_id.clone(),
user_config,
&config,
);
p.set_pdata_receiver(
node_id.clone(),
Receiver::Local(LocalReceiver::mpsc(input_rx)),
)
.expect("set pdata receiver");
p.set_pdata_sender(
node_id,
"default".into(),
Sender::Local(LocalSender::mpsc(out_tx)),
)
.expect("set pdata sender");
let (metrics_rx, metrics_reporter) =
otel_arrow_dfe_telemetry::reporter::MetricsReporter::create_new_and_receiver(4);
let (runtime_ctrl_tx, _runtime_ctrl_rx) = runtime_ctrl_msg_channel(4);
let (completion_tx, _completion_rx) = pipeline_completion_msg_channel(4);
input_tx
.send(FlowMetricTestPData::default())
.expect("pdata should enqueue");
drop(input_tx);
let _ctrl_keepalive = p.control_sender();
let (_, forced_shutdown_signal) = crate::forced_shutdown::ForcedShutdownTrigger::pair();
let result = p
.start_with_completion_metrics(
runtime_ctrl_tx,
completion_tx,
metrics_reporter,
crate::Interests::empty(),
None,
true,
false,
None, Some(input_metric), None, None, None, None, None, None, true,
false,
crate::terminal_state::TerminalMetricsDeadline::default(),
forced_shutdown_signal,
crate::testing::test_pipeline_runtime_services(),
)
.await;
drop(_ctrl_keepalive);
let flow_snapshot = metrics_rx
.try_recv()
.expect("flow input-items snapshot must still be delivered");
let [MetricValue::U64(input)] = flow_snapshot.get_metrics() else {
panic!(
"expected U64 input-items metric, got {:?}",
flow_snapshot.get_metrics()
);
};
assert_eq!(*input, 1, "one PData message entered before the error");
assert!(
metrics_rx.try_recv().is_err(),
"no processor-local snapshot expected when CollectTelemetry itself errors"
);
let err = result.expect_err("must return an error");
let Error::ProcessorError { error, .. } = err else {
panic!("expected ProcessorError, got {err:?}");
};
assert_eq!(
error, "original processing error",
"original error must take precedence over CollectTelemetry error"
);
}
}