use crate::channel_metrics::{
NodeCompletionMetrics, NodeInputItemMetrics, NodeInputMetrics, NodeInputSizeMetrics,
NodeOutputItemMetrics, NodeOutputMetrics, NodeOutputSizeMetrics,
};
use crate::clock;
use crate::completion_emission_metrics::CompletionEmissionMetricsHandle;
use crate::context::PipelineContext;
use crate::control::RouteData;
use crate::control::UnwindData;
use crate::control::{
AckMsg, ControlSenders, NackMsg, NodeControlMsg, PipelineCompletionMsg,
PipelineCompletionMsgReceiver, RuntimeControlMsg, RuntimeCtrlMsgReceiver,
};
use crate::control_plane_metrics::{PipelineCompletionMetricsState, RuntimeControlMetricsState};
use crate::error::Error;
use crate::forced_shutdown::ForcedShutdownTrigger;
use crate::memory_limiter::MemoryPressureChanged;
use crate::pipeline_metrics::PipelineMetricsMonitor;
use crate::terminal_state::TerminalMetricsDeadline;
use crate::{Interests, RequestOutcome, Unwindable};
use otel_arrow_dfe_config::DeployedPipelineKey;
use otel_arrow_dfe_config::MetricLevel;
use otel_arrow_dfe_config::SignalType;
use otel_arrow_dfe_config::policy::TelemetryPolicy;
use otel_arrow_dfe_telemetry::common_attributes::{Outcome, SignalOutcomeAttributes};
use otel_arrow_dfe_telemetry::error::Error as TelemetryError;
use otel_arrow_dfe_telemetry::event::{EngineEvent, ObservedEventReporter};
use otel_arrow_dfe_telemetry::metrics::{MeasurementMetricSet, MetricSetSnapshot};
use otel_arrow_dfe_telemetry::registry::TelemetryRegistryHandle;
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use otel_arrow_dfe_telemetry::{otel_debug, otel_warn};
use std::cell::RefCell;
use std::cmp::Reverse;
use std::collections::{BinaryHeap, HashMap, HashSet, VecDeque};
use std::rc::Rc;
use std::time::{Duration, Instant};
use tokio::sync::watch;
const PENDING_SENDS_WARN_THRESHOLD: usize = 100;
const RUNTIME_CTRL_BURST: usize = 64;
struct TimerState {
scheduled_time: Instant,
duration: Duration,
is_canceled: bool,
}
struct TimerSet {
timers: BinaryHeap<Reverse<(Instant, usize)>>,
timer_states: HashMap<usize, TimerState>,
}
impl TimerSet {
fn new() -> Self {
Self {
timers: BinaryHeap::new(),
timer_states: HashMap::new(),
}
}
fn start(&mut self, node_id: usize, duration: Duration) {
let when = clock::now() + duration;
self.timers.push(Reverse((when, node_id)));
let _ = self.timer_states.insert(
node_id,
TimerState {
scheduled_time: when,
duration,
is_canceled: false,
},
);
}
fn cancel(&mut self, node_id: usize) {
if let Some(timer_state) = self.timer_states.get_mut(&node_id) {
timer_state.is_canceled = true;
}
}
fn cancel_all(&mut self) {
self.timers.clear();
self.timer_states.clear();
}
fn next_expiry(&self) -> Option<Instant> {
self.timers.peek().map(|Reverse((when, _))| *when)
}
fn fire_due<F: FnMut(&usize)>(&mut self, now: Instant, mut on_fire: F) {
while let Some(Reverse((when, node_id))) = self.timers.peek().cloned() {
if when > now {
break;
}
let _ = self.timers.pop();
if let Some(timer_state) = self.timer_states.get_mut(&node_id) {
if !timer_state.is_canceled && timer_state.scheduled_time == when {
on_fire(&node_id);
let next_when = timer_state.scheduled_time + timer_state.duration;
self.timers.push(Reverse((next_when, node_id)));
timer_state.scheduled_time = next_when;
} else if timer_state.is_canceled {
let _ = self.timer_states.remove(&node_id);
}
}
}
}
}
fn opt_min<T: Ord>(a: Option<T>, b: Option<T>) -> Option<T> {
match (a, b) {
(Some(a), Some(b)) => Some(std::cmp::min(a, b)),
(Some(a), None) => Some(a),
(None, Some(b)) => Some(b),
(None, None) => None,
}
}
pub(crate) struct NodeMetricHandles {
pub(crate) registry: TelemetryRegistryHandle,
pub(crate) input: Option<MeasurementMetricSet<NodeInputMetrics>>,
pub(crate) input_completion: Option<MeasurementMetricSet<NodeCompletionMetrics>>,
pub(crate) input_items: Option<MeasurementMetricSet<NodeInputItemMetrics>>,
pub(crate) input_size: Option<MeasurementMetricSet<NodeInputSizeMetrics>>,
pub(crate) outputs: Vec<MeasurementMetricSet<NodeOutputMetrics>>,
pub(crate) output_completion: Vec<MeasurementMetricSet<NodeCompletionMetrics>>,
pub(crate) output_items: Vec<MeasurementMetricSet<NodeOutputItemMetrics>>,
pub(crate) output_size: Vec<MeasurementMetricSet<NodeOutputSizeMetrics>>,
pub(crate) completion_emission: Option<CompletionEmissionMetricsHandle>,
}
pub(crate) fn report_node_metrics_with_handles(
node_metric_handles: &Rc<RefCell<Vec<Option<NodeMetricHandles>>>>,
metrics_reporter: &mut MetricsReporter,
) -> Result<(), TelemetryError> {
let mut handles_guard = node_metric_handles.borrow_mut();
for handles in handles_guard.iter_mut().flatten() {
if let Some(input) = &mut handles.input {
metrics_reporter.report_measurement(input)?;
}
if let Some(input_completion) = &mut handles.input_completion {
metrics_reporter.report_measurement(input_completion)?;
}
if let Some(input_items) = &mut handles.input_items {
metrics_reporter.report_measurement(input_items)?;
}
if let Some(input_size) = &mut handles.input_size {
metrics_reporter.report_measurement(input_size)?;
}
for output in &mut handles.outputs {
metrics_reporter.report_measurement(output)?;
}
for output_completion in &mut handles.output_completion {
metrics_reporter.report_measurement(output_completion)?;
}
for output_items in &mut handles.output_items {
metrics_reporter.report_measurement(output_items)?;
}
for output_size in &mut handles.output_size {
metrics_reporter.report_measurement(output_size)?;
}
if let Some(completion_emission) = &handles.completion_emission {
let mut completion_emission = completion_emission
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
completion_emission.report(metrics_reporter)?;
}
}
Ok(())
}
pub(crate) fn snapshot_node_metrics_with_handles(
node_metric_handles: &Rc<RefCell<Vec<Option<NodeMetricHandles>>>>,
) -> Vec<MetricSetSnapshot> {
let mut handles_guard = node_metric_handles.borrow_mut();
let mut snapshots = Vec::new();
for handles in handles_guard.iter_mut().flatten() {
if let Some(input) = &mut handles.input {
snapshots.extend(input.terminal_snapshots());
}
if let Some(input_completion) = &mut handles.input_completion {
snapshots.extend(input_completion.terminal_snapshots());
}
if let Some(input_items) = &mut handles.input_items {
snapshots.extend(input_items.terminal_snapshots());
}
if let Some(input_size) = &mut handles.input_size {
snapshots.extend(input_size.terminal_snapshots());
}
for output in &mut handles.outputs {
snapshots.extend(output.terminal_snapshots());
}
for output_completion in &mut handles.output_completion {
snapshots.extend(output_completion.terminal_snapshots());
}
for output_items in &mut handles.output_items {
snapshots.extend(output_items.terminal_snapshots());
}
for output_size in &mut handles.output_size {
snapshots.extend(output_size.terminal_snapshots());
}
if let Some(completion_emission) = &handles.completion_emission
&& let Some(snapshot) = completion_emission
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.snapshot()
{
snapshots.push(snapshot);
}
}
snapshots
}
impl Drop for NodeMetricHandles {
fn drop(&mut self) {
if let Some(input) = self.input.take() {
let _ = self.registry.unregister_metric_set(input.metric_set_key());
}
if let Some(input_completion) = self.input_completion.take() {
let _ = self
.registry
.unregister_metric_set(input_completion.metric_set_key());
}
if let Some(input_items) = self.input_items.take() {
let _ = self
.registry
.unregister_metric_set(input_items.metric_set_key());
}
if let Some(input_size) = self.input_size.take() {
let _ = self
.registry
.unregister_metric_set(input_size.metric_set_key());
}
for output in self.outputs.drain(..) {
let _ = self.registry.unregister_metric_set(output.metric_set_key());
}
for output_completion in self.output_completion.drain(..) {
let _ = self
.registry
.unregister_metric_set(output_completion.metric_set_key());
}
for output_items in self.output_items.drain(..) {
let _ = self
.registry
.unregister_metric_set(output_items.metric_set_key());
}
for output_size in self.output_size.drain(..) {
let _ = self
.registry
.unregister_metric_set(output_size.metric_set_key());
}
}
}
pub struct RuntimeCtrlMsgManager<PData> {
pipeline_key: DeployedPipelineKey,
pipeline_context: PipelineContext,
runtime_ctrl_msg_receiver: RuntimeCtrlMsgReceiver<PData>,
memory_pressure_rx: watch::Receiver<MemoryPressureChanged>,
control_senders: ControlSenders<PData>,
tick_timers: TimerSet,
telemetry_timers: TimerSet,
event_reporter: ObservedEventReporter,
metrics_reporter: MetricsReporter,
control_plane_metrics_flush_interval: Duration,
channel_metrics: Vec<crate::channel_metrics::ChannelMetricsHandle>,
admission_metrics: Vec<crate::admission::metrics::AdmissionMetricsHandle>,
node_metric_handles: Rc<RefCell<Vec<Option<NodeMetricHandles>>>>,
telemetry: TelemetryPolicy,
pending_sends: VecDeque<(usize, NodeControlMsg<PData>)>,
runtime_control_metrics: RuntimeControlMetricsState,
terminal_metrics_deadline: TerminalMetricsDeadline,
forced_shutdown_trigger: ForcedShutdownTrigger,
}
impl<PData> RuntimeCtrlMsgManager<PData> {
#[must_use]
pub(crate) fn new(
pipeline_key: DeployedPipelineKey,
pipeline_context: PipelineContext,
runtime_ctrl_msg_receiver: RuntimeCtrlMsgReceiver<PData>,
memory_pressure_rx: watch::Receiver<MemoryPressureChanged>,
control_senders: ControlSenders<PData>,
event_reporter: ObservedEventReporter,
metrics_reporter: MetricsReporter,
control_plane_metrics_flush_interval: Duration,
telemetry_policy: TelemetryPolicy,
channel_metrics: Vec<crate::channel_metrics::ChannelMetricsHandle>,
admission_metrics: Vec<crate::admission::metrics::AdmissionMetricsHandle>,
node_metric_handles: Rc<RefCell<Vec<Option<NodeMetricHandles>>>>,
terminal_metrics_deadline: TerminalMetricsDeadline,
forced_shutdown_trigger: ForcedShutdownTrigger,
) -> Self {
let mut result = Self {
runtime_control_metrics: RuntimeControlMetricsState::new(
&pipeline_context,
metrics_reporter.clone(),
telemetry_policy.runtime_metrics,
0,
0,
),
pipeline_key,
pipeline_context,
runtime_ctrl_msg_receiver,
memory_pressure_rx,
control_senders,
tick_timers: TimerSet::new(),
telemetry_timers: TimerSet::new(),
event_reporter,
metrics_reporter,
control_plane_metrics_flush_interval,
channel_metrics,
admission_metrics,
node_metric_handles,
telemetry: telemetry_policy,
pending_sends: VecDeque::new(),
terminal_metrics_deadline,
forced_shutdown_trigger,
};
for node_id in result.control_senders.node_ids() {
result
.telemetry_timers
.start(node_id, result.control_plane_metrics_flush_interval);
}
result
}
pub async fn run(mut self) -> Result<(), Error> {
let internal_telemetry_enabled =
self.telemetry.pipeline_metrics || self.telemetry.tokio_metrics;
let mut pipeline_metrics_monitor = internal_telemetry_enabled
.then(|| PipelineMetricsMonitor::new(self.pipeline_context.clone()));
let mut is_draining_ingress = false;
let mut shutdown_deadline: Option<Instant> = None;
let mut shutdown_reason: Option<String> = None;
let mut pending_receivers: HashSet<usize> = HashSet::new();
let mut downstream_shutdown_sent = false;
let mut consecutive_runtime_ctrl = 0usize;
let mut retry_delay: Option<clock::Sleep> = None;
let mut metrics_flush_delay: Option<clock::Sleep> = None;
let mut shutdown_deadline_forced = false;
let mut memory_pressure_updates_open = true;
loop {
self.drain_pending_sends();
if self.runtime_control_metrics.is_dirty() && metrics_flush_delay.is_none() {
metrics_flush_delay = Some(clock::sleep(self.control_plane_metrics_flush_interval));
} else if !self.runtime_control_metrics.is_dirty() {
metrics_flush_delay = None;
}
if !self.pending_sends.is_empty() && retry_delay.is_none() {
retry_delay = Some(clock::sleep(Duration::from_millis(5)));
} else if self.pending_sends.is_empty() {
retry_delay = None;
}
let now = clock::now();
if let Some(deadline) = shutdown_deadline
&& now >= deadline
{
self.forced_shutdown_trigger.trigger();
shutdown_deadline_forced = true;
self.runtime_control_metrics
.record_shutdown_deadline_forced(now);
self.event_reporter
.report(EngineEvent::drain_deadline_reached(
self.pipeline_key.clone(),
shutdown_reason.clone(),
));
if let Some(reason) = shutdown_reason.as_ref() {
for node_id in self.control_senders.non_receiver_ids() {
self.send(
node_id,
NodeControlMsg::Shutdown {
deadline,
reason: reason.clone(),
},
);
}
for node_id in pending_receivers.iter().copied() {
self.send(
node_id,
NodeControlMsg::Shutdown {
deadline,
reason: reason.clone(),
},
);
}
}
self.report_runtime_control_metrics();
break;
}
let next_earliest = if is_draining_ingress {
shutdown_deadline
} else {
let next_expiry = self.tick_timers.next_expiry();
let next_tel_expiry = self.telemetry_timers.next_expiry();
opt_min(next_expiry, next_tel_expiry)
};
if consecutive_runtime_ctrl >= RUNTIME_CTRL_BURST
&& next_earliest.is_some_and(|when| when <= now)
{
self.handle_due_events(now, &mut pipeline_metrics_monitor);
consecutive_runtime_ctrl = 0;
continue;
}
tokio::select! {
biased;
msg = self.runtime_ctrl_msg_receiver.recv() => {
let Some(msg) = msg.ok() else { break; };
consecutive_runtime_ctrl += 1;
match msg {
RuntimeControlMsg::Shutdown { deadline, reason } => {
if is_draining_ingress {
continue;
}
self.terminal_metrics_deadline.record(deadline);
self.event_reporter.report(EngineEvent::shutdown_requested(
self.pipeline_key.clone(),
Some(reason.clone()),
));
is_draining_ingress = true;
shutdown_deadline = Some(deadline);
shutdown_reason = Some(reason.clone());
pending_receivers = self.control_senders.receiver_ids().into_iter().collect();
self.runtime_control_metrics
.record_shutdown_received(now, pending_receivers.len());
self.tick_timers.cancel_all();
self.telemetry_timers.cancel_all();
self.runtime_control_metrics.set_timer_counts(
self.tick_timers.timer_states.len(),
self.telemetry_timers.timer_states.len(),
);
for node_id in pending_receivers.iter().copied() {
self.send(
node_id,
NodeControlMsg::DrainIngress {
deadline,
reason: reason.clone(),
},
);
}
self.runtime_control_metrics.record_drain_ingress_sent();
self.event_reporter.report(EngineEvent::ingress_drain_started(
self.pipeline_key.clone(),
Some(reason.clone()),
));
if pending_receivers.is_empty() {
for node_id in self.control_senders.non_receiver_ids() {
self.send(
node_id,
NodeControlMsg::Shutdown {
deadline,
reason: reason.clone(),
},
);
}
downstream_shutdown_sent = true;
self.runtime_control_metrics.record_downstream_shutdown_sent();
self.event_reporter.report(
EngineEvent::downstream_shutdown_started(
self.pipeline_key.clone(),
Some(reason.clone()),
),
);
}
self.report_runtime_control_metrics();
},
RuntimeControlMsg::ReceiverDrained { node_id } => {
if !is_draining_ingress {
continue;
}
let removed = pending_receivers.remove(&node_id);
if removed {
self.runtime_control_metrics
.record_receiver_drained(now, pending_receivers.len());
}
if pending_receivers.is_empty() && !downstream_shutdown_sent {
self.event_reporter.report(EngineEvent::receivers_drained(
self.pipeline_key.clone(),
shutdown_reason.clone(),
));
let deadline = shutdown_deadline.unwrap_or(now);
let reason = shutdown_reason
.clone()
.unwrap_or_else(|| "pipeline shutting down".to_owned());
for node_id in self.control_senders.non_receiver_ids() {
self.send(
node_id,
NodeControlMsg::Shutdown {
deadline,
reason: reason.clone(),
},
);
}
downstream_shutdown_sent = true;
self.runtime_control_metrics.record_downstream_shutdown_sent();
self.event_reporter.report(
EngineEvent::downstream_shutdown_started(
self.pipeline_key.clone(),
Some(reason.clone()),
),
);
self.report_runtime_control_metrics();
}
}
RuntimeControlMsg::StartTimer { node_id, duration } => {
self.runtime_control_metrics.record_start_timer_received();
if is_draining_ingress {
otel_debug!(
"pipeline.draining.ignored_start_timer",
node_id = node_id,
);
} else {
self.tick_timers.start(node_id, duration);
self.runtime_control_metrics.set_timer_counts(
self.tick_timers.timer_states.len(),
self.telemetry_timers.timer_states.len(),
);
}
}
RuntimeControlMsg::CancelTimer { node_id } => {
self.runtime_control_metrics.record_cancel_timer_received();
if !is_draining_ingress {
self.tick_timers.cancel(node_id);
self.runtime_control_metrics.set_timer_counts(
self.tick_timers.timer_states.len(),
self.telemetry_timers.timer_states.len(),
);
}
}
RuntimeControlMsg::StartTelemetryTimer { node_id, duration } => {
self.runtime_control_metrics
.record_start_telemetry_timer_received();
if is_draining_ingress {
otel_debug!(
"pipeline.draining.ignored_start_telemetry_timer",
node_id = node_id,
"Ignoring StartTelemetryTimer during shutdown draining"
);
} else {
self.telemetry_timers.start(node_id, duration);
self.runtime_control_metrics.set_timer_counts(
self.tick_timers.timer_states.len(),
self.telemetry_timers.timer_states.len(),
);
}
}
RuntimeControlMsg::CancelTelemetryTimer { node_id, .. } => {
self.runtime_control_metrics
.record_cancel_telemetry_timer_received();
if !is_draining_ingress {
self.telemetry_timers.cancel(node_id);
self.runtime_control_metrics.set_timer_counts(
self.tick_timers.timer_states.len(),
self.telemetry_timers.timer_states.len(),
);
}
}
}
}
changed = self.memory_pressure_rx.changed(), if memory_pressure_updates_open => {
if changed.is_err() {
memory_pressure_updates_open = false;
continue;
}
let update = *self.memory_pressure_rx.borrow_and_update();
for node_id in self.control_senders.receiver_ids() {
self.send(node_id, NodeControlMsg::MemoryPressureChanged { update });
}
}
_ = async {
if let Some(when) = next_earliest
&& when > now
{
clock::sleep_until(when).await;
}
}, if next_earliest.is_some() => {
consecutive_runtime_ctrl = 0;
if !is_draining_ingress {
self.handle_due_events(clock::now(), &mut pipeline_metrics_monitor);
}
}
_ = async {
if let Some(delay) = retry_delay.as_mut() {
delay.await;
}
}, if retry_delay.is_some() => {
retry_delay = None;
continue;
}
_ = async {
if let Some(delay) = metrics_flush_delay.as_mut() {
delay.await;
}
}, if metrics_flush_delay.is_some() => {
metrics_flush_delay = None;
self.report_runtime_control_metrics();
}
}
}
if is_draining_ingress && !shutdown_deadline_forced {
self.runtime_control_metrics
.record_shutdown_completed(clock::now());
}
let terminal_metrics_deadline = self.terminal_metrics_deadline.get();
if let Some(pipeline_metrics_monitor) = pipeline_metrics_monitor.as_mut() {
if self.telemetry.pipeline_metrics {
pipeline_metrics_monitor.update_pipeline_metrics();
if let Err(err) = self
.metrics_reporter
.report_reliably_until(
pipeline_metrics_monitor.metrics_mut(),
terminal_metrics_deadline,
)
.await
{
otel_warn!(
"pipeline.metrics.final_reporting.fail",
error = err.to_string()
);
}
}
if self.telemetry.tokio_metrics {
pipeline_metrics_monitor.update_tokio_metrics();
if let Err(err) = self
.metrics_reporter
.report_reliably_until(
pipeline_metrics_monitor.tokio_metrics_mut(),
terminal_metrics_deadline,
)
.await
{
otel_warn!(
"tokio.metrics.final_reporting.fail",
error = err.to_string()
);
}
}
if let Err(err) = self
.metrics_reporter
.flush_until(terminal_metrics_deadline)
.await
{
otel_warn!(
"pipeline.metrics.final_collection.flush.fail",
error = err.to_string()
);
}
}
if let Err(err) = self
.runtime_control_metrics
.finish_reporting_until(terminal_metrics_deadline)
.await
{
otel_warn!(
"pipeline.runtime_control.metrics.final_reporting.fail",
error = err.to_string()
);
}
let _ = self.report_node_metrics();
Ok(())
}
fn handle_due_events(
&mut self,
now: Instant,
pipeline_metrics_monitor: &mut Option<PipelineMetricsMonitor>,
) {
let mut to_send: Vec<(usize, NodeControlMsg<PData>)> = Vec::new();
let mut timer_tick_count = 0usize;
let mut collect_telemetry_count = 0usize;
self.tick_timers.fire_due(now, |node_id| {
to_send.push((*node_id, NodeControlMsg::TimerTick {}));
timer_tick_count += 1;
});
let metrics_reporter = self.metrics_reporter.clone();
self.telemetry_timers.fire_due(now, |node_id| {
to_send.push((
*node_id,
NodeControlMsg::CollectTelemetry {
metrics_reporter: metrics_reporter.clone(),
},
));
collect_telemetry_count += 1;
});
self.runtime_control_metrics.set_timer_counts(
self.tick_timers.timer_states.len(),
self.telemetry_timers.timer_states.len(),
);
self.runtime_control_metrics
.record_due_events(timer_tick_count, collect_telemetry_count);
if let Some(pipeline_metrics_monitor) = pipeline_metrics_monitor.as_mut() {
if self.telemetry.pipeline_metrics {
pipeline_metrics_monitor.update_pipeline_metrics();
if let Err(err) = self
.metrics_reporter
.report(pipeline_metrics_monitor.metrics_mut())
{
otel_warn!("pipeline.metrics.reporting.fail", error = err.to_string());
}
}
if self.telemetry.tokio_metrics {
pipeline_metrics_monitor.update_tokio_metrics();
if let Err(err) = self
.metrics_reporter
.report(pipeline_metrics_monitor.tokio_metrics_mut())
{
otel_warn!("tokio.metrics.reporting.fail", error = err.to_string());
}
}
}
if self.telemetry.runtime_metrics >= MetricLevel::Basic {
for metrics in &self.channel_metrics {
if let Err(err) = metrics.report(&mut self.metrics_reporter) {
otel_warn!("channel.metrics.reporting.fail", error = err.to_string());
}
}
for metrics in &self.admission_metrics {
if let Err(err) = metrics.report(&mut self.metrics_reporter) {
otel_warn!("admission.metrics.reporting.fail", error = err.to_string());
}
}
}
if let Err(err) = self.report_node_metrics() {
otel_warn!("node.metrics.reporting.fail", error = err.to_string());
}
for (node_id, msg) in to_send {
self.send(node_id, msg);
}
}
fn report_node_metrics(&mut self) -> Result<(), TelemetryError> {
report_node_metrics_with_handles(&self.node_metric_handles, &mut self.metrics_reporter)
}
fn send(&mut self, node_id: usize, msg: NodeControlMsg<PData>) {
if let Some(sender) = self.control_senders.get(node_id) {
match sender.try_send(msg) {
Ok(()) => {}
Err(otel_arrow_dfe_channel::error::SendError::Full(msg)) => {
self.pending_sends.push_back((node_id, msg));
self.runtime_control_metrics
.set_pending_sends_buffered(self.pending_sends.len());
if self.pending_sends.len() == PENDING_SENDS_WARN_THRESHOLD {
otel_warn!(
"pipeline.ctrl.pending_sends.high",
count = self.pending_sends.len(),
"Pending sends buffer reached threshold; \
a node's control channel may be persistently full"
);
}
}
Err(otel_arrow_dfe_channel::error::SendError::Closed(_)) => {
}
}
}
}
fn drain_pending_sends(&mut self) {
let n = self.pending_sends.len();
for _ in 0..n {
let Some((node_id, msg)) = self.pending_sends.pop_front() else {
break;
};
if let Some(sender) = self.control_senders.get(node_id) {
match sender.try_send(msg) {
Ok(()) => {}
Err(otel_arrow_dfe_channel::error::SendError::Full(msg)) => {
self.pending_sends.push_back((node_id, msg));
}
Err(otel_arrow_dfe_channel::error::SendError::Closed(_)) => {
}
}
}
}
self.runtime_control_metrics
.set_pending_sends_buffered(self.pending_sends.len());
}
fn report_runtime_control_metrics(&mut self) {
if let Err(err) = self.runtime_control_metrics.report_if_needed() {
otel_warn!(
"pipeline.runtime_control.metrics.reporting.fail",
error = err.to_string()
);
}
}
}
pub struct PipelineCompletionMsgDispatcher<PData> {
pipeline_completion_msg_receiver: PipelineCompletionMsgReceiver<PData>,
control_senders: ControlSenders<PData>,
node_metric_handles: Rc<RefCell<Vec<Option<NodeMetricHandles>>>>,
pending_sends: VecDeque<(usize, NodeControlMsg<PData>)>,
completion_metrics: PipelineCompletionMetricsState,
control_plane_metrics_flush_interval: Duration,
terminal_metrics_deadline: TerminalMetricsDeadline,
}
impl<PData> PipelineCompletionMsgDispatcher<PData> {
#[must_use]
pub(crate) fn new(
pipeline_context: PipelineContext,
pipeline_completion_msg_receiver: PipelineCompletionMsgReceiver<PData>,
control_senders: ControlSenders<PData>,
node_metric_handles: Rc<RefCell<Vec<Option<NodeMetricHandles>>>>,
metrics_reporter: MetricsReporter,
control_plane_metrics_flush_interval: Duration,
telemetry_policy: TelemetryPolicy,
terminal_metrics_deadline: TerminalMetricsDeadline,
) -> Self {
Self {
pipeline_completion_msg_receiver,
control_senders,
node_metric_handles,
pending_sends: VecDeque::new(),
control_plane_metrics_flush_interval,
completion_metrics: PipelineCompletionMetricsState::new(
&pipeline_context,
metrics_reporter,
telemetry_policy.runtime_metrics,
),
terminal_metrics_deadline,
}
}
fn send(&mut self, node_id: usize, msg: NodeControlMsg<PData>) {
let is_ack = matches!(msg, NodeControlMsg::Ack(_));
let is_nack = matches!(msg, NodeControlMsg::Nack(_));
if let Some(sender) = self.control_senders.get(node_id) {
match sender.try_send(msg) {
Ok(()) => {
if is_ack {
self.completion_metrics.record_ack_delivered();
} else if is_nack {
self.completion_metrics.record_nack_delivered();
}
}
Err(otel_arrow_dfe_channel::error::SendError::Full(msg)) => {
self.pending_sends.push_back((node_id, msg));
self.completion_metrics
.set_pending_sends_buffered(self.pending_sends.len());
if self.pending_sends.len() == PENDING_SENDS_WARN_THRESHOLD {
otel_warn!(
"pipeline.return.pending_sends.high",
count = self.pending_sends.len(),
"Pending return sends buffer reached threshold"
);
}
}
Err(otel_arrow_dfe_channel::error::SendError::Closed(_)) => {}
}
}
}
fn drain_pending_sends(&mut self) {
let n = self.pending_sends.len();
for _ in 0..n {
let Some((node_id, msg)) = self.pending_sends.pop_front() else {
break;
};
let is_ack = matches!(msg, NodeControlMsg::Ack(_));
let is_nack = matches!(msg, NodeControlMsg::Nack(_));
if let Some(sender) = self.control_senders.get(node_id) {
match sender.try_send(msg) {
Ok(()) => {
if is_ack {
self.completion_metrics.record_ack_delivered();
} else if is_nack {
self.completion_metrics.record_nack_delivered();
}
}
Err(otel_arrow_dfe_channel::error::SendError::Full(msg)) => {
self.pending_sends.push_back((node_id, msg));
}
Err(otel_arrow_dfe_channel::error::SendError::Closed(_)) => {}
}
}
}
self.completion_metrics
.set_pending_sends_buffered(self.pending_sends.len());
}
fn record_frame_metrics(
&mut self,
node_id: usize,
interests: Interests,
route: &RouteData,
signal: Option<SignalType>,
output_items: u32,
input_items: u32,
output_size: u64,
input_size: u64,
outcome: RequestOutcome,
now_ns: u64,
) {
let mut handles_guard = self.node_metric_handles.borrow_mut();
if let Some(Some(handles)) = handles_guard.get_mut(node_id)
&& let Some(signal) = signal
{
let outcome = match outcome {
RequestOutcome::Success => Outcome::Success,
RequestOutcome::Failure => Outcome::Failure,
RequestOutcome::Refused => Outcome::Refused,
};
let attributes = SignalOutcomeAttributes { signal, outcome };
if interests.contains(Interests::NODE_INPUT_METRICS)
&& let Some(input) = &mut handles.input
{
input.with(attributes).messages.inc();
}
if input_items > 0
&& let Some(input_item_metrics) = &mut handles.input_items
{
input_item_metrics
.with(attributes)
.items
.add(input_items as u64);
}
if input_size > 0
&& let Some(input_size_metrics) = &mut handles.input_size
{
input_size_metrics.with(attributes).size.add(input_size);
}
let port = route.output_port_index as usize;
if interests.contains(Interests::NODE_OUTPUT_METRICS)
&& let Some(output) = handles.outputs.get_mut(port)
{
output.with(attributes).messages.inc();
}
if output_items > 0
&& let Some(output_item_metrics) = handles.output_items.get_mut(port)
{
output_item_metrics
.with(attributes)
.items
.add(output_items as u64);
}
if output_size > 0
&& let Some(output_size_metrics) = handles.output_size.get_mut(port)
{
output_size_metrics.with(attributes).size.add(output_size);
}
if interests.contains(Interests::NODE_COMPLETION_DURATION)
&& route.entry_time_ns > 0
&& now_ns > 0
{
let duration =
Duration::from_nanos(now_ns.saturating_sub(route.entry_time_ns)).as_secs_f64();
if let Some(input_completion) = &mut handles.input_completion {
input_completion.with(attributes).duration.record(duration);
} else if let Some(output_completion) = handles.output_completion.get_mut(port) {
output_completion.with(attributes).duration.record(duration);
}
}
}
self.completion_metrics
.set_pending_sends_buffered(self.pending_sends.len());
}
}
impl<PData: Unwindable> PipelineCompletionMsgDispatcher<PData> {
pub async fn run(mut self) -> Result<(), Error> {
let mut retry_delay: Option<clock::Sleep> = None;
let mut metrics_flush_delay: Option<clock::Sleep> = None;
loop {
self.drain_pending_sends();
if self.completion_metrics.is_dirty() && metrics_flush_delay.is_none() {
metrics_flush_delay = Some(clock::sleep(self.control_plane_metrics_flush_interval));
} else if !self.completion_metrics.is_dirty() {
metrics_flush_delay = None;
}
if !self.pending_sends.is_empty() && retry_delay.is_none() {
retry_delay = Some(clock::sleep(Duration::from_millis(5)));
} else if self.pending_sends.is_empty() {
retry_delay = None;
}
tokio::select! {
biased;
msg = self.pipeline_completion_msg_receiver.recv() => {
let Some(msg) = msg.ok() else { break; };
match msg {
PipelineCompletionMsg::DeliverAck { ack } => {
self.completion_metrics.record_deliver_ack_received();
self.unwind_ack(ack);
}
PipelineCompletionMsg::DeliverNack { nack } => {
self.completion_metrics.record_deliver_nack_received();
self.unwind_nack(nack);
}
}
}
_ = async {
if let Some(delay) = retry_delay.as_mut() {
delay.await;
}
}, if retry_delay.is_some() => {
retry_delay = None;
}
_ = async {
if let Some(delay) = metrics_flush_delay.as_mut() {
delay.await;
}
}, if metrics_flush_delay.is_some() => {
metrics_flush_delay = None;
self.report_completion_metrics();
}
}
}
if let Err(err) = self
.completion_metrics
.finish_reporting_until(self.terminal_metrics_deadline.get())
.await
{
otel_warn!(
"pipeline.completion.metrics.final_reporting.fail",
error = err.to_string()
);
}
Ok(())
}
fn unwind_frames(
&mut self,
pdata: &mut PData,
now_ns: u64,
outcome: RequestOutcome,
interest: Interests,
) -> (Option<(usize, RouteData)>, usize) {
let signal = pdata.signal();
let mut unwind_depth = 0usize;
loop {
match pdata.pop_frame() {
None => return (None, unwind_depth),
Some(frame) => {
unwind_depth += 1;
if frame.interests.intersects(
Interests::NODE_METRICS
| Interests::NODE_COMPLETION_DURATION
| Interests::NODE_ITEM_COUNTS
| Interests::NODE_SIZE,
) {
self.record_frame_metrics(
frame.node_id,
frame.interests,
&frame.route,
signal,
frame.output_items,
frame.input_items,
frame.output_size,
frame.input_size,
outcome,
now_ns,
);
}
if frame.interests.contains(interest) {
if !frame.interests.contains(Interests::RETURN_DATA) {
pdata.drop_payload();
}
return (Some((frame.node_id, frame.route)), unwind_depth);
}
}
}
}
}
fn unwind_ack(&mut self, mut ack: AckMsg<PData>) {
let now_ns = ack.unwind.return_time_ns;
let (next, unwind_depth) = self.unwind_frames(
&mut ack.accepted,
now_ns,
RequestOutcome::Success,
Interests::ACKS,
);
if let Some((node_id, route)) = next {
ack.unwind = UnwindData::new(route, now_ns);
self.completion_metrics.record_ack_attempted(unwind_depth);
self.send(node_id, NodeControlMsg::Ack(ack));
} else {
self.completion_metrics
.record_ack_dropped_no_interest(unwind_depth);
}
}
fn unwind_nack(&mut self, mut nack: NackMsg<PData>) {
let now_ns = nack.unwind.return_time_ns;
let outcome = if nack.permanent {
RequestOutcome::Refused
} else {
RequestOutcome::Failure
};
let (next, unwind_depth) =
self.unwind_frames(&mut nack.refused, now_ns, outcome, Interests::NACKS);
if let Some((node_id, route)) = next {
nack.unwind = UnwindData::new(route, now_ns);
self.completion_metrics.record_nack_attempted(unwind_depth);
self.send(node_id, NodeControlMsg::Nack(nack));
} else {
self.completion_metrics
.record_nack_dropped_no_interest(unwind_depth);
}
}
fn report_completion_metrics(&mut self) {
if let Err(err) = self.completion_metrics.report_if_needed() {
otel_warn!(
"pipeline.completion.metrics.reporting.fail",
error = err.to_string()
);
}
}
}
#[cfg(test)]
impl<PData> RuntimeCtrlMsgManager<PData> {
pub(crate) fn test_tick_count(&self) -> usize {
self.tick_timers.timers.len()
}
pub(crate) fn test_telemetry_count(&self) -> usize {
self.telemetry_timers.timers.len()
}
pub(crate) fn test_control_senders_len(&self) -> usize {
self.control_senders.len()
}
pub(crate) fn test_push_tick_heap(&mut self, when: Instant, node_id: usize) {
self.tick_timers.timers.push(Reverse((when, node_id)));
}
pub(crate) fn test_pop_tick_heap(&mut self) -> Option<(Instant, usize)> {
self.tick_timers
.timers
.pop()
.map(|Reverse((when, node))| (when, node))
}
pub(crate) fn test_tick_heap_len(&self) -> usize {
self.tick_timers.timers.len()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::attributes::{ChannelImplementation, ChannelKind, ChannelMode, ChannelType};
use crate::channel_metrics::{
NodeCompletionMetrics, NodeInputItemMetrics, NodeInputMetrics, NodeOutputItemMetrics,
NodeOutputMetrics,
};
use crate::context::{ControllerContext, PipelineContextParams};
use crate::control::{AckMsg, Frame, NackMsg, RouteData, nanos_since_birth};
use crate::control::{
NodeControlMsg, PipelineCompletionMsg, RuntimeControlMsg, pipeline_completion_msg_channel,
runtime_ctrl_msg_channel,
};
use crate::message::{Receiver, Sender};
use crate::node::{NodeId, NodeType};
use crate::shared::message::{SharedReceiver, SharedSender};
use crate::testing::test_nodes;
use otel_arrow_dfe_channel::error::RecvError;
use otel_arrow_dfe_config::observed_state::{ObservedStateSettings, SendPolicy};
use otel_arrow_dfe_config::{PipelineGroupId, PipelineId};
use otel_arrow_dfe_state::store::ObservedStateStore;
use otel_arrow_dfe_telemetry::event::{
EngineEvent, ErrorEvent as TelemetryErrorEvent, EventType, ObservedEventReporter,
RequestEvent as TelemetryRequestEvent, SuccessEvent as TelemetrySuccessEvent,
};
use otel_arrow_dfe_telemetry::metrics::{MetricSetSnapshot, MetricValue};
use otel_arrow_dfe_telemetry::registry::MetricSetKey;
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
use std::cell::RefCell;
use std::collections::HashMap;
use std::rc::Rc;
use std::time::{Duration, Instant};
use tokio::task::LocalSet;
use tokio::time::timeout;
const TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL: Duration = Duration::from_secs(3600);
fn empty_node_metric_handles() -> Rc<RefCell<Vec<Option<NodeMetricHandles>>>> {
Rc::new(RefCell::new(Vec::new()))
}
async fn yield_cycles(count: usize) {
for _ in 0..count {
tokio::task::yield_now().await;
}
}
async fn expect_timer_tick(receiver: &mut Receiver<NodeControlMsg<()>>, node_name: &str) {
let result = timeout(Duration::from_secs(1), receiver.recv()).await;
match result {
Ok(Ok(NodeControlMsg::TimerTick {})) => {}
Ok(Ok(other)) => panic!("Expected TimerTick for {node_name}, got {other:?}"),
Ok(Err(e)) => panic!("Failed to receive message for {node_name}: {e:?}"),
Err(_) => panic!("Timed out waiting for TimerTick for {node_name}"),
}
}
fn assert_no_control_msg(receiver: &mut Receiver<NodeControlMsg<()>>, node_name: &str) {
match receiver.try_recv() {
Err(RecvError::Empty) => {}
Ok(msg) => panic!("Expected no control message for {node_name}, got {msg:?}"),
Err(e) => panic!("Unexpected receive error for {node_name}: {e:?}"),
}
}
fn create_test_pipeline_context()
-> (PipelineContext, crate::entity_context::PipelineEntityScope) {
let metrics_system = otel_arrow_dfe_telemetry::InternalTelemetrySystem::default();
let controller_context = ControllerContext::new(metrics_system.registry());
let pipeline_context_params = PipelineContextParams {
pipeline_group_id: Default::default(),
pipeline_id: Default::default(),
core_id: 0,
num_cores: 1,
thread_id: 0,
numa_node_id: 0,
};
let pipeline_context = PipelineContext::new(controller_context, pipeline_context_params);
let pipeline_entity_key = pipeline_context.register_pipeline_entity();
let pipeline_entity_guard = crate::entity_context::set_pipeline_entity_key(
pipeline_context.metrics_registry(),
pipeline_entity_key,
);
(pipeline_context, pipeline_entity_guard)
}
fn create_mock_control_sender<PData>() -> (
Sender<NodeControlMsg<PData>>,
Receiver<NodeControlMsg<PData>>,
) {
create_mock_control_sender_with_capacity(10)
}
fn create_mock_control_sender_with_capacity<PData>(
capacity: usize,
) -> (
Sender<NodeControlMsg<PData>>,
Receiver<NodeControlMsg<PData>>,
) {
let (tx, rx) = tokio::sync::mpsc::channel(capacity);
(
Sender::Shared(SharedSender::mpsc(tx)),
Receiver::Shared(SharedReceiver::mpsc(rx)),
)
}
fn build_test_manager<PData>(
pipeline_capacity: usize,
control_senders: ControlSenders<PData>,
) -> (
RuntimeCtrlMsgManager<PData>,
crate::control::RuntimeCtrlMsgSender<PData>,
crate::entity_context::PipelineEntityScope,
) {
let (pipeline_tx, pipeline_rx) = runtime_ctrl_msg_channel(pipeline_capacity);
let (_memory_pressure_tx, memory_pressure_rx) =
watch::channel(MemoryPressureChanged::initial());
let metrics_system = otel_arrow_dfe_telemetry::InternalTelemetrySystem::default();
let metrics_reporter = metrics_system.reporter();
let observed_state_store =
ObservedStateStore::new(&ObservedStateSettings::default(), metrics_system.registry());
let pipeline_group_id: PipelineGroupId = Default::default();
let pipeline_id: PipelineId = Default::default();
let core_id = 0;
let controller_context = ControllerContext::new(metrics_system.registry());
let pipeline_context_params = PipelineContextParams {
pipeline_group_id: pipeline_group_id.clone(),
pipeline_id: pipeline_id.clone(),
core_id,
num_cores: 1,
thread_id: 0,
numa_node_id: 0,
};
let pipeline_context = PipelineContext::new(controller_context, pipeline_context_params);
let pipeline_entity_key = pipeline_context.register_pipeline_entity();
let pipeline_entity_guard = crate::entity_context::set_pipeline_entity_key(
pipeline_context.metrics_registry(),
pipeline_entity_key,
);
let (forced_shutdown_trigger, _) = ForcedShutdownTrigger::pair();
let manager = RuntimeCtrlMsgManager::new(
DeployedPipelineKey {
pipeline_group_id,
pipeline_id,
core_id,
deployment_generation: 0,
},
pipeline_context,
pipeline_rx,
memory_pressure_rx,
control_senders,
observed_state_store.reporter(SendPolicy::default()),
metrics_reporter,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
TelemetryPolicy::default(),
Vec::new(),
Vec::new(),
empty_node_metric_handles(),
TerminalMetricsDeadline::default(),
forced_shutdown_trigger,
);
(manager, pipeline_tx, pipeline_entity_guard)
}
fn setup_test_manager_with_capacities<PData: Clone>(
pipeline_capacity: usize,
control_capacity: usize,
) -> (
RuntimeCtrlMsgManager<PData>,
crate::control::RuntimeCtrlMsgSender<PData>,
ControlSenders<PData>,
HashMap<usize, Receiver<NodeControlMsg<PData>>>,
Vec<NodeId>,
crate::entity_context::PipelineEntityScope,
) {
let mut control_senders = ControlSenders::new();
let mut control_receivers = HashMap::new();
let nodes = test_nodes(vec!["node1", "node2", "node3"]);
for node in &nodes {
let (sender, receiver) = create_mock_control_sender_with_capacity(control_capacity);
control_senders.register(node.clone(), NodeType::Processor, sender);
let _ = control_receivers.insert(node.index, receiver);
}
let (manager, pipeline_tx, pipeline_entity_guard) =
build_test_manager(pipeline_capacity, control_senders.clone());
(
manager,
pipeline_tx,
control_senders,
control_receivers,
nodes,
pipeline_entity_guard,
)
}
fn setup_test_manager<PData: Clone>() -> (
RuntimeCtrlMsgManager<PData>,
crate::control::RuntimeCtrlMsgSender<PData>,
HashMap<usize, Receiver<NodeControlMsg<PData>>>,
Vec<NodeId>,
crate::entity_context::PipelineEntityScope,
) {
let (
manager,
pipeline_tx,
_control_senders,
control_receivers,
nodes,
pipeline_entity_guard,
) = setup_test_manager_with_capacities(10, 10);
(
manager,
pipeline_tx,
control_receivers,
nodes,
pipeline_entity_guard,
)
}
#[tokio::test]
async fn test_run_start_timer_integration() {
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, mut control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let node = nodes.first().expect("ok");
let duration = Duration::from_millis(100);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let start_msg = RuntimeControlMsg::StartTimer {
node_id: node.index,
duration,
};
pipeline_tx.send(start_msg).await.unwrap();
let mut receiver = control_receivers.remove(&node.index).unwrap();
let tick_result =
timeout(Duration::from_millis(200), async { receiver.recv().await }).await;
assert!(
tick_result.is_ok(),
"Should receive TimerTick within timeout"
);
match tick_result.unwrap() {
Ok(NodeControlMsg::TimerTick {}) => {
}
Ok(other) => panic!("Expected TimerTick, got {other:?}"),
Err(e) => panic!("Failed to receive message: {e:?}"),
}
let second_tick_result =
timeout(Duration::from_millis(150), async { receiver.recv().await }).await;
assert!(
second_tick_result.is_ok(),
"Should receive second TimerTick for recurring timer"
);
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "".to_owned(),
})
.await
.unwrap();
drop(pipeline_tx);
let _ = timeout(Duration::from_millis(100), manager_handle).await;
})
.await;
}
#[tokio::test]
async fn test_run_cancel_timer_integration() {
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, mut control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let node = nodes.first().expect("ok");
let duration = Duration::from_millis(100);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let start_msg = RuntimeControlMsg::StartTimer {
node_id: node.index,
duration,
};
pipeline_tx.send(start_msg).await.unwrap();
let cancel_msg = RuntimeControlMsg::CancelTimer {
node_id: node.index,
};
pipeline_tx.send(cancel_msg).await.unwrap();
let mut receiver = control_receivers.remove(&node.index).unwrap();
let tick_result =
timeout(Duration::from_millis(200), async { receiver.recv().await }).await;
assert!(
tick_result.is_err(),
"Should not receive TimerTick for canceled timer"
);
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "".to_owned(),
})
.await
.unwrap();
drop(pipeline_tx);
let _ = timeout(Duration::from_millis(100), manager_handle).await;
})
.await;
}
#[tokio::test]
async fn test_run_multiple_timers_integration() {
let local = LocalSet::new();
let clock = clock::SimClock::new();
let _clock_guard = clock.install();
local
.run_until(async {
let (manager, pipeline_tx, mut control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let node1 = nodes.first().expect("ok");
let node2 = nodes.get(1).expect("ok");
let duration1 = Duration::from_millis(80); let duration2 = Duration::from_millis(120);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let start_msg1 = RuntimeControlMsg::StartTimer {
node_id: node1.index,
duration: duration1,
};
let start_msg2 = RuntimeControlMsg::StartTimer {
node_id: node2.index,
duration: duration2,
};
pipeline_tx.send(start_msg1).await.unwrap();
pipeline_tx.send(start_msg2).await.unwrap();
yield_cycles(2).await;
let mut receiver1 = control_receivers.remove(&node1.index).unwrap();
let mut receiver2 = control_receivers.remove(&node2.index).unwrap();
clock.advance(Duration::from_millis(79));
yield_cycles(2).await;
assert_no_control_msg(&mut receiver1, "node1");
assert_no_control_msg(&mut receiver2, "node2");
clock.advance(Duration::from_millis(1));
yield_cycles(2).await;
expect_timer_tick(&mut receiver1, "node1").await;
assert_no_control_msg(&mut receiver2, "node2");
clock.advance(Duration::from_millis(39));
yield_cycles(2).await;
assert_no_control_msg(&mut receiver2, "node2");
clock.advance(Duration::from_millis(1));
yield_cycles(2).await;
expect_timer_tick(&mut receiver2, "node2").await;
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: clock.now() + Duration::from_secs(1),
reason: "".to_owned(),
})
.await
.unwrap();
drop(pipeline_tx);
let _ = timeout(Duration::from_millis(100), manager_handle).await;
})
.await;
}
#[tokio::test]
async fn test_run_timer_replacement_integration() {
let local = LocalSet::new();
let clock = clock::SimClock::new();
let _clock_guard = clock.install();
local
.run_until(async {
let (manager, pipeline_tx, mut control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let node = nodes.first().expect("ok");
let first_duration = Duration::from_millis(150); let second_duration = Duration::from_millis(80);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let start_msg1 = RuntimeControlMsg::StartTimer {
node_id: node.index,
duration: first_duration,
};
pipeline_tx.send(start_msg1).await.unwrap();
yield_cycles(2).await;
clock.advance(Duration::from_millis(20));
yield_cycles(2).await;
let start_msg2 = RuntimeControlMsg::StartTimer {
node_id: node.index,
duration: second_duration,
};
pipeline_tx.send(start_msg2).await.unwrap();
yield_cycles(2).await;
let mut receiver = control_receivers.remove(&node.index).unwrap();
clock.advance(Duration::from_millis(79));
yield_cycles(2).await;
assert_no_control_msg(&mut receiver, "node");
clock.advance(Duration::from_millis(1));
yield_cycles(2).await;
expect_timer_tick(&mut receiver, "node").await;
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: clock.now() + Duration::from_secs(1),
reason: "".to_owned(),
})
.await
.unwrap();
drop(pipeline_tx);
let _ = timeout(Duration::from_millis(100), manager_handle).await;
})
.await;
}
#[tokio::test]
async fn test_run_shutdown_integration() {
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, _control_receivers, _, _pipeline_entity_guard) =
setup_test_manager::<()>();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "".to_owned(),
})
.await
.unwrap();
drop(pipeline_tx);
let shutdown_result = timeout(Duration::from_millis(100), manager_handle).await;
assert!(shutdown_result.is_ok(), "Manager should shutdown cleanly");
})
.await;
}
#[tokio::test]
async fn test_run_start_telemetry_timer_integration() {
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, mut control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let node = nodes.first().expect("ok");
let duration = Duration::from_millis(60);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let start_msg = RuntimeControlMsg::StartTelemetryTimer {
node_id: node.index,
duration,
};
pipeline_tx.send(start_msg).await.unwrap();
let mut receiver = control_receivers.remove(&node.index).unwrap();
let telemetry_result =
timeout(Duration::from_millis(200), async { receiver.recv().await }).await;
assert!(
telemetry_result.is_ok(),
"Should receive CollectTelemetry within timeout"
);
match telemetry_result.unwrap() {
Ok(NodeControlMsg::CollectTelemetry { .. }) => {
}
Ok(other) => panic!("Expected CollectTelemetry, got {other:?}"),
Err(e) => panic!("Failed to receive message: {e:?}"),
}
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "".to_owned(),
})
.await
.unwrap();
drop(pipeline_tx);
let _ = timeout(Duration::from_millis(100), manager_handle).await;
})
.await;
}
#[tokio::test]
async fn test_run_no_control_sender_integration() {
let local = LocalSet::new();
local
.run_until(async {
let (pipeline_tx, pipeline_rx) = runtime_ctrl_msg_channel(10);
let metrics_system = otel_arrow_dfe_telemetry::InternalTelemetrySystem::default();
let metrics_reporter = metrics_system.reporter();
let observed_state_store = ObservedStateStore::new(
&ObservedStateSettings::default(),
metrics_system.registry(),
);
let pipeline_group_id: PipelineGroupId = Default::default();
let pipeline_id: PipelineId = Default::default();
let core_id = 0;
let pipeline_key = DeployedPipelineKey {
pipeline_group_id: pipeline_group_id.clone(),
pipeline_id: pipeline_id.clone(),
core_id,
deployment_generation: 0,
};
let controller_context = ControllerContext::new(metrics_system.registry());
let pipeline_context_params = PipelineContextParams {
pipeline_group_id: pipeline_group_id.clone(),
pipeline_id: pipeline_id.clone(),
core_id,
num_cores: 1,
thread_id: 0,
numa_node_id: 0,
};
let pipeline_context =
PipelineContext::new(controller_context, pipeline_context_params);
let pipeline_entity_key = pipeline_context.register_pipeline_entity();
let _pipeline_entity_guard = crate::entity_context::set_pipeline_entity_key(
pipeline_context.metrics_registry(),
pipeline_entity_key,
);
let (_memory_pressure_tx, memory_pressure_rx) =
watch::channel(MemoryPressureChanged::initial());
let (forced_shutdown_trigger, _) = ForcedShutdownTrigger::pair();
let manager = RuntimeCtrlMsgManager::<()>::new(
pipeline_key,
pipeline_context,
pipeline_rx,
memory_pressure_rx,
ControlSenders::new(),
observed_state_store.reporter(SendPolicy::default()),
metrics_reporter,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
TelemetryPolicy::default(),
Vec::new(),
Vec::new(),
empty_node_metric_handles(),
TerminalMetricsDeadline::default(),
forced_shutdown_trigger,
);
let duration = Duration::from_millis(50);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let start_msg = RuntimeControlMsg::StartTimer {
node_id: 1234,
duration,
};
pipeline_tx.send(start_msg).await.unwrap();
tokio::time::sleep(Duration::from_millis(100)).await;
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "".to_owned(),
})
.await
.unwrap();
drop(pipeline_tx);
let shutdown_result = timeout(Duration::from_millis(100), manager_handle).await;
assert!(
shutdown_result.is_ok(),
"Manager should handle missing control sender gracefully"
);
})
.await;
}
#[tokio::test]
async fn test_run_timer_ordering_integration() {
let local = LocalSet::new();
let clock = clock::SimClock::new();
let _clock_guard = clock.install();
local
.run_until(async {
let (manager, pipeline_tx, mut control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let node1 = nodes.first().expect("ok");
let node2 = nodes.get(1).expect("ok");
let node3 = nodes.get(2).expect("ok");
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let start_msg1 = RuntimeControlMsg::StartTimer {
node_id: node1.index,
duration: Duration::from_millis(120), };
let start_msg2 = RuntimeControlMsg::StartTimer {
node_id: node2.index,
duration: Duration::from_millis(60), };
let start_msg3 = RuntimeControlMsg::StartTimer {
node_id: node3.index,
duration: Duration::from_millis(90), };
pipeline_tx.send(start_msg1).await.unwrap();
pipeline_tx.send(start_msg2).await.unwrap();
pipeline_tx.send(start_msg3).await.unwrap();
yield_cycles(2).await;
let mut receiver1 = control_receivers.remove(&node1.index).unwrap();
let mut receiver2 = control_receivers.remove(&node2.index).unwrap();
let mut receiver3 = control_receivers.remove(&node3.index).unwrap();
clock.advance(Duration::from_millis(59));
yield_cycles(2).await;
assert_no_control_msg(&mut receiver1, "node1");
assert_no_control_msg(&mut receiver2, "node2");
assert_no_control_msg(&mut receiver3, "node3");
clock.advance(Duration::from_millis(1));
yield_cycles(2).await;
expect_timer_tick(&mut receiver2, "node2").await;
assert_no_control_msg(&mut receiver1, "node1");
assert_no_control_msg(&mut receiver3, "node3");
clock.advance(Duration::from_millis(29));
yield_cycles(2).await;
assert_no_control_msg(&mut receiver1, "node1");
assert_no_control_msg(&mut receiver3, "node3");
clock.advance(Duration::from_millis(1));
yield_cycles(2).await;
expect_timer_tick(&mut receiver3, "node3").await;
assert_no_control_msg(&mut receiver1, "node1");
clock.advance(Duration::from_millis(29));
yield_cycles(2).await;
assert_no_control_msg(&mut receiver1, "node1");
clock.advance(Duration::from_millis(1));
yield_cycles(2).await;
expect_timer_tick(&mut receiver1, "node1").await;
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: clock.now() + Duration::from_secs(1),
reason: "".to_owned(),
})
.await
.unwrap();
drop(pipeline_tx);
let _ = timeout(Duration::from_millis(100), manager_handle).await;
})
.await;
}
#[tokio::test]
async fn test_manager_creation() {
let (manager, _pipeline_tx, _control_receivers, _, _pipeline_entity_guard) =
setup_test_manager::<()>();
assert_eq!(
manager.tick_timers.timers.len(),
0,
"Timer queue should be empty initially"
);
assert_eq!(
manager.tick_timers.timer_states.len(),
0,
"Timer states map should be empty initially"
);
let tick_count = manager.test_tick_count();
assert_eq!(tick_count, 0, "Tick timer queue should be empty initially");
let telemetry_count = manager.test_telemetry_count();
assert_eq!(
telemetry_count, 3,
"Telemetry timers should be pre-registered for all 3 nodes"
);
assert_eq!(
manager.test_control_senders_len(),
3,
"Should have 3 mock control senders"
);
}
#[tokio::test]
async fn test_timer_heap_ordering() {
let (mut manager, _pipeline_tx, _control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let node1 = nodes.first().expect("ok");
let node2 = nodes.get(1).expect("ok");
let node3 = nodes.get(2).expect("ok");
let now = Instant::now();
let when1 = now + Duration::from_millis(300); let when2 = now + Duration::from_millis(100); let when3 = now + Duration::from_millis(200);
manager.test_push_tick_heap(when1, node1.index);
manager.test_push_tick_heap(when2, node2.index);
manager.test_push_tick_heap(when3, node3.index);
assert_eq!(
manager.test_tick_heap_len(),
3,
"All timers should be in the heap"
);
if let Some(Reverse((first_when, first_node))) = manager.tick_timers.timers.pop() {
assert_eq!(first_when, when2, "Earliest timer should be popped first");
assert_eq!(
first_node, node2.index,
"Correct node should be associated with earliest timer"
);
}
if let Some((second_when, second_node)) = manager.test_pop_tick_heap() {
assert_eq!(second_when, when3, "Middle timer should be popped second");
assert_eq!(
second_node, node3.index,
"Correct node should be associated with middle timer"
);
}
if let Some((third_when, third_node)) = manager.test_pop_tick_heap() {
assert_eq!(third_when, when1, "Latest timer should be popped last");
assert_eq!(
third_node, node1.index,
"Correct node should be associated with latest timer"
);
}
}
#[tokio::test]
async fn test_due_timer_tick_progress_under_runtime_ctrl_burst() {
let local = LocalSet::new();
local
.run_until(async {
let (
mut manager,
pipeline_tx,
_control_senders,
mut control_receivers,
nodes,
_pipeline_entity_guard,
) = setup_test_manager_with_capacities::<String>(128, 10);
let noisy_node = nodes[0].clone();
let target = nodes[1].clone();
manager
.tick_timers
.start(target.index, Duration::from_millis(1));
tokio::time::sleep(Duration::from_millis(5)).await;
for _ in 0..96 {
pipeline_tx
.send(RuntimeControlMsg::StartTimer {
node_id: noisy_node.index,
duration: Duration::from_secs(60),
})
.await
.unwrap();
}
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let mut receiver = control_receivers.remove(&target.index).unwrap();
let msg = timeout(Duration::from_millis(500), receiver.recv())
.await
.expect("TimerTick should make progress under runtime control burst")
.expect("target control channel should stay open");
assert!(matches!(msg, NodeControlMsg::TimerTick {}));
drop(pipeline_tx);
let shutdown_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(shutdown_result.is_ok(), "Manager should shutdown cleanly");
})
.await;
}
#[tokio::test]
async fn test_due_collect_telemetry_progress_under_runtime_ctrl_burst() {
let local = LocalSet::new();
local
.run_until(async {
let (
mut manager,
pipeline_tx,
_control_senders,
mut control_receivers,
nodes,
_pipeline_entity_guard,
) = setup_test_manager_with_capacities::<String>(128, 10);
let noisy_node = nodes[0].clone();
let target = nodes[1].clone();
manager
.telemetry_timers
.start(target.index, Duration::from_millis(1));
tokio::time::sleep(Duration::from_millis(5)).await;
for _ in 0..96 {
pipeline_tx
.send(RuntimeControlMsg::StartTimer {
node_id: noisy_node.index,
duration: Duration::from_secs(60),
})
.await
.unwrap();
}
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let mut receiver = control_receivers.remove(&target.index).unwrap();
let msg = timeout(Duration::from_millis(500), receiver.recv())
.await
.expect("CollectTelemetry should make progress under runtime control burst")
.expect("target control channel should stay open");
assert!(matches!(msg, NodeControlMsg::CollectTelemetry { .. }));
drop(pipeline_tx);
let shutdown_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(shutdown_result.is_ok(), "Manager should shutdown cleanly");
})
.await;
}
#[tokio::test]
async fn test_shutdown_allows_nodes_to_send_cleanup_messages() {
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, mut control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let node = nodes.first().expect("ok");
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let start_msg = RuntimeControlMsg::StartTimer {
node_id: node.index,
duration: Duration::from_secs(1), };
pipeline_tx.send(start_msg).await.unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "test shutdown".to_owned(),
})
.await
.unwrap();
let cancel_result = pipeline_tx
.send(RuntimeControlMsg::CancelTimer {
node_id: node.index,
})
.await;
assert!(
cancel_result.is_ok(),
"Nodes should be able to send control messages during cleanup, \
but the channel was closed prematurely"
);
drop(pipeline_tx);
let shutdown_result = timeout(Duration::from_millis(100), manager_handle).await;
assert!(
shutdown_result.is_ok(),
"Manager should shutdown cleanly after cleanup"
);
let mut receiver = control_receivers.remove(&node.index).unwrap();
while receiver.recv().await.is_ok() {}
})
.await;
}
#[tokio::test]
async fn test_duplicate_shutdown_ignored_during_draining() {
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, _control_receivers, _nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "first shutdown".to_owned(),
})
.await
.unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
let duplicate_result = pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "duplicate shutdown".to_owned(),
})
.await;
assert!(
duplicate_result.is_ok(),
"Duplicate shutdown should be accepted (and ignored)"
);
drop(pipeline_tx);
let shutdown_result = timeout(Duration::from_millis(100), manager_handle).await;
assert!(shutdown_result.is_ok(), "Manager should shutdown cleanly");
})
.await;
}
#[tokio::test(start_paused = true)]
async fn duplicate_shutdown_preserves_original_deadline() {
LocalSet::new()
.run_until(async {
let (manager, pipeline_tx, _control_receivers, _nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let deadline = manager.terminal_metrics_deadline.clone();
let forced_shutdown_signal = manager.forced_shutdown_trigger.subscribe();
let original = tokio::time::Instant::now() + Duration::from_millis(100);
let duplicate = original - Duration::from_millis(80);
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: original.into_std(),
reason: "first shutdown".to_owned(),
})
.await
.unwrap();
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: duplicate.into_std(),
reason: "duplicate shutdown".to_owned(),
})
.await
.unwrap();
let observer = tokio::task::spawn_local(async move {
let pending_through_duplicate = forced_shutdown_signal.clone();
tokio::select! {
biased;
() = tokio::time::sleep_until(duplicate) => {}
() = pending_through_duplicate.triggered() => {
panic!("forced shutdown fired at the duplicate deadline")
}
}
forced_shutdown_signal.triggered().await;
tokio::time::Instant::now()
});
manager.run().await.unwrap();
drop(pipeline_tx);
let fired_at = observer.await.expect("observer completes");
assert_eq!(deadline.get(), original.into_std());
assert_eq!(fired_at, original);
assert_eq!(tokio::time::Instant::now(), original);
})
.await;
}
#[tokio::test]
async fn test_deliver_ack_during_draining() {
use crate::control::AckMsg;
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, _control_receivers, _nodes, _pipeline_entity_guard) =
setup_test_manager::<String>();
let (return_tx, return_rx) = pipeline_completion_msg_channel(10);
let (dispatcher_context, dispatcher_guard) = create_test_pipeline_context();
let dispatcher = PipelineCompletionMsgDispatcher::new(
dispatcher_context,
return_rx,
ControlSenders::new(),
empty_node_metric_handles(),
MetricsReporter::create_new_and_receiver(16).1,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
TelemetryPolicy::default(),
TerminalMetricsDeadline::default(),
);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "test shutdown".to_owned(),
})
.await
.unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
let ack = AckMsg::new("ack_data".to_owned());
return_tx
.send(PipelineCompletionMsg::DeliverAck { ack })
.await
.unwrap();
drop(return_tx);
drop(pipeline_tx);
let dispatcher_result =
timeout(Duration::from_millis(100), dispatcher_handle).await;
assert!(
dispatcher_result.is_ok(),
"Return dispatcher should shutdown cleanly"
);
let shutdown_result = timeout(Duration::from_millis(100), manager_handle).await;
assert!(shutdown_result.is_ok(), "Manager should shutdown cleanly");
drop(dispatcher_guard);
})
.await;
}
#[tokio::test]
async fn test_deliver_nack_during_draining() {
use crate::control::NackMsg;
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, _control_receivers, _nodes, _pipeline_entity_guard) =
setup_test_manager::<String>();
let (return_tx, return_rx) = pipeline_completion_msg_channel(10);
let (dispatcher_context, dispatcher_guard) = create_test_pipeline_context();
let dispatcher = PipelineCompletionMsgDispatcher::new(
dispatcher_context,
return_rx,
ControlSenders::new(),
empty_node_metric_handles(),
MetricsReporter::create_new_and_receiver(16).1,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
TelemetryPolicy::default(),
TerminalMetricsDeadline::default(),
);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "test shutdown".to_owned(),
})
.await
.unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
let nack = NackMsg::new("test failure", "nack_data".to_owned());
return_tx
.send(PipelineCompletionMsg::DeliverNack { nack })
.await
.unwrap();
drop(return_tx);
drop(pipeline_tx);
let dispatcher_result =
timeout(Duration::from_millis(100), dispatcher_handle).await;
assert!(
dispatcher_result.is_ok(),
"Return dispatcher should shutdown cleanly"
);
let shutdown_result = timeout(Duration::from_millis(100), manager_handle).await;
assert!(shutdown_result.is_ok(), "Manager should shutdown cleanly");
drop(dispatcher_guard);
})
.await;
}
#[tokio::test]
async fn test_draining_deadline_forces_shutdown() {
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, _control_receivers, _nodes, _pipeline_entity_guard) =
setup_test_manager::<()>();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_millis(50),
reason: "test shutdown with short deadline".to_owned(),
})
.await
.unwrap();
let shutdown_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(
shutdown_result.is_ok(),
"Manager should shutdown when draining deadline is exceeded, even with active senders"
);
})
.await;
}
#[tokio::test]
async fn test_timer_tick_does_not_fire_during_draining() {
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, mut control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<String>();
let node = nodes.first().expect("ok");
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::StartTimer {
node_id: node.index,
duration: Duration::from_millis(10),
})
.await
.unwrap();
tokio::time::sleep(Duration::from_millis(5)).await;
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_millis(500),
reason: "test shutdown before timer fires".to_owned(),
})
.await
.unwrap();
let mut receiver = control_receivers.remove(&node.index).unwrap();
let shutdown = timeout(Duration::from_millis(100), receiver.recv())
.await
.expect("processor should be shut down during draining")
.expect("processor control channel should stay open");
assert!(
matches!(shutdown, NodeControlMsg::Shutdown { .. }),
"Processors should receive Shutdown immediately during draining"
);
tokio::time::sleep(Duration::from_millis(50)).await;
let msg = timeout(Duration::from_millis(100), receiver.recv()).await;
assert!(
msg.is_err(),
"Should NOT receive TimerTick during draining - timer ticks are suppressed"
);
drop(pipeline_tx);
let shutdown_result = timeout(Duration::from_millis(100), manager_handle).await;
assert!(shutdown_result.is_ok(), "Manager should shutdown cleanly");
})
.await;
}
#[tokio::test]
async fn test_start_timer_ignored_during_draining() {
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, mut control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<String>();
let node = nodes.first().expect("ok");
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "test shutdown".to_owned(),
})
.await
.unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
pipeline_tx
.send(RuntimeControlMsg::StartTimer {
node_id: node.index,
duration: Duration::from_millis(10),
})
.await
.unwrap();
let mut receiver = control_receivers.remove(&node.index).unwrap();
let shutdown = timeout(Duration::from_millis(100), receiver.recv())
.await
.expect("processor should be shut down during draining")
.expect("processor control channel should stay open");
assert!(
matches!(shutdown, NodeControlMsg::Shutdown { .. }),
"Processors should receive Shutdown immediately during draining"
);
tokio::time::sleep(Duration::from_millis(50)).await;
let msg = timeout(Duration::from_millis(100), receiver.recv()).await;
assert!(
msg.is_err(),
"Should NOT receive TimerTick - StartTimer should be ignored during draining"
);
drop(pipeline_tx);
let shutdown_result = timeout(Duration::from_millis(100), manager_handle).await;
assert!(shutdown_result.is_ok(), "Manager should shutdown cleanly");
})
.await;
}
#[tokio::test]
async fn test_start_telemetry_timer_ignored_during_draining() {
let local = LocalSet::new();
local
.run_until(async {
let (manager, pipeline_tx, mut control_receivers, nodes, _pipeline_entity_guard) =
setup_test_manager::<String>();
let node = nodes.first().expect("ok");
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "test shutdown".to_owned(),
})
.await
.unwrap();
tokio::time::sleep(Duration::from_millis(10)).await;
pipeline_tx
.send(RuntimeControlMsg::StartTelemetryTimer {
node_id: node.index,
duration: Duration::from_millis(10),
})
.await
.unwrap();
let mut receiver = control_receivers.remove(&node.index).unwrap();
let shutdown = timeout(Duration::from_millis(100), receiver.recv())
.await
.expect("processor should be shut down during draining")
.expect("processor control channel should stay open");
assert!(matches!(shutdown, NodeControlMsg::Shutdown { .. }));
let msg = timeout(Duration::from_millis(100), receiver.recv()).await;
assert!(
msg.is_err(),
"StartTelemetryTimer should be ignored during draining"
);
drop(pipeline_tx);
let shutdown_result = timeout(Duration::from_millis(100), manager_handle).await;
assert!(shutdown_result.is_ok(), "Manager should shutdown cleanly");
})
.await;
}
#[derive(Debug, Clone)]
struct TestPData {
frames: Vec<Frame>,
signal: Option<SignalType>,
}
impl TestPData {
fn new() -> Self {
Self {
frames: Vec::new(),
signal: None,
}
}
fn push_frame(&mut self, frame: Frame) {
self.frames.push(frame);
}
}
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> {
self.signal
}
fn drop_payload(&mut self) {}
}
impl crate::ReceivedAtNode for TestPData {
fn received_at_node(&mut self, _node_id: usize, _node_interests: Interests) {}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
enum MetricLabel {
RecvOutput,
RecvCompletion,
RecvOutputItems,
RecvOutputSize,
ProcInput,
ProcCompletion,
ProcInputItems,
ProcInputSize,
ProcOutput,
ProcOutputItems,
ProcOutputSize,
ExpInput,
ExpCompletion,
ExpInputItems,
ExpInputSize,
}
struct MetricsTestHarness {
manager: RuntimeCtrlMsgManager<TestPData>,
metrics_reporter: MetricsReporter,
pipeline_context: PipelineContext,
telemetry_policy: TelemetryPolicy,
pipeline_tx: crate::control::RuntimeCtrlMsgSender<TestPData>,
nodes: Vec<NodeId>,
_guard: crate::entity_context::PipelineEntityScope,
snapshot_rx: flume::Receiver<MetricSetSnapshot>,
key_labels: HashMap<MetricSetKey, MetricLabel>,
node_metric_handles: Rc<RefCell<Vec<Option<NodeMetricHandles>>>>,
}
fn setup_test_manager_with_metrics() -> MetricsTestHarness {
let (pipeline_tx, pipeline_rx) = runtime_ctrl_msg_channel(10);
let mut control_senders = ControlSenders::new();
let nodes = test_nodes(vec!["receiver", "processor", "exporter"]);
let node_types = [NodeType::Receiver, NodeType::Processor, NodeType::Exporter];
for (node, nt) in nodes.iter().zip(node_types.iter()) {
let (sender, _receiver) = create_mock_control_sender();
control_senders.register(node.clone(), *nt, sender);
}
let metrics_system = otel_arrow_dfe_telemetry::InternalTelemetrySystem::default();
let (snapshot_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(64);
let observed_state_store =
ObservedStateStore::new(&ObservedStateSettings::default(), metrics_system.registry());
let pipeline_group_id: PipelineGroupId = Default::default();
let pipeline_id: PipelineId = Default::default();
let controller_context = ControllerContext::new(metrics_system.registry());
let pipeline_context_params = PipelineContextParams {
pipeline_group_id: pipeline_group_id.clone(),
pipeline_id: pipeline_id.clone(),
core_id: 0,
num_cores: 1,
thread_id: 0,
numa_node_id: 0,
};
let pipeline_context = PipelineContext::new(controller_context, pipeline_context_params);
let pipeline_entity_key = pipeline_context.register_pipeline_entity();
let pipeline_entity_guard = crate::entity_context::set_pipeline_entity_key(
pipeline_context.metrics_registry(),
pipeline_entity_key,
);
let registry = pipeline_context.metrics_registry();
let recv_out_key = pipeline_context.register_node_channel_entity(
"recv:out".into(),
"output".into(),
ChannelKind::Pdata,
ChannelMode::Local,
ChannelType::Mpsc,
ChannelImplementation::Internal,
);
let proc_in_key = pipeline_context.register_node_channel_entity(
"proc:in".into(),
"input".into(),
ChannelKind::Pdata,
ChannelMode::Local,
ChannelType::Mpsc,
ChannelImplementation::Internal,
);
let proc_out_key = pipeline_context.register_node_channel_entity(
"proc:out".into(),
"output".into(),
ChannelKind::Pdata,
ChannelMode::Local,
ChannelType::Mpsc,
ChannelImplementation::Internal,
);
let exp_in_key = pipeline_context.register_node_channel_entity(
"exp:in".into(),
"input".into(),
ChannelKind::Pdata,
ChannelMode::Local,
ChannelType::Mpsc,
ChannelImplementation::Internal,
);
let recv_output: MeasurementMetricSet<NodeOutputMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(recv_out_key);
let recv_completion: MeasurementMetricSet<NodeCompletionMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(recv_out_key);
let recv_output_items: MeasurementMetricSet<NodeOutputItemMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(recv_out_key);
let recv_output_size: MeasurementMetricSet<NodeOutputSizeMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(recv_out_key);
let proc_input: MeasurementMetricSet<NodeInputMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(proc_in_key);
let proc_completion: MeasurementMetricSet<NodeCompletionMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(proc_in_key);
let proc_input_items: MeasurementMetricSet<NodeInputItemMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(proc_in_key);
let proc_input_size: MeasurementMetricSet<NodeInputSizeMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(proc_in_key);
let proc_output: MeasurementMetricSet<NodeOutputMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(proc_out_key);
let proc_output_items: MeasurementMetricSet<NodeOutputItemMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(proc_out_key);
let proc_output_size: MeasurementMetricSet<NodeOutputSizeMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(proc_out_key);
let exp_input: MeasurementMetricSet<NodeInputMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(exp_in_key);
let exp_completion: MeasurementMetricSet<NodeCompletionMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(exp_in_key);
let exp_input_items: MeasurementMetricSet<NodeInputItemMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(exp_in_key);
let exp_input_size: MeasurementMetricSet<NodeInputSizeMetrics> =
registry.register_metric_set_with_measurement_attributes_for_entity(exp_in_key);
let mut key_labels = HashMap::new();
let _ = key_labels.insert(recv_output.metric_set_key(), MetricLabel::RecvOutput);
let _ = key_labels.insert(
recv_completion.metric_set_key(),
MetricLabel::RecvCompletion,
);
let _ = key_labels.insert(
recv_output_items.metric_set_key(),
MetricLabel::RecvOutputItems,
);
let _ = key_labels.insert(
recv_output_size.metric_set_key(),
MetricLabel::RecvOutputSize,
);
let _ = key_labels.insert(proc_input.metric_set_key(), MetricLabel::ProcInput);
let _ = key_labels.insert(
proc_completion.metric_set_key(),
MetricLabel::ProcCompletion,
);
let _ = key_labels.insert(
proc_input_items.metric_set_key(),
MetricLabel::ProcInputItems,
);
let _ = key_labels.insert(proc_input_size.metric_set_key(), MetricLabel::ProcInputSize);
let _ = key_labels.insert(proc_output.metric_set_key(), MetricLabel::ProcOutput);
let _ = key_labels.insert(
proc_output_items.metric_set_key(),
MetricLabel::ProcOutputItems,
);
let _ = key_labels.insert(
proc_output_size.metric_set_key(),
MetricLabel::ProcOutputSize,
);
let _ = key_labels.insert(exp_input.metric_set_key(), MetricLabel::ExpInput);
let _ = key_labels.insert(exp_completion.metric_set_key(), MetricLabel::ExpCompletion);
let _ = key_labels.insert(exp_input_items.metric_set_key(), MetricLabel::ExpInputItems);
let _ = key_labels.insert(exp_input_size.metric_set_key(), MetricLabel::ExpInputSize);
let mut node_metric_handles: Vec<Option<NodeMetricHandles>> = Vec::new();
let max_idx = nodes.iter().map(|n| n.index).max().unwrap_or(0);
for _ in 0..=max_idx {
node_metric_handles.push(None);
}
node_metric_handles[nodes[0].index] = Some(NodeMetricHandles {
registry: registry.clone(),
input: None,
input_completion: None,
input_items: None,
input_size: None,
outputs: vec![recv_output],
output_completion: vec![recv_completion],
output_items: vec![recv_output_items],
output_size: vec![recv_output_size],
completion_emission: None,
});
node_metric_handles[nodes[1].index] = Some(NodeMetricHandles {
registry: registry.clone(),
input: Some(proc_input),
input_completion: Some(proc_completion),
input_items: Some(proc_input_items),
input_size: Some(proc_input_size),
outputs: vec![proc_output],
output_completion: Vec::new(),
output_items: vec![proc_output_items],
output_size: vec![proc_output_size],
completion_emission: None,
});
node_metric_handles[nodes[2].index] = Some(NodeMetricHandles {
registry: registry.clone(),
input: Some(exp_input),
input_completion: Some(exp_completion),
input_items: Some(exp_input_items),
input_size: Some(exp_input_size),
outputs: Vec::new(),
output_completion: Vec::new(),
output_items: Vec::new(),
output_size: Vec::new(),
completion_emission: None,
});
let telemetry_policy = TelemetryPolicy {
runtime_metrics: MetricLevel::Detailed,
..Default::default()
};
let node_metric_handles = Rc::new(RefCell::new(node_metric_handles));
let (_memory_pressure_tx, memory_pressure_rx) =
watch::channel(MemoryPressureChanged::initial());
let (forced_shutdown_trigger, _) = ForcedShutdownTrigger::pair();
let manager = RuntimeCtrlMsgManager::new(
DeployedPipelineKey {
pipeline_group_id,
pipeline_id,
core_id: 0,
deployment_generation: 0,
},
pipeline_context.clone(),
pipeline_rx,
memory_pressure_rx,
control_senders,
observed_state_store.reporter(SendPolicy::default()),
metrics_reporter.clone(),
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
telemetry_policy.clone(),
Vec::new(),
Vec::new(),
node_metric_handles.clone(),
TerminalMetricsDeadline::default(),
forced_shutdown_trigger,
);
MetricsTestHarness {
manager,
metrics_reporter: metrics_reporter.clone(),
pipeline_context,
telemetry_policy,
pipeline_tx,
nodes,
_guard: pipeline_entity_guard,
snapshot_rx,
key_labels,
node_metric_handles,
}
}
fn collect_snapshots(
rx: &flume::Receiver<MetricSetSnapshot>,
key_labels: &HashMap<MetricSetKey, MetricLabel>,
) -> HashMap<MetricLabel, Vec<MetricValue>> {
let mut result: HashMap<MetricLabel, Vec<MetricValue>> = HashMap::new();
while let Ok(snapshot) = rx.try_recv() {
let Some(&label) = key_labels.get(&snapshot.key()) else {
continue;
};
let values = snapshot.get_metrics().to_vec();
let _ = result
.entry(label)
.and_modify(|existing| {
for (dst, src) in existing.iter_mut().zip(values.iter()) {
dst.add_in_place(src);
}
})
.or_insert(values);
}
result
}
fn assert_u64(values: &[MetricValue], index: usize, expected: u64, msg: &str) {
match &values[index] {
MetricValue::U64(v) => assert_eq!(*v, expected, "{msg}"),
other => panic!("{msg}: expected U64, got {other:?}"),
}
}
fn assert_u64_gte(values: &[MetricValue], index: usize, min: u64, msg: &str) {
match &values[index] {
MetricValue::U64(v) => assert!(*v >= min, "{msg}: expected >= {min}, got {v}"),
other => panic!("{msg}: expected U64, got {other:?}"),
}
}
fn assert_dist_is_mmsc(
values: &[MetricValue],
index: usize,
msg: &str,
) -> otel_arrow_dfe_telemetry::instrument::Mmsc {
match &values[index] {
MetricValue::Distribution(
otel_arrow_dfe_telemetry::instrument::DistributionValue::Basic(mmsc),
) => **mmsc,
other => panic!("{msg}: expected a basic-tier distribution, got {other:?}"),
}
}
fn assert_dist_is_normal_histogram(
values: &[MetricValue],
index: usize,
msg: &str,
) -> (u64, f64, f64, f64) {
match &values[index] {
MetricValue::Distribution(distribution) => {
assert_eq!(
distribution.tier_name(),
"normal",
"{msg}: expected a normal-tier histogram"
);
distribution.summary()
}
other => panic!("{msg}: expected a distribution, got {other:?}"),
}
}
fn assert_duration_seconds(values: &[MetricValue], index: usize, expected: f64, msg: &str) {
let (count, sum, min, max) = assert_dist_is_normal_histogram(values, index, msg);
assert_eq!(count, 1, "{msg}: expected one duration observation");
for (name, actual) in [("sum", sum), ("min", min), ("max", max)] {
assert!(
(actual - expected).abs() < f64::EPSILON,
"{msg}: expected {name}={expected}, got {actual}"
);
}
}
const COMPLETION_DURATION: usize = 0;
const INPUT_MESSAGES: usize = 0;
const OUTPUT_MESSAGES: usize = 0;
const ITEMS: usize = 0;
const SIZE: usize = 0;
const RUNTIME_DRAIN_ACTIVE: usize = 0;
const RUNTIME_DRAIN_PENDING_RECEIVERS: usize = 1;
const RUNTIME_PENDING_SENDS_BUFFERED: usize = 2;
const RUNTIME_TIMERS_ACTIVE: usize = 3;
const RUNTIME_TELEMETRY_TIMERS_ACTIVE: usize = 4;
const RUNTIME_SHUTDOWN_RECEIVED: usize = 5;
const RUNTIME_DRAIN_INGRESS_SENT: usize = 6;
const RUNTIME_RECEIVER_DRAINED_RECEIVED: usize = 7;
const RUNTIME_DOWNSTREAM_SHUTDOWN_SENT: usize = 8;
const RUNTIME_SHUTDOWN_DEADLINE_FORCED: usize = 9;
const RUNTIME_START_TIMER_RECEIVED: usize = 10;
const RUNTIME_START_TELEMETRY_TIMER_RECEIVED: usize = 12;
const RUNTIME_TIMER_TICK_SENT: usize = 14;
const RUNTIME_COLLECT_TELEMETRY_SENT: usize = 15;
const RUNTIME_DRAIN_RECEIVER_PHASE_DURATION_NS: usize = 16;
const RUNTIME_DRAIN_TOTAL_DURATION_NS: usize = 17;
const COMPLETION_PENDING_SENDS_BUFFERED: usize = 0;
const COMPLETION_DELIVER_ACK_RECEIVED: usize = 1;
const COMPLETION_DELIVER_NACK_RECEIVED: usize = 2;
const COMPLETION_ACK_ATTEMPTED: usize = 3;
const COMPLETION_NACK_ATTEMPTED: usize = 4;
const COMPLETION_ACK_DELIVERED: usize = 5;
const COMPLETION_NACK_DELIVERED: usize = 6;
const COMPLETION_ACK_DROPPED_NO_INTEREST: usize = 7;
const COMPLETION_NACK_DROPPED_NO_INTEREST: usize = 8;
const COMPLETION_UNWIND_DEPTH: usize = 9;
fn merge_metric_snapshots(
snapshots: Vec<Vec<MetricValue>>,
gauge_indices: &[usize],
) -> Option<Vec<MetricValue>> {
let mut iter = snapshots.into_iter();
let mut merged = iter.next()?;
for values in iter {
for (index, value) in values.iter().enumerate() {
if gauge_indices.contains(&index) {
merged[index] = value.clone();
} else {
merged[index].add_in_place(value);
}
}
}
Some(merged)
}
fn collect_metric_set_snapshots(
rx: &flume::Receiver<MetricSetSnapshot>,
key: Option<MetricSetKey>,
gauge_indices: &[usize],
) -> Option<Vec<MetricValue>> {
let key = key?;
let mut matching = Vec::new();
while let Ok(snapshot) = rx.try_recv() {
if snapshot.key() == key {
matching.push(snapshot.get_metrics().to_vec());
}
}
merge_metric_snapshots(matching, gauge_indices)
}
fn collect_engine_event_types(rx: &flume::Receiver<EngineEvent>) -> Vec<EventType> {
let mut events = Vec::new();
while let Ok(event) = rx.try_recv() {
events.push(event.r#type);
}
events
}
struct RuntimeControlTelemetryHarness<PData> {
manager: RuntimeCtrlMsgManager<PData>,
pipeline_tx: crate::control::RuntimeCtrlMsgSender<PData>,
control_receivers: HashMap<usize, Receiver<NodeControlMsg<PData>>>,
nodes: Vec<NodeId>,
runtime_metrics_key: Option<MetricSetKey>,
snapshot_rx: flume::Receiver<MetricSetSnapshot>,
engine_rx: flume::Receiver<EngineEvent>,
_guard: crate::entity_context::PipelineEntityScope,
}
struct MemoryPressureFanoutHarness<PData> {
manager: RuntimeCtrlMsgManager<PData>,
_pipeline_tx: crate::control::RuntimeCtrlMsgSender<PData>,
memory_pressure_tx: watch::Sender<MemoryPressureChanged>,
control_receivers: HashMap<usize, Receiver<NodeControlMsg<PData>>>,
nodes: Vec<NodeId>,
_guard: crate::entity_context::PipelineEntityScope,
}
fn setup_runtime_control_telemetry_harness<PData: Clone>(
node_specs: Vec<(&'static str, NodeType, usize)>,
metric_level: MetricLevel,
) -> RuntimeControlTelemetryHarness<PData> {
setup_runtime_control_telemetry_harness_with_flush_interval(
node_specs,
metric_level,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
)
}
fn setup_runtime_control_telemetry_harness_with_flush_interval<PData: Clone>(
node_specs: Vec<(&'static str, NodeType, usize)>,
metric_level: MetricLevel,
flush_interval: Duration,
) -> RuntimeControlTelemetryHarness<PData> {
let (pipeline_tx, pipeline_rx) = runtime_ctrl_msg_channel(16);
let metrics_system = otel_arrow_dfe_telemetry::InternalTelemetrySystem::default();
let (snapshot_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(64);
let controller_context = ControllerContext::new(metrics_system.registry());
let pipeline_group_id: PipelineGroupId = Default::default();
let pipeline_id: PipelineId = Default::default();
let pipeline_context_params = PipelineContextParams {
pipeline_group_id: pipeline_group_id.clone(),
pipeline_id: pipeline_id.clone(),
core_id: 0,
num_cores: 1,
thread_id: 0,
numa_node_id: 0,
};
let pipeline_context = PipelineContext::new(controller_context, pipeline_context_params);
let pipeline_entity_key = pipeline_context.register_pipeline_entity();
let pipeline_entity_guard = crate::entity_context::set_pipeline_entity_key(
pipeline_context.metrics_registry(),
pipeline_entity_key,
);
let nodes = test_nodes(node_specs.iter().map(|(name, _, _)| *name).collect());
let mut control_senders = ControlSenders::new();
let mut control_receivers = HashMap::new();
for (node, (_, node_type, capacity)) in nodes.iter().zip(node_specs.iter()) {
let (sender, receiver) = create_mock_control_sender_with_capacity(*capacity);
control_senders.register(node.clone(), *node_type, sender);
let _ = control_receivers.insert(node.index, receiver);
}
let (_log_tx, _log_rx) = flume::bounded(1);
let (engine_tx, engine_rx) = flume::unbounded();
let event_reporter = ObservedEventReporter::new_with_engine_sender(
SendPolicy::default(),
_log_tx,
engine_tx,
);
let (_memory_pressure_tx, memory_pressure_rx) =
watch::channel(MemoryPressureChanged::initial());
let (forced_shutdown_trigger, _) = ForcedShutdownTrigger::pair();
let manager = RuntimeCtrlMsgManager::new(
DeployedPipelineKey {
pipeline_group_id,
pipeline_id,
core_id: 0,
deployment_generation: 0,
},
pipeline_context,
pipeline_rx,
memory_pressure_rx,
control_senders,
event_reporter,
metrics_reporter,
flush_interval,
TelemetryPolicy {
pipeline_metrics: false,
tokio_metrics: false,
runtime_metrics: metric_level,
flow_metrics: Vec::new(),
},
Vec::new(),
Vec::new(),
empty_node_metric_handles(),
TerminalMetricsDeadline::default(),
forced_shutdown_trigger,
);
let runtime_metrics_key = manager.runtime_control_metrics.metric_set_key();
RuntimeControlTelemetryHarness {
manager,
pipeline_tx,
control_receivers,
nodes,
runtime_metrics_key,
snapshot_rx,
engine_rx,
_guard: pipeline_entity_guard,
}
}
fn setup_memory_pressure_fanout_harness<PData: Clone>(
node_specs: Vec<(&'static str, NodeType, usize)>,
) -> MemoryPressureFanoutHarness<PData> {
let (pipeline_tx, pipeline_rx) = runtime_ctrl_msg_channel(16);
let metrics_system = otel_arrow_dfe_telemetry::InternalTelemetrySystem::default();
let metrics_reporter = metrics_system.reporter();
let controller_context = ControllerContext::new(metrics_system.registry());
let pipeline_group_id: PipelineGroupId = Default::default();
let pipeline_id: PipelineId = Default::default();
let pipeline_context_params = PipelineContextParams {
pipeline_group_id: pipeline_group_id.clone(),
pipeline_id: pipeline_id.clone(),
core_id: 0,
num_cores: 1,
thread_id: 0,
numa_node_id: 0,
};
let pipeline_context = PipelineContext::new(controller_context, pipeline_context_params);
let pipeline_entity_key = pipeline_context.register_pipeline_entity();
let pipeline_entity_guard = crate::entity_context::set_pipeline_entity_key(
pipeline_context.metrics_registry(),
pipeline_entity_key,
);
let nodes = test_nodes(node_specs.iter().map(|(name, _, _)| *name).collect());
let mut control_senders = ControlSenders::new();
let mut control_receivers = HashMap::new();
for (node, (_, node_type, capacity)) in nodes.iter().zip(node_specs.iter()) {
let (sender, receiver) = create_mock_control_sender_with_capacity(*capacity);
control_senders.register(node.clone(), *node_type, sender);
let _ = control_receivers.insert(node.index, receiver);
}
let (_log_tx, _log_rx) = flume::bounded(1);
let (engine_tx, _engine_rx) = flume::unbounded();
let event_reporter = ObservedEventReporter::new_with_engine_sender(
SendPolicy::default(),
_log_tx,
engine_tx,
);
let (memory_pressure_tx, memory_pressure_rx) =
watch::channel(MemoryPressureChanged::initial());
let (forced_shutdown_trigger, _) = ForcedShutdownTrigger::pair();
let manager = RuntimeCtrlMsgManager::new(
DeployedPipelineKey {
pipeline_group_id,
pipeline_id,
core_id: 0,
deployment_generation: 0,
},
pipeline_context,
pipeline_rx,
memory_pressure_rx,
control_senders,
event_reporter,
metrics_reporter,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
TelemetryPolicy {
pipeline_metrics: false,
tokio_metrics: false,
runtime_metrics: MetricLevel::None,
flow_metrics: Vec::new(),
},
Vec::new(),
Vec::new(),
empty_node_metric_handles(),
TerminalMetricsDeadline::default(),
forced_shutdown_trigger,
);
MemoryPressureFanoutHarness {
manager,
_pipeline_tx: pipeline_tx,
memory_pressure_tx,
control_receivers,
nodes,
_guard: pipeline_entity_guard,
}
}
struct CompletionTelemetryHarness {
dispatcher: PipelineCompletionMsgDispatcher<TestPData>,
completion_tx: crate::control::PipelineCompletionMsgSender<TestPData>,
control_receivers: HashMap<usize, Receiver<NodeControlMsg<TestPData>>>,
nodes: Vec<NodeId>,
completion_metrics_key: Option<MetricSetKey>,
snapshot_rx: flume::Receiver<MetricSetSnapshot>,
_guard: crate::entity_context::PipelineEntityScope,
}
fn setup_completion_telemetry_harness(metric_level: MetricLevel) -> CompletionTelemetryHarness {
setup_completion_telemetry_harness_with_options(
metric_level,
16,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
64,
)
}
fn setup_completion_telemetry_harness_with_capacity(
metric_level: MetricLevel,
control_capacity: usize,
) -> CompletionTelemetryHarness {
setup_completion_telemetry_harness_with_options(
metric_level,
control_capacity,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
64,
)
}
fn setup_completion_telemetry_harness_with_capacity_and_flush_interval(
metric_level: MetricLevel,
control_capacity: usize,
flush_interval: Duration,
) -> CompletionTelemetryHarness {
setup_completion_telemetry_harness_with_options(
metric_level,
control_capacity,
flush_interval,
64,
)
}
fn setup_completion_telemetry_harness_with_options(
metric_level: MetricLevel,
control_capacity: usize,
flush_interval: Duration,
reporter_channel_size: usize,
) -> CompletionTelemetryHarness {
let (completion_tx, completion_rx) = pipeline_completion_msg_channel(16);
let metrics_system = otel_arrow_dfe_telemetry::InternalTelemetrySystem::default();
let (snapshot_rx, metrics_reporter) =
MetricsReporter::create_new_and_receiver(reporter_channel_size);
let controller_context = ControllerContext::new(metrics_system.registry());
let pipeline_context_params = PipelineContextParams {
pipeline_group_id: Default::default(),
pipeline_id: Default::default(),
core_id: 0,
num_cores: 1,
thread_id: 0,
numa_node_id: 0,
};
let pipeline_context = PipelineContext::new(controller_context, pipeline_context_params);
let pipeline_entity_key = pipeline_context.register_pipeline_entity();
let pipeline_entity_guard = crate::entity_context::set_pipeline_entity_key(
pipeline_context.metrics_registry(),
pipeline_entity_key,
);
let nodes = test_nodes(vec!["receiver", "processor", "exporter"]);
let mut control_senders = ControlSenders::new();
let mut control_receivers = HashMap::new();
for (node, node_type) in [
(&nodes[0], NodeType::Receiver),
(&nodes[1], NodeType::Processor),
(&nodes[2], NodeType::Exporter),
] {
let (sender, receiver) = create_mock_control_sender_with_capacity(control_capacity);
control_senders.register(node.clone(), node_type, sender);
let _ = control_receivers.insert(node.index, receiver);
}
let dispatcher = PipelineCompletionMsgDispatcher::new(
pipeline_context,
completion_rx,
control_senders,
empty_node_metric_handles(),
metrics_reporter,
flush_interval,
TelemetryPolicy {
pipeline_metrics: false,
tokio_metrics: false,
runtime_metrics: metric_level,
flow_metrics: Vec::new(),
},
TerminalMetricsDeadline::default(),
);
let completion_metrics_key = dispatcher.completion_metrics.metric_set_key();
CompletionTelemetryHarness {
dispatcher,
completion_tx,
control_receivers,
nodes,
completion_metrics_key,
snapshot_rx,
_guard: pipeline_entity_guard,
}
}
fn build_3node_pdata(nodes: &[NodeId], with_timestamp: bool) -> TestPData {
let mut pdata = TestPData::new();
pdata.signal = Some(SignalType::Logs);
pdata.push_frame(Frame {
node_id: nodes[0].index,
interests: Interests::NODE_OUTPUT_METRICS | Interests::ACKS | Interests::NACKS,
route: RouteData {
calldata: Default::default(),
entry_time_ns: 0,
output_port_index: 0,
},
output_items: 10,
input_items: 0,
output_size: 100,
input_size: 0,
});
let entry_time_ns = if with_timestamp {
nanos_since_birth()
} else {
0
};
pdata.push_frame(Frame {
node_id: nodes[1].index,
interests: Interests::NODE_INPUT_METRICS
| Interests::NODE_OUTPUT_METRICS
| Interests::NODE_COMPLETION_DURATION
| Interests::ACKS
| Interests::NACKS,
route: RouteData {
calldata: Default::default(),
entry_time_ns,
output_port_index: 0,
},
output_items: 7,
input_items: 10,
output_size: 70,
input_size: 100,
});
let entry_time_ns = if with_timestamp {
nanos_since_birth()
} else {
0
};
pdata.push_frame(Frame {
node_id: nodes[2].index,
interests: Interests::NODE_INPUT_METRICS | Interests::NODE_COMPLETION_DURATION,
route: RouteData {
calldata: Default::default(),
entry_time_ns,
output_port_index: 0,
},
output_items: 0,
input_items: 7,
output_size: 0,
input_size: 70,
});
pdata
}
fn build_3node_pdata_no_subscribers(nodes: &[NodeId], with_timestamp: bool) -> TestPData {
let mut pdata = TestPData::new();
pdata.signal = Some(SignalType::Logs);
let entry_time_ns = if with_timestamp {
nanos_since_birth()
} else {
0
};
pdata.push_frame(Frame {
node_id: nodes[0].index,
interests: Interests::NODE_OUTPUT_METRICS | Interests::NODE_COMPLETION_DURATION,
route: RouteData {
calldata: Default::default(),
entry_time_ns,
output_port_index: 0,
},
output_items: 10,
input_items: 0,
output_size: 100,
input_size: 0,
});
let entry_time_ns = if with_timestamp {
nanos_since_birth()
} else {
0
};
pdata.push_frame(Frame {
node_id: nodes[1].index,
interests: Interests::NODE_INPUT_METRICS
| Interests::NODE_OUTPUT_METRICS
| Interests::NODE_COMPLETION_DURATION,
route: RouteData {
calldata: Default::default(),
entry_time_ns,
output_port_index: 0,
},
output_items: 7,
input_items: 10,
output_size: 70,
input_size: 100,
});
let entry_time_ns = if with_timestamp {
nanos_since_birth()
} else {
0
};
pdata.push_frame(Frame {
node_id: nodes[2].index,
interests: Interests::NODE_INPUT_METRICS | Interests::NODE_COMPLETION_DURATION,
route: RouteData {
calldata: Default::default(),
entry_time_ns,
output_port_index: 0,
},
output_items: 0,
input_items: 7,
output_size: 0,
input_size: 70,
});
pdata
}
async fn run_and_collect(
harness: MetricsTestHarness,
send_fn: impl FnOnce(&[NodeId]) -> Vec<PipelineCompletionMsg<TestPData>>,
) -> HashMap<MetricLabel, Vec<MetricValue>> {
let MetricsTestHarness {
manager,
mut metrics_reporter,
pipeline_context,
telemetry_policy,
pipeline_tx,
nodes,
_guard,
snapshot_rx,
key_labels,
node_metric_handles,
} = harness;
let local = LocalSet::new();
local
.run_until(async {
let msgs = send_fn(&nodes);
let (return_tx, return_rx) = pipeline_completion_msg_channel(32);
let return_dispatcher = PipelineCompletionMsgDispatcher::new(
pipeline_context.clone(),
return_rx,
ControlSenders::new(),
node_metric_handles.clone(),
metrics_reporter.clone(),
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
telemetry_policy.clone(),
TerminalMetricsDeadline::default(),
);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let dispatcher_handle =
tokio::task::spawn_local(async move { return_dispatcher.run().await });
for msg in msgs {
return_tx.send(msg).await.unwrap();
}
drop(return_tx);
drop(pipeline_tx);
let dispatcher_result =
timeout(Duration::from_millis(500), dispatcher_handle).await;
assert!(
dispatcher_result.is_ok(),
"Return dispatcher should shut down cleanly"
);
let result = timeout(Duration::from_millis(500), manager_handle).await;
assert!(result.is_ok(), "Manager should shut down cleanly");
report_node_metrics_with_handles(&node_metric_handles, &mut metrics_reporter)
.expect("Final node metrics flush should succeed");
drop(_guard);
collect_snapshots(&snapshot_rx, &key_labels)
})
.await
}
#[tokio::test]
async fn test_ack_lifecycle_input_output_metrics() {
let harness = setup_test_manager_with_metrics();
let nodes_clone = harness.nodes.clone();
let snapshots = run_and_collect(harness, |nodes| {
let pdata = build_3node_pdata(nodes, false);
vec![PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata),
}]
})
.await;
let exp = &snapshots[&MetricLabel::ExpInput];
assert_u64(exp, INPUT_MESSAGES, 1, "Exporter input messages");
let proc_c = &snapshots[&MetricLabel::ProcInput];
assert_u64(proc_c, INPUT_MESSAGES, 1, "Processor input messages");
let proc_p = &snapshots[&MetricLabel::ProcOutput];
assert_u64(proc_p, OUTPUT_MESSAGES, 1, "Processor output messages");
assert!(
!snapshots.contains_key(&MetricLabel::RecvOutput),
"Receiver output should have no metrics (ack delivered at processor)"
);
drop(nodes_clone);
}
#[tokio::test]
async fn test_nack_lifecycle_failure_metrics() {
let harness = setup_test_manager_with_metrics();
let snapshots = run_and_collect(harness, |nodes| {
let pdata = build_3node_pdata(nodes, false);
vec![PipelineCompletionMsg::DeliverNack {
nack: NackMsg::new("transient error", pdata),
}]
})
.await;
let exp = &snapshots[&MetricLabel::ExpInput];
assert_u64(exp, INPUT_MESSAGES, 1, "Exporter input messages");
let proc_c = &snapshots[&MetricLabel::ProcInput];
assert_u64(proc_c, INPUT_MESSAGES, 1, "Processor input messages");
let proc_p = &snapshots[&MetricLabel::ProcOutput];
assert_u64(proc_p, OUTPUT_MESSAGES, 1, "Processor output messages");
}
#[tokio::test]
async fn test_permanent_nack_lifecycle_refused_metrics() {
let harness = setup_test_manager_with_metrics();
let snapshots = run_and_collect(harness, |nodes| {
let pdata = build_3node_pdata(nodes, false);
vec![PipelineCompletionMsg::DeliverNack {
nack: NackMsg::new_permanent("permanent refusal", pdata),
}]
})
.await;
let exp = &snapshots[&MetricLabel::ExpInput];
assert_u64(exp, INPUT_MESSAGES, 1, "Exporter input messages");
let proc_c = &snapshots[&MetricLabel::ProcInput];
assert_u64(proc_c, INPUT_MESSAGES, 1, "Processor input messages");
let proc_p = &snapshots[&MetricLabel::ProcOutput];
assert_u64(proc_p, OUTPUT_MESSAGES, 1, "Processor output messages");
}
#[tokio::test]
async fn test_ack_lifecycle_duration_histogram() {
const ENTRY_TIME_NS: u64 = 1_000_000_000;
const RETURN_TIME_NS: u64 = 1_250_000_000;
const EXPECTED_DURATION_SECONDS: f64 = 0.25;
let harness = setup_test_manager_with_metrics();
let snapshots = run_and_collect(harness, |nodes| {
let mut pdata = build_3node_pdata(nodes, true);
for frame in &mut pdata.frames {
if frame.route.entry_time_ns > 0 {
frame.route.entry_time_ns = ENTRY_TIME_NS;
}
}
let mut ack = AckMsg::new(pdata);
ack.unwind.return_time_ns = RETURN_TIME_NS;
vec![PipelineCompletionMsg::DeliverAck { ack }]
})
.await;
let exp = &snapshots[&MetricLabel::ExpCompletion];
assert_duration_seconds(
exp,
COMPLETION_DURATION,
EXPECTED_DURATION_SECONDS,
"Exporter completion duration",
);
let proc_c = &snapshots[&MetricLabel::ProcCompletion];
assert_duration_seconds(
proc_c,
COMPLETION_DURATION,
EXPECTED_DURATION_SECONDS,
"Processor completion duration",
);
assert!(!snapshots.contains_key(&MetricLabel::RecvCompletion));
}
#[tokio::test]
async fn test_completion_duration_records_without_message_metrics() {
const ENTRY_TIME_NS: u64 = 1_000_000_000;
const RETURN_TIME_NS: u64 = 1_250_000_000;
const EXPECTED_DURATION_SECONDS: f64 = 0.25;
let harness = setup_test_manager_with_metrics();
let snapshots = run_and_collect(harness, |nodes| {
let mut pdata = TestPData::new();
pdata.signal = Some(SignalType::Logs);
pdata.push_frame(Frame {
node_id: nodes[2].index,
interests: Interests::NODE_COMPLETION_DURATION,
route: RouteData {
entry_time_ns: ENTRY_TIME_NS,
..Default::default()
},
output_items: 0,
input_items: 0,
output_size: 0,
input_size: 0,
});
let mut ack = AckMsg::new(pdata);
ack.unwind.return_time_ns = RETURN_TIME_NS;
vec![PipelineCompletionMsg::DeliverAck { ack }]
})
.await;
let completion = &snapshots[&MetricLabel::ExpCompletion];
assert_duration_seconds(
completion,
COMPLETION_DURATION,
EXPECTED_DURATION_SECONDS,
"Exporter completion duration",
);
assert!(!snapshots.contains_key(&MetricLabel::ExpInput));
}
#[tokio::test]
async fn test_ack_lifecycle_receiver_completion_duration() {
const ENTRY_TIME_NS: u64 = 1_000_000_000;
const RETURN_TIME_NS: u64 = 1_250_000_000;
const EXPECTED_DURATION_SECONDS: f64 = 0.25;
let harness = setup_test_manager_with_metrics();
let snapshots = run_and_collect(harness, |nodes| {
let mut pdata = build_3node_pdata_no_subscribers(nodes, true);
for frame in &mut pdata.frames {
frame.route.entry_time_ns = ENTRY_TIME_NS;
}
let mut ack = AckMsg::new(pdata);
ack.unwind.return_time_ns = RETURN_TIME_NS;
vec![PipelineCompletionMsg::DeliverAck { ack }]
})
.await;
let recv_p = &snapshots[&MetricLabel::RecvCompletion];
assert_duration_seconds(
recv_p,
COMPLETION_DURATION,
EXPECTED_DURATION_SECONDS,
"Receiver completion duration",
);
let proc_c = &snapshots[&MetricLabel::ProcCompletion];
let (count, _, _, _) = assert_dist_is_normal_histogram(
proc_c,
COMPLETION_DURATION,
"Processor completion duration",
);
assert_eq!(
count, 1,
"Processor should have 1 completion duration observation"
);
}
#[tokio::test]
async fn test_per_signal_item_counts() {
let harness = setup_test_manager_with_metrics();
for handles in harness
.node_metric_handles
.borrow_mut()
.iter_mut()
.flatten()
{
let _ = handles.input.take();
let _ = handles.input_completion.take();
let _ = handles.input_size.take();
handles.outputs.clear();
handles.output_completion.clear();
handles.output_size.clear();
}
let snapshots = run_and_collect(harness, |nodes| {
let mut pdata = build_3node_pdata_no_subscribers(nodes, false);
for frame in &mut pdata.frames {
frame.interests = Interests::NODE_ITEM_COUNTS;
frame.input_size = 0;
frame.output_size = 0;
}
vec![PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata),
}]
})
.await;
assert!(!snapshots.contains_key(&MetricLabel::RecvOutput));
let recv_p = &snapshots[&MetricLabel::RecvOutputItems];
assert_u64(recv_p, ITEMS, 10, "Receiver output logs");
let proc_c = &snapshots[&MetricLabel::ProcInputItems];
assert_u64(proc_c, ITEMS, 10, "Processor input logs");
let proc_p = &snapshots[&MetricLabel::ProcOutputItems];
assert_u64(proc_p, ITEMS, 7, "Processor output logs");
let exp = &snapshots[&MetricLabel::ExpInputItems];
assert_u64(exp, ITEMS, 7, "Exporter input logs");
}
#[test]
fn node_metrics_report_at_basic_runtime_level() {
let mut harness = setup_test_manager_with_metrics();
harness.manager.telemetry.runtime_metrics = MetricLevel::Basic;
let processor_id = harness.nodes[1].index;
{
let mut handles = harness.node_metric_handles.borrow_mut();
let input = handles[processor_id]
.as_mut()
.and_then(|handles| handles.input.as_mut())
.expect("processor input metrics");
input
.with(SignalOutcomeAttributes {
signal: SignalType::Logs,
outcome: Outcome::Success,
})
.messages
.inc();
}
let mut pipeline_metrics_monitor = None;
harness
.manager
.handle_due_events(Instant::now(), &mut pipeline_metrics_monitor);
let snapshots = collect_snapshots(&harness.snapshot_rx, &harness.key_labels);
let processor_input = &snapshots[&MetricLabel::ProcInput];
assert_u64(
processor_input,
INPUT_MESSAGES,
1,
"basic-level node opt-in should be reported",
);
}
#[tokio::test]
async fn test_per_signal_payload_size() {
let harness = setup_test_manager_with_metrics();
for handles in harness
.node_metric_handles
.borrow_mut()
.iter_mut()
.flatten()
{
let _ = handles.input.take();
let _ = handles.input_completion.take();
let _ = handles.input_items.take();
handles.outputs.clear();
handles.output_completion.clear();
handles.output_items.clear();
}
let snapshots = run_and_collect(harness, |nodes| {
let mut pdata = build_3node_pdata_no_subscribers(nodes, false);
for frame in &mut pdata.frames {
frame.interests = Interests::NODE_SIZE;
frame.input_items = 0;
frame.output_items = 0;
}
vec![PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata),
}]
})
.await;
assert!(!snapshots.contains_key(&MetricLabel::RecvOutput));
let recv_p = &snapshots[&MetricLabel::RecvOutputSize];
assert_u64(recv_p, SIZE, 100, "Receiver output size");
let proc_c = &snapshots[&MetricLabel::ProcInputSize];
assert_u64(proc_c, SIZE, 100, "Processor input size");
let proc_p = &snapshots[&MetricLabel::ProcOutputSize];
assert_u64(proc_p, SIZE, 70, "Processor output size");
let exp = &snapshots[&MetricLabel::ExpInputSize];
assert_u64(exp, SIZE, 70, "Exporter input size");
}
#[tokio::test]
async fn node_metrics_omit_items_without_optin() {
let harness = setup_test_manager_with_metrics();
for handles in harness
.node_metric_handles
.borrow_mut()
.iter_mut()
.flatten()
{
let _ = handles.input_items.take();
handles.output_items.clear();
}
let snapshots = run_and_collect(harness, |nodes| {
let pdata = build_3node_pdata_no_subscribers(nodes, false);
vec![PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata),
}]
})
.await;
assert!(
snapshots.contains_key(&MetricLabel::RecvOutput),
"Request metrics should remain enabled"
);
assert!(
!snapshots.contains_key(&MetricLabel::RecvOutputItems)
&& !snapshots.contains_key(&MetricLabel::ProcInputItems)
&& !snapshots.contains_key(&MetricLabel::ProcOutputItems)
&& !snapshots.contains_key(&MetricLabel::ExpInputItems),
"Item metrics should be absent when item counting is disabled"
);
}
#[tokio::test]
async fn node_metrics_omit_size_without_optin() {
let harness = setup_test_manager_with_metrics();
for handles in harness
.node_metric_handles
.borrow_mut()
.iter_mut()
.flatten()
{
let _ = handles.input_size.take();
handles.output_size.clear();
}
let snapshots = run_and_collect(harness, |nodes| {
let pdata = build_3node_pdata_no_subscribers(nodes, false);
vec![PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata),
}]
})
.await;
assert!(
snapshots.contains_key(&MetricLabel::RecvOutput),
"Request metrics should remain enabled"
);
assert!(
!snapshots.contains_key(&MetricLabel::RecvOutputSize)
&& !snapshots.contains_key(&MetricLabel::ProcInputSize)
&& !snapshots.contains_key(&MetricLabel::ProcOutputSize)
&& !snapshots.contains_key(&MetricLabel::ExpInputSize),
"Size metrics should be absent when payload sizing is disabled"
);
}
#[test]
fn node_metrics_export_signal_and_outcome_attributes() {
let mut harness = setup_test_manager_with_metrics();
let processor_id = harness.nodes[1].index;
let interests = Interests::NODE_INPUT_METRICS | Interests::NODE_OUTPUT_METRICS;
let route = RouteData::default();
let (_completion_tx, completion_rx) = pipeline_completion_msg_channel::<TestPData>(1);
let mut dispatcher = PipelineCompletionMsgDispatcher::new(
harness.pipeline_context.clone(),
completion_rx,
ControlSenders::new(),
harness.node_metric_handles.clone(),
harness.metrics_reporter.clone(),
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
harness.telemetry_policy.clone(),
TerminalMetricsDeadline::default(),
);
dispatcher.record_frame_metrics(
processor_id,
interests,
&route,
Some(SignalType::Logs),
4,
5,
0,
0,
RequestOutcome::Success,
0,
);
dispatcher.record_frame_metrics(
processor_id,
interests,
&route,
Some(SignalType::Traces),
6,
7,
0,
0,
RequestOutcome::Failure,
0,
);
report_node_metrics_with_handles(
&harness.node_metric_handles,
&mut harness.metrics_reporter,
)
.expect("node metric reporting should succeed");
let mut request_snapshots = Vec::new();
let mut item_snapshots = Vec::new();
while let Ok(snapshot) = harness.snapshot_rx.try_recv() {
let Some(&label) = harness.key_labels.get(&snapshot.key()) else {
continue;
};
let snapshot = (
label,
snapshot.measurement_attribute_value("signal"),
snapshot.measurement_attribute_value("outcome"),
snapshot.get_metrics().to_vec(),
);
match label {
MetricLabel::ProcInput | MetricLabel::ProcOutput => {
request_snapshots.push(snapshot);
}
MetricLabel::ProcInputItems | MetricLabel::ProcOutputItems => {
item_snapshots.push(snapshot);
}
_ => {}
}
}
assert_eq!(
request_snapshots.len(),
4,
"two signal/outcome request buckets per direction"
);
assert_eq!(
item_snapshots.len(),
4,
"two signal/outcome item buckets per direction"
);
let find =
|snapshots: &Vec<(MetricLabel, Option<&str>, Option<&str>, Vec<MetricValue>)>,
label,
signal,
outcome| {
snapshots
.iter()
.find(|(actual_label, actual_signal, actual_outcome, _)| {
*actual_label == label
&& *actual_signal == Some(signal)
&& *actual_outcome == Some(outcome)
})
.map(|(_, _, _, values)| values.clone())
.expect("expected datapoint attribute bucket")
};
let input_logs = find(
&request_snapshots,
MetricLabel::ProcInput,
"logs",
"success",
);
assert_u64(
&input_logs,
INPUT_MESSAGES,
1,
"logs success input messages",
);
let input_traces = find(
&request_snapshots,
MetricLabel::ProcInput,
"traces",
"failure",
);
assert_u64(
&input_traces,
INPUT_MESSAGES,
1,
"traces failure input messages",
);
let output_logs = find(
&request_snapshots,
MetricLabel::ProcOutput,
"logs",
"success",
);
assert_u64(
&output_logs,
OUTPUT_MESSAGES,
1,
"logs success output messages",
);
let output_traces = find(
&request_snapshots,
MetricLabel::ProcOutput,
"traces",
"failure",
);
assert_u64(
&output_traces,
OUTPUT_MESSAGES,
1,
"traces failure output messages",
);
let input_logs = find(
&item_snapshots,
MetricLabel::ProcInputItems,
"logs",
"success",
);
assert_u64(&input_logs, ITEMS, 5, "logs success input items");
let input_traces = find(
&item_snapshots,
MetricLabel::ProcInputItems,
"traces",
"failure",
);
assert_u64(&input_traces, ITEMS, 7, "traces failure input items");
let output_logs = find(
&item_snapshots,
MetricLabel::ProcOutputItems,
"logs",
"success",
);
assert_u64(&output_logs, ITEMS, 4, "logs success output items");
let output_traces = find(
&item_snapshots,
MetricLabel::ProcOutputItems,
"traces",
"failure",
);
assert_u64(&output_traces, ITEMS, 6, "traces failure output items");
}
#[tokio::test]
async fn test_completion_duration_not_recorded_without_timestamp() {
let harness = setup_test_manager_with_metrics();
let snapshots = run_and_collect(harness, |nodes| {
let pdata = build_3node_pdata_no_subscribers(nodes, false);
vec![PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata),
}]
})
.await;
assert!(!snapshots.contains_key(&MetricLabel::RecvCompletion));
}
#[tokio::test]
async fn test_ack_lifecycle_no_duration_without_timestamp() {
let harness = setup_test_manager_with_metrics();
let snapshots = run_and_collect(harness, |nodes| {
let pdata = build_3node_pdata(nodes, false);
vec![PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata),
}]
})
.await;
assert!(!snapshots.contains_key(&MetricLabel::ExpCompletion));
assert!(!snapshots.contains_key(&MetricLabel::ProcCompletion));
}
#[tokio::test]
async fn test_multiple_acks_lifecycle_accumulate() {
let harness = setup_test_manager_with_metrics();
let snapshots = run_and_collect(harness, |nodes| {
(0..3)
.map(|_| {
let pdata = build_3node_pdata(nodes, false);
PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata),
}
})
.collect()
})
.await;
let exp = &snapshots[&MetricLabel::ExpInput];
assert_u64(
exp,
INPUT_MESSAGES,
3,
"Exporter 3 input messages after 3 acks",
);
let proc_c = &snapshots[&MetricLabel::ProcInput];
assert_u64(proc_c, INPUT_MESSAGES, 3, "Processor 3 input messages");
let proc_p = &snapshots[&MetricLabel::ProcOutput];
assert_u64(proc_p, OUTPUT_MESSAGES, 3, "Processor 3 output messages");
}
#[tokio::test]
async fn test_mixed_ack_nack_lifecycle() {
let harness = setup_test_manager_with_metrics();
let snapshots = run_and_collect(harness, |nodes| {
vec![
PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(build_3node_pdata(nodes, false)),
},
PipelineCompletionMsg::DeliverNack {
nack: NackMsg::new("transient", build_3node_pdata(nodes, false)),
},
PipelineCompletionMsg::DeliverNack {
nack: NackMsg::new_permanent("refused", build_3node_pdata(nodes, false)),
},
]
})
.await;
let exp = &snapshots[&MetricLabel::ExpInput];
assert_u64(exp, INPUT_MESSAGES, 3, "Exporter input messages");
let proc_c = &snapshots[&MetricLabel::ProcInput];
assert_u64(proc_c, INPUT_MESSAGES, 3, "Processor input messages");
let proc_p = &snapshots[&MetricLabel::ProcOutput];
assert_u64(proc_p, OUTPUT_MESSAGES, 3, "Processor output messages");
}
#[tokio::test]
async fn test_two_pass_unwind_receiver_completion_duration() {
let harness = setup_test_manager_with_metrics();
let snapshots = run_and_collect(harness, |nodes| {
let pdata_full = build_3node_pdata(nodes, true);
let mut ack1 = AckMsg::new(pdata_full);
ack1.unwind.return_time_ns = nanos_since_birth();
let mut pdata_recv_only = TestPData::new();
pdata_recv_only.signal = Some(SignalType::Logs);
pdata_recv_only.push_frame(Frame {
node_id: nodes[0].index,
interests: Interests::NODE_OUTPUT_METRICS | Interests::NODE_COMPLETION_DURATION,
route: RouteData {
calldata: Default::default(),
entry_time_ns: nanos_since_birth(),
output_port_index: 0,
},
output_items: 0,
input_items: 0,
output_size: 0,
input_size: 0,
});
let mut ack2 = AckMsg::new(pdata_recv_only);
ack2.unwind.return_time_ns = nanos_since_birth();
vec![
PipelineCompletionMsg::DeliverAck { ack: ack1 },
PipelineCompletionMsg::DeliverAck { ack: ack2 },
]
})
.await;
let exp = &snapshots[&MetricLabel::ExpInput];
assert_u64(exp, INPUT_MESSAGES, 1, "Exporter input messages");
let exp_completion = &snapshots[&MetricLabel::ExpCompletion];
let (count, _, min, _) = assert_dist_is_normal_histogram(
exp_completion,
COMPLETION_DURATION,
"Exporter completion duration",
);
assert_eq!(count, 1, "Exporter should have 1 completion duration");
assert!(min > 0.0, "Exporter completion duration > 0");
let proc_c = &snapshots[&MetricLabel::ProcInput];
assert_u64(proc_c, INPUT_MESSAGES, 1, "Processor input messages");
let proc_completion = &snapshots[&MetricLabel::ProcCompletion];
let (count, _, _, _) = assert_dist_is_normal_histogram(
proc_completion,
COMPLETION_DURATION,
"Processor completion duration",
);
assert_eq!(count, 1, "Processor should have 1 completion duration");
let proc_p = &snapshots[&MetricLabel::ProcOutput];
assert_u64(proc_p, OUTPUT_MESSAGES, 1, "Processor output messages");
let recv_p = &snapshots[&MetricLabel::RecvOutput];
assert_u64(recv_p, OUTPUT_MESSAGES, 1, "Receiver output messages");
let recv_completion = &snapshots[&MetricLabel::RecvCompletion];
let (count, _, min, _) = assert_dist_is_normal_histogram(
recv_completion,
COMPLETION_DURATION,
"Receiver completion duration",
);
assert_eq!(
count, 1,
"Receiver should have 1 completion duration observation from two-pass unwind"
);
assert!(min > 0.0, "Receiver completion duration should be > 0");
}
#[tokio::test]
async fn test_memory_pressure_updates_are_fanned_out_only_to_receivers() {
let local = LocalSet::new();
local
.run_until(async {
let MemoryPressureFanoutHarness {
manager,
_pipeline_tx,
memory_pressure_tx,
mut control_receivers,
nodes,
_guard: _,
} = setup_memory_pressure_fanout_harness::<String>(vec![
("receiver", NodeType::Receiver, 16),
("processor", NodeType::Processor, 16),
]);
let receiver = nodes[0].clone();
let processor = nodes[1].clone();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
memory_pressure_tx
.send(MemoryPressureChanged {
generation: 1,
level: crate::memory_limiter::MemoryPressureLevel::Hard,
retry_after_secs: 5,
usage_bytes: 123,
})
.expect("watch send should succeed");
let mut receiver_ctrl = control_receivers.remove(&receiver.index).unwrap();
let receiver_msg = timeout(Duration::from_millis(100), receiver_ctrl.recv())
.await
.expect("receiver should get memory pressure update")
.expect("receiver control channel should stay open");
assert!(matches!(
receiver_msg,
NodeControlMsg::MemoryPressureChanged {
update: MemoryPressureChanged {
generation: 1,
level: crate::memory_limiter::MemoryPressureLevel::Hard,
retry_after_secs: 5,
usage_bytes: 123,
}
}
));
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
assert!(
timeout(Duration::from_millis(50), processor_ctrl.recv())
.await
.is_err(),
"non-receiver nodes should not get memory pressure updates"
);
manager_handle.abort();
})
.await;
}
#[tokio::test]
async fn test_runtime_control_metrics_track_receiver_first_drain() {
let local = LocalSet::new();
local
.run_until(async {
let RuntimeControlTelemetryHarness {
manager,
pipeline_tx,
mut control_receivers,
nodes,
runtime_metrics_key,
snapshot_rx,
engine_rx,
_guard: _,
} = setup_runtime_control_telemetry_harness::<String>(
vec![
("receiver", NodeType::Receiver, 16),
("processor", NodeType::Processor, 16),
],
MetricLevel::Detailed,
);
let receiver = nodes[0].clone();
let processor = nodes[1].clone();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "test shutdown".to_owned(),
})
.await
.unwrap();
let mut receiver_ctrl = control_receivers.remove(&receiver.index).unwrap();
let drain_msg = timeout(Duration::from_millis(100), receiver_ctrl.recv())
.await
.expect("receiver should get DrainIngress")
.expect("receiver control channel should stay open");
assert!(matches!(drain_msg, NodeControlMsg::DrainIngress { .. }));
let shutdown_start_metrics = collect_metric_set_snapshots(
&snapshot_rx,
runtime_metrics_key,
&[
RUNTIME_DRAIN_ACTIVE,
RUNTIME_DRAIN_PENDING_RECEIVERS,
RUNTIME_PENDING_SENDS_BUFFERED,
RUNTIME_TIMERS_ACTIVE,
RUNTIME_TELEMETRY_TIMERS_ACTIVE,
],
)
.expect("runtime-control metrics should be exported");
assert_u64(
&shutdown_start_metrics,
RUNTIME_DRAIN_ACTIVE,
1,
"drain.active should latch on shutdown",
);
assert_u64(
&shutdown_start_metrics,
RUNTIME_DRAIN_PENDING_RECEIVERS,
1,
"drain.pending_receivers should reflect the single receiver",
);
assert_u64(
&shutdown_start_metrics,
RUNTIME_SHUTDOWN_RECEIVED,
1,
"shutdown.received should increment once per shutdown request",
);
assert_u64(
&shutdown_start_metrics,
RUNTIME_DRAIN_INGRESS_SENT,
1,
"drain_ingress.sent should increment once when ingress drain starts",
);
pipeline_tx
.send(RuntimeControlMsg::ReceiverDrained {
node_id: receiver.index,
})
.await
.unwrap();
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
let shutdown_msg = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("processor should get Shutdown after receivers drain")
.expect("processor control channel should stay open");
assert!(matches!(shutdown_msg, NodeControlMsg::Shutdown { .. }));
drop(pipeline_tx);
let manager_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(manager_result.is_ok(), "manager should stop cleanly");
let drain_finish_metrics = collect_metric_set_snapshots(
&snapshot_rx,
runtime_metrics_key,
&[
RUNTIME_DRAIN_ACTIVE,
RUNTIME_DRAIN_PENDING_RECEIVERS,
RUNTIME_PENDING_SENDS_BUFFERED,
RUNTIME_TIMERS_ACTIVE,
RUNTIME_TELEMETRY_TIMERS_ACTIVE,
],
)
.expect("runtime-control metrics should export receiver-drained transition");
assert_u64(
&drain_finish_metrics,
RUNTIME_RECEIVER_DRAINED_RECEIVED,
1,
"receiver_drained.received should increment once",
);
assert_u64(
&drain_finish_metrics,
RUNTIME_DOWNSTREAM_SHUTDOWN_SENT,
1,
"downstream_shutdown.sent should increment once receivers are drained",
);
let receiver_phase = assert_dist_is_mmsc(
&drain_finish_metrics,
RUNTIME_DRAIN_RECEIVER_PHASE_DURATION_NS,
"receiver-phase duration",
);
assert_eq!(
receiver_phase.count, 1,
"receiver phase duration should record once"
);
let total_drain = assert_dist_is_mmsc(
&drain_finish_metrics,
RUNTIME_DRAIN_TOTAL_DURATION_NS,
"total drain duration",
);
assert_eq!(
total_drain.count, 1,
"total drain duration should record once"
);
let event_types = collect_engine_event_types(&engine_rx);
assert!(matches!(
event_types.as_slice(),
[
EventType::Request(TelemetryRequestEvent::ShutdownRequested),
EventType::Success(TelemetrySuccessEvent::IngressDrainStarted),
EventType::Success(TelemetrySuccessEvent::ReceiversDrained),
EventType::Success(TelemetrySuccessEvent::DownstreamShutdownStarted),
]
));
})
.await;
}
#[tokio::test]
async fn test_runtime_control_metrics_record_forced_deadline() {
let local = LocalSet::new();
local
.run_until(async {
let RuntimeControlTelemetryHarness {
manager,
pipeline_tx,
mut control_receivers,
nodes,
runtime_metrics_key,
snapshot_rx,
engine_rx,
_guard: _,
} = setup_runtime_control_telemetry_harness::<String>(
vec![("receiver", NodeType::Receiver, 16)],
MetricLevel::Detailed,
);
let receiver = nodes[0].clone();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_millis(25),
reason: "deadline test".to_owned(),
})
.await
.unwrap();
let mut receiver_ctrl = control_receivers.remove(&receiver.index).unwrap();
let drain_msg = timeout(Duration::from_millis(100), receiver_ctrl.recv())
.await
.expect("receiver should get DrainIngress")
.expect("receiver control channel should stay open");
assert!(matches!(drain_msg, NodeControlMsg::DrainIngress { .. }));
let manager_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(manager_result.is_ok(), "manager should stop on deadline");
let forced_metrics = collect_metric_set_snapshots(
&snapshot_rx,
runtime_metrics_key,
&[
RUNTIME_DRAIN_ACTIVE,
RUNTIME_DRAIN_PENDING_RECEIVERS,
RUNTIME_PENDING_SENDS_BUFFERED,
RUNTIME_TIMERS_ACTIVE,
RUNTIME_TELEMETRY_TIMERS_ACTIVE,
],
)
.expect("runtime-control metrics should export forced-deadline snapshot");
assert_u64(
&forced_metrics,
RUNTIME_SHUTDOWN_DEADLINE_FORCED,
1,
"shutdown.deadline_forced should increment once",
);
let total_drain = assert_dist_is_mmsc(
&forced_metrics,
RUNTIME_DRAIN_TOTAL_DURATION_NS,
"forced drain duration",
);
assert_eq!(
total_drain.count, 1,
"forced drain should still record total duration"
);
let event_types = collect_engine_event_types(&engine_rx);
assert!(matches!(
event_types.as_slice(),
[
EventType::Request(TelemetryRequestEvent::ShutdownRequested),
EventType::Success(TelemetrySuccessEvent::IngressDrainStarted),
EventType::Error(TelemetryErrorEvent::DrainDeadlineReached),
]
));
})
.await;
}
#[tokio::test]
async fn test_no_receiver_shutdown_emits_downstream_event_immediately() {
let local = LocalSet::new();
local
.run_until(async {
let RuntimeControlTelemetryHarness {
manager,
pipeline_tx,
mut control_receivers,
nodes,
engine_rx,
..
} = setup_runtime_control_telemetry_harness::<String>(
vec![("processor", NodeType::Processor, 16)],
MetricLevel::Normal,
);
let processor = nodes[0].clone();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "no receiver".to_owned(),
})
.await
.unwrap();
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
let msg = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("processor should get immediate Shutdown")
.expect("processor control channel should stay open");
assert!(matches!(msg, NodeControlMsg::Shutdown { .. }));
drop(pipeline_tx);
let manager_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(manager_result.is_ok(), "manager should stop cleanly");
let event_types = collect_engine_event_types(&engine_rx);
assert!(matches!(
event_types.as_slice(),
[
EventType::Request(TelemetryRequestEvent::ShutdownRequested),
EventType::Success(TelemetrySuccessEvent::IngressDrainStarted),
EventType::Success(TelemetrySuccessEvent::DownstreamShutdownStarted),
]
));
})
.await;
}
#[tokio::test]
async fn test_runtime_control_metrics_track_due_work_dispatch() {
let local = LocalSet::new();
local
.run_until(async {
let RuntimeControlTelemetryHarness {
manager,
pipeline_tx,
mut control_receivers,
nodes,
runtime_metrics_key,
snapshot_rx,
_guard: _,
engine_rx: _,
} = setup_runtime_control_telemetry_harness::<String>(
vec![("processor", NodeType::Processor, 16)],
MetricLevel::Detailed,
);
let processor = nodes[0].clone();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::StartTimer {
node_id: processor.index,
duration: Duration::from_millis(5),
})
.await
.unwrap();
pipeline_tx
.send(RuntimeControlMsg::StartTelemetryTimer {
node_id: processor.index,
duration: Duration::from_millis(5),
})
.await
.unwrap();
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
let mut timer_tick = false;
let mut collect_telemetry = false;
let deadline = tokio::time::Instant::now() + Duration::from_millis(250);
while !(timer_tick && collect_telemetry) && tokio::time::Instant::now() < deadline {
match timeout(Duration::from_millis(100), processor_ctrl.recv()).await {
Ok(Ok(NodeControlMsg::TimerTick {})) => timer_tick = true,
Ok(Ok(NodeControlMsg::CollectTelemetry { .. })) => collect_telemetry = true,
Ok(Ok(_)) => {}
Ok(Err(_)) | Err(_) => break,
}
}
assert!(timer_tick, "due timer should reach the processor");
assert!(
collect_telemetry,
"due telemetry timer should reach the processor"
);
drop(pipeline_tx);
let manager_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(manager_result.is_ok(), "manager should stop cleanly");
let due_metrics = collect_metric_set_snapshots(
&snapshot_rx,
runtime_metrics_key,
&[
RUNTIME_DRAIN_ACTIVE,
RUNTIME_DRAIN_PENDING_RECEIVERS,
RUNTIME_PENDING_SENDS_BUFFERED,
RUNTIME_TIMERS_ACTIVE,
RUNTIME_TELEMETRY_TIMERS_ACTIVE,
],
)
.expect("runtime-control metrics should export due-work counters");
assert_u64(
&due_metrics,
RUNTIME_START_TIMER_RECEIVED,
1,
"start_timer.received should count timer requests",
);
assert_u64(
&due_metrics,
RUNTIME_START_TELEMETRY_TIMER_RECEIVED,
1,
"start_telemetry_timer.received should count telemetry timer requests",
);
assert_u64_gte(
&due_metrics,
RUNTIME_TIMER_TICK_SENT,
1,
"timer_tick.sent should count due timer dispatches",
);
assert_u64_gte(
&due_metrics,
RUNTIME_COLLECT_TELEMETRY_SENT,
1,
"collect_telemetry.sent should count due telemetry dispatches",
);
})
.await;
}
#[tokio::test]
async fn test_runtime_control_metrics_flush_on_interval_without_repeated_snapshots() {
let local = LocalSet::new();
local
.run_until(async {
let flush_interval = Duration::from_millis(200);
let RuntimeControlTelemetryHarness {
manager,
pipeline_tx,
nodes,
runtime_metrics_key,
snapshot_rx,
..
} = setup_runtime_control_telemetry_harness_with_flush_interval::<String>(
vec![("processor", NodeType::Processor, 16)],
MetricLevel::Normal,
flush_interval,
);
let processor = nodes[0].clone();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::StartTimer {
node_id: processor.index,
duration: Duration::from_secs(60),
})
.await
.unwrap();
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(
snapshot_rx.try_recv().is_err(),
"dirty runtime-control state should not flush before the configured interval"
);
tokio::time::sleep(flush_interval + Duration::from_millis(50)).await;
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
runtime_metrics_key,
&[
RUNTIME_DRAIN_ACTIVE,
RUNTIME_DRAIN_PENDING_RECEIVERS,
RUNTIME_PENDING_SENDS_BUFFERED,
RUNTIME_TIMERS_ACTIVE,
RUNTIME_TELEMETRY_TIMERS_ACTIVE,
],
)
.expect("dirty runtime-control state should flush on the configured interval");
assert_u64(
&metrics,
RUNTIME_START_TIMER_RECEIVED,
1,
"periodic flush should include the timer-start counter",
);
assert_u64(
&metrics,
RUNTIME_TIMERS_ACTIVE,
1,
"periodic flush should include the updated timer gauge",
);
tokio::time::sleep(flush_interval + Duration::from_millis(50)).await;
assert!(
collect_metric_set_snapshots(
&snapshot_rx,
runtime_metrics_key,
&[
RUNTIME_DRAIN_ACTIVE,
RUNTIME_DRAIN_PENDING_RECEIVERS,
RUNTIME_PENDING_SENDS_BUFFERED,
RUNTIME_TIMERS_ACTIVE,
RUNTIME_TELEMETRY_TIMERS_ACTIVE,
],
)
.is_none(),
"unchanged runtime-control state should not emit on later intervals"
);
drop(pipeline_tx);
let manager_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(manager_result.is_ok(), "manager should stop cleanly");
})
.await;
}
#[tokio::test]
async fn test_runtime_control_metrics_shutdown_phase_flushes_immediately() {
let local = LocalSet::new();
local
.run_until(async {
let RuntimeControlTelemetryHarness {
manager,
pipeline_tx,
mut control_receivers,
nodes,
runtime_metrics_key,
snapshot_rx,
..
} = setup_runtime_control_telemetry_harness_with_flush_interval::<String>(
vec![
("receiver", NodeType::Receiver, 16),
("processor", NodeType::Processor, 16),
],
MetricLevel::Normal,
Duration::from_secs(1),
);
let receiver = nodes[0].clone();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(5),
reason: "immediate-flush".to_owned(),
})
.await
.unwrap();
let mut receiver_ctrl = control_receivers.remove(&receiver.index).unwrap();
let drain_msg = timeout(Duration::from_millis(100), receiver_ctrl.recv())
.await
.expect("receiver should get DrainIngress")
.expect("receiver control channel should stay open");
assert!(matches!(drain_msg, NodeControlMsg::DrainIngress { .. }));
tokio::time::sleep(Duration::from_millis(50)).await;
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
runtime_metrics_key,
&[
RUNTIME_DRAIN_ACTIVE,
RUNTIME_DRAIN_PENDING_RECEIVERS,
RUNTIME_PENDING_SENDS_BUFFERED,
RUNTIME_TIMERS_ACTIVE,
RUNTIME_TELEMETRY_TIMERS_ACTIVE,
],
)
.expect("shutdown latch should flush without waiting for the full interval");
assert_u64(
&metrics,
RUNTIME_SHUTDOWN_RECEIVED,
1,
"shutdown transition should flush immediately",
);
assert_u64(
&metrics,
RUNTIME_DRAIN_INGRESS_SENT,
1,
"ingress-drain transition should flush immediately",
);
pipeline_tx
.send(RuntimeControlMsg::ReceiverDrained {
node_id: receiver.index,
})
.await
.unwrap();
drop(pipeline_tx);
let manager_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(manager_result.is_ok(), "manager should stop cleanly");
})
.await;
}
#[tokio::test]
async fn test_runtime_control_metrics_gated_off_at_none() {
let local = LocalSet::new();
local
.run_until(async {
let RuntimeControlTelemetryHarness {
manager,
pipeline_tx,
runtime_metrics_key,
snapshot_rx,
engine_rx: _engine_rx,
..
} = setup_runtime_control_telemetry_harness::<String>(
vec![("processor", NodeType::Processor, 16)],
MetricLevel::None,
);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "none".to_owned(),
})
.await
.unwrap();
drop(pipeline_tx);
let manager_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(manager_result.is_ok(), "manager should stop cleanly");
assert!(
runtime_metrics_key.is_none(),
"runtime-control metrics should not be registered at level none"
);
assert!(
collect_metric_set_snapshots(
&snapshot_rx,
runtime_metrics_key,
&[RUNTIME_DRAIN_ACTIVE]
)
.is_none(),
"no runtime-control snapshots should be emitted at level none"
);
})
.await;
}
#[tokio::test]
async fn test_runtime_control_metrics_basic_only_exports_gauges() {
let local = LocalSet::new();
local
.run_until(async {
let RuntimeControlTelemetryHarness {
manager,
pipeline_tx,
mut control_receivers,
nodes,
runtime_metrics_key,
snapshot_rx,
engine_rx: _engine_rx,
..
} = setup_runtime_control_telemetry_harness::<String>(
vec![
("receiver", NodeType::Receiver, 16),
("processor", NodeType::Processor, 16),
],
MetricLevel::Basic,
);
let receiver = nodes[0].clone();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "basic".to_owned(),
})
.await
.unwrap();
let mut receiver_ctrl = control_receivers.remove(&receiver.index).unwrap();
let _ = timeout(Duration::from_millis(100), receiver_ctrl.recv())
.await
.expect("receiver should get DrainIngress")
.expect("receiver control channel should stay open");
drop(pipeline_tx);
let manager_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(manager_result.is_ok(), "manager should stop cleanly");
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
runtime_metrics_key,
&[
RUNTIME_DRAIN_ACTIVE,
RUNTIME_DRAIN_PENDING_RECEIVERS,
RUNTIME_PENDING_SENDS_BUFFERED,
RUNTIME_TIMERS_ACTIVE,
RUNTIME_TELEMETRY_TIMERS_ACTIVE,
],
)
.expect("basic runtime-control metrics should be exported");
assert_u64(
&metrics,
RUNTIME_DRAIN_ACTIVE,
1,
"basic should export gauges",
);
assert_u64(
&metrics,
RUNTIME_DRAIN_PENDING_RECEIVERS,
1,
"basic should export pending receiver gauge",
);
assert_u64(
&metrics,
RUNTIME_SHUTDOWN_RECEIVED,
0,
"basic should suppress normal counters",
);
let receiver_phase = assert_dist_is_mmsc(
&metrics,
RUNTIME_DRAIN_RECEIVER_PHASE_DURATION_NS,
"basic receiver-phase duration",
);
assert_eq!(
receiver_phase.count, 0,
"basic should suppress detailed durations"
);
})
.await;
}
#[tokio::test]
async fn test_runtime_control_metrics_normal_exports_counters_without_durations() {
let local = LocalSet::new();
local
.run_until(async {
let RuntimeControlTelemetryHarness {
manager,
pipeline_tx,
mut control_receivers,
nodes,
runtime_metrics_key,
snapshot_rx,
engine_rx: _engine_rx,
..
} = setup_runtime_control_telemetry_harness::<String>(
vec![
("receiver", NodeType::Receiver, 16),
("processor", NodeType::Processor, 16),
],
MetricLevel::Normal,
);
let receiver = nodes[0].clone();
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
pipeline_tx
.send(RuntimeControlMsg::Shutdown {
deadline: Instant::now() + Duration::from_secs(1),
reason: "normal".to_owned(),
})
.await
.unwrap();
let mut receiver_ctrl = control_receivers.remove(&receiver.index).unwrap();
let _ = timeout(Duration::from_millis(100), receiver_ctrl.recv())
.await
.expect("receiver should get DrainIngress")
.expect("receiver control channel should stay open");
pipeline_tx
.send(RuntimeControlMsg::ReceiverDrained {
node_id: receiver.index,
})
.await
.unwrap();
drop(pipeline_tx);
let manager_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(manager_result.is_ok(), "manager should stop cleanly");
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
runtime_metrics_key,
&[
RUNTIME_DRAIN_ACTIVE,
RUNTIME_DRAIN_PENDING_RECEIVERS,
RUNTIME_PENDING_SENDS_BUFFERED,
RUNTIME_TIMERS_ACTIVE,
RUNTIME_TELEMETRY_TIMERS_ACTIVE,
],
)
.expect("normal runtime-control metrics should be exported");
assert_u64(
&metrics,
RUNTIME_SHUTDOWN_RECEIVED,
1,
"normal should export shutdown counter",
);
assert_u64(
&metrics,
RUNTIME_RECEIVER_DRAINED_RECEIVED,
1,
"normal should export receiver_drained counter",
);
assert_u64(
&metrics,
RUNTIME_DOWNSTREAM_SHUTDOWN_SENT,
1,
"normal should export downstream shutdown counter",
);
let total_drain = assert_dist_is_mmsc(
&metrics,
RUNTIME_DRAIN_TOTAL_DURATION_NS,
"normal total drain duration",
);
assert_eq!(
total_drain.count, 0,
"normal should suppress detailed durations"
);
})
.await;
}
#[tokio::test]
async fn test_completion_metrics_track_ack_delivery_and_unwind_depth() {
let local = LocalSet::new();
local
.run_until(async {
let CompletionTelemetryHarness {
dispatcher,
completion_tx,
mut control_receivers,
nodes,
completion_metrics_key,
snapshot_rx,
_guard: _,
} = setup_completion_telemetry_harness(MetricLevel::Detailed);
let processor = nodes[1].clone();
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
completion_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(build_3node_pdata(&nodes, false)),
})
.await
.unwrap();
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
let ack_msg = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("processor should receive Ack")
.expect("processor control channel should stay open");
assert!(matches!(ack_msg, NodeControlMsg::Ack(_)));
drop(completion_tx);
let dispatcher_result =
timeout(Duration::from_millis(200), dispatcher_handle).await;
assert!(dispatcher_result.is_ok(), "dispatcher should stop cleanly");
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
completion_metrics_key,
&[COMPLETION_PENDING_SENDS_BUFFERED],
)
.expect("completion metrics should be exported");
assert_u64(
&metrics,
COMPLETION_DELIVER_ACK_RECEIVED,
1,
"deliver_ack.received should increment once",
);
assert_u64(
&metrics,
COMPLETION_ACK_ATTEMPTED,
1,
"ack.attempted should increment once",
);
assert_u64(
&metrics,
COMPLETION_ACK_DELIVERED,
1,
"ack.delivered should increment once for immediate delivery",
);
assert_u64(
&metrics,
COMPLETION_ACK_DROPPED_NO_INTEREST,
0,
"ack.dropped_no_interest should stay at zero for interested unwind",
);
let unwind = assert_dist_is_mmsc(
&metrics,
COMPLETION_UNWIND_DEPTH,
"completion unwind depth",
);
assert_eq!(unwind.count, 1, "unwind depth should record one Ack unwind");
assert_eq!(
unwind.min, 2.0,
"Ack should unwind exporter+processor frames"
);
assert_eq!(unwind.max, 2.0, "single unwind depth should be exact");
})
.await;
}
#[tokio::test]
async fn test_completion_metrics_track_nack_delivery() {
let local = LocalSet::new();
local
.run_until(async {
let CompletionTelemetryHarness {
dispatcher,
completion_tx,
mut control_receivers,
nodes,
completion_metrics_key,
snapshot_rx,
_guard: _,
} = setup_completion_telemetry_harness(MetricLevel::Normal);
let processor = nodes[1].clone();
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
completion_tx
.send(PipelineCompletionMsg::DeliverNack {
nack: NackMsg::new("transient", build_3node_pdata(&nodes, false)),
})
.await
.unwrap();
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
let nack_msg = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("processor should receive Nack")
.expect("processor control channel should stay open");
assert!(matches!(nack_msg, NodeControlMsg::Nack(_)));
drop(completion_tx);
let dispatcher_result =
timeout(Duration::from_millis(200), dispatcher_handle).await;
assert!(dispatcher_result.is_ok(), "dispatcher should stop cleanly");
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
completion_metrics_key,
&[COMPLETION_PENDING_SENDS_BUFFERED],
)
.expect("completion metrics should be exported");
assert_u64(
&metrics,
COMPLETION_DELIVER_NACK_RECEIVED,
1,
"deliver_nack.received should increment once",
);
assert_u64(
&metrics,
COMPLETION_NACK_ATTEMPTED,
1,
"nack.attempted should increment once",
);
assert_u64(
&metrics,
COMPLETION_NACK_DELIVERED,
1,
"nack.delivered should increment once for immediate delivery",
);
assert_u64(
&metrics,
COMPLETION_NACK_DROPPED_NO_INTEREST,
0,
"nack.dropped_no_interest should stay at zero for interested unwind",
);
let unwind = assert_dist_is_mmsc(
&metrics,
COMPLETION_UNWIND_DEPTH,
"normal completion unwind depth",
);
assert_eq!(
unwind.count, 0,
"normal should suppress detailed unwind depth"
);
})
.await;
}
#[tokio::test]
async fn test_completion_metrics_track_buffered_retry_delivery() {
let local = LocalSet::new();
local
.run_until(async {
let CompletionTelemetryHarness {
dispatcher,
completion_tx,
mut control_receivers,
nodes,
completion_metrics_key,
snapshot_rx,
_guard: _,
} = setup_completion_telemetry_harness_with_capacity(MetricLevel::Normal, 1);
let processor = nodes[1].clone();
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
completion_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(build_3node_pdata(&nodes, false)),
})
.await
.unwrap();
completion_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(build_3node_pdata(&nodes, false)),
})
.await
.unwrap();
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
let first_ack = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("processor should receive the first Ack")
.expect("processor control channel should stay open");
assert!(matches!(first_ack, NodeControlMsg::Ack(_)));
let second_ack = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("buffered Ack should be retried and eventually delivered")
.expect("processor control channel should stay open");
assert!(matches!(second_ack, NodeControlMsg::Ack(_)));
drop(completion_tx);
let dispatcher_result =
timeout(Duration::from_millis(200), dispatcher_handle).await;
assert!(dispatcher_result.is_ok(), "dispatcher should stop cleanly");
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
completion_metrics_key,
&[COMPLETION_PENDING_SENDS_BUFFERED],
)
.expect("completion metrics should be exported");
assert_u64(
&metrics,
COMPLETION_DELIVER_ACK_RECEIVED,
2,
"two Ack completions should be received",
);
assert_u64(
&metrics,
COMPLETION_ACK_ATTEMPTED,
2,
"both Ack completions should be counted as attempted",
);
assert_u64(
&metrics,
COMPLETION_ACK_DELIVERED,
2,
"both Ack completions should be counted as delivered once the retry succeeds",
);
assert_u64(
&metrics,
COMPLETION_PENDING_SENDS_BUFFERED,
0,
"pending_sends.buffered should return to zero after retry delivery",
);
})
.await;
}
#[tokio::test]
async fn test_completion_metrics_track_dropped_no_interest() {
let local = LocalSet::new();
local
.run_until(async {
let CompletionTelemetryHarness {
dispatcher,
completion_tx,
completion_metrics_key,
snapshot_rx,
nodes,
_guard: _,
control_receivers: _,
} = setup_completion_telemetry_harness(MetricLevel::Detailed);
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
completion_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(build_3node_pdata_no_subscribers(&nodes, false)),
})
.await
.unwrap();
drop(completion_tx);
let dispatcher_result =
timeout(Duration::from_millis(200), dispatcher_handle).await;
assert!(dispatcher_result.is_ok(), "dispatcher should stop cleanly");
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
completion_metrics_key,
&[COMPLETION_PENDING_SENDS_BUFFERED],
)
.expect("completion metrics should be exported");
assert_u64(
&metrics,
COMPLETION_DELIVER_ACK_RECEIVED,
1,
"deliver_ack.received should increment once",
);
assert_u64(
&metrics,
COMPLETION_ACK_ATTEMPTED,
0,
"ack.attempted should stay at zero when no frame is interested",
);
assert_u64(
&metrics,
COMPLETION_ACK_DELIVERED,
0,
"ack.delivered should stay at zero when no frame is interested",
);
assert_u64(
&metrics,
COMPLETION_ACK_DROPPED_NO_INTEREST,
1,
"ack.dropped_no_interest should count uninterested unwinds",
);
let unwind = assert_dist_is_mmsc(
&metrics,
COMPLETION_UNWIND_DEPTH,
"dropped-no-interest unwind depth",
);
assert_eq!(
unwind.count, 1,
"unwind depth should record dropped completion depth"
);
assert_eq!(unwind.min, 3.0, "all three frames should be popped");
assert_eq!(
unwind.max, 3.0,
"single dropped unwind depth should be exact"
);
})
.await;
}
#[tokio::test]
async fn test_completion_metrics_gated_off_at_none() {
let local = LocalSet::new();
local
.run_until(async {
let CompletionTelemetryHarness {
dispatcher,
completion_tx,
completion_metrics_key,
snapshot_rx,
..
} = setup_completion_telemetry_harness(MetricLevel::None);
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
completion_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(TestPData::new()),
})
.await
.unwrap();
drop(completion_tx);
let dispatcher_result =
timeout(Duration::from_millis(200), dispatcher_handle).await;
assert!(dispatcher_result.is_ok(), "dispatcher should stop cleanly");
assert!(
completion_metrics_key.is_none(),
"completion metrics should not be registered at level none"
);
assert!(
collect_metric_set_snapshots(
&snapshot_rx,
completion_metrics_key,
&[COMPLETION_PENDING_SENDS_BUFFERED],
)
.is_none(),
"no completion snapshots should be emitted at level none"
);
})
.await;
}
#[tokio::test]
async fn test_completion_metrics_basic_only_exports_gauge() {
let local = LocalSet::new();
local
.run_until(async {
let CompletionTelemetryHarness {
dispatcher,
completion_tx,
mut control_receivers,
nodes,
completion_metrics_key,
snapshot_rx,
_guard: _,
} = setup_completion_telemetry_harness(MetricLevel::Basic);
let processor = nodes[1].clone();
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
completion_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(build_3node_pdata(&nodes, false)),
})
.await
.unwrap();
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
let ack_msg = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("processor should receive Ack")
.expect("processor control channel should stay open");
assert!(matches!(ack_msg, NodeControlMsg::Ack(_)));
drop(completion_tx);
let dispatcher_result =
timeout(Duration::from_millis(200), dispatcher_handle).await;
assert!(dispatcher_result.is_ok(), "dispatcher should stop cleanly");
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
completion_metrics_key,
&[COMPLETION_PENDING_SENDS_BUFFERED],
)
.expect("basic completion metrics should be exported");
assert_u64(
&metrics,
COMPLETION_PENDING_SENDS_BUFFERED,
0,
"basic should export the backlog gauge",
);
assert_u64(
&metrics,
COMPLETION_DELIVER_ACK_RECEIVED,
0,
"basic should suppress completion counters",
);
let unwind = assert_dist_is_mmsc(
&metrics,
COMPLETION_UNWIND_DEPTH,
"basic completion unwind depth",
);
assert_eq!(
unwind.count, 0,
"basic should suppress unwind-depth details"
);
})
.await;
}
#[tokio::test]
async fn test_completion_metrics_normal_exports_counters_without_unwind_depth() {
let local = LocalSet::new();
local
.run_until(async {
let CompletionTelemetryHarness {
dispatcher,
completion_tx,
mut control_receivers,
nodes,
completion_metrics_key,
snapshot_rx,
_guard: _,
} = setup_completion_telemetry_harness(MetricLevel::Normal);
let processor = nodes[1].clone();
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
completion_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(build_3node_pdata(&nodes, false)),
})
.await
.unwrap();
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
let ack_msg = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("processor should receive Ack")
.expect("processor control channel should stay open");
assert!(matches!(ack_msg, NodeControlMsg::Ack(_)));
drop(completion_tx);
let dispatcher_result =
timeout(Duration::from_millis(200), dispatcher_handle).await;
assert!(dispatcher_result.is_ok(), "dispatcher should stop cleanly");
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
completion_metrics_key,
&[COMPLETION_PENDING_SENDS_BUFFERED],
)
.expect("normal completion metrics should be exported");
assert_u64(
&metrics,
COMPLETION_DELIVER_ACK_RECEIVED,
1,
"normal should export completion receive counters",
);
assert_u64(
&metrics,
COMPLETION_ACK_ATTEMPTED,
1,
"normal should export completion attempted counters",
);
assert_u64(
&metrics,
COMPLETION_ACK_DELIVERED,
1,
"normal should export completion delivered counters",
);
let unwind = assert_dist_is_mmsc(
&metrics,
COMPLETION_UNWIND_DEPTH,
"normal completion unwind depth",
);
assert_eq!(
unwind.count, 0,
"normal should suppress unwind-depth details"
);
})
.await;
}
#[tokio::test]
async fn test_completion_metrics_flush_on_interval_without_repeated_snapshots() {
let local = LocalSet::new();
local
.run_until(async {
let flush_interval = Duration::from_millis(200);
let CompletionTelemetryHarness {
dispatcher,
completion_tx,
mut control_receivers,
nodes,
completion_metrics_key,
snapshot_rx,
..
} = setup_completion_telemetry_harness_with_capacity_and_flush_interval(
MetricLevel::Normal,
16,
flush_interval,
);
let processor = nodes[1].clone();
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
completion_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(build_3node_pdata(&nodes, false)),
})
.await
.unwrap();
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
let ack_msg = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("processor should receive Ack")
.expect("processor control channel should stay open");
assert!(matches!(ack_msg, NodeControlMsg::Ack(_)));
tokio::time::sleep(Duration::from_millis(50)).await;
assert!(
snapshot_rx.try_recv().is_err(),
"completion metrics should not flush before the configured interval"
);
tokio::time::sleep(flush_interval + Duration::from_millis(50)).await;
let metrics = collect_metric_set_snapshots(
&snapshot_rx,
completion_metrics_key,
&[COMPLETION_PENDING_SENDS_BUFFERED],
)
.expect("completion metrics should flush on the configured interval");
assert_u64(
&metrics,
COMPLETION_DELIVER_ACK_RECEIVED,
1,
"deliver_ack.received should flush on interval",
);
assert_u64(
&metrics,
COMPLETION_ACK_ATTEMPTED,
1,
"ack.attempted should flush on interval",
);
assert_u64(
&metrics,
COMPLETION_ACK_DELIVERED,
1,
"ack.delivered should flush on interval",
);
tokio::time::sleep(flush_interval + Duration::from_millis(50)).await;
assert!(
collect_metric_set_snapshots(
&snapshot_rx,
completion_metrics_key,
&[COMPLETION_PENDING_SENDS_BUFFERED],
)
.is_none(),
"unchanged completion state should not emit on later intervals"
);
drop(completion_tx);
let dispatcher_result =
timeout(Duration::from_millis(200), dispatcher_handle).await;
assert!(dispatcher_result.is_ok(), "dispatcher should stop cleanly");
})
.await;
}
#[tokio::test]
async fn test_completion_metrics_retry_after_deferred_flush() {
let local = LocalSet::new();
local
.run_until(async {
let flush_interval = Duration::from_millis(60);
let CompletionTelemetryHarness {
dispatcher,
completion_tx,
mut control_receivers,
nodes,
completion_metrics_key,
snapshot_rx,
..
} = setup_completion_telemetry_harness_with_options(
MetricLevel::Normal,
16,
flush_interval,
1,
);
let processor = nodes[1].clone();
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
completion_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(build_3node_pdata(&nodes, false)),
})
.await
.unwrap();
let mut processor_ctrl = control_receivers.remove(&processor.index).unwrap();
let first_ack = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("processor should receive first Ack")
.expect("processor control channel should stay open");
assert!(matches!(first_ack, NodeControlMsg::Ack(_)));
tokio::time::sleep(flush_interval + Duration::from_millis(30)).await;
assert_eq!(
snapshot_rx.len(),
1,
"first periodic flush should occupy the single-slot reporter channel"
);
completion_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(build_3node_pdata(&nodes, false)),
})
.await
.unwrap();
let second_ack = timeout(Duration::from_millis(100), processor_ctrl.recv())
.await
.expect("processor should receive second Ack")
.expect("processor control channel should stay open");
assert!(matches!(second_ack, NodeControlMsg::Ack(_)));
tokio::time::sleep(flush_interval + Duration::from_millis(30)).await;
assert_eq!(
snapshot_rx.len(),
1,
"deferred flush should keep the dirty snapshot pending instead of dropping it"
);
let _first_snapshot = snapshot_rx
.try_recv()
.expect("first snapshot should be buffered");
tokio::time::sleep(flush_interval + Duration::from_millis(30)).await;
let retried_metrics = collect_metric_set_snapshots(
&snapshot_rx,
completion_metrics_key,
&[COMPLETION_PENDING_SENDS_BUFFERED],
)
.expect("dirty completion state should retry on the next interval");
assert_u64(
&retried_metrics,
COMPLETION_DELIVER_ACK_RECEIVED,
1,
"retried snapshot should preserve the second ack delta",
);
assert_u64(
&retried_metrics,
COMPLETION_ACK_ATTEMPTED,
1,
"retried snapshot should preserve the second attempted delta",
);
assert_u64(
&retried_metrics,
COMPLETION_ACK_DELIVERED,
1,
"retried snapshot should preserve the second delivered delta",
);
drop(completion_tx);
let dispatcher_result =
timeout(Duration::from_millis(200), dispatcher_handle).await;
assert!(dispatcher_result.is_ok(), "dispatcher should stop cleanly");
})
.await;
}
#[tokio::test]
async fn test_return_lane_progress_while_runtime_ctrl_lane_is_busy() {
use crate::control::AckMsg;
fn pdata_for_node(node_id: usize) -> TestPData {
let mut pdata = TestPData::new();
pdata.push_frame(Frame {
node_id,
interests: Interests::ACKS,
route: RouteData {
calldata: Default::default(),
entry_time_ns: 0,
output_port_index: 0,
},
output_items: 0,
input_items: 0,
output_size: 0,
input_size: 0,
});
pdata
}
let local = LocalSet::new();
local
.run_until(async {
let (
manager,
pipeline_tx,
control_senders,
mut control_receivers,
nodes,
_pipeline_entity_guard,
) = setup_test_manager_with_capacities::<TestPData>(128, 10);
let noisy_node = nodes[0].clone();
let target = nodes[1].clone();
for _ in 0..96 {
pipeline_tx
.send(RuntimeControlMsg::StartTimer {
node_id: noisy_node.index,
duration: Duration::from_secs(60),
})
.await
.unwrap();
}
let (return_tx, return_rx) = pipeline_completion_msg_channel(8);
return_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata_for_node(target.index)),
})
.await
.unwrap();
let (dispatcher_context, dispatcher_guard) = create_test_pipeline_context();
let dispatcher = PipelineCompletionMsgDispatcher::new(
dispatcher_context,
return_rx,
control_senders,
empty_node_metric_handles(),
MetricsReporter::create_new_and_receiver(16).1,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
TelemetryPolicy::default(),
TerminalMetricsDeadline::default(),
);
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
let mut receiver = control_receivers.remove(&target.index).unwrap();
let msg = timeout(Duration::from_millis(500), receiver.recv())
.await
.expect("Ack should make progress while the runtime control lane is busy")
.expect("target control channel should stay open");
assert!(matches!(msg, NodeControlMsg::Ack(_)));
drop(return_tx);
drop(pipeline_tx);
let dispatcher_result =
timeout(Duration::from_millis(200), dispatcher_handle).await;
assert!(
dispatcher_result.is_ok(),
"Return dispatcher should shut down cleanly"
);
let manager_result = timeout(Duration::from_millis(200), manager_handle).await;
assert!(manager_result.is_ok(), "Manager should shut down cleanly");
drop(dispatcher_guard);
})
.await;
}
#[tokio::test]
async fn test_circular_wait_between_node_and_manager() {
use crate::control::AckMsg;
fn pdata_for_node(node_id: usize) -> TestPData {
let mut pdata = TestPData::new();
pdata.push_frame(Frame {
node_id,
interests: Interests::ACKS,
route: RouteData {
calldata: Default::default(),
entry_time_ns: 0,
output_port_index: 0,
},
output_items: 0,
input_items: 0,
output_size: 0,
input_size: 0,
});
pdata
}
let local = LocalSet::new();
local
.run_until(async {
let (return_tx, return_rx) = pipeline_completion_msg_channel(3);
let mut control_senders = ControlSenders::new();
let nodes = test_nodes(vec!["node_a", "node_b"]);
let node_a = nodes[0].clone();
let node_b = nodes[1].clone();
let (tx_a, rx_a) = tokio::sync::mpsc::channel::<NodeControlMsg<TestPData>>(1);
control_senders.register(
node_a.clone(),
NodeType::Processor,
Sender::Shared(SharedSender::mpsc(tx_a)),
);
let (tx_b, rx_b) = tokio::sync::mpsc::channel::<NodeControlMsg<TestPData>>(10);
control_senders.register(
node_b.clone(),
NodeType::Processor,
Sender::Shared(SharedSender::mpsc(tx_b)),
);
let (dispatcher_context, dispatcher_guard) = create_test_pipeline_context();
let dispatcher = PipelineCompletionMsgDispatcher::new(
dispatcher_context,
return_rx,
control_senders,
empty_node_metric_handles(),
MetricsReporter::create_new_and_receiver(16).1,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
TelemetryPolicy::default(),
TerminalMetricsDeadline::default(),
);
return_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata_for_node(node_a.index)),
})
.await
.unwrap();
return_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata_for_node(node_a.index)),
})
.await
.unwrap();
return_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata_for_node(node_b.index)),
})
.await
.unwrap();
let node_a_tx = return_tx.clone();
let node_a_index = node_a.index;
let _node_a_handle = tokio::task::spawn_local(async move {
let _rx_a = rx_a; loop {
if node_a_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata_for_node(node_a_index)),
})
.await
.is_err()
{
break;
}
}
});
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
let mut receiver_b = Receiver::Shared(SharedReceiver::mpsc(rx_b));
let received = timeout(Duration::from_millis(500), receiver_b.recv()).await;
assert!(
received.is_ok(),
"Node B should receive its Ack within 500 ms, but the \
dispatcher is stuck in a circular wait with Node A: the \
dispatcher is blocked sending to Node A's full control \
channel, while Node A is blocked sending to the full \
shared return channel"
);
drop(return_tx);
dispatcher_handle.abort();
drop(dispatcher_guard);
})
.await;
}
#[tokio::test]
async fn test_runtime_ctrl_progress_while_return_lane_is_busy() {
use crate::control::AckMsg;
fn pdata_for_node(node_id: usize) -> TestPData {
let mut pdata = TestPData::new();
pdata.push_frame(Frame {
node_id,
interests: Interests::ACKS,
route: RouteData {
calldata: Default::default(),
entry_time_ns: 0,
output_port_index: 0,
},
output_items: 0,
input_items: 0,
output_size: 0,
input_size: 0,
});
pdata
}
let local = LocalSet::new();
local
.run_until(async {
let mut control_senders = ControlSenders::new();
let nodes = test_nodes(vec!["node_a", "node_timer"]);
let node_a = nodes[0].clone();
let node_timer = nodes[1].clone();
let (tx_a, rx_a) = tokio::sync::mpsc::channel::<NodeControlMsg<TestPData>>(1);
control_senders.register(
node_a.clone(),
NodeType::Processor,
Sender::Shared(SharedSender::mpsc(tx_a)),
);
let (tx_timer, rx_timer) =
tokio::sync::mpsc::channel::<NodeControlMsg<TestPData>>(10);
control_senders.register(
node_timer.clone(),
NodeType::Processor,
Sender::Shared(SharedSender::mpsc(tx_timer)),
);
let (mut manager, pipeline_tx, _pipeline_entity_guard) =
build_test_manager(16, control_senders.clone());
manager.tick_timers.start(node_timer.index, Duration::from_millis(1));
tokio::time::sleep(Duration::from_millis(5)).await;
let (return_tx, return_rx) = pipeline_completion_msg_channel(3);
for _ in 0..3 {
return_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata_for_node(node_a.index)),
})
.await
.unwrap();
}
let (dispatcher_context, dispatcher_guard) = create_test_pipeline_context();
let dispatcher = PipelineCompletionMsgDispatcher::new(
dispatcher_context,
return_rx,
control_senders,
empty_node_metric_handles(),
MetricsReporter::create_new_and_receiver(16).1,
TEST_CONTROL_PLANE_METRICS_FLUSH_INTERVAL,
TelemetryPolicy::default(),
TerminalMetricsDeadline::default(),
);
let node_a_tx = return_tx.clone();
let node_a_index = node_a.index;
let node_a_handle = tokio::task::spawn_local(async move {
let _rx_a = rx_a;
for _ in 0..200 {
if node_a_tx
.send(PipelineCompletionMsg::DeliverAck {
ack: AckMsg::new(pdata_for_node(node_a_index)),
})
.await
.is_err()
{
break;
}
}
});
let manager_handle = tokio::task::spawn_local(async move { manager.run().await });
let dispatcher_handle =
tokio::task::spawn_local(async move { dispatcher.run().await });
let mut receiver_timer = Receiver::Shared(SharedReceiver::mpsc(rx_timer));
let received = timeout(Duration::from_millis(500), receiver_timer.recv()).await;
assert!(
received.is_ok(),
"TimerTick should make progress even while the return lane is under sustained Ack load"
);
let msg = received.unwrap().expect("timer control channel should stay open");
assert!(matches!(msg, NodeControlMsg::TimerTick {}));
drop(return_tx);
drop(pipeline_tx);
let _ = timeout(Duration::from_millis(500), node_a_handle).await;
let dispatcher_result =
timeout(Duration::from_millis(500), dispatcher_handle).await;
assert!(
dispatcher_result.is_ok(),
"Return dispatcher should shut down cleanly"
);
let manager_result = timeout(Duration::from_millis(500), manager_handle).await;
assert!(manager_result.is_ok(), "Manager should shut down cleanly");
drop(dispatcher_guard);
})
.await;
}
}