use crate::{
admission::AdmissionBinder,
attributes::{ChannelImplementation, ChannelKind, ChannelMode, ChannelType},
channel_metrics::{
ChannelMetricsRegistry, ChannelReceiverMetricSets, ChannelReceiverMetrics,
ChannelReceiverStateMetrics, ChannelSenderFailureMetrics, ChannelSenderMetricSets,
ChannelSenderMetrics, LocalChannelQueueDepth, PdataChannelReceiverMetricSets,
PdataChannelSenderMetricSets, SharedChannelQueueDepth,
},
config::{ExporterConfig, ExtensionConfig, ProcessorConfig, ReceiverConfig},
control::{AckMsg, CallData, NackMsg},
effect_handler::SourceTagging,
entity_context::{NodeTelemetryGuard, NodeTelemetryHandle, with_node_telemetry_handle},
error::{Error, TypedError},
exporter::ExporterWrapper,
extension::ExtensionBundle,
local::message::{LocalReceiver, LocalSender},
message::{Receiver, Sender},
node::{Node, NodeDefs, NodeId, NodeName, NodeType},
processor::{ProcessorWrapper, validate_local_wakeup_requirements},
receiver::ReceiverWrapper,
runtime_pipeline::{PipeNode, RuntimePipeline},
shared::message::{SharedReceiver, SharedSender},
};
use async_trait::async_trait;
pub use channel_metrics::RequestOutcome;
use context::ExtensionContext;
use context::NodeNameIndex;
use context::PipelineContext;
pub use linkme::distributed_slice;
use otel_arrow_dfe_config::MetricLevel;
use otel_arrow_dfe_config::SignalType;
use otel_arrow_dfe_config::{
PipelineGroupId, PipelineId, PortName,
engine::INTERNAL_TELEMETRY_RECEIVER_URN,
node::NodeUserConfig,
pipeline::{DispatchPolicy, PipelineConfig},
policy::{
ChannelCapacityPolicy, RateLimiterDeclarationScope, RateLimiterPolicy, TelemetryPolicy,
},
};
use otel_arrow_dfe_telemetry::InternalTelemetrySettings;
use otel_arrow_dfe_telemetry::{otel_debug, otel_debug_span, otel_info, otel_warn};
use std::borrow::Cow;
use std::fmt::Debug;
use std::num::NonZeroUsize;
use std::sync::Arc;
use std::{
collections::{BTreeMap, HashMap, HashSet},
sync::OnceLock,
};
pub mod admission;
pub mod capability;
#[doc(hidden)]
pub mod clock;
pub mod context_declaration;
pub mod error;
pub mod exporter;
pub mod extension;
mod extension_lifecycle;
mod extension_monitor;
mod forced_shutdown;
pub mod inventory;
pub use otel_arrow_dfe_engine_macros::component_inventory;
pub mod message;
pub mod processor;
pub mod receiver;
pub mod retained_work;
pub mod runtime_services;
mod attributes;
mod channel_metrics;
mod channel_mode;
mod completion_emission_metrics;
pub mod config;
pub mod context;
pub mod control;
mod control_plane_metrics;
pub mod effect_handler;
pub mod engine_metrics;
pub mod entity_context;
pub mod flow_metrics;
pub(crate) mod indexed_min_heap;
pub mod listener_group;
pub mod local;
pub mod memory_limiter;
pub mod node;
mod node_local_scheduler;
pub mod output_router;
pub mod pipeline_ctrl;
mod pipeline_metrics;
pub mod process_duration;
mod route_admission;
pub mod runtime_pipeline;
pub mod shared;
pub mod state_dir;
pub mod terminal_state;
pub mod testing;
pub mod topic;
pub mod topology;
pub mod wiring_contract;
pub use node_local_scheduler::{WakeupError, WakeupSetOutcome};
pub use processor::{LocalWakeupRequirements, ProcessorRuntimeRequirements};
pub use route_admission::RouteAdmission;
fn resolve_admission_binding(
node_config: &NodeUserConfig,
policies: &BTreeMap<String, RateLimiterPolicy>,
declaration_scope: Option<RateLimiterDeclarationScope>,
) -> Result<AdmissionBinder, String> {
match node_config.rate_limiters.as_deref() {
Some([]) => Ok(AdmissionBinder::none()),
Some([limiter_name]) => {
let policy = policies.get(limiter_name).copied().ok_or_else(|| {
format!("rate limiter binding '{limiter_name}' does not name an effective limiter")
})?;
Ok(AdmissionBinder::configured_at_scope(
limiter_name.clone(),
declaration_scope,
policy,
))
}
Some(limiter_names) => Err(format!(
"V1 supports at most one rate limiter binding per node; found {}",
limiter_names.len()
)),
None => Ok(AdmissionBinder::none()),
}
}
pub trait NamedFactory {
fn name(&self) -> &'static str;
}
pub struct ReceiverFactory<PData> {
pub name: &'static str,
pub create: fn(
pipeline_ctx: PipelineContext,
node: NodeId,
node_config: Arc<NodeUserConfig>,
receiver_config: &ReceiverConfig,
capabilities: &capability::registry::Capabilities,
) -> Result<ReceiverWrapper<PData>, otel_arrow_dfe_config::error::Error>,
pub context_declarations: Option<context_declaration::ContextDeclarationProvider>,
pub wiring_contract: wiring_contract::WiringContract,
pub validate_config:
fn(config: &serde_json::Value) -> Result<(), otel_arrow_dfe_config::error::Error>,
}
impl<PData> Clone for ReceiverFactory<PData> {
fn clone(&self) -> Self {
ReceiverFactory {
name: self.name,
create: self.create,
context_declarations: self.context_declarations,
wiring_contract: self.wiring_contract,
validate_config: self.validate_config,
}
}
}
impl<PData> NamedFactory for ReceiverFactory<PData> {
fn name(&self) -> &'static str {
self.name
}
}
pub struct ProcessorFactory<PData> {
pub name: &'static str,
pub create: fn(
pipeline: PipelineContext,
node: NodeId,
node_config: Arc<NodeUserConfig>,
processor_config: &ProcessorConfig,
capabilities: &capability::registry::Capabilities,
) -> Result<ProcessorWrapper<PData>, otel_arrow_dfe_config::error::Error>,
pub context_declarations: Option<context_declaration::ContextDeclarationProvider>,
pub wiring_contract: wiring_contract::WiringContract,
pub validate_config:
fn(config: &serde_json::Value) -> Result<(), otel_arrow_dfe_config::error::Error>,
}
impl<PData> Clone for ProcessorFactory<PData> {
fn clone(&self) -> Self {
ProcessorFactory {
name: self.name,
create: self.create,
context_declarations: self.context_declarations,
wiring_contract: self.wiring_contract,
validate_config: self.validate_config,
}
}
}
impl<PData> NamedFactory for ProcessorFactory<PData> {
fn name(&self) -> &'static str {
self.name
}
}
pub struct ExporterFactory<PData> {
pub name: &'static str,
pub create: fn(
pipeline: PipelineContext,
node: NodeId,
node_config: Arc<NodeUserConfig>,
exporter_config: &ExporterConfig,
capabilities: &capability::registry::Capabilities,
) -> Result<ExporterWrapper<PData>, otel_arrow_dfe_config::error::Error>,
pub context_declarations: Option<context_declaration::ContextDeclarationProvider>,
pub wiring_contract: wiring_contract::WiringContract,
pub validate_config:
fn(config: &serde_json::Value) -> Result<(), otel_arrow_dfe_config::error::Error>,
}
impl<PData> Clone for ExporterFactory<PData> {
fn clone(&self) -> Self {
ExporterFactory {
name: self.name,
create: self.create,
context_declarations: self.context_declarations,
wiring_contract: self.wiring_contract,
validate_config: self.validate_config,
}
}
}
impl<PData> NamedFactory for ExporterFactory<PData> {
fn name(&self) -> &'static str {
self.name
}
}
#[derive(Clone)]
pub struct ExtensionFactory {
pub name: &'static str,
pub description: &'static str,
pub documentation_url: &'static str,
pub capabilities: Option<capability::ExtensionCapabilities>,
pub create: fn(
ext_ctx: &ExtensionContext,
name: otel_arrow_dfe_config::ExtensionId,
ext_config: Arc<otel_arrow_dfe_config::extension::ExtensionUserConfig>,
extension_config: &ExtensionConfig,
) -> Result<ExtensionBundle, otel_arrow_dfe_config::error::Error>,
pub validate_config:
fn(config: &serde_json::Value) -> Result<(), otel_arrow_dfe_config::error::Error>,
}
impl NamedFactory for ExtensionFactory {
fn name(&self) -> &'static str {
self.name
}
}
pub fn get_factory_map<T>(
factory_map: &'static OnceLock<HashMap<&'static str, T>>,
factory_slice: &'static [T],
) -> &'static HashMap<&'static str, T>
where
T: NamedFactory + Clone,
{
factory_map.get_or_init(|| {
factory_slice
.iter()
.map(|f| (f.name(), f.clone()))
.collect::<HashMap<&'static str, T>>()
})
}
bitflags::bitflags! {
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct Interests: u16 {
const ACKS = 1 << 0;
const NACKS = 1 << 1;
const ACKS_OR_NACKS = Self::ACKS.bits() | Self::NACKS.bits();
const RETURN_DATA = 1 << 2;
const NODE_COMPLETION_DURATION = 1 << 3;
const NODE_INPUT_METRICS = 1 << 4;
const NODE_OUTPUT_METRICS = 1 << 5;
const SOURCE_TAGGING = 1 << 6;
const NODE_ITEM_COUNTS = 1 << 8;
const NODE_SIZE = 1 << 9;
const NODE_LOCAL_DURATION = 1 << 7;
const NODE_METRICS = Self::NODE_INPUT_METRICS.bits() | Self::NODE_OUTPUT_METRICS.bits();
}
}
impl Interests {
#[must_use]
pub fn from_metric_level(level: MetricLevel) -> Self {
match level {
MetricLevel::None | MetricLevel::Basic => Self::empty(),
MetricLevel::Normal => Self::NODE_METRICS,
MetricLevel::Detailed => {
Self::NODE_METRICS
| Self::NODE_COMPLETION_DURATION
| Self::NODE_ITEM_COUNTS
| Self::NODE_SIZE
| Self::NODE_LOCAL_DURATION
}
}
}
#[must_use]
pub fn for_node(level: MetricLevel, node_config: &NodeUserConfig) -> Self {
let mut interests = Self::from_metric_level(level);
if let Some(telemetry) = node_config
.policies
.as_ref()
.and_then(|policies| policies.telemetry.as_ref())
{
if telemetry.messages {
interests |= Self::NODE_METRICS;
}
if telemetry.completion_duration {
interests |= Self::NODE_COMPLETION_DURATION;
}
if telemetry.item_counts {
interests |= Self::NODE_ITEM_COUNTS;
}
if telemetry.size {
interests |= Self::NODE_SIZE;
}
if telemetry.duration {
interests |= Self::NODE_LOCAL_DURATION;
}
}
interests
}
}
pub trait Unwindable {
fn has_frames(&self) -> bool;
fn pop_frame(&mut self) -> Option<control::Frame>;
fn signal(&self) -> Option<SignalType>;
fn drop_payload(&mut self);
}
impl Unwindable for () {
fn has_frames(&self) -> bool {
false
}
fn pop_frame(&mut self) -> Option<control::Frame> {
None
}
fn signal(&self) -> Option<SignalType> {
None
}
fn drop_payload(&mut self) {}
}
impl Unwindable for String {
fn has_frames(&self) -> bool {
false
}
fn pop_frame(&mut self) -> Option<control::Frame> {
None
}
fn signal(&self) -> Option<SignalType> {
None
}
fn drop_payload(&mut self) {}
}
pub trait ReceivedAtNode {
fn received_at_node(&mut self, node_id: usize, node_interests: Interests);
}
impl ReceivedAtNode for () {
fn received_at_node(&mut self, _node_id: usize, _node_interests: Interests) {}
}
impl ReceivedAtNode for String {
fn received_at_node(&mut self, _node_id: usize, _node_interests: Interests) {}
}
impl processor::FlowMetricHook for () {}
impl processor::FlowMetricHook for String {}
pub trait StampOutputPort {
fn stamp_output_port_index(&mut self, node_id: usize, index: u16);
}
impl StampOutputPort for () {
fn stamp_output_port_index(&mut self, _node_id: usize, _index: u16) {}
}
impl StampOutputPort for String {
fn stamp_output_port_index(&mut self, _node_id: usize, _index: u16) {}
}
pub trait FlowMetricAccumulation {
fn start_flow_metric(&mut self);
fn add_flow_compute(&mut self, ns: u64);
fn take_flow_compute(&mut self) -> Option<u64>;
}
#[async_trait(?Send)]
pub trait ProducerEffectHandlerExtension<PData> {
fn subscribe_to(&self, int: Interests, ctx: CallData, data: &mut PData);
}
#[async_trait(?Send)]
pub trait ConsumerEffectHandlerExtension<PData> {
async fn notify_ack(&self, ack: AckMsg<PData>) -> Result<(), Error>;
async fn notify_nack(&self, nack: NackMsg<PData>) -> Result<(), Error>;
}
#[doc(hidden)]
pub mod _private {
use super::*;
#[async_trait(?Send)]
pub trait AckNackRouting<PData> {
async fn route_ack(&self, ack: AckMsg<PData>) -> Result<(), Error>;
async fn route_nack(&self, nack: NackMsg<PData>) -> Result<(), Error>;
}
}
#[async_trait(?Send)]
pub trait MessageSourceLocalEffectHandlerExtension<PData> {
async fn send_message_with_source_node(&self, data: PData) -> Result<(), TypedError<PData>>;
fn try_send_message_with_source_node(&self, data: PData) -> Result<(), TypedError<PData>>;
fn try_admit_message_with_source_node(
&self,
data: PData,
) -> Result<RouteAdmission<PData>, TypedError<PData>> {
route_admission::classify_route_admission(self.try_send_message_with_source_node(data))
}
async fn send_message_with_source_node_to<P>(
&self,
port: P,
data: PData,
) -> Result<(), TypedError<PData>>
where
P: Into<PortName> + Send + 'static;
fn try_send_message_with_source_node_to<P>(
&self,
port: P,
data: PData,
) -> Result<(), TypedError<PData>>
where
P: Into<PortName> + Send + 'static;
fn try_admit_message_with_source_node_to<P>(
&self,
port: P,
data: PData,
) -> Result<RouteAdmission<PData>, TypedError<PData>>
where
P: Into<PortName> + Send + 'static,
{
route_admission::classify_route_admission(
self.try_send_message_with_source_node_to(port, data),
)
}
}
#[async_trait]
pub trait MessageSourceSharedEffectHandlerExtension<PData: Send + 'static> {
async fn send_message_with_source_node(&self, data: PData) -> Result<(), TypedError<PData>>;
fn try_send_message_with_source_node(&self, data: PData) -> Result<(), TypedError<PData>>;
fn try_admit_message_with_source_node(
&self,
data: PData,
) -> Result<RouteAdmission<PData>, TypedError<PData>> {
route_admission::classify_route_admission(self.try_send_message_with_source_node(data))
}
async fn send_message_with_source_node_to<P>(
&self,
port: P,
data: PData,
) -> Result<(), TypedError<PData>>
where
P: Into<PortName> + Send + 'static;
fn try_send_message_with_source_node_to<P>(
&self,
port: P,
data: PData,
) -> Result<(), TypedError<PData>>
where
P: Into<PortName> + Send + 'static;
fn try_admit_message_with_source_node_to<P>(
&self,
port: P,
data: PData,
) -> Result<RouteAdmission<PData>, TypedError<PData>>
where
P: Into<PortName> + Send + 'static,
{
route_admission::classify_route_admission(
self.try_send_message_with_source_node_to(port, data),
)
}
}
#[must_use]
pub const fn build_factory<PData: 'static + Clone>() -> PipelineFactory<PData> {
panic!(
"build_registry() should never be called - it's replaced by the #[factory_registry] macro"
)
}
pub struct PipelineFactory<PData: 'static + Clone> {
receiver_factory_map: OnceLock<HashMap<&'static str, ReceiverFactory<PData>>>,
processor_factory_map: OnceLock<HashMap<&'static str, ProcessorFactory<PData>>>,
exporter_factory_map: OnceLock<HashMap<&'static str, ExporterFactory<PData>>>,
extension_factory_map: OnceLock<HashMap<&'static str, ExtensionFactory>>,
receiver_factories: &'static [ReceiverFactory<PData>],
processor_factories: &'static [ProcessorFactory<PData>],
exporter_factories: &'static [ExporterFactory<PData>],
extension_factories: &'static [ExtensionFactory],
}
impl<PData: 'static + Clone + Debug> PipelineFactory<PData> {
#[must_use]
pub const fn new(
receiver_factories: &'static [ReceiverFactory<PData>],
processor_factories: &'static [ProcessorFactory<PData>],
exporter_factories: &'static [ExporterFactory<PData>],
extension_factories: &'static [ExtensionFactory],
) -> Self {
Self {
receiver_factory_map: OnceLock::new(),
processor_factory_map: OnceLock::new(),
exporter_factory_map: OnceLock::new(),
extension_factory_map: OnceLock::new(),
receiver_factories,
processor_factories,
exporter_factories,
extension_factories,
}
}
pub fn get_receiver_factory_map(&self) -> &HashMap<&'static str, ReceiverFactory<PData>> {
self.receiver_factory_map.get_or_init(|| {
self.receiver_factories
.iter()
.map(|f| (f.name(), f.clone()))
.collect::<HashMap<&'static str, ReceiverFactory<PData>>>()
})
}
pub fn get_processor_factory_map(&self) -> &HashMap<&'static str, ProcessorFactory<PData>> {
self.processor_factory_map.get_or_init(|| {
self.processor_factories
.iter()
.map(|f| (f.name(), f.clone()))
.collect::<HashMap<&'static str, ProcessorFactory<PData>>>()
})
}
pub fn get_exporter_factory_map(&self) -> &HashMap<&'static str, ExporterFactory<PData>> {
self.exporter_factory_map.get_or_init(|| {
self.exporter_factories
.iter()
.map(|f| (f.name(), f.clone()))
.collect::<HashMap<&'static str, ExporterFactory<PData>>>()
})
}
pub fn get_extension_factory_map(&self) -> &HashMap<&'static str, ExtensionFactory> {
self.extension_factory_map.get_or_init(|| {
self.extension_factories
.iter()
.map(|f| (f.name(), f.clone()))
.collect::<HashMap<&'static str, ExtensionFactory>>()
})
}
pub fn build(
self: &PipelineFactory<PData>,
mut pipeline_ctx: PipelineContext,
mut config: PipelineConfig,
channel_capacity_policy: ChannelCapacityPolicy,
telemetry_policy: TelemetryPolicy,
rate_limiter_policies: BTreeMap<String, RateLimiterPolicy>,
rate_limiter_scope: Option<RateLimiterDeclarationScope>,
internal_telemetry: Option<InternalTelemetrySettings>,
) -> Result<RuntimePipeline<PData>, Error>
where
PData: Unwindable,
{
let mut receivers = Vec::new();
let mut processors = Vec::new();
let mut exporters = Vec::new();
let mut build_state = BuildState::new();
let pipeline_group_id = pipeline_ctx.pipeline_group_id();
let pipeline_id = pipeline_ctx.pipeline_id();
let core_id = pipeline_ctx.core_id();
let span = otel_debug_span!(
"pipeline.build",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id
);
let _enter = span.enter();
let unconnected = config.remove_unconnected_nodes();
for (node_id, node_kind) in &unconnected {
let kind: Cow<'static, str> = (*node_kind).into();
otel_info!(
"pipeline.build.unconnected_node.removed",
message = "Removed unconnected node from pipeline.",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
node_id = node_id.as_ref(),
node_kind = kind.as_ref(),
);
}
if !unconnected.is_empty() {
otel_warn!(
"pipeline.build.unconnected_nodes",
message = "Some pipeline nodes were removed because they had no active incoming or outgoing edges. These nodes will not participate in data processing. Check pipeline configuration if this is unintentional.",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
removed_count = unconnected.len(),
);
}
if config.nodes().is_empty() {
return Err(Error::EmptyPipeline);
}
self.validate_connection_wiring_contracts(&config)?;
let basic_runtime_metrics_enabled = telemetry_policy.runtime_metrics >= MetricLevel::Basic;
let mut receiver_count = 0usize;
let mut processor_count = 0usize;
let mut exporter_count = 0usize;
let mut node_ids: HashMap<NodeName, NodeId> = HashMap::new();
for (name, node_config) in config.node_iter() {
let (node_type, pipe_node) = match node_config.kind() {
otel_arrow_dfe_config::node::NodeKind::Receiver => {
let pn = PipeNode::new(receiver_count);
receiver_count += 1;
(NodeType::Receiver, pn)
}
otel_arrow_dfe_config::node::NodeKind::Processor => {
let pn = PipeNode::new(processor_count);
processor_count += 1;
(NodeType::Processor, pn)
}
otel_arrow_dfe_config::node::NodeKind::Exporter => {
let pn = PipeNode::new(exporter_count);
exporter_count += 1;
(NodeType::Exporter, pn)
}
};
let node_id = build_state.next_node_id(name.clone(), node_type, pipe_node)?;
let _ = node_ids.insert(name.clone(), node_id);
}
let node_names: NodeNameIndex = Arc::new(
node_ids
.iter()
.map(|(name, id)| (name.clone(), id.clone()))
.collect(),
);
pipeline_ctx.set_node_names(node_names);
let known_extensions: HashSet<otel_arrow_dfe_config::ExtensionId> =
config.extensions().keys().cloned().collect();
let mut capability_registry = capability::registry::CapabilityRegistry::new();
let mut extension_bundles: Vec<(
otel_arrow_dfe_config::ExtensionId,
ExtensionBundle,
bool,
extension::wrapper::ExtensionEntityKeys,
)> = Vec::new();
for (ext_id, ext_user_config) in config.extension_iter() {
let raw_urn = ext_user_config.r#type.as_str();
let factory = self
.get_extension_factory_map()
.get(raw_urn)
.ok_or_else(|| Error::UnknownExtension {
plugin_urn: raw_urn.to_string(),
})?;
let runtime_config = ExtensionConfig::with_control_channel_capacity(
ext_id.clone(),
channel_capacity_policy.control.node,
);
let ext_ctx = pipeline_ctx.extension_context();
let bundle = (factory.create)(
&ext_ctx,
ext_id.clone(),
ext_user_config.clone(),
&runtime_config,
)
.map_err(|e| Error::ConfigError(Box::new(e)))?;
let mut bundle = bundle;
let entity_keys = bundle.wire_telemetry(
ext_id.clone(),
&ext_ctx,
&mut build_state.channel_metrics,
basic_runtime_metrics_enabled,
);
bundle
.register_into(factory.capabilities.as_ref(), &mut capability_registry)
.map_err(|e| Error::CapabilityRegistrationFailed {
extension: ext_id.clone(),
message: format!("{e}"),
})?;
let is_background = factory.capabilities.is_none();
extension_bundles.push((ext_id.clone(), bundle, is_background, entity_keys));
}
let mut consumed_tracker = capability::registry::ConsumedTracker::new();
let mut per_node_capabilities: HashMap<NodeName, capability::registry::Capabilities> =
HashMap::new();
for (name, node_config) in config.node_iter() {
let caps = capability::registry::resolve_bindings(
&node_config.capabilities,
&capability_registry,
&known_extensions,
&mut consumed_tracker,
)
.map_err(|e| Error::CapabilityResolutionFailed {
node: name.clone(),
message: format!("{e}"),
})?;
let _ = per_node_capabilities.insert(name.clone(), caps);
}
let empty_capabilities = capability::registry::Capabilities::empty();
let mut admission_bound_nodes = Vec::new();
let mut admission_explicitly_opted_out_nodes = Vec::new();
for (name, node_config) in config.node_iter() {
let node_kind = node_config.kind();
let node_id = node_ids.get(name).expect("allocated in first pass").clone();
let mut base_ctx = pipeline_ctx.with_node_context(
name.clone(),
node_config.r#type.clone(),
node_kind,
node_config.identity_attributes(),
);
base_ctx.set_node_interests(Interests::for_node(
telemetry_policy.runtime_metrics,
node_config,
));
base_ctx.set_node_duration_distribution(node_config.duration_distribution());
let invalid_binding = |error: String| {
Error::ConfigError(Box::new(
otel_arrow_dfe_config::error::Error::InvalidUserConfig {
error: format!(
"Component `{}` in pipeline_group={} pipeline={} node={}: {error}",
node_config.r#type.as_ref(),
pipeline_ctx.pipeline_group_id().as_ref(),
pipeline_ctx.pipeline_id().as_ref(),
name.as_ref(),
),
},
))
};
let admission = resolve_admission_binding(
node_config,
&rate_limiter_policies,
rate_limiter_scope.clone(),
)
.map_err(invalid_binding)?;
base_ctx.set_admission(admission);
let node_capabilities = per_node_capabilities
.get(name)
.unwrap_or(&empty_capabilities);
match node_kind {
otel_arrow_dfe_config::node::NodeKind::Receiver => {
if node_config.r#type.as_ref() == INTERNAL_TELEMETRY_RECEIVER_URN
&& let Some(ref settings) = internal_telemetry
{
base_ctx.set_internal_telemetry(settings.clone());
}
let wrapper = self.build_node_wrapper(
&mut build_state,
&base_ctx,
NodeType::Receiver,
node_id.clone(),
basic_runtime_metrics_enabled,
|| {
self.create_receiver(
&base_ctx,
node_id.clone(),
node_config.clone(),
channel_capacity_policy.control.node,
channel_capacity_policy.pdata,
node_capabilities,
)
},
)?;
receivers.push(wrapper);
}
otel_arrow_dfe_config::node::NodeKind::Processor => {
let wrapper = self.build_node_wrapper(
&mut build_state,
&base_ctx,
NodeType::Processor,
node_id.clone(),
basic_runtime_metrics_enabled,
|| {
self.create_processor(
&base_ctx,
node_id.clone(),
node_config.clone(),
channel_capacity_policy.control.node,
channel_capacity_policy.pdata,
node_capabilities,
)
},
)?;
processors.push(wrapper);
}
otel_arrow_dfe_config::node::NodeKind::Exporter => {
let wrapper = self.build_node_wrapper(
&mut build_state,
&base_ctx,
NodeType::Exporter,
node_id.clone(),
basic_runtime_metrics_enabled,
|| {
self.create_exporter(
&base_ctx,
node_id.clone(),
node_config.clone(),
channel_capacity_policy.control.node,
channel_capacity_policy.pdata,
node_capabilities,
)
},
)?;
exporters.push(wrapper);
}
}
if base_ctx.admission().was_bound() {
admission_bound_nodes.push(name.as_ref().to_owned());
} else if matches!(node_config.rate_limiters.as_deref(), Some([])) {
admission_explicitly_opted_out_nodes.push(name.as_ref().to_owned());
}
}
if !admission_bound_nodes.is_empty() || !admission_explicitly_opted_out_nodes.is_empty() {
otel_info!(
"admission.binding.summary",
bound_nodes = admission_bound_nodes.join(","),
explicitly_opted_out_nodes = admission_explicitly_opted_out_nodes.join(","),
message = "Resolved pipeline admission bindings"
);
}
let bound_extensions: HashSet<otel_arrow_dfe_config::ExtensionId> = config
.node_iter()
.flat_map(|(_, node_config)| node_config.capabilities.values().cloned())
.collect();
let consumed_local: HashSet<otel_arrow_dfe_config::ExtensionId> =
consumed_tracker.consumed_local();
let consumed_shared: HashSet<otel_arrow_dfe_config::ExtensionId> =
consumed_tracker.consumed_shared();
let extension_wrappers: Vec<(
extension::ExtensionWrapper,
otel_arrow_dfe_telemetry::registry::EntityKey,
)> = extension_bundles
.into_iter()
.flat_map(|(ext_id, mut bundle, is_background, entity_keys)| {
let mut kept: Vec<(
extension::ExtensionWrapper,
otel_arrow_dfe_telemetry::registry::EntityKey,
)> = Vec::new();
if is_background {
if let Some(local) = bundle.take_local() {
kept.push((
local,
entity_keys
.local
.expect("wire_telemetry mints a key for every present variant"),
));
}
if let Some(shared) = bundle.take_shared() {
kept.push((
shared,
entity_keys
.shared
.expect("wire_telemetry mints a key for every present variant"),
));
}
return kept;
}
if !bound_extensions.contains(&ext_id) {
otel_warn!(
"extension.unbound",
message = "extension defined in pipeline config but no node binds to any of its capabilities; dropping",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
extension = ext_id.as_ref(),
);
return kept;
}
let local_present = bundle.local().is_some();
let shared_present = bundle.shared().is_some();
let local_consumed = local_present && consumed_local.contains(&ext_id);
let shared_consumed = shared_present && consumed_shared.contains(&ext_id);
if !local_consumed && !shared_consumed {
otel_warn!(
"extension.unconsumed",
message = "node bindings reference this extension but no node called require_*/optional_* for any of its variants; dropping",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
extension = ext_id.as_ref(),
);
return kept;
}
if let Some(local) = bundle.take_local()
&& local_consumed
{
kept.push((
local,
entity_keys
.local
.expect("wire_telemetry mints a key for every present variant"),
));
}
if let Some(shared) = bundle.take_shared()
&& shared_consumed
{
kept.push((
shared,
entity_keys
.shared
.expect("wire_telemetry mints a key for every present variant"),
));
}
kept
})
.collect();
let edges = collect_hyper_edges_runtime_from_connections(&config, &build_state)?;
let buffer_size = NonZeroUsize::new(channel_capacity_policy.pdata)
.expect("channel_capacity.pdata must be non-zero");
let nodes = std::mem::take(&mut build_state.nodes);
let mut pipeline = RuntimePipeline::new(
config,
receivers,
processors,
exporters,
extension_wrappers,
nodes,
telemetry_policy,
);
let wirings = edges
.into_iter()
.map(|hyper_edge| {
let resolved = hyper_edge.resolve(&build_state)?;
resolved.into_wiring(
&pipeline,
&mut build_state,
buffer_size,
basic_runtime_metrics_enabled,
&pipeline_group_id,
&pipeline_id,
core_id,
)
})
.collect::<Result<Vec<_>, Error>>()?;
for wiring in wirings {
wiring.apply(&mut pipeline, &pipeline_group_id, &pipeline_id, core_id)?;
}
pipeline.set_channel_metrics(build_state.channel_metrics.into_handles());
pipeline.set_admission_metrics(build_state.admission_metrics.into_handles());
Ok(pipeline)
}
fn validate_connection_wiring_contracts(&self, config: &PipelineConfig) -> Result<(), Error> {
let mut contracts_by_node: HashMap<NodeName, wiring_contract::WiringContract> =
HashMap::new();
for (node_name, node_config) in config.node_iter() {
let contract = match node_config.kind() {
otel_arrow_dfe_config::node::NodeKind::Receiver => {
let normalized = otel_arrow_dfe_config::node_urn::validate_plugin_urn(
node_config.r#type.as_ref(),
otel_arrow_dfe_config::node::NodeKind::Receiver,
)
.map_err(|e| Error::ConfigError(Box::new(e)))?;
self.get_receiver_factory_map()
.get(normalized.as_str())
.ok_or(Error::UnknownReceiver {
plugin_urn: normalized,
})?
.wiring_contract
}
otel_arrow_dfe_config::node::NodeKind::Processor => {
let normalized = otel_arrow_dfe_config::node_urn::validate_plugin_urn(
node_config.r#type.as_ref(),
otel_arrow_dfe_config::node::NodeKind::Processor,
)
.map_err(|e| Error::ConfigError(Box::new(e)))?;
self.get_processor_factory_map()
.get(normalized.as_str())
.ok_or(Error::UnknownProcessor {
plugin_urn: normalized,
})?
.wiring_contract
}
otel_arrow_dfe_config::node::NodeKind::Exporter => {
let normalized = otel_arrow_dfe_config::node_urn::validate_plugin_urn(
node_config.r#type.as_ref(),
otel_arrow_dfe_config::node::NodeKind::Exporter,
)
.map_err(|e| Error::ConfigError(Box::new(e)))?;
self.get_exporter_factory_map()
.get(normalized.as_str())
.ok_or(Error::UnknownExporter {
plugin_urn: normalized,
})?
.wiring_contract
}
};
_ = contracts_by_node.insert(node_name.as_ref().to_string().into(), contract);
}
let mut destinations_by_source_output: HashMap<(NodeName, PortName), HashSet<NodeName>> =
HashMap::new();
for connection in config.connection_iter() {
let mut destinations: Vec<NodeName> = connection
.to_nodes()
.into_iter()
.map(|node_id| node_id.as_ref().to_string().into())
.collect();
if destinations.is_empty() {
continue;
}
destinations.sort_unstable_by(|left, right| left.as_ref().cmp(right.as_ref()));
destinations.dedup_by(|left, right| left.as_ref() == right.as_ref());
for source in connection.from_sources() {
let source_name: NodeName = source.node_id().as_ref().to_string().into();
let source_port = source.resolved_output_port();
let entry = destinations_by_source_output
.entry((source_name, source_port))
.or_default();
entry.extend(destinations.iter().cloned());
}
}
for ((source, output), destination_set) in destinations_by_source_output {
let Some(contract) = contracts_by_node.get(&source) else {
return Err(Error::UnknownNode { node: source });
};
let mut destinations: Vec<NodeName> = destination_set.into_iter().collect();
destinations.sort_unstable_by(|left, right| left.as_ref().cmp(right.as_ref()));
contract.validate_output_destinations(&source, &output, &destinations)?;
}
Ok(())
}
fn build_node_wrapper<W, F>(
&self,
build_state: &mut BuildState<PData>,
base_ctx: &PipelineContext,
node_type: NodeType,
node_id: NodeId,
basic_runtime_metrics_enabled: bool,
create_wrapper: F,
) -> Result<W, Error>
where
W: TelemetryWrapped,
F: FnOnce() -> Result<W, Error>,
{
let node_entity_key = base_ctx.register_node_entity();
let node_telemetry_handle =
NodeTelemetryHandle::new(base_ctx.metrics_registry(), node_entity_key);
let mut node_guard = Some(NodeTelemetryGuard::new(node_telemetry_handle.clone()));
build_state.register_node(
node_type,
node_id,
base_ctx.clone(),
node_telemetry_handle.clone(),
)?;
let wrapper =
with_node_telemetry_handle(node_telemetry_handle.clone(), || -> Result<W, Error> {
let wrapper = create_wrapper()?;
let wrapper = wrapper.with_control_channel_metrics(
base_ctx,
&mut build_state.channel_metrics,
basic_runtime_metrics_enabled,
);
build_state
.admission_metrics
.register_if_enabled(basic_runtime_metrics_enabled, || {
base_ctx.admission().metrics_handle(base_ctx)
});
Ok(wrapper)
})?;
Ok(wrapper
.with_node_telemetry_guard(node_guard.take().expect("node telemetry guard missing")))
}
fn select_channel_type(
src_nodes: &[&dyn Node<PData>],
dest_nodes: &[&dyn Node<PData>],
buffer_size: NonZeroUsize,
channel_id: Cow<'static, str>,
source_ports: &[PortName],
source_contexts: &[PipelineContext],
source_telemetries: &[NodeTelemetryHandle],
dest_contexts: &[PipelineContext],
dest_telemetries: &[NodeTelemetryHandle],
channel_metrics: &mut ChannelMetricsRegistry,
channel_metrics_enabled: bool,
) -> Result<(Vec<Sender<PData>>, Vec<Receiver<PData>>), Error>
where
PData: Unwindable,
{
let any_source_is_shared = src_nodes.iter().any(|source| source.is_shared());
let any_dest_is_shared = dest_nodes.iter().any(|dest| dest.is_shared());
let use_shared_channels = any_source_is_shared || any_dest_is_shared;
let num_sources = src_nodes.len();
let num_destinations = dest_nodes.len();
debug_assert_eq!(num_sources, source_ports.len());
debug_assert_eq!(num_sources, source_contexts.len());
debug_assert_eq!(num_sources, source_telemetries.len());
debug_assert_eq!(num_destinations, dest_contexts.len());
debug_assert_eq!(num_destinations, dest_telemetries.len());
let channel_kind = ChannelKind::Pdata;
let capacity = buffer_size.get() as u64;
let register_sender_metrics = |ctx: &PipelineContext,
telemetry: &NodeTelemetryHandle,
entity_key| {
let metrics = PdataChannelSenderMetricSets {
messages: ctx
.register_measurement_metric_set_for_entity::<ChannelSenderMetrics>(entity_key),
failures: ctx
.register_measurement_metric_set_for_entity::<ChannelSenderFailureMetrics>(
entity_key,
),
};
for key in metrics.metric_set_keys() {
telemetry.track_metric_set(key);
}
ChannelSenderMetricSets::Pdata(metrics)
};
let register_receiver_metrics = |ctx: &PipelineContext,
telemetry: &NodeTelemetryHandle,
entity_key| {
let metrics = PdataChannelReceiverMetricSets {
messages: ctx.register_measurement_metric_set_for_entity::<ChannelReceiverMetrics>(
entity_key,
),
state: ctx
.register_metric_set_for_entity::<ChannelReceiverStateMetrics>(entity_key),
};
for key in metrics.metric_set_keys() {
telemetry.track_metric_set(key);
}
ChannelReceiverMetricSets::Pdata(metrics)
};
if channel_metrics_enabled {
match (use_shared_channels, num_destinations > 1) {
(true, true) => {
let channel_mode = ChannelMode::Shared;
let channel_type = ChannelType::Mpmc;
let channel_impl = ChannelImplementation::Flume;
let (pdata_sender, pdata_receiver) = flume::bounded(buffer_size.get());
let queue_depth = SharedChannelQueueDepth::default();
let mut pdata_senders = Vec::with_capacity(num_sources);
for ((ctx, telemetry), port) in source_contexts
.iter()
.zip(source_telemetries.iter())
.zip(source_ports.iter())
{
let sender_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
port.clone(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.add_output_channel_key(port.clone(), sender_entity_key);
let sender_metrics =
register_sender_metrics(ctx, telemetry, sender_entity_key);
let sender = SharedSender::mpmc_with_metrics(
pdata_sender.clone(),
channel_metrics,
sender_metrics,
queue_depth.clone(),
Some(<PData as Unwindable>::signal),
);
pdata_senders.push(Sender::Shared(sender));
}
let pdata_receivers = dest_contexts
.iter()
.zip(dest_telemetries.iter())
.map(|(ctx, telemetry)| {
let receiver_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
"input".into(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.set_input_channel_key(receiver_entity_key);
let receiver_metrics =
register_receiver_metrics(ctx, telemetry, receiver_entity_key);
let receiver = SharedReceiver::mpmc_with_metrics(
pdata_receiver.clone(),
channel_metrics,
receiver_metrics,
capacity,
queue_depth.clone(),
Some(<PData as Unwindable>::signal),
);
Receiver::Shared(receiver)
})
.collect::<Vec<_>>();
Ok((pdata_senders, pdata_receivers))
}
(true, false) => {
let channel_mode = ChannelMode::Shared;
let channel_type = ChannelType::Mpsc;
let channel_impl = ChannelImplementation::Tokio;
let (pdata_sender, pdata_receiver) =
tokio::sync::mpsc::channel::<PData>(buffer_size.get());
let queue_depth = SharedChannelQueueDepth::default();
let mut pdata_senders = Vec::with_capacity(num_sources);
for ((ctx, telemetry), port) in source_contexts
.iter()
.zip(source_telemetries.iter())
.zip(source_ports.iter())
{
let sender_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
port.clone(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.add_output_channel_key(port.clone(), sender_entity_key);
let sender_metrics =
register_sender_metrics(ctx, telemetry, sender_entity_key);
let sender = SharedSender::mpsc_with_metrics(
pdata_sender.clone(),
channel_metrics,
sender_metrics,
queue_depth.clone(),
Some(<PData as Unwindable>::signal),
);
pdata_senders.push(Sender::Shared(sender));
}
let ctx = dest_contexts.first().expect("dest_contexts is empty");
let telemetry = dest_telemetries.first().expect("dest_telemetries is empty");
let receiver_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
"input".into(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.set_input_channel_key(receiver_entity_key);
let receiver_metrics =
register_receiver_metrics(ctx, telemetry, receiver_entity_key);
let pdata_receiver = SharedReceiver::mpsc_with_metrics(
pdata_receiver,
channel_metrics,
receiver_metrics,
capacity,
queue_depth,
Some(<PData as Unwindable>::signal),
);
Ok((pdata_senders, vec![Receiver::Shared(pdata_receiver)]))
}
(false, true) => {
let channel_mode = ChannelMode::Local;
let channel_type = ChannelType::Mpmc;
let channel_impl = ChannelImplementation::Internal;
let (pdata_sender, pdata_receiver) =
otel_arrow_dfe_channel::mpmc::Channel::new(buffer_size);
let queue_depth = LocalChannelQueueDepth::default();
let mut pdata_senders = Vec::with_capacity(num_sources);
for ((ctx, telemetry), port) in source_contexts
.iter()
.zip(source_telemetries.iter())
.zip(source_ports.iter())
{
let sender_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
port.clone(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.add_output_channel_key(port.clone(), sender_entity_key);
let sender_metrics =
register_sender_metrics(ctx, telemetry, sender_entity_key);
let sender = LocalSender::mpmc_with_metrics(
pdata_sender.clone(),
channel_metrics,
sender_metrics,
queue_depth.clone(),
Some(<PData as Unwindable>::signal),
);
pdata_senders.push(Sender::Local(sender));
}
let pdata_receivers = dest_contexts
.iter()
.zip(dest_telemetries.iter())
.map(|(ctx, telemetry)| {
let receiver_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
"input".into(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.set_input_channel_key(receiver_entity_key);
let receiver_metrics =
register_receiver_metrics(ctx, telemetry, receiver_entity_key);
let receiver = LocalReceiver::mpmc_with_metrics(
pdata_receiver.clone(),
channel_metrics,
receiver_metrics,
capacity,
queue_depth.clone(),
Some(<PData as Unwindable>::signal),
);
Receiver::Local(receiver)
})
.collect::<Vec<_>>();
Ok((pdata_senders, pdata_receivers))
}
(false, false) => {
let channel_mode = ChannelMode::Local;
let channel_type = ChannelType::Mpsc;
let channel_impl = ChannelImplementation::Internal;
let (pdata_sender, pdata_receiver) =
otel_arrow_dfe_channel::mpsc::Channel::new(buffer_size.get());
let queue_depth = LocalChannelQueueDepth::default();
let mut pdata_senders = Vec::with_capacity(num_sources);
for ((ctx, telemetry), port) in source_contexts
.iter()
.zip(source_telemetries.iter())
.zip(source_ports.iter())
{
let sender_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
port.clone(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.add_output_channel_key(port.clone(), sender_entity_key);
let sender_metrics =
register_sender_metrics(ctx, telemetry, sender_entity_key);
let sender = LocalSender::mpsc_with_metrics(
pdata_sender.clone(),
channel_metrics,
sender_metrics,
queue_depth.clone(),
Some(<PData as Unwindable>::signal),
);
pdata_senders.push(Sender::Local(sender));
}
let ctx = dest_contexts.first().expect("dest_contexts is empty");
let telemetry = dest_telemetries.first().expect("dest_telemetries is empty");
let receiver_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
"input".into(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.set_input_channel_key(receiver_entity_key);
let receiver_metrics =
register_receiver_metrics(ctx, telemetry, receiver_entity_key);
let pdata_receiver = LocalReceiver::mpsc_with_metrics(
pdata_receiver,
channel_metrics,
receiver_metrics,
capacity,
queue_depth,
Some(<PData as Unwindable>::signal),
);
Ok((pdata_senders, vec![Receiver::Local(pdata_receiver)]))
}
}
} else {
match (use_shared_channels, num_destinations > 1) {
(true, true) => {
let channel_mode = ChannelMode::Shared;
let channel_type = ChannelType::Mpmc;
let channel_impl = ChannelImplementation::Flume;
let (pdata_sender, pdata_receiver) = flume::bounded(buffer_size.get());
let mut pdata_senders = Vec::with_capacity(num_sources);
for ((ctx, telemetry), port) in source_contexts
.iter()
.zip(source_telemetries.iter())
.zip(source_ports.iter())
{
let sender_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
port.clone(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.add_output_channel_key(port.clone(), sender_entity_key);
let sender = SharedSender::mpmc(pdata_sender.clone());
pdata_senders.push(Sender::Shared(sender));
}
let pdata_receivers = dest_contexts
.iter()
.zip(dest_telemetries.iter())
.map(|(ctx, telemetry)| {
let receiver_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
"input".into(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.set_input_channel_key(receiver_entity_key);
Receiver::Shared(SharedReceiver::mpmc(pdata_receiver.clone()))
})
.collect::<Vec<_>>();
Ok((pdata_senders, pdata_receivers))
}
(true, false) => {
let channel_mode = ChannelMode::Shared;
let channel_type = ChannelType::Mpsc;
let channel_impl = ChannelImplementation::Tokio;
let (pdata_sender, pdata_receiver) =
tokio::sync::mpsc::channel::<PData>(buffer_size.get());
let mut pdata_senders = Vec::with_capacity(num_sources);
for ((ctx, telemetry), port) in source_contexts
.iter()
.zip(source_telemetries.iter())
.zip(source_ports.iter())
{
let sender_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
port.clone(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.add_output_channel_key(port.clone(), sender_entity_key);
let sender = SharedSender::mpsc(pdata_sender.clone());
pdata_senders.push(Sender::Shared(sender));
}
let ctx = dest_contexts.first().expect("dest_contexts is empty");
let telemetry = dest_telemetries.first().expect("dest_telemetries is empty");
let receiver_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
"input".into(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.set_input_channel_key(receiver_entity_key);
Ok((
pdata_senders,
vec![Receiver::Shared(SharedReceiver::mpsc(pdata_receiver))],
))
}
(false, true) => {
let channel_mode = ChannelMode::Local;
let channel_type = ChannelType::Mpmc;
let channel_impl = ChannelImplementation::Internal;
let (pdata_sender, pdata_receiver) =
otel_arrow_dfe_channel::mpmc::Channel::new(buffer_size);
let mut pdata_senders = Vec::with_capacity(num_sources);
for ((ctx, telemetry), port) in source_contexts
.iter()
.zip(source_telemetries.iter())
.zip(source_ports.iter())
{
let sender_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
port.clone(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.add_output_channel_key(port.clone(), sender_entity_key);
let sender = LocalSender::mpmc(pdata_sender.clone());
pdata_senders.push(Sender::Local(sender));
}
let pdata_receivers = dest_contexts
.iter()
.zip(dest_telemetries.iter())
.map(|(ctx, telemetry)| {
let receiver_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
"input".into(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.set_input_channel_key(receiver_entity_key);
Receiver::Local(LocalReceiver::mpmc(pdata_receiver.clone()))
})
.collect::<Vec<_>>();
Ok((pdata_senders, pdata_receivers))
}
(false, false) => {
let channel_mode = ChannelMode::Local;
let channel_type = ChannelType::Mpsc;
let channel_impl = ChannelImplementation::Internal;
let (pdata_sender, pdata_receiver) =
otel_arrow_dfe_channel::mpsc::Channel::new(buffer_size.get());
let mut pdata_senders = Vec::with_capacity(num_sources);
for ((ctx, telemetry), port) in source_contexts
.iter()
.zip(source_telemetries.iter())
.zip(source_ports.iter())
{
let sender_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
port.clone(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.add_output_channel_key(port.clone(), sender_entity_key);
let sender = LocalSender::mpsc(pdata_sender.clone());
pdata_senders.push(Sender::Local(sender));
}
let ctx = dest_contexts.first().expect("dest_contexts is empty");
let telemetry = dest_telemetries.first().expect("dest_telemetries is empty");
let receiver_entity_key = ctx.register_node_channel_entity(
channel_id.clone(),
"input".into(),
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
telemetry.set_input_channel_key(receiver_entity_key);
Ok((
pdata_senders,
vec![Receiver::Local(LocalReceiver::mpsc(pdata_receiver))],
))
}
}
}
}
fn create_receiver(
&self,
pipeline_ctx: &PipelineContext,
node_id: NodeId,
node_config: Arc<NodeUserConfig>,
control_channel_capacity: usize,
pdata_channel_capacity: usize,
capabilities: &capability::registry::Capabilities,
) -> Result<ReceiverWrapper<PData>, Error> {
let pipeline_group_id = pipeline_ctx.pipeline_group_id();
let pipeline_id = pipeline_ctx.pipeline_id();
let core_id = pipeline_ctx.core_id();
let name = node_id.name.clone();
otel_debug!(
"receiver.create.start",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
node_id = name.as_ref(),
);
let normalized = otel_arrow_dfe_config::node_urn::validate_plugin_urn(
node_config.r#type.as_ref(),
otel_arrow_dfe_config::node::NodeKind::Receiver,
)
.map_err(|e| Error::ConfigError(Box::new(e)))?;
let factory = self
.get_receiver_factory_map()
.get(normalized.as_str())
.ok_or_else(|| Error::UnknownReceiver {
plugin_urn: normalized.clone(),
})?;
let runtime_config = ReceiverConfig::with_channel_capacities(
name.clone(),
control_channel_capacity,
pdata_channel_capacity,
);
let create = factory.create;
let capture_policy = pipeline_ctx
.compiled_context_bindings()
.header_capture_policy(&pipeline_ctx.pipeline_key(), &pipeline_ctx.node_id())
.cloned();
let authorized_identity_policy = pipeline_ctx
.compiled_context_bindings()
.authorized_identity_policy(&pipeline_ctx.pipeline_key(), &pipeline_ctx.node_id())
.cloned();
let receiver = create(
(*pipeline_ctx).clone(),
node_id.clone(),
node_config,
&runtime_config,
capabilities,
)
.map_err(|e| Error::ConfigError(Box::new(e)))?
.with_capture_policy(capture_policy)
.with_authorized_identity_policy(authorized_identity_policy);
pipeline_ctx
.admission()
.validate_factory_consumption(normalized.as_str())
.map_err(|error| {
Error::ConfigError(Box::new(
otel_arrow_dfe_config::error::Error::InvalidUserConfig {
error: error.to_string(),
},
))
})?;
otel_debug!(
"receiver.create.complete",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
node_id = name.as_ref(),
);
Ok(receiver)
}
fn create_processor(
&self,
pipeline_ctx: &PipelineContext,
node_id: NodeId,
node_config: Arc<NodeUserConfig>,
control_channel_capacity: usize,
pdata_channel_capacity: usize,
capabilities: &capability::registry::Capabilities,
) -> Result<ProcessorWrapper<PData>, Error> {
let pipeline_group_id = pipeline_ctx.pipeline_group_id();
let pipeline_id = pipeline_ctx.pipeline_id();
let core_id = pipeline_ctx.core_id();
let name = node_id.name.clone();
otel_debug!(
"processor.create.start",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
node_id = name.as_ref(),
);
let normalized = otel_arrow_dfe_config::node_urn::validate_plugin_urn(
node_config.r#type.as_ref(),
otel_arrow_dfe_config::node::NodeKind::Processor,
)
.map_err(|e| Error::ConfigError(Box::new(e)))?;
let factory = self
.get_processor_factory_map()
.get(normalized.as_str())
.ok_or(Error::UnknownProcessor {
plugin_urn: normalized.clone(),
})?;
let processor_config = ProcessorConfig::with_channel_capacities(
name.clone(),
control_channel_capacity,
pdata_channel_capacity,
);
let create = factory.create;
let processor = create(
(*pipeline_ctx).clone(),
node_id.clone(),
node_config.clone(),
&processor_config,
capabilities,
)
.map_err(|e| Error::ConfigError(Box::new(e)))?;
pipeline_ctx
.admission()
.validate_factory_consumption(normalized.as_str())
.map_err(|error| {
Error::ConfigError(Box::new(
otel_arrow_dfe_config::error::Error::InvalidUserConfig {
error: error.to_string(),
},
))
})?;
validate_local_wakeup_requirements(&node_id, processor.runtime_requirements())?;
otel_debug!(
"processor.create.complete",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
node_id = name.as_ref(),
);
Ok(processor)
}
fn create_exporter(
&self,
pipeline_ctx: &PipelineContext,
node_id: NodeId,
node_config: Arc<NodeUserConfig>,
control_channel_capacity: usize,
pdata_channel_capacity: usize,
capabilities: &capability::registry::Capabilities,
) -> Result<ExporterWrapper<PData>, Error> {
let pipeline_group_id = pipeline_ctx.pipeline_group_id();
let pipeline_id = pipeline_ctx.pipeline_id();
let core_id = pipeline_ctx.core_id();
let name = node_id.name.clone();
otel_debug!(
"exporter.create.start",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
node_id = name.as_ref(),
);
let normalized = otel_arrow_dfe_config::node_urn::validate_plugin_urn(
node_config.r#type.as_ref(),
otel_arrow_dfe_config::node::NodeKind::Exporter,
)
.map_err(|e| Error::ConfigError(Box::new(e)))?;
let factory = self
.get_exporter_factory_map()
.get(normalized.as_str())
.ok_or(Error::UnknownExporter {
plugin_urn: normalized.clone(),
})?;
let exporter_config = ExporterConfig::with_channel_capacities(
name.clone(),
control_channel_capacity,
pdata_channel_capacity,
);
let create = factory.create;
let propagation_policy = pipeline_ctx
.compiled_context_bindings()
.header_propagation_policy(&pipeline_ctx.pipeline_key(), &pipeline_ctx.node_id())
.cloned();
let exporter = create(
(*pipeline_ctx).clone(),
node_id.clone(),
node_config,
&exporter_config,
capabilities,
)
.map_err(|e| Error::ConfigError(Box::new(e)))?
.with_propagation_policy(propagation_policy);
pipeline_ctx
.admission()
.validate_factory_consumption(normalized.as_str())
.map_err(|error| {
Error::ConfigError(Box::new(
otel_arrow_dfe_config::error::Error::InvalidUserConfig {
error: error.to_string(),
},
))
})?;
otel_debug!(
"exporter.create.complete",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
node_id = name.as_ref(),
);
Ok(exporter)
}
}
trait TelemetryWrapped: Sized {
fn with_control_channel_metrics(
self,
pipeline_ctx: &PipelineContext,
channel_metrics: &mut ChannelMetricsRegistry,
channel_metrics_enabled: bool,
) -> Self;
fn with_node_telemetry_guard(self, guard: NodeTelemetryGuard) -> Self;
}
impl<PData> TelemetryWrapped for ReceiverWrapper<PData> {
fn with_control_channel_metrics(
self,
pipeline_ctx: &PipelineContext,
channel_metrics: &mut ChannelMetricsRegistry,
channel_metrics_enabled: bool,
) -> Self {
ReceiverWrapper::with_control_channel_metrics(
self,
pipeline_ctx,
channel_metrics,
channel_metrics_enabled,
)
}
fn with_node_telemetry_guard(self, guard: NodeTelemetryGuard) -> Self {
ReceiverWrapper::with_node_telemetry_guard(self, guard)
}
}
impl<PData> TelemetryWrapped for ProcessorWrapper<PData> {
fn with_control_channel_metrics(
self,
pipeline_ctx: &PipelineContext,
channel_metrics: &mut ChannelMetricsRegistry,
channel_metrics_enabled: bool,
) -> Self {
ProcessorWrapper::with_control_channel_metrics(
self,
pipeline_ctx,
channel_metrics,
channel_metrics_enabled,
)
}
fn with_node_telemetry_guard(self, guard: NodeTelemetryGuard) -> Self {
ProcessorWrapper::with_node_telemetry_guard(self, guard)
}
}
impl<PData> TelemetryWrapped for ExporterWrapper<PData> {
fn with_control_channel_metrics(
self,
pipeline_ctx: &PipelineContext,
channel_metrics: &mut ChannelMetricsRegistry,
channel_metrics_enabled: bool,
) -> Self {
ExporterWrapper::with_control_channel_metrics(
self,
pipeline_ctx,
channel_metrics,
channel_metrics_enabled,
)
}
fn with_node_telemetry_guard(self, guard: NodeTelemetryGuard) -> Self {
ExporterWrapper::with_node_telemetry_guard(self, guard)
}
}
struct NodeRegistration {
node_id: NodeId,
node_type: NodeType,
context: PipelineContext,
telemetry: NodeTelemetryHandle,
}
struct BuildState<PData> {
nodes: NodeDefs<PData, PipeNode>,
registry: HashMap<NodeName, NodeRegistration>,
channel_metrics: ChannelMetricsRegistry,
admission_metrics: admission::metrics::AdmissionMetricsRegistry,
}
impl<PData> BuildState<PData> {
fn new() -> Self {
Self {
nodes: NodeDefs::default(),
registry: HashMap::new(),
channel_metrics: ChannelMetricsRegistry::default(),
admission_metrics: admission::metrics::AdmissionMetricsRegistry::default(),
}
}
fn next_node_id(
&mut self,
name: NodeName,
node_type: NodeType,
inner: PipeNode,
) -> Result<NodeId, Error> {
self.nodes.next(name, node_type, inner)
}
fn register_node(
&mut self,
node_type: NodeType,
node_id: NodeId,
context: PipelineContext,
telemetry: NodeTelemetryHandle,
) -> Result<(), Error> {
if self.registry.contains_key(&node_id.name) {
return Err(match node_type {
NodeType::Receiver => Error::ReceiverAlreadyExists { receiver: node_id },
NodeType::Processor => Error::ProcessorAlreadyExists { processor: node_id },
NodeType::Exporter => Error::ExporterAlreadyExists { exporter: node_id },
});
}
let _ = self.registry.insert(
node_id.name.clone(),
NodeRegistration {
node_id,
node_type,
context,
telemetry,
},
);
Ok(())
}
fn registration(&self, name: &NodeName) -> Result<&NodeRegistration, Error> {
self.registry
.get(name)
.ok_or_else(|| Error::UnknownNode { node: name.clone() })
}
fn node_context(&self, name: &NodeName) -> Result<PipelineContext, Error> {
Ok(self.registration(name)?.context.clone())
}
fn node_telemetry(&self, name: &NodeName) -> Result<NodeTelemetryHandle, Error> {
Ok(self.registration(name)?.telemetry.clone())
}
fn resolve_destination_id(&self, name: &NodeName) -> Result<NodeId, Error> {
let registration = self.registration(name)?;
match registration.node_type {
NodeType::Processor | NodeType::Exporter => Ok(registration.node_id.clone()),
NodeType::Receiver => Err(Error::UnknownNode { node: name.clone() }),
}
}
}
struct NodeIdPortName {
node_id: NodeId,
port: PortName,
}
struct HyperEdgeWiring<PData> {
sources: Vec<NodeIdPortName>,
senders: Vec<Sender<PData>>,
destinations: Vec<(NodeId, Receiver<PData>)>,
}
impl<PData> HyperEdgeWiring<PData>
where
PData: 'static + Clone + Debug,
{
fn apply(
self,
pipeline: &mut RuntimePipeline<PData>,
pipeline_group_id: &PipelineGroupId,
pipeline_id: &PipelineId,
core_id: usize,
) -> Result<(), Error> {
debug_assert_eq!(self.sources.len(), self.senders.len());
let multi_source = self.sources.len() > 1;
for (source, sender) in self.sources.into_iter().zip(self.senders) {
let src_node = pipeline
.get_mut_node_with_pdata_sender(source.node_id.index)
.ok_or_else(|| Error::UnknownNode {
node: source.node_id.name.clone(),
})?;
if multi_source {
src_node.set_source_tagging(SourceTagging::Enabled);
}
otel_debug!(
"pdata.sender.set",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
node_id = source.node_id.name.as_ref(),
port = source.port.as_ref(),
);
src_node.set_pdata_sender(source.node_id, source.port, sender)?;
}
for (dest, receiver) in self.destinations {
let dest_node = pipeline
.get_mut_node_with_pdata_receiver(dest.index)
.ok_or_else(|| Error::UnknownNode {
node: dest.name.clone(),
})?;
otel_debug!(
"pdata.receiver.set",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
node_id = dest.name.as_ref(),
);
dest_node.set_pdata_receiver(dest, receiver)?;
}
Ok(())
}
}
struct HyperEdgeRuntime {
sources: Vec<NodeIdPortName>,
dispatch_policy: DispatchPolicy,
destinations: Vec<NodeName>,
}
struct ResolvedHyperEdgeRuntime {
sources: Vec<NodeIdPortName>,
destinations: Vec<NodeId>,
dispatch_policy: DispatchPolicy,
source_ids_display: String,
destination_ids_display: String,
}
#[derive(Hash, PartialEq, Eq)]
struct HyperEdgeKey {
dispatch_policy: std::mem::Discriminant<DispatchPolicy>,
destinations: Vec<NodeName>,
}
impl HyperEdgeRuntime {
fn resolve<PData>(
self,
build_state: &BuildState<PData>,
) -> Result<ResolvedHyperEdgeRuntime, Error> {
let destinations = self
.destinations
.iter()
.map(|name| build_state.resolve_destination_id(name))
.collect::<Result<Vec<_>, Error>>()?;
let source_ids_display = self
.sources
.iter()
.map(|source| format!("{}:{}", source.node_id.name, source.port))
.collect::<Vec<_>>()
.join(", ");
let destination_ids_display = destinations
.iter()
.map(|dest| dest.name.as_ref().to_string())
.collect::<Vec<_>>()
.join(", ");
Ok(ResolvedHyperEdgeRuntime {
sources: self.sources,
destinations,
dispatch_policy: self.dispatch_policy,
source_ids_display,
destination_ids_display,
})
}
}
impl ResolvedHyperEdgeRuntime {
fn channel_id(&self) -> Cow<'static, str> {
let sources = self
.sources
.iter()
.map(|source| format!("{}:{}", source.node_id.name, source.port))
.collect::<Vec<_>>();
let destinations = self
.destinations
.iter()
.map(|dest| dest.name.as_ref().to_string())
.collect::<Vec<_>>();
let signature = format!(
"src:[{}]|dst:[{}]|dispatch:{}",
sources.join(","),
destinations.join(","),
dispatch_policy_label(&self.dispatch_policy)
);
let hash = stable_hash64(&signature);
format!("hyperedge:{:016x}", hash).into()
}
fn into_wiring<PData>(
self,
pipeline: &RuntimePipeline<PData>,
build_state: &mut BuildState<PData>,
buffer_size: NonZeroUsize,
channel_metrics_enabled: bool,
pipeline_group_id: &PipelineGroupId,
pipeline_id: &PipelineId,
core_id: usize,
) -> Result<HyperEdgeWiring<PData>, Error>
where
PData: 'static + Clone + Debug + Unwindable,
{
let channel_id = self.channel_id();
let ResolvedHyperEdgeRuntime {
sources,
destinations,
dispatch_policy: _,
source_ids_display,
destination_ids_display,
} = self;
let span = otel_debug_span!(
"hyper_edge.wireup",
pipeline_group_id = pipeline_group_id.as_ref(),
pipeline_id = pipeline_id.as_ref(),
core_id = core_id,
source_ids = source_ids_display,
dest_ids = destination_ids_display
);
let _enter = span.enter();
let mut source_nodes = Vec::with_capacity(sources.len());
let mut source_ports = Vec::with_capacity(sources.len());
let mut source_contexts = Vec::with_capacity(sources.len());
let mut source_telemetries = Vec::with_capacity(sources.len());
for source in &sources {
let src_node =
pipeline
.get_node(source.node_id.index)
.ok_or_else(|| Error::UnknownNode {
node: source.node_id.name.clone(),
})?;
source_nodes.push(src_node);
source_ports.push(source.port.clone());
source_contexts.push(build_state.node_context(&source.node_id.name)?);
source_telemetries.push(build_state.node_telemetry(&source.node_id.name)?);
}
let mut dest_nodes = Vec::with_capacity(destinations.len());
let mut dest_contexts = Vec::with_capacity(destinations.len());
let mut dest_telemetries = Vec::with_capacity(destinations.len());
for node_id in &destinations {
let node = pipeline
.get_node(node_id.index)
.ok_or_else(|| Error::UnknownNode {
node: node_id.name.clone(),
})?;
dest_nodes.push(node);
dest_contexts.push(build_state.node_context(&node_id.name)?);
dest_telemetries.push(build_state.node_telemetry(&node_id.name)?);
}
let (pdata_senders, pdata_receivers) = PipelineFactory::<PData>::select_channel_type(
&source_nodes,
&dest_nodes,
buffer_size,
channel_id,
&source_ports,
&source_contexts,
&source_telemetries,
&dest_contexts,
&dest_telemetries,
&mut build_state.channel_metrics,
channel_metrics_enabled,
)?;
let destinations = destinations.into_iter().zip(pdata_receivers).collect();
Ok(HyperEdgeWiring {
sources,
senders: pdata_senders,
destinations,
})
}
}
fn collect_hyper_edges_runtime_from_connections<PData>(
config: &PipelineConfig,
build_state: &BuildState<PData>,
) -> Result<Vec<HyperEdgeRuntime>, Error> {
let mut edges: Vec<HyperEdgeRuntime> = Vec::new();
let mut edge_index: HashMap<HyperEdgeKey, Vec<usize>> = HashMap::new();
for connection in config.connection_iter() {
let mut destinations: Vec<NodeName> = connection
.to_nodes()
.into_iter()
.map(|node_id| node_id.as_ref().to_string().into())
.collect();
if destinations.is_empty() {
continue;
}
destinations.sort_unstable_by(|a, b| a.as_ref().cmp(b.as_ref()));
destinations.dedup_by(|a, b| a.as_ref() == b.as_ref());
let mut sources = Vec::new();
for source in connection.from_sources() {
let source_name: NodeName = source.node_id().as_ref().to_string().into();
let registration = build_state.registration(&source_name)?;
if !matches!(
registration.node_type,
NodeType::Receiver | NodeType::Processor
) {
return Err(Error::UnknownNode { node: source_name });
}
sources.push(NodeIdPortName {
node_id: registration.node_id.clone(),
port: source.resolved_output_port(),
});
}
if sources.is_empty() {
continue;
}
sources.sort_by(|left, right| {
let left_key = (left.node_id.name.as_ref(), left.port.as_ref());
let right_key = (right.node_id.name.as_ref(), right.port.as_ref());
left_key.cmp(&right_key)
});
sources.dedup_by(|left, right| {
left.node_id.index == right.node_id.index && left.port.as_ref() == right.port.as_ref()
});
let dispatch_policy = connection.effective_dispatch_policy();
let key = HyperEdgeKey {
dispatch_policy: std::mem::discriminant(&dispatch_policy),
destinations: destinations.clone(),
};
let mut match_index = None;
if let Some(indexes) = edge_index.get(&key) {
'candidate: for &index in indexes {
let edge = &edges[index];
for source in &sources {
if edge.sources.iter().any(|existing| {
existing.node_id.index == source.node_id.index
&& existing.port.as_ref() != source.port.as_ref()
}) {
continue 'candidate;
}
}
match_index = Some(index);
break;
}
}
if let Some(index) = match_index {
edges[index].sources.extend(sources);
} else {
edges.push(HyperEdgeRuntime {
sources,
dispatch_policy,
destinations,
});
edge_index.entry(key).or_default().push(edges.len() - 1);
}
}
for edge in &mut edges {
edge.sources.sort_by(|left, right| {
let left_key = (left.node_id.name.as_ref(), left.port.as_ref());
let right_key = (right.node_id.name.as_ref(), right.port.as_ref());
left_key.cmp(&right_key)
});
edge.sources.dedup_by(|left, right| {
left.node_id.index == right.node_id.index && left.port.as_ref() == right.port.as_ref()
});
}
Ok(edges)
}
const fn dispatch_policy_label(policy: &DispatchPolicy) -> &'static str {
match policy {
DispatchPolicy::OneOf => "one_of",
DispatchPolicy::Broadcast => "broadcast",
}
}
fn stable_hash64(value: &str) -> u64 {
let mut hash = 0xcbf29ce484222325u64;
for byte in value.as_bytes() {
hash ^= u64::from(*byte);
hash = hash.wrapping_mul(0x100000001b3);
}
hash
}
#[cfg(test)]
mod test {
use super::*;
use otel_arrow_dfe_config::policy::{
RateLimitAggregation, RateLimitEnforcement, RateLimitPressure, RateLimitUnit,
TokenBucketPolicy,
};
use std::time::Duration;
#[test]
fn detailed_runtime_metrics_enable_optional_data_path_measurements() {
let normal = Interests::from_metric_level(MetricLevel::Normal);
assert!(normal.contains(Interests::NODE_METRICS));
assert!(!normal.contains(Interests::NODE_COMPLETION_DURATION));
assert!(!normal.contains(Interests::NODE_LOCAL_DURATION));
assert!(!normal.contains(Interests::NODE_ITEM_COUNTS));
assert!(!normal.contains(Interests::NODE_SIZE));
let detailed = Interests::from_metric_level(MetricLevel::Detailed);
assert!(detailed.contains(Interests::NODE_METRICS));
assert!(detailed.contains(Interests::NODE_COMPLETION_DURATION));
assert!(detailed.contains(Interests::NODE_LOCAL_DURATION));
assert!(detailed.contains(Interests::NODE_ITEM_COUNTS));
assert!(detailed.contains(Interests::NODE_SIZE));
}
#[test]
fn node_telemetry_policy_extends_effective_interests() {
let mut node_config = NodeUserConfig::new_exporter_config("console");
node_config.policies = Some(otel_arrow_dfe_config::node::NodePolicies {
telemetry: Some(otel_arrow_dfe_config::node::NodeTelemetryPolicy {
messages: true,
completion_duration: true,
duration: true,
duration_distribution: otel_arrow_dfe_config::policy::DistributionTier::Detailed,
item_counts: true,
size: true,
}),
});
let interests = Interests::for_node(MetricLevel::Normal, &node_config);
assert!(interests.contains(Interests::NODE_METRICS));
assert!(interests.contains(Interests::NODE_COMPLETION_DURATION));
assert!(interests.contains(Interests::NODE_LOCAL_DURATION));
assert!(interests.contains(Interests::NODE_ITEM_COUNTS));
assert!(interests.contains(Interests::NODE_SIZE));
}
#[test]
fn node_telemetry_policy_enables_interests_at_basic_level() {
let mut node_config = NodeUserConfig::new_exporter_config("console");
node_config.policies = Some(otel_arrow_dfe_config::node::NodePolicies {
telemetry: Some(otel_arrow_dfe_config::node::NodeTelemetryPolicy {
messages: true,
completion_duration: true,
duration: true,
duration_distribution: otel_arrow_dfe_config::policy::DistributionTier::Detailed,
item_counts: true,
size: true,
}),
});
let interests = Interests::for_node(MetricLevel::Basic, &node_config);
assert!(interests.contains(Interests::NODE_METRICS));
assert!(interests.contains(Interests::NODE_COMPLETION_DURATION));
assert!(interests.contains(Interests::NODE_LOCAL_DURATION));
assert!(interests.contains(Interests::NODE_ITEM_COUNTS));
assert!(interests.contains(Interests::NODE_SIZE));
}
fn admission_policy(unit: RateLimitUnit) -> RateLimiterPolicy {
RateLimiterPolicy {
enforcement: RateLimitEnforcement::Enforce,
aggregation: RateLimitAggregation::ReceiverInstance,
unit,
pressure: RateLimitPressure::Soft,
token_bucket: TokenBucketPolicy {
allow: 10,
interval: Duration::from_secs(1),
burst: Some(10),
},
}
}
#[test]
fn admission_resolution_honors_explicit_empty_opt_out() {
let mut node = NodeUserConfig::new_receiver_config("urn:test:receiver:example");
node.rate_limiters = Some(Vec::new());
let policies = BTreeMap::from([(
"ingress".to_owned(),
admission_policy(RateLimitUnit::RequestBytes),
)]);
let binder = resolve_admission_binding(&node, &policies, None).expect("opt out resolves");
assert!(!binder.is_configured());
}
#[test]
fn admission_resolution_rejects_unknown_explicit_name() {
let mut node = NodeUserConfig::new_receiver_config("urn:test:receiver:example");
node.rate_limiters = Some(vec!["missing".to_owned()]);
let error = resolve_admission_binding(&node, &BTreeMap::new(), None)
.expect_err("unknown limiter must fail");
assert!(error.contains("does not name an effective limiter"));
}
#[test]
fn admission_resolution_rejects_multiple_explicit_names() {
let mut node = NodeUserConfig::new_receiver_config("urn:test:receiver:example");
node.rate_limiters = Some(vec!["first".to_owned(), "second".to_owned()]);
let policy = admission_policy(RateLimitUnit::RequestBytes);
let policies =
BTreeMap::from([("first".to_owned(), policy), ("second".to_owned(), policy)]);
let error = resolve_admission_binding(&node, &policies, None)
.expect_err("multiple bindings must fail");
assert!(error.contains("at most one rate limiter binding"));
}
#[test]
fn admission_resolution_requires_explicit_binding() {
let node = NodeUserConfig::new_receiver_config("urn:test:receiver:example");
let policies = BTreeMap::from([(
"ingress".to_owned(),
admission_policy(RateLimitUnit::RequestBytes),
)]);
let binder =
resolve_admission_binding(&node, &policies, Some(RateLimiterDeclarationScope::Engine))
.expect("omitted binding resolves");
assert!(!binder.is_configured());
}
#[test]
fn admission_resolution_ignores_unselected_limiters() {
let node = NodeUserConfig::new_receiver_config("urn:test:receiver:example");
let policy = admission_policy(RateLimitUnit::RequestBytes);
let policies =
BTreeMap::from([("first".to_owned(), policy), ("second".to_owned(), policy)]);
let binder =
resolve_admission_binding(&node, &policies, None).expect("omitted binding resolves");
assert!(!binder.is_configured());
}
#[test]
fn test_interests() {
assert_eq!(Interests::ACKS | Interests::NACKS, Interests::ACKS_OR_NACKS);
assert_eq!(
Interests::NODE_INPUT_METRICS | Interests::NODE_OUTPUT_METRICS,
Interests::NODE_METRICS
);
}
#[test]
fn test_extension_factory_named_factory() {
fn dummy_create(
_: &ExtensionContext,
_: otel_arrow_dfe_config::ExtensionId,
_: Arc<otel_arrow_dfe_config::extension::ExtensionUserConfig>,
_: &ExtensionConfig,
) -> Result<ExtensionBundle, otel_arrow_dfe_config::error::Error> {
unimplemented!()
}
fn dummy_validate(
_: &serde_json::Value,
) -> Result<(), otel_arrow_dfe_config::error::Error> {
Ok(())
}
let factory = ExtensionFactory {
name: "urn:test:example",
description: "test extension",
documentation_url: "",
capabilities: None,
create: dummy_create,
validate_config: dummy_validate,
};
assert_eq!(factory.name(), "urn:test:example");
let cloned = factory.clone();
assert_eq!(cloned.name(), "urn:test:example");
assert_eq!(cloned.description, "test extension");
}
#[test]
fn test_extension_already_exists_variant_name() {
let err = Error::ExtensionAlreadyExists {
extension: "dup_ext".into(),
};
assert_eq!(err.variant_name(), "ExtensionAlreadyExists");
}
#[test]
fn test_extension_factory_validate_config() {
fn dummy_create(
_: &ExtensionContext,
_: otel_arrow_dfe_config::ExtensionId,
_: Arc<otel_arrow_dfe_config::extension::ExtensionUserConfig>,
_: &ExtensionConfig,
) -> Result<ExtensionBundle, otel_arrow_dfe_config::error::Error> {
unimplemented!()
}
fn dummy_validate(
config: &serde_json::Value,
) -> Result<(), otel_arrow_dfe_config::error::Error> {
if config.is_null() {
Ok(())
} else {
Err(otel_arrow_dfe_config::error::Error::InvalidUserConfig {
error: "expected null".into(),
})
}
}
let factory = ExtensionFactory {
name: "urn:test:example",
description: "test",
documentation_url: "",
capabilities: None,
create: dummy_create,
validate_config: dummy_validate,
};
assert!((factory.validate_config)(&serde_json::Value::Null).is_ok());
assert!((factory.validate_config)(&serde_json::json!({"key": "val"})).is_err());
}
#[otel_arrow_dfe_telemetry_macros::metric_set(name = "test.extension.factory")]
#[derive(Debug, Default, Clone)]
struct FactoryTestMetrics {
#[metric(unit = "{tick}")]
ticks: otel_arrow_dfe_telemetry::instrument::Counter<u64>,
}
#[test]
fn test_extension_factory_create_receives_extension_context() {
use crate::extension::wrapper::ExtensionVariant;
use crate::testing::test_extension_ctx;
use otel_arrow_dfe_config::extension::{ExtensionUrn, ExtensionUserConfig};
use otel_arrow_dfe_telemetry::registry::EntityKey;
use std::cell::Cell;
thread_local! {
static REGISTERED_ENTITY: Cell<Option<EntityKey>> = const { Cell::new(None) };
}
fn entity_registering_create(
ext_ctx: &ExtensionContext,
name: otel_arrow_dfe_config::ExtensionId,
_: Arc<ExtensionUserConfig>,
_: &ExtensionConfig,
) -> Result<ExtensionBundle, otel_arrow_dfe_config::error::Error> {
let entity = ext_ctx.register_extension_entity(name, ExtensionVariant::Local);
REGISTERED_ENTITY.with(|cell| cell.set(Some(entity)));
Err(otel_arrow_dfe_config::error::Error::InvalidUserConfig {
error: "no-op factory".into(),
})
}
fn dummy_validate(
_: &serde_json::Value,
) -> Result<(), otel_arrow_dfe_config::error::Error> {
Ok(())
}
let factory = ExtensionFactory {
name: "urn:otel:extension:test_ext",
description: "test extension that registers an entity via ext_ctx",
documentation_url: "",
capabilities: None,
create: entity_registering_create,
validate_config: dummy_validate,
};
let (ctx, registry) = test_extension_ctx();
let entities_before = registry.entity_count();
let metrics_before = registry.metric_set_count();
let user_config = Arc::new(ExtensionUserConfig::with_type(ExtensionUrn::from(
"urn:otel:extension:test_ext",
)));
let ext_config = ExtensionConfig::with_control_channel_capacity("test_ext", 16);
let result = (factory.create)(&ctx, "test_ext".into(), user_config, &ext_config);
assert!(result.is_err());
assert_eq!(registry.entity_count(), entities_before + 1);
let entity_key = REGISTERED_ENTITY
.with(|cell| cell.take())
.expect("factory should have registered an entity via ext_ctx");
let _metrics = ctx.register_metric_set_for_entity::<FactoryTestMetrics>(entity_key);
assert_eq!(registry.metric_set_count(), metrics_before + 1);
}
}