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
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
//! Set of system configuration structures used by the engine, for example, to define channel sizes.
//!
//! Note: This type of system configuration is distinct from the pipeline configuration, which
//! focuses instead on defining the interconnection of nodes within the DAG and each node's specific
//! settings.
use otel_arrow_dfe_config::ExtensionId;
use otel_arrow_dfe_config::NodeId;
/// Configuration for the process-wide console output service.
///
/// The type is defined by the telemetry crate, which owns the writer threads,
/// and re-exported here so the engine's system configuration surface stays in
/// one place. The defaults work with no user configuration.
pub use otel_arrow_dfe_telemetry::output_service::OutputServiceConfig;
/// Default control channel capacity used by legacy constructor paths.
const DEFAULT_CONTROL_CHANNEL_CAPACITY: usize = 32;
/// Default pdata channel capacity used by legacy constructor paths.
const DEFAULT_PDATA_CHANNEL_CAPACITY: usize = 256;
/// Generic configuration for a control channel.
#[derive(Clone, Debug)]
pub struct ControlChannelConfig {
/// Max capacity of the channel.
pub capacity: usize,
}
/// Generic configuration for a pdata channel.
#[derive(Clone, Debug)]
pub struct PdataChannelConfig {
/// Max capacity of the channel.
pub capacity: usize,
}
/// Runtime configuration for a receiver.
#[derive(Clone, Debug)]
pub struct ReceiverConfig {
/// Name of the receiver.
pub name: NodeId,
/// Configuration for control channel.
pub control_channel: ControlChannelConfig,
/// Configuration for output pdata channel.
pub output_pdata_channel: PdataChannelConfig,
}
/// Generic configuration for a processor.
#[derive(Clone, Debug)]
pub struct ProcessorConfig {
/// Name of the processor.
pub name: NodeId,
/// Configuration for control channel.
pub control_channel: ControlChannelConfig,
/// Configuration for input pdata channel.
pub input_pdata_channel: PdataChannelConfig,
/// Configuration for output pdata channel.
pub output_pdata_channel: PdataChannelConfig,
}
/// Generic configuration for an exporter.
#[derive(Clone, Debug)]
pub struct ExporterConfig {
/// Name of the exporter.
pub name: NodeId,
/// Configuration for control channel.
pub control_channel: ControlChannelConfig,
/// Configuration for input pdata channel.
pub input_pdata_channel: PdataChannelConfig,
}
/// Runtime configuration for an extension.
///
/// Extensions only have a control channel -- they do not process pipeline data.
#[derive(Clone, Debug)]
pub struct ExtensionConfig {
/// Name of the extension.
pub name: ExtensionId,
/// Configuration for control channel.
pub control_channel: ControlChannelConfig,
}
impl ReceiverConfig {
/// Creates a new receiver configuration with default channel capacities.
pub fn new<T>(name: T) -> Self
where
T: Into<NodeId>,
{
Self::with_channel_capacities(
name,
DEFAULT_CONTROL_CHANNEL_CAPACITY,
DEFAULT_PDATA_CHANNEL_CAPACITY,
)
}
/// Creates a new receiver configuration with explicit channel capacities.
pub fn with_channel_capacities<T>(
name: T,
control_channel_capacity: usize,
pdata_channel_capacity: usize,
) -> Self
where
T: Into<NodeId>,
{
ReceiverConfig {
name: name.into(),
control_channel: ControlChannelConfig {
capacity: control_channel_capacity,
},
output_pdata_channel: PdataChannelConfig {
capacity: pdata_channel_capacity,
},
}
}
}
impl ProcessorConfig {
/// Creates a new processor configuration with default channel capacities.
#[must_use]
pub fn new<T>(name: T) -> Self
where
T: Into<NodeId>,
{
Self::with_channel_capacities(
name,
DEFAULT_CONTROL_CHANNEL_CAPACITY,
DEFAULT_PDATA_CHANNEL_CAPACITY,
)
}
/// Creates a new processor configuration with explicit channel capacities.
#[must_use]
pub fn with_channel_capacities<T>(
name: T,
control_channel_capacity: usize,
pdata_channel_capacity: usize,
) -> Self
where
T: Into<NodeId>,
{
ProcessorConfig {
name: name.into(),
control_channel: ControlChannelConfig {
capacity: control_channel_capacity,
},
input_pdata_channel: PdataChannelConfig {
capacity: pdata_channel_capacity,
},
output_pdata_channel: PdataChannelConfig {
capacity: pdata_channel_capacity,
},
}
}
}
impl ExporterConfig {
/// Creates a new exporter configuration with default channel capacities.
#[must_use]
pub fn new<T>(name: T) -> Self
where
T: Into<NodeId>,
{
Self::with_channel_capacities(
name,
DEFAULT_CONTROL_CHANNEL_CAPACITY,
DEFAULT_PDATA_CHANNEL_CAPACITY,
)
}
/// Creates a new exporter configuration with explicit channel capacities.
#[must_use]
pub fn with_channel_capacities<T>(
name: T,
control_channel_capacity: usize,
pdata_channel_capacity: usize,
) -> Self
where
T: Into<NodeId>,
{
ExporterConfig {
name: name.into(),
control_channel: ControlChannelConfig {
capacity: control_channel_capacity,
},
input_pdata_channel: PdataChannelConfig {
capacity: pdata_channel_capacity,
},
}
}
}
impl ExtensionConfig {
/// Creates a new extension configuration with default channel capacities.
#[must_use]
pub fn new<T>(name: T) -> Self
where
T: Into<ExtensionId>,
{
Self::with_control_channel_capacity(name, DEFAULT_CONTROL_CHANNEL_CAPACITY)
}
/// Creates a new extension configuration with explicit control channel capacity.
#[must_use]
pub fn with_control_channel_capacity<T>(name: T, control_channel_capacity: usize) -> Self
where
T: Into<ExtensionId>,
{
ExtensionConfig {
name: name.into(),
control_channel: ControlChannelConfig {
capacity: control_channel_capacity,
},
}
}
}