1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
//! Set of traits defining the common properties between all types of nodes in the pipeline engine.
//!
//! Receivers, processors, and exporters implement the [`Node`] trait.
//! Receivers and processors implement the [`NodeWithPDataSender`] trait.
//! Processors and exporters implement the [`NodeWithPDataReceiver`] trait.
use crate::control::NodeControlMsg;
use crate::effect_handler::SourceTagging;
use crate::error::Error;
use crate::message::{Receiver, Sender};
use otel_arrow_dfe_channel::error::SendError;
use otel_arrow_dfe_config::PortName;
use otel_arrow_dfe_config::node::NodeUserConfig;
use std::marker::PhantomData;
use std::sync::Arc;
pub use otel_arrow_dfe_config::NodeId as NodeName;
/// Common trait for nodes in the pipeline.
#[async_trait::async_trait(?Send)]
pub trait Node<PData> {
/// Flag indicating whether the node is shared (true) or local (false).
#[must_use]
fn is_shared(&self) -> bool;
/// Node identifier.
fn node_id(&self) -> NodeId;
/// Returns a reference to the node's user configuration.
#[must_use]
fn user_config(&self) -> Arc<NodeUserConfig>;
/// Sends a control message to the node.
async fn send_control_msg(
&self,
msg: NodeControlMsg<PData>,
) -> Result<(), SendError<NodeControlMsg<PData>>>;
}
/// NodeId consists of a unique integer index and a name.
#[derive(Clone, Debug)]
pub struct NodeId {
/// A unique integer.
pub index: usize,
/// A unique name as defined by otel_arrow_dfe_config.
pub name: NodeName,
}
/// Enum to identify the type of a node for registry lookups
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum NodeType {
/// Represents a node that acts as a receiver, receiving data from an external source.
Receiver,
/// Represents a node that processes data, transforming or analyzing it.
Processor,
/// Represents a node that exports data to an external destination.
Exporter,
}
/// Trait for nodes that can send pdata to a specific port.
pub trait NodeWithPDataSender<PData>: Node<PData> {
/// Sets the sender for pdata messages on the node.
fn set_pdata_sender(
&mut self,
node: NodeId,
port: PortName,
sender: Sender<PData>,
) -> Result<(), Error>;
/// Marks this node as needing to tag outgoing messages with the source node id.
/// Called by the pipeline wiring when the destination node has multiple input sources.
fn set_source_tagging(&mut self, source_tag: SourceTagging);
}
/// Trait for nodes that can receive pdata.
pub trait NodeWithPDataReceiver<PData>: Node<PData> {
/// Sets the receiver for pdata messages on the node.
fn set_pdata_receiver(&mut self, node: NodeId, receiver: Receiver<PData>) -> Result<(), Error>;
}
/// NodeDefinition is an entry in NodeDefs, indexed by the corresponding NodeIndex assignment.
pub(crate) struct NodeDefinition<Inner> {
/// Type of node.
pub(crate) ntype: NodeType,
// Node name.
pub(crate) name: NodeName,
/// Inner data.
pub(crate) inner: Inner,
}
/// NodeDefs is a NodeIndex-indexed set of node definitions.
pub(crate) struct NodeDefs<PData, Inner> {
/// Entries have an implicit index equal to their NodeIndex value.
entries: Vec<NodeDefinition<Inner>>,
_data: PhantomData<PData>,
}
impl<PData, Inner> Default for NodeDefs<PData, Inner> {
fn default() -> Self {
Self {
entries: Vec::new(),
_data: PhantomData,
}
}
}
impl<PData, Inner> NodeDefs<PData, Inner> {
/// Gets a the node definition
#[must_use]
pub(crate) fn get(&self, index: usize) -> Option<&NodeDefinition<Inner>> {
self.entries.get(index)
}
/// Gets the next unique node identifier. Returns an error when
/// the node limit (65,535) is exceeded.
pub fn next(&mut self, name: NodeName, ntype: NodeType, inner: Inner) -> Result<NodeId, Error> {
let index = self.entries.len();
if index > u16::MAX as usize {
return Err(Error::TooManyNodes {});
}
let uniq = NodeId::build(index, name.clone());
self.entries.push(NodeDefinition { ntype, name, inner });
Ok(uniq)
}
/// Returns an iterator over NodeId values for this set.
pub(crate) fn iter(&self) -> impl Iterator<Item = (NodeId, &NodeDefinition<Inner>)> {
self.entries.iter().enumerate().map(|(idx, val)| {
(
NodeId {
name: val.name.clone(),
index: idx,
},
val,
)
})
}
}
impl NodeId {
pub(crate) const fn build(index: usize, name: NodeName) -> NodeId {
NodeId { index, name }
}
}
impl std::fmt::Display for NodeId {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{}", self.name)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn test_too_many_nodes_error() {
let mut node_defs: NodeDefs<(), ()> = NodeDefs::default();
let name: NodeName = "test_node".into();
const LIMIT: usize = u16::MAX as usize + 1;
for i in 0..=LIMIT {
let result = node_defs.next(name.clone(), NodeType::Processor, ());
if i == LIMIT {
// This should fail with TooManyNodes error
assert!(matches!(result, Err(Error::TooManyNodes {})));
break;
} else {
assert!(result.is_ok());
}
}
}
}