use crate::Interests;
#[cfg(any(test, feature = "test-utils"))]
use crate::control::WakeupRevision;
use crate::control::{AckMsg, NackMsg, RuntimeCtrlMsgSender, WakeupSlot};
use crate::effect_handler::{
EffectHandlerCore, SourceTagging, TelemetryTimerCancelHandle, TimerCancelHandle,
};
use crate::error::{Error, TypedError};
use crate::flow_metrics::{
DecisionFlowMetrics, EndFlowMetrics, FLOW_SIGNALS, FlowDroppedItemsMetrics,
FlowDurationMetricSet, FlowInputItemsMetrics, FlowInputMessageMetrics, FlowInputSizeMetrics,
FlowOutputItemsMetrics, FlowOutputMessageMetrics, FlowOutputSizeMetrics, InputFlowMetrics,
LocalFlowMetricState, flow_signal_index, nanos_u64,
};
use crate::message::{Message, Sender};
use crate::node::NodeId;
use crate::output_router::OutputRouter;
use crate::process_duration::ComputeDuration;
use crate::processor::ProcessorRuntimeRequirements;
use crate::runtime_services::{CodecEffectHandler, PipelineRuntimeServices};
use crate::{WakeupError, WakeupSetOutcome};
use async_trait::async_trait;
use otel_arrow_dfe_config::{PortName, SignalType};
use otel_arrow_dfe_pdata_codec::CodecService;
use otel_arrow_dfe_telemetry::common_attributes::SignalAttributes;
use otel_arrow_dfe_telemetry::error::Error as TelemetryError;
use otel_arrow_dfe_telemetry::metrics::{MeasurementMetricSet, MetricSet, MetricSetHandler};
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use std::cell::{Cell, RefCell};
use std::collections::HashMap;
use std::time::{Duration, Instant};
#[async_trait(?Send)]
pub trait Processor<PData> {
async fn process(
&mut self,
msg: Message<PData>,
effect_handler: &mut EffectHandler<PData>,
) -> Result<(), Error>;
fn accept_pdata(&self) -> bool {
true
}
fn runtime_requirements(&self) -> ProcessorRuntimeRequirements {
ProcessorRuntimeRequirements::none()
}
}
#[derive(Clone)]
pub struct EffectHandler<PData> {
pub(crate) core: EffectHandlerCore<PData>,
pub router: OutputRouter<Sender<PData>>,
pub(crate) flow: LocalFlowMetricState,
}
impl<PData> EffectHandler<PData> {
#[must_use]
pub fn new(
node_id: NodeId,
msg_senders: HashMap<PortName, Sender<PData>>,
default_port: Option<PortName>,
metrics_reporter: MetricsReporter,
runtime_services: PipelineRuntimeServices,
) -> Self {
let core = EffectHandlerCore::new(node_id.clone(), metrics_reporter, runtime_services);
let router = OutputRouter::new(node_id, msg_senders, default_port);
EffectHandler {
core,
router,
flow: LocalFlowMetricState::default(),
}
}
#[must_use]
pub fn processor_id(&self) -> NodeId {
self.core.node_id()
}
pub fn set_source_tagging(&mut self, value: SourceTagging) {
self.core.set_source_tagging(value);
}
#[must_use]
pub const fn source_tagging(&self) -> SourceTagging {
self.core.source_tagging()
}
#[must_use]
pub fn connected_ports(&self) -> Vec<PortName> {
self.router.connected_ports()
}
#[must_use]
pub fn default_port(&self) -> Option<PortName> {
self.router.default_port()
}
#[must_use]
pub fn node_interests(&self) -> Interests {
self.core.node_interests()
}
pub(crate) fn begin_process_timing(&self) {
if self.flow.needs_timing {
self.flow.last_send_marker.set(Some(Instant::now()));
}
}
#[must_use]
pub fn take_elapsed_since_send_marker_ns(&self) -> u64 {
let Some(prev) = self.flow.last_send_marker.get() else {
return 0;
};
let now = Instant::now();
self.flow.last_send_marker.set(Some(now));
nanos_u64(now.duration_since(prev).as_nanos())
}
#[must_use]
pub fn flow_metrics_active(&self) -> bool {
self.flow.active
}
pub(crate) fn set_flow_roles(
&mut self,
is_start: bool,
is_end: bool,
input_message_metric: Option<MeasurementMetricSet<FlowInputMessageMetrics>>,
input_items_metric: Option<MeasurementMetricSet<FlowInputItemsMetrics>>,
input_size_metric: Option<MeasurementMetricSet<FlowInputSizeMetrics>>,
duration_metric: Option<FlowDurationMetricSet>,
output_items_metric: Option<MeasurementMetricSet<FlowOutputItemsMetrics>>,
output_message_metric: Option<MeasurementMetricSet<FlowOutputMessageMetrics>>,
output_size_metric: Option<MeasurementMetricSet<FlowOutputSizeMetrics>>,
dropped_items_metric: Option<MeasurementMetricSet<FlowDroppedItemsMetrics>>,
flow_metrics_active: bool,
flow_needs_timing: bool,
) {
self.flow.is_start = is_start;
self.flow.is_end = is_end;
self.flow.is_decision = dropped_items_metric.is_some();
self.flow.interests.set(
crate::flow_metrics::FlowMetricInterests::INPUT_MESSAGES,
input_message_metric.is_some(),
);
self.flow.interests.set(
crate::flow_metrics::FlowMetricInterests::INPUT_ITEMS,
input_items_metric.is_some(),
);
self.flow.interests.set(
crate::flow_metrics::FlowMetricInterests::INPUT_SIZE,
input_size_metric.is_some(),
);
self.flow.interests.set(
crate::flow_metrics::FlowMetricInterests::COMPUTE_DURATION,
duration_metric.is_some(),
);
self.flow.interests.set(
crate::flow_metrics::FlowMetricInterests::OUTPUT_ITEMS,
output_items_metric.is_some(),
);
self.flow.interests.set(
crate::flow_metrics::FlowMetricInterests::OUTPUT_MESSAGES,
output_message_metric.is_some(),
);
self.flow.interests.set(
crate::flow_metrics::FlowMetricInterests::OUTPUT_SIZE,
output_size_metric.is_some(),
);
self.flow.interests.set(
crate::flow_metrics::FlowMetricInterests::DROPPED_ITEMS,
dropped_items_metric.is_some(),
);
self.flow.active = flow_metrics_active;
self.flow.needs_timing = flow_needs_timing;
self.flow.input = InputFlowMetrics {
input_messages: input_message_metric.map(|metrics| (metrics, Cell::new([0; 3]))),
input_items: input_items_metric.map(|metrics| (metrics, Cell::new([0; 3]))),
input_size: input_size_metric.map(|metrics| (metrics, Cell::new([0; 3]))),
};
self.flow.end = EndFlowMetrics {
duration: duration_metric
.map(FlowDurationMetricSet::into_measurement)
.map(RefCell::new),
output_messages: output_message_metric.map(|metrics| (metrics, Cell::new([0; 3]))),
output_items: output_items_metric.map(|metrics| (metrics, Cell::new([0; 3]))),
output_size: output_size_metric.map(|metrics| (metrics, Cell::new([0; 3]))),
};
self.flow.decision = DecisionFlowMetrics {
dropped_items: dropped_items_metric.map(|metrics| (metrics, Cell::new([0; 3]))),
};
}
#[must_use]
pub fn is_flow_start(&self) -> bool {
self.flow.is_start
}
#[must_use]
pub fn is_flow_end(&self) -> bool {
self.flow.is_end
}
#[must_use]
pub fn is_flow_decision(&self) -> bool {
self.flow.is_decision
}
pub fn record_flow_duration(&self, signal: SignalType, total: u64) {
let Some(measurement) = self.flow.end.duration.as_ref() else {
return;
};
measurement
.borrow_mut()
.record(signal, total as f64 / 1_000_000_000.0);
}
pub fn record_flow_input_items(&self, signal: SignalType, items: u64) {
let Some((_, acc_cell)) = self.flow.input.input_items.as_ref() else {
return;
};
let mut acc = acc_cell.get();
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(items);
acc_cell.set(acc);
}
pub fn record_flow_input_message(&self, signal: SignalType) {
let Some((_, acc_cell)) = self.flow.input.input_messages.as_ref() else {
return;
};
let mut acc = acc_cell.get();
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(1);
acc_cell.set(acc);
}
pub fn record_flow_input_size(&self, signal: SignalType, size: u64) {
let Some((_, acc_cell)) = self.flow.input.input_size.as_ref() else {
return;
};
let mut acc = acc_cell.get();
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(size);
acc_cell.set(acc);
}
pub fn record_flow_output_items(&self, signal: SignalType, items: u64) {
let Some((_, acc_cell)) = self.flow.end.output_items.as_ref() else {
return;
};
let mut acc = acc_cell.get();
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(items);
acc_cell.set(acc);
}
pub fn record_flow_output_message(&self, signal: SignalType) {
let Some((_, acc_cell)) = self.flow.end.output_messages.as_ref() else {
return;
};
let mut acc = acc_cell.get();
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(1);
acc_cell.set(acc);
}
pub fn record_flow_output_size(&self, signal: SignalType, size: u64) {
let Some((_, acc_cell)) = self.flow.end.output_size.as_ref() else {
return;
};
let mut acc = acc_cell.get();
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(size);
acc_cell.set(acc);
}
pub fn record_flow_dropped_items(&self, signal: SignalType, items: u64) {
let Some((_, acc_cell)) = self.flow.decision.dropped_items.as_ref() else {
return;
};
let mut acc = acc_cell.get();
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(items);
acc_cell.set(acc);
}
pub(crate) fn report_flow_metrics(&mut self) {
if let Some((metrics, acc_cell)) = self.flow.input.input_messages.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let count = drained[flow_signal_index(signal)];
if count != 0 {
metrics
.with(SignalAttributes { signal })
.messages
.add(count);
}
}
let _ = self.core.metrics_reporter.report_measurement(metrics);
}
if let Some((metrics, acc_cell)) = self.flow.input.input_items.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let count = drained[flow_signal_index(signal)];
if count != 0 {
metrics.with(SignalAttributes { signal }).items.add(count);
}
}
let _ = self.core.metrics_reporter.report_measurement(metrics);
}
if let Some((metrics, acc_cell)) = self.flow.input.input_size.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let size = drained[flow_signal_index(signal)];
if size != 0 {
metrics.with(SignalAttributes { signal }).size.add(size);
}
}
let _ = self.core.metrics_reporter.report_measurement(metrics);
}
if let Some(measurement) = self.flow.end.duration.as_mut() {
measurement
.borrow_mut()
.report(&mut self.core.metrics_reporter);
}
if let Some((metrics, acc_cell)) = self.flow.end.output_messages.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let count = drained[flow_signal_index(signal)];
if count != 0 {
metrics
.with(SignalAttributes { signal })
.messages
.add(count);
}
}
let _ = self.core.metrics_reporter.report_measurement(metrics);
}
if let Some((metrics, acc_cell)) = self.flow.end.output_items.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let count = drained[flow_signal_index(signal)];
if count != 0 {
metrics.with(SignalAttributes { signal }).items.add(count);
}
}
let _ = self.core.metrics_reporter.report_measurement(metrics);
}
if let Some((metrics, acc_cell)) = self.flow.end.output_size.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let size = drained[flow_signal_index(signal)];
if size != 0 {
metrics.with(SignalAttributes { signal }).size.add(size);
}
}
let _ = self.core.metrics_reporter.report_measurement(metrics);
}
if let Some((metrics, acc_cell)) = self.flow.decision.dropped_items.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let count = drained[flow_signal_index(signal)];
if count != 0 {
metrics.with(SignalAttributes { signal }).items.add(count);
}
}
let _ = self.core.metrics_reporter.report_measurement(metrics);
}
}
pub(crate) async fn report_flow_metrics_reliably(
&mut self,
deadline: Instant,
) -> Result<(), TelemetryError> {
let reporter = self.core.metrics_reporter.clone();
if let Some((metrics, acc_cell)) = self.flow.input.input_messages.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let count = drained[flow_signal_index(signal)];
if count != 0 {
metrics
.with(SignalAttributes { signal })
.messages
.add(count);
}
}
let _ = reporter
.report_measurement_reliably_until(metrics, deadline)
.await?;
}
if let Some((metrics, acc_cell)) = self.flow.input.input_items.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let count = drained[flow_signal_index(signal)];
if count != 0 {
metrics.with(SignalAttributes { signal }).items.add(count);
}
}
let _ = reporter
.report_measurement_reliably_until(metrics, deadline)
.await?;
}
if let Some((metrics, acc_cell)) = self.flow.input.input_size.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let size = drained[flow_signal_index(signal)];
if size != 0 {
metrics.with(SignalAttributes { signal }).size.add(size);
}
}
let _ = reporter
.report_measurement_reliably_until(metrics, deadline)
.await?;
}
if let Some(measurement) = self.flow.end.duration.as_mut() {
let snapshots = measurement.borrow_mut().terminal_snapshots();
for snapshot in snapshots {
let _ = reporter
.report_snapshot_reliably_until(snapshot, deadline)
.await?;
}
}
if let Some((metrics, acc_cell)) = self.flow.end.output_messages.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let count = drained[flow_signal_index(signal)];
if count != 0 {
metrics
.with(SignalAttributes { signal })
.messages
.add(count);
}
}
let _ = reporter
.report_measurement_reliably_until(metrics, deadline)
.await?;
}
if let Some((metrics, acc_cell)) = self.flow.end.output_items.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let count = drained[flow_signal_index(signal)];
if count != 0 {
metrics.with(SignalAttributes { signal }).items.add(count);
}
}
let _ = reporter
.report_measurement_reliably_until(metrics, deadline)
.await?;
}
if let Some((metrics, acc_cell)) = self.flow.end.output_size.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let size = drained[flow_signal_index(signal)];
if size != 0 {
metrics.with(SignalAttributes { signal }).size.add(size);
}
}
let _ = reporter
.report_measurement_reliably_until(metrics, deadline)
.await?;
}
if let Some((metrics, acc_cell)) = self.flow.decision.dropped_items.as_mut() {
let drained = acc_cell.replace([0; 3]);
for signal in FLOW_SIGNALS {
let count = drained[flow_signal_index(signal)];
if count != 0 {
metrics.with(SignalAttributes { signal }).items.add(count);
}
}
let _ = reporter
.report_measurement_reliably_until(metrics, deadline)
.await?;
}
Ok(())
}
#[inline]
pub fn timed<T, E>(
&self,
cd: &ComputeDuration,
f: impl FnOnce() -> Result<T, E>,
) -> Result<T, E> {
cd.timed(self.core.node_interests(), f)
}
#[inline]
pub async fn send_message(&self, mut data: PData) -> Result<(), TypedError<PData>>
where
PData: crate::processor::FlowMetricHook,
{
data.before_processor_send(self);
self.router.send_default(data).await
}
#[inline]
pub fn try_send_message(&self, mut data: PData) -> Result<(), TypedError<PData>>
where
PData: crate::processor::FlowMetricHook,
{
data.before_processor_send(self);
self.router.try_send_default(data)
}
#[inline]
pub async fn send_message_to<P>(
&self,
port: P,
mut data: PData,
) -> Result<(), TypedError<PData>>
where
P: Into<PortName>,
PData: crate::processor::FlowMetricHook,
{
data.before_processor_send(self);
self.router.send_to(port, data).await
}
#[inline]
pub fn try_send_message_to<P>(&self, port: P, mut data: PData) -> Result<(), TypedError<PData>>
where
P: Into<PortName>,
PData: crate::processor::FlowMetricHook,
{
data.before_processor_send(self);
self.router.try_send_to(port, data)
}
pub async fn info(&self, message: &str) {
self.core.info(message).await;
}
pub async fn start_periodic_timer(
&self,
duration: Duration,
) -> Result<TimerCancelHandle<PData>, Error> {
self.core.start_periodic_timer(duration).await
}
pub async fn start_periodic_telemetry(
&self,
duration: Duration,
) -> Result<TelemetryTimerCancelHandle<PData>, Error> {
self.core.start_periodic_telemetry(duration).await
}
pub fn requeue_later(&self, when: Instant, data: Box<PData>) -> Result<(), PData> {
self.core.requeue_later(when, data)
}
pub fn set_wakeup(
&self,
slot: WakeupSlot,
when: Instant,
) -> Result<WakeupSetOutcome, WakeupError> {
self.core.set_wakeup(slot, when)
}
#[must_use]
pub fn cancel_wakeup(&self, slot: WakeupSlot) -> bool {
self.core.cancel_wakeup(slot)
}
#[cfg(any(test, feature = "test-utils"))]
#[must_use]
pub fn pop_wakeup(&self) -> Option<(WakeupSlot, Instant, WakeupRevision)> {
self.core.pop_wakeup()
}
#[allow(dead_code)] pub(crate) fn report_metrics<M: MetricSetHandler + 'static>(
&mut self,
metrics: &mut MetricSet<M>,
) -> Result<(), TelemetryError> {
self.core.report_metrics(metrics)
}
pub fn report_local_scheduler_metrics(
&self,
metrics_reporter: &mut MetricsReporter,
) -> Result<(), TelemetryError> {
self.core.report_local_scheduler_metrics(metrics_reporter)
}
pub fn set_runtime_ctrl_msg_sender(
&mut self,
runtime_ctrl_msg_sender: RuntimeCtrlMsgSender<PData>,
) {
self.core
.set_runtime_ctrl_msg_sender(runtime_ctrl_msg_sender);
}
pub fn set_pipeline_completion_msg_sender(
&mut self,
pipeline_completion_msg_sender: crate::control::PipelineCompletionMsgSender<PData>,
) {
self.core
.set_pipeline_completion_msg_sender(pipeline_completion_msg_sender);
}
}
impl<PData> CodecEffectHandler for EffectHandler<PData> {
fn codec_service(&self) -> &CodecService {
self.core.runtime_services.codecs()
}
}
impl<PData> crate::processor::FlowMetricEffectHandler for EffectHandler<PData> {
#[inline]
fn is_flow_start(&self) -> bool {
EffectHandler::is_flow_start(self)
}
#[inline]
fn is_flow_end(&self) -> bool {
EffectHandler::is_flow_end(self)
}
#[inline]
fn flow_metric_interests(&self) -> crate::flow_metrics::FlowMetricInterests {
self.flow.interests
}
#[inline]
fn take_elapsed_since_send_marker_ns(&self) -> u64 {
EffectHandler::take_elapsed_since_send_marker_ns(self)
}
#[inline]
fn record_flow_duration(&self, signal: SignalType, total: u64) {
EffectHandler::record_flow_duration(self, signal, total);
}
#[inline]
fn record_flow_input_items(&self, signal: SignalType, items: u64) {
EffectHandler::record_flow_input_items(self, signal, items);
}
#[inline]
fn record_flow_input_message(&self, signal: SignalType) {
EffectHandler::record_flow_input_message(self, signal);
}
#[inline]
fn record_flow_input_size(&self, signal: SignalType, size: u64) {
EffectHandler::record_flow_input_size(self, signal, size);
}
#[inline]
fn record_flow_output_items(&self, signal: SignalType, items: u64) {
EffectHandler::record_flow_output_items(self, signal, items);
}
#[inline]
fn record_flow_output_message(&self, signal: SignalType) {
EffectHandler::record_flow_output_message(self, signal);
}
#[inline]
fn record_flow_output_size(&self, signal: SignalType, size: u64) {
EffectHandler::record_flow_output_size(self, signal, size);
}
}
#[async_trait(?Send)]
impl<PData: crate::Unwindable> crate::_private::AckNackRouting<PData> for EffectHandler<PData> {
async fn route_ack(&self, ack: AckMsg<PData>) -> Result<(), Error> {
self.core.route_ack(ack).await
}
async fn route_nack(&self, nack: NackMsg<PData>) -> Result<(), Error> {
self.core.route_nack(nack).await
}
}
#[cfg(test)]
mod tests {
#![allow(missing_docs)]
use super::*;
use crate::_private::AckNackRouting;
use crate::completion_emission_metrics::make_completion_emission_metrics;
use crate::context::ControllerContext;
use crate::control::{
AckMsg, Frame, NackMsg, PipelineCompletionMsg, RouteData, WakeupSlot,
pipeline_completion_msg_channel,
};
use crate::entity_context::NodeTelemetryHandle;
use crate::flow_metrics::FlowAttributeSet;
use crate::local::message::LocalSender;
use crate::testing::{test_node, test_pipeline_ctx};
use crate::{Interests, Unwindable, WakeupError};
use otel_arrow_dfe_channel::error::SendError;
use otel_arrow_dfe_channel::mpsc;
use otel_arrow_dfe_config::{MetricLevel, node::NodeKind};
use otel_arrow_dfe_telemetry::registry::TelemetryRegistryHandle;
use std::borrow::Cow;
use std::collections::{HashMap, HashSet};
use tokio::time::{Duration, timeout};
fn channel<T>(capacity: usize) -> (mpsc::Sender<T>, mpsc::Receiver<T>) {
mpsc::Channel::new(capacity)
}
impl crate::processor::FlowMetricHook for u64 {}
#[derive(Clone, Debug)]
struct TestPData {
frames: Vec<Frame>,
}
impl TestPData {
fn with_ack_frame(node_id: usize) -> Self {
Self {
frames: vec![Frame {
node_id,
interests: Interests::ACKS,
route: RouteData::default(),
output_items: 0,
input_items: 0,
output_size: 0,
input_size: 0,
}],
}
}
fn with_nack_frame(node_id: usize) -> Self {
Self {
frames: vec![Frame {
node_id,
interests: Interests::NACKS,
route: RouteData::default(),
output_items: 0,
input_items: 0,
output_size: 0,
input_size: 0,
}],
}
}
}
impl Unwindable for TestPData {
fn has_frames(&self) -> bool {
!self.frames.is_empty()
}
fn pop_frame(&mut self) -> Option<Frame> {
self.frames.pop()
}
fn signal(&self) -> Option<SignalType> {
None
}
fn drop_payload(&mut self) {}
}
fn test_node_telemetry() -> (TelemetryRegistryHandle, NodeTelemetryHandle) {
let registry = TelemetryRegistryHandle::new();
let controller = ControllerContext::new(registry.clone());
let pipeline_ctx = controller
.pipeline_context_with("test_grp".into(), "test_pipeline".into(), 0, 1, 0)
.with_node_context(
"test_node".into(),
"urn:test:processor:example".into(),
NodeKind::Processor,
HashMap::new(),
);
let entity_key = pipeline_ctx.register_node_entity();
(
registry,
NodeTelemetryHandle::new(pipeline_ctx.metrics_registry(), entity_key),
)
}
#[tokio::test]
async fn effect_handler_send_message_to_named_port() {
let (a_tx, a_rx) = channel::<u64>(10);
let (b_tx, b_rx) = channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("a".into(), Sender::Local(LocalSender::mpsc(a_tx)));
let _ = senders.insert("b".into(), Sender::Local(LocalSender::mpsc(b_tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
eh.send_message_to("b", 42).await.unwrap();
assert!(
timeout(Duration::from_millis(50), a_rx.recv())
.await
.is_err()
);
assert_eq!(b_rx.recv().await.unwrap(), 42);
}
#[tokio::test]
async fn effect_handler_send_message_single_port_fallback() {
let (tx, rx) = channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("only".into(), Sender::Local(LocalSender::mpsc(tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
eh.send_message(7).await.unwrap();
assert_eq!(rx.recv().await.unwrap(), 7);
}
#[tokio::test]
async fn effect_handler_send_message_uses_default_port() {
let (a_tx, a_rx) = channel::<u64>(10);
let (b_tx, b_rx) = channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("a".into(), Sender::Local(LocalSender::mpsc(a_tx)));
let _ = senders.insert("b".into(), Sender::Local(LocalSender::mpsc(b_tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
Some("a".into()),
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
eh.send_message(11).await.unwrap();
assert_eq!(a_rx.recv().await.unwrap(), 11);
assert!(
timeout(Duration::from_millis(50), b_rx.recv())
.await
.is_err()
);
}
#[test]
fn effect_handler_set_wakeup_without_runtime_support_returns_unsupported() {
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::<u64>::new(
test_node("proc"),
HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
assert_eq!(
eh.set_wakeup(WakeupSlot(0), Instant::now()),
Err(WakeupError::Unsupported)
);
assert!(!eh.cancel_wakeup(WakeupSlot(0)));
}
#[tokio::test]
async fn effect_handler_send_message_ambiguous_without_default() {
let (a_tx, a_rx) = channel::<u64>(10);
let (b_tx, b_rx) = channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("a".into(), Sender::Local(LocalSender::mpsc(a_tx)));
let _ = senders.insert("b".into(), Sender::Local(LocalSender::mpsc(b_tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
let res = eh.send_message(5).await;
assert!(res.is_err());
assert!(
timeout(Duration::from_millis(50), a_rx.recv())
.await
.is_err()
);
assert!(
timeout(Duration::from_millis(50), b_rx.recv())
.await
.is_err()
);
}
#[tokio::test]
async fn effect_handler_connected_ports_lists_all() {
let (a_tx, _a_rx) = channel::<u64>(1);
let (b_tx, _b_rx) = channel::<u64>(1);
let mut senders = HashMap::new();
let _ = senders.insert("a".into(), Sender::Local(LocalSender::mpsc(a_tx)));
let _ = senders.insert("b".into(), Sender::Local(LocalSender::mpsc(b_tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
let ports: HashSet<_> = eh.connected_ports().into_iter().collect();
let expected: HashSet<_> = [Cow::from("a"), Cow::from("b")].into_iter().collect();
assert_eq!(ports, expected);
}
#[test]
fn effect_handler_try_send_message_success() {
let (tx, rx) = channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("out".into(), Sender::Local(LocalSender::mpsc(tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
Some("out".into()),
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
assert!(eh.try_send_message(42).is_ok());
assert_eq!(rx.try_recv().unwrap(), 42);
}
#[test]
fn effect_handler_try_send_message_inbox_full() {
let (tx, _rx) = channel::<u64>(1);
let mut senders = HashMap::new();
let _ = senders.insert("out".into(), Sender::Local(LocalSender::mpsc(tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
Some("out".into()),
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
assert!(eh.try_send_message(1).is_ok());
let result = eh.try_send_message(2);
assert!(matches!(
result,
Err(TypedError::ChannelSendError(SendError::Full(2)))
));
}
#[test]
fn effect_handler_try_send_message_no_default_sender() {
let (a_tx, _a_rx) = channel::<u64>(10);
let (b_tx, _b_rx) = channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("a".into(), Sender::Local(LocalSender::mpsc(a_tx)));
let _ = senders.insert("b".into(), Sender::Local(LocalSender::mpsc(b_tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
let result = eh.try_send_message(99);
assert!(matches!(result, Err(TypedError::Error(_))));
}
#[test]
fn effect_handler_try_send_message_to_success() {
let (a_tx, a_rx) = channel::<u64>(10);
let (b_tx, b_rx) = channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("a".into(), Sender::Local(LocalSender::mpsc(a_tx)));
let _ = senders.insert("b".into(), Sender::Local(LocalSender::mpsc(b_tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
assert!(eh.try_send_message_to("b", 42).is_ok());
assert_eq!(b_rx.try_recv().unwrap(), 42);
assert!(a_rx.try_recv().is_err());
}
#[test]
fn effect_handler_try_send_message_to_channel_full() {
let (tx, _rx) = channel::<u64>(1);
let mut senders = HashMap::new();
let _ = senders.insert("out".into(), Sender::Local(LocalSender::mpsc(tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
assert!(eh.try_send_message_to("out", 1).is_ok());
let result = eh.try_send_message_to("out", 2);
assert!(matches!(
result,
Err(TypedError::ChannelSendError(SendError::Full(2)))
));
}
#[test]
fn effect_handler_try_send_message_to_unknown_port() {
let (tx, _rx) = channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("out".into(), Sender::Local(LocalSender::mpsc(tx)));
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let eh = EffectHandler::new(
test_node("proc"),
senders,
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
let result = eh.try_send_message_to("unknown", 99);
assert!(matches!(result, Err(TypedError::Error(_))));
}
#[tokio::test]
async fn effect_handler_route_ack_records_completion_emission_metrics() {
let (_registry, telemetry_handle) = test_node_telemetry();
let completion_metrics =
make_completion_emission_metrics(&Some(telemetry_handle), MetricLevel::Normal)
.expect("completion emission metrics should be registered");
let (completion_tx, mut completion_rx) = pipeline_completion_msg_channel(1);
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let mut eh = EffectHandler::<TestPData>::new(
test_node("proc"),
HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
eh.set_pipeline_completion_msg_sender(completion_tx);
eh.core
.set_completion_emission_metrics(Some(completion_metrics.clone()));
eh.route_ack(AckMsg::new(TestPData::with_ack_frame(1)))
.await
.expect("route_ack should succeed");
assert!(matches!(
completion_rx.recv().await.expect("completion message"),
PipelineCompletionMsg::DeliverAck { .. }
));
let counts = completion_metrics
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.counts();
assert_eq!(counts, (1, 0));
}
#[tokio::test]
async fn effect_handler_route_nack_records_completion_emission_metrics() {
let (_registry, telemetry_handle) = test_node_telemetry();
let completion_metrics =
make_completion_emission_metrics(&Some(telemetry_handle), MetricLevel::Normal)
.expect("completion emission metrics should be registered");
let (completion_tx, mut completion_rx) = pipeline_completion_msg_channel(1);
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let mut eh = EffectHandler::<TestPData>::new(
test_node("proc"),
HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
eh.set_pipeline_completion_msg_sender(completion_tx);
eh.core
.set_completion_emission_metrics(Some(completion_metrics.clone()));
eh.route_nack(NackMsg::new("test nack", TestPData::with_nack_frame(1)))
.await
.expect("route_nack should succeed");
assert!(matches!(
completion_rx.recv().await.expect("completion message"),
PipelineCompletionMsg::DeliverNack { .. }
));
let counts = completion_metrics
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.counts();
assert_eq!(counts, (0, 1));
}
#[tokio::test]
async fn effect_handler_route_ack_without_frames_does_not_record_completion_emission_metrics() {
let (_registry, telemetry_handle) = test_node_telemetry();
let completion_metrics =
make_completion_emission_metrics(&Some(telemetry_handle), MetricLevel::Normal)
.expect("completion emission metrics should be registered");
let (completion_tx, mut completion_rx) = pipeline_completion_msg_channel::<TestPData>(1);
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let mut eh = EffectHandler::<TestPData>::new(
test_node("proc"),
HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
eh.set_pipeline_completion_msg_sender(completion_tx);
eh.core
.set_completion_emission_metrics(Some(completion_metrics.clone()));
eh.route_ack(AckMsg::new(TestPData { frames: Vec::new() }))
.await
.expect("route_ack without frames should be a no-op");
assert!(
timeout(Duration::from_millis(50), completion_rx.recv())
.await
.is_err()
);
let counts = completion_metrics
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.counts();
assert_eq!(counts, (0, 0));
}
#[test]
fn flow_accumulate_then_report() {
use crate::context::ControllerContext;
use crate::flow_metrics::{
FlowAttributeSet, FlowDurationNormalMetrics, FlowInputItemsMetrics,
FlowOutputItemsMetrics,
};
use otel_arrow_dfe_config::node::NodeKind;
use otel_arrow_dfe_telemetry::registry::TelemetryRegistryHandle;
let registry = TelemetryRegistryHandle::new();
let controller = ControllerContext::new(registry.clone());
let pipeline_ctx = controller
.pipeline_context_with("g".into(), "p".into(), 0, 1, 0)
.with_node_context(
"end_node".into(),
"urn:test:processor:example".into(),
NodeKind::Processor,
HashMap::new(),
);
let entity_key = pipeline_ctx
.metrics_registry()
.register_entity(FlowAttributeSet::default());
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 (_snapshot_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(64);
let mut eh = EffectHandler::<u64>::new(
test_node("end_node"),
HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
eh.set_flow_roles(
true,
true,
None,
Some(start_metric_set),
None,
Some(duration_metric_set.into()),
Some(outgoing_metric_set),
None,
None,
None,
true,
true,
);
eh.record_flow_input_items(SignalType::Logs, 10);
eh.record_flow_input_items(SignalType::Metrics, 20);
eh.record_flow_duration(SignalType::Logs, 1000);
eh.record_flow_duration(SignalType::Metrics, 2000);
eh.record_flow_output_items(SignalType::Logs, 7);
eh.record_flow_output_items(SignalType::Metrics, 8);
let start_acc = eh.flow.input.input_items.as_ref().unwrap().1.get();
assert_eq!(
start_acc,
[0, 20, 10],
"start items should accumulate by signal"
);
let duration = eh.flow.end.duration.as_ref().unwrap().borrow();
assert!(duration.is_empty(SignalType::Traces));
assert_eq!(duration.pending_summary(SignalType::Metrics).0, 1);
assert_eq!(duration.pending_summary(SignalType::Logs).0, 1);
drop(duration);
let output_acc = eh.flow.end.output_items.as_ref().unwrap().1.get();
assert_eq!(
output_acc,
[0, 8, 7],
"output items should accumulate by signal"
);
let ms_snap = eh
.flow
.end
.duration
.as_ref()
.unwrap()
.borrow()
.reported_summary(SignalType::Logs);
assert_eq!(ms_snap.0, 0, "MetricSet should be empty before report");
eh.report_flow_metrics();
let start_acc_after = eh.flow.input.input_items.as_ref().unwrap().1.get();
assert_eq!(
start_acc_after, [0; 3],
"start accumulator should be drained"
);
let acc_after = eh.flow.end.duration.as_ref().unwrap().borrow();
assert_eq!(
acc_after.pending_summary(SignalType::Metrics).0,
0,
"duration accumulator should be drained"
);
drop(acc_after);
let output_acc_after = eh.flow.end.output_items.as_ref().unwrap().1.get();
assert_eq!(
output_acc_after, [0; 3],
"stop item accumulator should be drained"
);
eh.record_flow_duration(SignalType::Logs, 500);
eh.report_flow_metrics();
let acc_final = eh.flow.end.duration.as_ref().unwrap().borrow();
assert_eq!(
acc_final.pending_summary(SignalType::Logs).0,
0,
"accumulator drained after second report"
);
}
#[test]
fn flow_metric_marker_accumulates_after_begin_process_timing() {
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let mut eh = EffectHandler::<u64>::new(
test_node("proc"),
HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
eh.set_flow_roles(
true, false, None, None, None, None, None, None, None, None, true, true,
);
assert!(eh.flow.active);
assert!(eh.flow.needs_timing);
eh.begin_process_timing();
let mut value = 0u64;
for i in 0..10_000 {
value = value.wrapping_add(std::hint::black_box(i));
}
let _ = std::hint::black_box(value);
let ns = eh.take_elapsed_since_send_marker_ns();
assert!(
ns > 0,
"take_elapsed_since_send_marker_ns should be non-zero after begin_process_timing, got {ns}"
);
}
#[test]
fn flow_message_and_size_accumulators_drain_to_metric_sets() {
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_messages = FlowInputMessageMetrics::register(®istrar);
let input_size = FlowInputSizeMetrics::register(®istrar);
let output_messages = FlowOutputMessageMetrics::register(®istrar);
let output_size = FlowOutputSizeMetrics::register(®istrar);
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(16);
let mut handler = EffectHandler::<u64>::new(
test_node("proc"),
HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
handler.set_flow_roles(
true,
true,
Some(input_messages),
None,
Some(input_size),
None,
None,
Some(output_messages),
Some(output_size),
None,
true,
false,
);
handler.record_flow_input_message(SignalType::Logs);
handler.record_flow_input_message(SignalType::Metrics);
handler.record_flow_input_size(SignalType::Logs, 10);
handler.record_flow_input_size(SignalType::Metrics, 20);
handler.record_flow_output_message(SignalType::Logs);
handler.record_flow_output_message(SignalType::Logs);
handler.record_flow_output_size(SignalType::Logs, 30);
handler.record_flow_output_size(SignalType::Metrics, 40);
assert_eq!(
handler.flow.input.input_messages.as_ref().unwrap().1.get(),
[0, 1, 1]
);
assert_eq!(
handler.flow.input.input_size.as_ref().unwrap().1.get(),
[0, 20, 10]
);
assert_eq!(
handler.flow.end.output_messages.as_ref().unwrap().1.get(),
[0, 0, 2]
);
assert_eq!(
handler.flow.end.output_size.as_ref().unwrap().1.get(),
[0, 40, 30]
);
handler.report_flow_metrics();
assert_eq!(
handler.flow.input.input_messages.as_ref().unwrap().1.get(),
[0; 3]
);
assert_eq!(
handler.flow.input.input_size.as_ref().unwrap().1.get(),
[0; 3]
);
assert_eq!(
handler.flow.end.output_messages.as_ref().unwrap().1.get(),
[0; 3]
);
assert_eq!(
handler.flow.end.output_size.as_ref().unwrap().1.get(),
[0; 3]
);
}
#[test]
fn flow_metric_marker_returns_zero_when_unarmed() {
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let mut eh = EffectHandler::<u64>::new(
test_node("proc"),
HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
eh.set_flow_roles(
true, false, None, None, None, None, None, None, None, None, true, true,
);
assert_eq!(eh.take_elapsed_since_send_marker_ns(), 0);
}
#[test]
fn flow_metric_marker_not_armed_when_timing_disabled() {
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let mut eh = EffectHandler::<u64>::new(
test_node("proc"),
HashMap::new(),
None,
metrics_reporter,
crate::testing::test_pipeline_runtime_services(),
);
eh.set_flow_roles(
true, false, None, None, None, None, None, None, None, None, true, false,
);
assert!(eh.flow.active);
assert!(!eh.flow.needs_timing);
eh.begin_process_timing();
let mut value = 0u64;
for i in 0..10_000 {
value = value.wrapping_add(std::hint::black_box(i));
}
let _ = std::hint::black_box(value);
assert_eq!(
eh.take_elapsed_since_send_marker_ns(),
0,
"send marker must stay unarmed when timing is disabled"
);
}
}