use otel_arrow_dfe_telemetry::attributes::{
AttributeKeySchema, AttributeSetHandler, AttributeSetKeySchema, AttributeValue,
};
use otel_arrow_dfe_telemetry::descriptor::{
AttributeField, AttributeValueType, AttributesDescriptor,
};
use otel_arrow_dfe_telemetry_macros::{AttributeEnum, attribute_set};
use std::borrow::Cow;
use std::collections::BTreeMap;
use std::hash::Hash;
#[must_use]
pub fn config_to_telemetry_attr(
value: &otel_arrow_dfe_config::pipeline::telemetry::AttributeValue,
) -> AttributeValue {
use otel_arrow_dfe_config::pipeline::telemetry::AttributeValue as ConfigValue;
match value {
ConfigValue::String(s) => AttributeValue::String(s.clone()),
ConfigValue::Bool(b) => AttributeValue::Boolean(*b),
ConfigValue::I64(i) => AttributeValue::Int(*i),
ConfigValue::F64(f) => AttributeValue::Double(*f),
ConfigValue::Array(arr) => {
AttributeValue::String(format!("{:?}", arr))
}
}
}
#[must_use]
pub fn config_map_to_telemetry(
map: &std::collections::HashMap<
String,
otel_arrow_dfe_config::pipeline::telemetry::TelemetryAttribute,
>,
) -> BTreeMap<String, AttributeValue> {
map.iter()
.map(|(k, attr)| (k.clone(), config_to_telemetry_attr(attr.value())))
.collect()
}
#[attribute_set(scope, name = "controller.attrs")]
#[derive(Debug, Clone, Default, Hash)]
pub struct EngineAttributeSet {
pub core_id: usize,
pub numa_node_id: usize,
}
static ENGINE_ENTITY_DESCRIPTOR: AttributesDescriptor = AttributesDescriptor {
name: "engine",
fields: &[],
};
#[derive(Debug, Clone, Default, Hash)]
pub struct EngineEntityAttributeSet;
impl AttributeSetHandler for EngineEntityAttributeSet {
fn descriptor(&self) -> &'static AttributesDescriptor {
&ENGINE_ENTITY_DESCRIPTOR
}
fn attribute_values(&self) -> &[AttributeValue] {
&[]
}
}
#[attribute_set(scope, name = "pipeline.attrs")]
#[derive(Debug, Clone, Default, Hash)]
pub struct PipelineAttributeSet {
pub pipeline_id: Cow<'static, str>,
#[compose]
pub engine_attrs: EngineAttributeSet,
pub pipeline_group_id: Cow<'static, str>,
pub deployment_generation: u64,
}
#[attribute_set(scope, name = "extension.scope.attrs")]
#[derive(Debug, Clone, Hash)]
pub struct ExtensionScopeAttributeSet {
#[attribute_key = "scope.kind"]
pub(crate) kind: Cow<'static, str>,
#[compose]
pub(crate) pipeline: PipelineAttributeSet,
}
impl Default for ExtensionScopeAttributeSet {
fn default() -> Self {
Self {
kind: Cow::Borrowed(""),
pipeline: PipelineAttributeSet::default(),
}
}
}
impl ExtensionScopeAttributeSet {
#[must_use]
pub fn pipeline(pipeline: PipelineAttributeSet) -> Self {
Self {
kind: Cow::Borrowed("pipeline"),
pipeline,
}
}
}
#[attribute_set(scope, name = "extension.attrs")]
#[derive(Debug, Clone, Default, Hash)]
pub struct ExtensionAttributeSet {
pub extension_id: Cow<'static, str>,
#[compose]
pub extension_scope: ExtensionScopeAttributeSet,
#[attribute_key = "extension.variant"]
pub extension_variant: Cow<'static, str>,
}
#[attribute_set(scope, name = "node.attrs")]
#[derive(Debug, Clone, Default, Hash)]
pub struct NodeAttributeSet {
pub node_id: Cow<'static, str>,
#[compose]
pub pipeline_attrs: PipelineAttributeSet,
#[attribute_key = "node.urn"]
pub node_urn: Cow<'static, str>,
pub node_type: Cow<'static, str>,
}
#[attribute_set(scope, name = "node.custom.attrs")]
#[derive(Debug, Clone, Default, Hash)]
pub struct NodeWithCustomAttributeSet {
#[compose]
pub node_attrs: NodeAttributeSet,
#[compose]
pub custom_attrs: CustomAttributeSet,
}
#[attribute_set(scope, name = "node.topic.attrs")]
#[derive(Debug, Clone, Default, Hash)]
pub struct NodeWithTopicAttributeSet {
#[compose]
pub node_attrs: NodeAttributeSet,
pub topic: Cow<'static, str>,
}
#[attribute_set(scope, name = "node.custom.topic.attrs")]
#[derive(Debug, Clone, Default, Hash)]
pub struct NodeWithCustomTopicAttributeSet {
#[compose]
pub node_custom_attrs: NodeWithCustomAttributeSet,
pub topic: Cow<'static, str>,
}
#[derive(Debug, Clone)]
pub struct CustomAttributeSet {
values: Vec<AttributeValue>,
}
impl Default for CustomAttributeSet {
fn default() -> Self {
Self {
values: vec![AttributeValue::Map(BTreeMap::new())],
}
}
}
impl Hash for CustomAttributeSet {
fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
self.values.len().hash(state);
for v in &self.values {
v.to_string_value().hash(state);
}
}
}
static CUSTOM_ATTRIBUTES_DESCRIPTOR: AttributesDescriptor = AttributesDescriptor {
name: "custom.attrs",
fields: &[AttributeField {
key: "custom",
brief: "Custom user-defined attributes",
r#type: AttributeValueType::Map,
}],
};
impl CustomAttributeSet {
#[must_use]
pub fn new(custom_attrs: BTreeMap<String, AttributeValue>) -> Self {
Self {
values: vec![AttributeValue::Map(custom_attrs)],
}
}
}
impl AttributeSetHandler for CustomAttributeSet {
fn descriptor(&self) -> &'static AttributesDescriptor {
&CUSTOM_ATTRIBUTES_DESCRIPTOR
}
fn attribute_values(&self) -> &[AttributeValue] {
&self.values
}
}
impl AttributeSetKeySchema for CustomAttributeSet {
const KEY_SCHEMA: &'static [AttributeKeySchema] = &[AttributeKeySchema::Key("custom")];
}
#[cfg(test)]
mod tests {
use super::*;
use otel_arrow_dfe_telemetry::attributes::{AttributeEnum, AttributeSetHandler};
#[test]
fn pipeline_scope_ids_are_unambiguous_across_group_pipeline_splits() {
let a = ExtensionScopeAttributeSet::pipeline(PipelineAttributeSet {
pipeline_group_id: "a/b".into(),
pipeline_id: "c".into(),
..PipelineAttributeSet::default()
});
let b = ExtensionScopeAttributeSet::pipeline(PipelineAttributeSet {
pipeline_group_id: "a".into(),
pipeline_id: "b/c".into(),
..PipelineAttributeSet::default()
});
let a_values = a.attribute_values().to_vec();
let b_values = b.attribute_values().to_vec();
assert_ne!(
a_values, b_values,
"distinct (group, pipeline) pairs must not collide on attribute values; \
flattening `{{group}}/{{pipeline}}` into one opaque string allows \
two real scopes to register the same telemetry entity"
);
}
#[test]
fn channel_attribute_enums_have_stable_values() {
assert_eq!(ChannelKind::CARDINALITY, 2);
assert_eq!(ChannelKind::VARIANTS, &["control", "pdata"]);
assert_eq!(ChannelMode::CARDINALITY, 2);
assert_eq!(ChannelMode::VARIANTS, &["local", "shared"]);
assert_eq!(ChannelType::CARDINALITY, 2);
assert_eq!(ChannelType::VARIANTS, &["mpsc", "mpmc"]);
assert_eq!(ChannelImplementation::CARDINALITY, 3);
assert_eq!(
ChannelImplementation::VARIANTS,
&["internal", "tokio", "flume"]
);
}
#[test]
fn node_channel_attribute_enums_serialize_as_scope_strings() {
let attrs = NodeChannelAttributeSet {
channel_id: "channel-a".into(),
node_attrs: NodeAttributeSet::default(),
node_port: "output".into(),
channel_kind: ChannelKind::Pdata,
channel_mode: ChannelMode::Shared,
channel_type: ChannelType::Mpmc,
channel_impl: ChannelImplementation::Flume,
};
let attr_map: BTreeMap<&'static str, String> = attrs
.iter_attributes()
.map(|(key, value)| (key, value.to_string_value()))
.collect();
assert_eq!(
attr_map.get("channel.kind").map(String::as_str),
Some("pdata")
);
assert_eq!(
attr_map.get("channel.mode").map(String::as_str),
Some("shared")
);
assert_eq!(
attr_map.get("channel.type").map(String::as_str),
Some("mpmc")
);
assert_eq!(
attr_map.get("channel.impl").map(String::as_str),
Some("flume")
);
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, AttributeEnum)]
pub enum ChannelKind {
Control,
Pdata,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, AttributeEnum)]
pub enum ChannelMode {
Local,
Shared,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, AttributeEnum)]
pub enum ChannelType {
Mpsc,
Mpmc,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, AttributeEnum)]
pub enum ChannelImplementation {
Internal,
Tokio,
Flume,
}
#[attribute_set(scope, name = "node.channel.attrs")]
#[derive(Debug, Clone, Hash)]
pub struct NodeChannelAttributeSet {
#[attribute_key = "channel.id"]
pub channel_id: Cow<'static, str>,
#[compose]
pub node_attrs: NodeAttributeSet,
#[attribute_key = "node.port"]
pub node_port: Cow<'static, str>,
#[attribute_key = "channel.kind"]
pub channel_kind: ChannelKind,
#[attribute_key = "channel.mode"]
pub channel_mode: ChannelMode,
#[attribute_key = "channel.type"]
pub channel_type: ChannelType,
#[attribute_key = "channel.impl"]
pub channel_impl: ChannelImplementation,
}
#[attribute_set(scope, name = "node.channel.custom.attrs")]
#[derive(Debug, Clone, Hash)]
pub struct NodeWithCustomChannelAttributeSet {
#[compose]
pub channel_attrs: NodeChannelAttributeSet,
#[compose]
pub custom_attrs: CustomAttributeSet,
}
#[attribute_set(scope, name = "extension.channel.attrs")]
#[derive(Debug, Clone, Hash)]
pub struct ExtensionChannelAttributeSet {
#[attribute_key = "channel.id"]
pub channel_id: Cow<'static, str>,
#[compose]
pub extension_attrs: ExtensionAttributeSet,
#[attribute_key = "channel.mode"]
pub channel_mode: ChannelMode,
#[attribute_key = "channel.impl"]
pub channel_impl: ChannelImplementation,
}