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,
SharedFlowMetricState, flow_signal_index, nanos_u64,
};
use crate::message::Message;
use crate::node::NodeId;
use crate::output_router::OutputRouter;
use crate::processor::ProcessorRuntimeRequirements;
use crate::runtime_services::{CodecEffectHandler, PipelineRuntimeServices};
use crate::shared::message::SharedSender;
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::collections::HashMap;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
#[async_trait]
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<SharedSender<PData>>,
pub(crate) flow: SharedFlowMetricState,
}
impl<PData> EffectHandler<PData> {
#[must_use]
pub fn new(
node_id: NodeId,
msg_senders: HashMap<PortName, SharedSender<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: SharedFlowMetricState::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 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, Arc::new(Mutex::new([0; 3])))),
input_items: input_items_metric.map(|metrics| (metrics, Arc::new(Mutex::new([0; 3])))),
input_size: input_size_metric.map(|metrics| (metrics, Arc::new(Mutex::new([0; 3])))),
};
self.flow.end = EndFlowMetrics {
duration: duration_metric
.map(FlowDurationMetricSet::into_measurement)
.map(|measurement| Arc::new(Mutex::new(measurement))),
output_messages: output_message_metric
.map(|metrics| (metrics, Arc::new(Mutex::new([0; 3])))),
output_items: output_items_metric
.map(|metrics| (metrics, Arc::new(Mutex::new([0; 3])))),
output_size: output_size_metric.map(|metrics| (metrics, Arc::new(Mutex::new([0; 3])))),
};
self.flow.decision = DecisionFlowMetrics {
dropped_items: dropped_items_metric
.map(|metrics| (metrics, Arc::new(Mutex::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
}
#[must_use]
pub fn flow_metrics_active(&self) -> bool {
self.flow.active
}
pub(crate) fn begin_process_timing(&self) {
if self.flow.needs_timing {
*self
.flow
.last_send_marker
.lock()
.expect("last_send_marker poisoned") = Some(Instant::now());
}
}
#[must_use]
pub fn take_elapsed_since_send_marker_ns(&self) -> u64 {
let mut guard = self
.flow
.last_send_marker
.lock()
.expect("last_send_marker poisoned");
let Some(prev) = *guard else {
return 0;
};
let now = Instant::now();
*guard = Some(now);
nanos_u64(now.duration_since(prev).as_nanos())
}
pub fn record_flow_duration(&self, signal: SignalType, total: u64) {
let Some(measurement) = self.flow.end.duration.as_ref() else {
return;
};
let mut measurement = measurement
.lock()
.expect("flow duration accumulator poisoned");
measurement.record(signal, total as f64 / 1_000_000_000.0);
}
pub fn record_flow_input_items(&self, signal: SignalType, items: u64) {
let Some((_, acc_mutex)) = self.flow.input.input_items.as_ref() else {
return;
};
let mut acc = acc_mutex
.lock()
.expect("flow input_items accumulator poisoned");
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(items);
}
pub fn record_flow_input_message(&self, signal: SignalType) {
let Some((_, acc_mutex)) = self.flow.input.input_messages.as_ref() else {
return;
};
let mut acc = acc_mutex
.lock()
.expect("flow input_messages accumulator poisoned");
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(1);
}
pub fn record_flow_input_size(&self, signal: SignalType, size: u64) {
let Some((_, acc_mutex)) = self.flow.input.input_size.as_ref() else {
return;
};
let mut acc = acc_mutex
.lock()
.expect("flow input_size accumulator poisoned");
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(size);
}
pub fn record_flow_output_items(&self, signal: SignalType, items: u64) {
let Some((_, acc_mutex)) = self.flow.end.output_items.as_ref() else {
return;
};
let mut acc = acc_mutex
.lock()
.expect("flow output_items accumulator poisoned");
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(items);
}
pub fn record_flow_output_message(&self, signal: SignalType) {
let Some((_, acc_mutex)) = self.flow.end.output_messages.as_ref() else {
return;
};
let mut acc = acc_mutex
.lock()
.expect("flow output_messages accumulator poisoned");
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(1);
}
pub fn record_flow_output_size(&self, signal: SignalType, size: u64) {
let Some((_, acc_mutex)) = self.flow.end.output_size.as_ref() else {
return;
};
let mut acc = acc_mutex
.lock()
.expect("flow output_size accumulator poisoned");
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(size);
}
pub fn record_flow_dropped_items(&self, signal: SignalType, items: u64) {
let Some((_, acc_mutex)) = self.flow.decision.dropped_items.as_ref() else {
return;
};
let mut acc = acc_mutex
.lock()
.expect("flow dropped_items accumulator poisoned");
let index = flow_signal_index(signal);
acc[index] = acc[index].saturating_add(items);
}
pub(crate) fn report_flow_metrics(&mut self) {
if let Some((metrics, acc_mutex)) = self.flow.input.input_messages.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow input_messages accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.input.input_items.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow input_items accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.input.input_size.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow input_size accumulator poisoned");
std::mem::take(&mut *guard)
};
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
.lock()
.expect("flow duration accumulator poisoned")
.report(&mut self.core.metrics_reporter);
}
if let Some((metrics, acc_mutex)) = self.flow.end.output_messages.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow output_messages accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.end.output_items.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow output_items accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.end.output_size.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow output_size accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.decision.dropped_items.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow dropped_items accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.input.input_messages.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow input_messages accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.input.input_items.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow input_items accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.input.input_size.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow input_size accumulator poisoned");
std::mem::take(&mut *guard)
};
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 = {
let mut measurement = measurement
.lock()
.expect("flow duration accumulator poisoned");
measurement.terminal_snapshots()
};
for snapshot in snapshots {
let _ = reporter
.report_snapshot_reliably_until(snapshot, deadline)
.await?;
}
}
if let Some((metrics, acc_mutex)) = self.flow.end.output_messages.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow output_messages accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.end.output_items.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow output_items accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.end.output_size.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow output_size accumulator poisoned");
std::mem::take(&mut *guard)
};
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_mutex)) = self.flow.decision.dropped_items.as_mut() {
let drained = {
let mut guard = acc_mutex
.lock()
.expect("flow dropped_items accumulator poisoned");
std::mem::take(&mut *guard)
};
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 async fn send_message(&self, mut data: PData) -> Result<(), TypedError<PData>>
where
PData: crate::processor::FlowMetricHook + Send,
{
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 + Send,
{
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 + Send,
{
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 + Send,
{
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::flow_metrics::{
FlowAttributeSet, FlowDurationNormalMetrics, FlowOutputItemsMetrics,
};
use crate::shared::message::SharedSender;
use crate::testing::{test_node, test_pipeline_ctx};
use otel_arrow_dfe_channel::error::SendError;
use otel_arrow_dfe_telemetry::metrics::MetricValue;
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use std::collections::HashMap;
#[test]
fn effect_handler_try_send_message_success() {
let (tx, mut rx) = tokio::sync::mpsc::channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("out".into(), SharedSender::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) = tokio::sync::mpsc::channel::<u64>(1);
let mut senders = HashMap::new();
let _ = senders.insert("out".into(), SharedSender::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) = tokio::sync::mpsc::channel::<u64>(10);
let (b_tx, _b_rx) = tokio::sync::mpsc::channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("a".into(), SharedSender::mpsc(a_tx));
let _ = senders.insert("b".into(), SharedSender::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, mut a_rx) = tokio::sync::mpsc::channel::<u64>(10);
let (b_tx, mut b_rx) = tokio::sync::mpsc::channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("a".into(), SharedSender::mpsc(a_tx));
let _ = senders.insert("b".into(), SharedSender::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) = tokio::sync::mpsc::channel::<u64>(1);
let mut senders = HashMap::new();
let _ = senders.insert("out".into(), SharedSender::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) = tokio::sync::mpsc::channel::<u64>(10);
let mut senders = HashMap::new();
let _ = senders.insert("out".into(), SharedSender::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(_))));
}
#[test]
fn flow_metric_marker_accumulates_after_begin_process_timing_shared() {
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.is_flow_start());
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_metric_marker_not_armed_when_timing_disabled_shared() {
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"
);
}
#[test]
fn shared_handler_record_flow_duration_drains_to_metric_set() {
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 start_metric_set = FlowInputItemsMetrics::register(®istrar);
let duration_metric_set = FlowDurationNormalMetrics::register(®istrar);
let outgoing_metric_set = FlowOutputItemsMetrics::register(®istrar);
let (metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(5);
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,
true,
None,
Some(start_metric_set),
None,
Some(duration_metric_set.into()),
Some(outgoing_metric_set),
None,
None,
None,
true,
true,
);
assert!(eh.is_flow_start());
assert!(eh.is_flow_end());
eh.record_flow_input_items(SignalType::Logs, 10);
eh.record_flow_input_items(SignalType::Metrics, 20);
for ns in [1000, 2000, 3000] {
eh.record_flow_duration(SignalType::Logs, ns);
}
eh.record_flow_output_items(SignalType::Logs, 7);
eh.record_flow_output_items(SignalType::Metrics, 8);
let start_before_report = eh
.flow
.input
.input_items
.as_ref()
.unwrap()
.1
.lock()
.unwrap();
assert_eq!(*start_before_report, [0, 20, 10]);
let before_report = eh.flow.end.duration.as_ref().unwrap().lock().unwrap();
let (count, sum, _, _) = before_report.pending_summary(SignalType::Logs);
assert_eq!(count, 3);
assert!((sum - 0.000_006).abs() < f64::EPSILON);
drop(before_report);
let output_before_report = eh.flow.end.output_items.as_ref().unwrap().1.lock().unwrap();
assert_eq!(*output_before_report, [0, 8, 7]);
drop(start_before_report);
drop(output_before_report);
eh.report_flow_metrics();
let start_drained = eh
.flow
.input
.input_items
.as_ref()
.unwrap()
.1
.lock()
.unwrap();
assert_eq!(
*start_drained, [0; 3],
"start accumulator should be drained"
);
let drained = eh.flow.end.duration.as_ref().unwrap().lock().unwrap();
assert_eq!(
drained.pending_summary(SignalType::Logs).0,
0,
"duration accumulator should be drained"
);
let output_drained = eh.flow.end.output_items.as_ref().unwrap().1.lock().unwrap();
assert_eq!(
*output_drained, [0; 3],
"stop item accumulator should be drained"
);
let snapshot = metrics_rx
.try_recv()
.expect("start flow_metric metric should be reported");
let [MetricValue::U64(consumed_snapshot)] = snapshot.get_metrics() else {
panic!("expected one flow input-item metric");
};
assert_eq!(
snapshot.measurement_attribute_value("signal"),
Some("metrics")
);
assert_eq!(*consumed_snapshot, 20);
let snapshot = metrics_rx
.try_recv()
.expect("logs flow input-item metric should be reported");
let [MetricValue::U64(consumed_snapshot)] = snapshot.get_metrics() else {
panic!("expected one flow input-item metric");
};
assert_eq!(snapshot.measurement_attribute_value("signal"), Some("logs"));
assert_eq!(*consumed_snapshot, 10);
let snapshot = metrics_rx
.try_recv()
.expect("flow duration metric should be reported");
let [MetricValue::Distribution(duration_snapshot)] = snapshot.get_metrics() else {
panic!("expected flow duration histogram");
};
let (count, sum, _, _) = duration_snapshot.summary();
assert_eq!(count, 3);
assert!((sum - 0.000_006).abs() < f64::EPSILON);
assert_eq!(snapshot.measurement_attribute_value("signal"), Some("logs"));
let snapshot = metrics_rx
.try_recv()
.expect("metrics flow output-item metric should be reported");
let [MetricValue::U64(produced_snapshot)] = snapshot.get_metrics() else {
panic!("expected flow output-item metric");
};
assert_eq!(
snapshot.measurement_attribute_value("signal"),
Some("metrics")
);
assert_eq!(*produced_snapshot, 8);
let snapshot = metrics_rx
.try_recv()
.expect("logs flow output-item metric should be reported");
let [MetricValue::U64(produced_snapshot)] = snapshot.get_metrics() else {
panic!("expected flow output-item metric");
};
assert_eq!(snapshot.measurement_attribute_value("signal"), Some("logs"));
assert_eq!(*produced_snapshot, 7);
}
#[test]
fn shared_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_size(SignalType::Logs, 10);
handler.record_flow_output_message(SignalType::Metrics);
handler.record_flow_output_size(SignalType::Metrics, 20);
assert_eq!(
*handler
.flow
.input
.input_messages
.as_ref()
.unwrap()
.1
.lock()
.unwrap(),
[0, 0, 1]
);
assert_eq!(
*handler
.flow
.end
.output_size
.as_ref()
.unwrap()
.1
.lock()
.unwrap(),
[0, 20, 0]
);
handler.report_flow_metrics();
assert_eq!(
*handler
.flow
.input
.input_messages
.as_ref()
.unwrap()
.1
.lock()
.unwrap(),
[0; 3]
);
assert_eq!(
*handler
.flow
.input
.input_size
.as_ref()
.unwrap()
.1
.lock()
.unwrap(),
[0; 3]
);
assert_eq!(
*handler
.flow
.end
.output_messages
.as_ref()
.unwrap()
.1
.lock()
.unwrap(),
[0; 3]
);
assert_eq!(
*handler
.flow
.end
.output_size
.as_ref()
.unwrap()
.1
.lock()
.unwrap(),
[0; 3]
);
}
}