use crate::Interests;
use crate::attributes::{
ChannelImplementation, ChannelKind, ChannelMode, ChannelType, CustomAttributeSet,
EngineAttributeSet, EngineEntityAttributeSet, ExtensionAttributeSet,
ExtensionChannelAttributeSet, ExtensionScopeAttributeSet, NodeChannelAttributeSet,
NodeWithCustomChannelAttributeSet, NodeWithCustomTopicAttributeSet, NodeWithTopicAttributeSet,
PipelineAttributeSet, config_map_to_telemetry,
};
pub use crate::attributes::{NodeAttributeSet, NodeWithCustomAttributeSet};
use crate::context_declaration::CompiledContextBindings;
use crate::entity_context::{current_node_telemetry_handle, node_entity_key};
use crate::listener_group::ListenerGroupSnapshot;
use crate::memory_limiter::MemoryPressureState;
use crate::node::NodeId as EngineNodeId;
use data_encoding::BASE32_NOPAD;
use otel_arrow_dfe_config::node::NodeKind;
use otel_arrow_dfe_config::pipeline::telemetry::TelemetryAttribute;
use otel_arrow_dfe_config::policy::DistributionTier;
use otel_arrow_dfe_config::{
NodeId as ConfigNodeId, NodeUrn, PipelineGroupId, PipelineId, PipelineKey,
};
use otel_arrow_dfe_telemetry::InternalTelemetrySettings;
use otel_arrow_dfe_telemetry::attributes::AttributeSetHandler;
use otel_arrow_dfe_telemetry::metrics::MetricSetRegistrar;
use otel_arrow_dfe_telemetry::metrics::{
MeasurementMetricSet, MeasurementMetricSetHandler, MetricSet, MetricSetHandler,
RegistrationMetricSetHandler,
};
use otel_arrow_dfe_telemetry::registry::{EntityKey, MetricSetKey, TelemetryRegistryHandle};
use std::any::Any;
use std::borrow::Cow;
use std::collections::HashMap;
use std::fmt::Debug;
use std::sync::{Arc, LazyLock};
use uuid::Uuid;
pub type NodeNameIndex = Arc<HashMap<ConfigNodeId, EngineNodeId>>;
static PROCESS_INSTANCE_ID: LazyLock<Cow<'static, str>> = LazyLock::new(|| {
let uuid = Uuid::now_v7();
Cow::Owned(BASE32_NOPAD.encode(uuid.as_bytes()))
});
fn detect_host_id() -> Option<String> {
if let Ok(host) = std::env::var("HOSTNAME")
&& !host.is_empty()
{
return Some(host);
}
std::fs::read_to_string("/etc/hostname")
.ok()
.map(|host| host.trim().to_owned())
.filter(|host| !host.is_empty())
}
fn detect_container_id() -> Option<String> {
let cgroup = std::fs::read_to_string("/proc/self/cgroup").ok()?;
for line in cgroup.lines() {
let path = line.split(':').nth(2).unwrap_or("");
for part in path.split('/') {
let token = part.trim();
if (32..=128).contains(&token.len())
&& token
.chars()
.all(|c| c.is_ascii_hexdigit() || matches!(c, '.' | '-' | '_'))
&& token.chars().filter(|c| c.is_ascii_hexdigit()).count() >= 32
{
return Some(token.to_owned());
}
}
}
None
}
static HOST_ID: LazyLock<Cow<'static, str>> =
LazyLock::new(|| detect_host_id().map_or(Cow::Borrowed(""), Cow::Owned));
static CONTAINER_ID: LazyLock<Cow<'static, str>> =
LazyLock::new(|| detect_container_id().map_or(Cow::Borrowed(""), Cow::Owned));
#[derive(Clone, Debug)]
pub struct ControllerContext {
telemetry_registry_handle: TelemetryRegistryHandle,
process_instance_id: Cow<'static, str>,
host_id: Cow<'static, str>,
container_id: Cow<'static, str>,
memory_pressure_state: MemoryPressureState,
state_directory: Option<crate::state_dir::StateDirectory>,
}
#[derive(Clone, Debug)]
pub struct PipelineContextParams {
pub pipeline_group_id: PipelineGroupId,
pub pipeline_id: PipelineId,
pub core_id: usize,
pub num_cores: usize,
pub thread_id: usize,
pub numa_node_id: usize,
}
#[derive(Clone, Debug)]
pub struct PipelineContext {
controller_context: ControllerContext,
pipeline_context_params: PipelineContextParams,
deployment_generation: u64,
pipeline_telemetry_attrs: HashMap<String, TelemetryAttribute>,
node_id: ConfigNodeId,
node_urn: NodeUrn,
node_kind: NodeKind,
node_interests: Interests,
node_duration_distribution: DistributionTier,
node_telemetry_attrs: HashMap<String, TelemetryAttribute>,
admission: crate::admission::AdmissionBinder,
internal_telemetry: Option<InternalTelemetrySettings>,
node_names: NodeNameIndex,
topic_set: Option<Arc<dyn Any + Send + Sync>>,
listener_group_snapshot: Arc<ListenerGroupSnapshot>,
compiled_context_bindings: Arc<CompiledContextBindings>,
}
pub struct EntityMetricSetRegistrar<'a> {
pipeline_context: &'a PipelineContext,
entity_key: EntityKey,
}
impl ControllerContext {
#[must_use]
pub fn with_state_directory(mut self, root: crate::state_dir::StateDirectory) -> Self {
self.state_directory = Some(root);
self
}
#[must_use]
pub fn state_directory(&self) -> Option<&crate::state_dir::StateDirectory> {
self.state_directory.as_ref()
}
#[must_use]
pub fn new(telemetry_registry_handle: TelemetryRegistryHandle) -> Self {
Self {
telemetry_registry_handle,
process_instance_id: PROCESS_INSTANCE_ID.clone(),
host_id: HOST_ID.clone(),
container_id: CONTAINER_ID.clone(),
memory_pressure_state: MemoryPressureState::default(),
state_directory: None,
}
}
#[cfg(test)]
fn new_with_identity(
telemetry_registry_handle: TelemetryRegistryHandle,
process_instance_id: impl Into<Cow<'static, str>>,
host_id: impl Into<Cow<'static, str>>,
container_id: impl Into<Cow<'static, str>>,
) -> Self {
Self {
telemetry_registry_handle,
process_instance_id: process_instance_id.into(),
host_id: host_id.into(),
container_id: container_id.into(),
memory_pressure_state: MemoryPressureState::default(),
state_directory: None,
}
}
#[must_use]
pub fn pipeline_context_with(
&self,
pipeline_group_id: PipelineGroupId,
pipeline_id: PipelineId,
core_id: usize,
num_cores: usize,
thread_id: usize,
) -> PipelineContext {
self.pipeline_context_with_generation(
pipeline_group_id,
pipeline_id,
core_id,
num_cores,
thread_id,
0,
)
}
#[must_use]
pub fn pipeline_context_with_generation(
&self,
pipeline_group_id: PipelineGroupId,
pipeline_id: PipelineId,
core_id: usize,
num_cores: usize,
thread_id: usize,
deployment_generation: u64,
) -> PipelineContext {
PipelineContext::new_with_generation(
self.clone(),
PipelineContextParams {
pipeline_group_id,
pipeline_id,
core_id,
num_cores,
thread_id,
numa_node_id: 0,
},
deployment_generation,
)
}
#[must_use]
pub fn pipeline_context_with_placement(
&self,
pipeline_group_id: PipelineGroupId,
pipeline_id: PipelineId,
core_id: usize,
num_cores: usize,
thread_id: usize,
deployment_generation: u64,
numa_node_id: usize,
) -> PipelineContext {
PipelineContext::new_with_generation(
self.clone(),
PipelineContextParams {
pipeline_group_id,
pipeline_id,
core_id,
num_cores,
thread_id,
numa_node_id,
},
deployment_generation,
)
}
#[must_use]
pub fn register_engine_entity(&self) -> EntityKey {
self.telemetry_registry_handle
.register_entity(EngineEntityAttributeSet)
}
#[must_use]
pub fn resource_attributes(&self) -> Vec<(String, String)> {
let mut attributes = Vec::new();
if !self.host_id.is_empty() {
attributes.push(("host.id".to_owned(), self.host_id.to_string()));
}
if !self.container_id.is_empty() {
attributes.push(("container.id".to_owned(), self.container_id.to_string()));
}
if !self.process_instance_id.is_empty() {
attributes.push((
"service.instance.id".to_owned(),
self.process_instance_id.to_string(),
));
}
attributes
}
#[must_use]
pub fn telemetry_registry(&self) -> TelemetryRegistryHandle {
self.telemetry_registry_handle.clone()
}
#[must_use]
pub fn memory_pressure_state(&self) -> MemoryPressureState {
self.memory_pressure_state.clone()
}
}
impl From<&PipelineContextParams> for PipelineKey {
fn from(params: &PipelineContextParams) -> Self {
PipelineKey::new(params.pipeline_group_id.clone(), params.pipeline_id.clone())
}
}
impl PipelineContext {
#[must_use]
pub fn state_directory(&self) -> Option<&crate::state_dir::StateDirectory> {
self.controller_context.state_directory()
}
#[allow(dead_code)]
pub(crate) fn new(
parent_ctx: ControllerContext,
pipeline_context_params: PipelineContextParams,
) -> Self {
Self::new_with_generation(parent_ctx, pipeline_context_params, 0)
}
pub(crate) fn new_with_generation(
parent_ctx: ControllerContext,
pipeline_context_params: PipelineContextParams,
deployment_generation: u64,
) -> Self {
Self {
controller_context: parent_ctx,
pipeline_context_params,
deployment_generation,
node_id: Default::default(),
node_urn: Default::default(),
node_kind: Default::default(),
node_interests: Interests::empty(),
node_duration_distribution: DistributionTier::Normal,
node_telemetry_attrs: HashMap::new(),
admission: crate::admission::AdmissionBinder::none(),
pipeline_telemetry_attrs: HashMap::new(),
internal_telemetry: None,
node_names: Arc::new(HashMap::new()),
topic_set: None,
listener_group_snapshot: Arc::new(ListenerGroupSnapshot::empty()),
compiled_context_bindings: Arc::new(CompiledContextBindings::empty()),
}
}
#[must_use]
pub fn pipeline_group_id(&self) -> PipelineGroupId {
self.pipeline_context_params.pipeline_group_id.clone()
}
#[must_use]
pub fn pipeline_id(&self) -> PipelineId {
self.pipeline_context_params.pipeline_id.clone()
}
#[must_use]
pub fn pipeline_key(&self) -> PipelineKey {
PipelineKey::from(&self.pipeline_context_params)
}
#[must_use]
pub fn node_id(&self) -> ConfigNodeId {
self.node_id.clone()
}
#[must_use]
pub const fn core_id(&self) -> usize {
self.pipeline_context_params.core_id
}
#[must_use]
pub const fn deployment_generation(&self) -> u64 {
self.deployment_generation
}
#[must_use]
pub const fn num_cores(&self) -> usize {
self.pipeline_context_params.num_cores
}
pub fn set_internal_telemetry(&mut self, settings: InternalTelemetrySettings) {
self.internal_telemetry = Some(settings);
}
#[must_use]
pub const fn internal_telemetry(&self) -> Option<&InternalTelemetrySettings> {
self.internal_telemetry.as_ref()
}
#[must_use]
pub fn memory_pressure_state(&self) -> MemoryPressureState {
self.controller_context.memory_pressure_state()
}
#[must_use]
pub const fn admission(&self) -> &crate::admission::AdmissionBinder {
&self.admission
}
pub(crate) fn set_admission(&mut self, admission: crate::admission::AdmissionBinder) {
self.admission = admission;
}
pub fn set_node_names(&mut self, node_names: NodeNameIndex) {
self.node_names = node_names;
}
pub fn set_topic_set<T: Send + Sync + 'static>(
&mut self,
topic_set: crate::topic::TopicSet<T>,
) {
self.topic_set = Some(Arc::new(topic_set));
}
pub fn set_listener_group_snapshot(&mut self, snapshot: ListenerGroupSnapshot) {
self.listener_group_snapshot = Arc::new(snapshot);
}
pub fn set_listener_group_snapshot_arc(&mut self, snapshot: Arc<ListenerGroupSnapshot>) {
self.listener_group_snapshot = snapshot;
}
#[must_use]
pub fn listener_group_snapshot(&self) -> Arc<ListenerGroupSnapshot> {
Arc::clone(&self.listener_group_snapshot)
}
pub fn set_compiled_context_bindings(&mut self, bindings: Arc<CompiledContextBindings>) {
self.compiled_context_bindings = bindings;
}
#[must_use]
pub fn compiled_context_bindings(&self) -> &Arc<CompiledContextBindings> {
&self.compiled_context_bindings
}
#[must_use]
pub fn topic_set<T: Send + Sync + 'static>(&self) -> Option<crate::topic::TopicSet<T>> {
self.topic_set
.as_ref()
.and_then(|resource| resource.downcast_ref::<crate::topic::TopicSet<T>>())
.cloned()
}
#[must_use]
pub fn node_by_name(&self, name: &str) -> Option<EngineNodeId> {
self.node_names.get(name).cloned()
}
#[must_use]
pub const fn take_internal_telemetry(&mut self) -> Option<InternalTelemetrySettings> {
self.internal_telemetry.take()
}
#[must_use]
pub fn register_metric_set_for_entity<T: MetricSetHandler + Default + Debug + Send + Sync>(
&self,
entity_key: EntityKey,
) -> MetricSet<T> {
let metrics = self
.controller_context
.telemetry_registry_handle
.register_metric_set_for_entity::<T>(entity_key);
if let Some(telemetry) = current_node_telemetry_handle() {
telemetry.track_metric_set(metrics.metric_set_key());
}
metrics
}
#[must_use]
pub fn register_measurement_metric_set_for_entity<
T: MeasurementMetricSetHandler + Debug + Send + Sync,
>(
&self,
entity_key: EntityKey,
) -> MeasurementMetricSet<T> {
let metrics = self
.controller_context
.telemetry_registry_handle
.register_metric_set_with_measurement_attributes_for_entity::<T>(entity_key);
if let Some(telemetry) = current_node_telemetry_handle() {
telemetry.track_metric_set(metrics.metric_set_key());
}
metrics
}
#[must_use]
pub const fn metric_set_registrar_for_entity(
&self,
entity_key: EntityKey,
) -> EntityMetricSetRegistrar<'_> {
EntityMetricSetRegistrar {
pipeline_context: self,
entity_key,
}
}
#[must_use]
pub fn metric_set_registrar_with_topic(
&self,
topic: Cow<'static, str>,
) -> EntityMetricSetRegistrar<'_> {
let entity_key = self.register_topic_entity(topic);
if let Some(telemetry) = current_node_telemetry_handle() {
telemetry.track_entity(entity_key);
}
self.metric_set_registrar_for_entity(entity_key)
}
#[must_use]
pub fn register_entity(
&self,
attributes: impl AttributeSetHandler + Send + Sync + 'static,
) -> EntityKey {
let entity_key = self
.controller_context
.telemetry_registry_handle
.register_entity(attributes);
if let Some(telemetry) = current_node_telemetry_handle() {
telemetry.track_entity(entity_key);
}
entity_key
}
fn register_scoped_metrics<R>(
&self,
for_entity: impl FnOnce(&TelemetryRegistryHandle, EntityKey) -> R,
metric_set_key: impl FnOnce(&R) -> MetricSetKey,
with_scope: impl FnOnce(&Self, &TelemetryRegistryHandle) -> R,
) -> R {
let handle = &self.controller_context.telemetry_registry_handle;
if let Some(telemetry) = current_node_telemetry_handle() {
let metrics = for_entity(handle, telemetry.entity_key());
telemetry.track_metric_set(metric_set_key(&metrics));
metrics
} else if let Some(entity_key) = node_entity_key() {
for_entity(handle, entity_key)
} else {
#[cfg(any(test, feature = "test-utils"))]
{
with_scope(self, handle)
}
#[cfg(not(any(test, feature = "test-utils")))]
{
let _ = with_scope;
panic!(
"node entity key not set; ensure node entity is registered and instrumented"
);
}
}
}
fn register_topic_entity(&self, topic: Cow<'static, str>) -> EntityKey {
if self.node_telemetry_attrs.is_empty() {
self.controller_context
.telemetry_registry_handle
.register_entity(NodeWithTopicAttributeSet {
node_attrs: self.node_attribute_set(),
topic,
})
} else {
self.controller_context
.telemetry_registry_handle
.register_entity(NodeWithCustomTopicAttributeSet {
node_custom_attrs: self.node_with_custom_attribute_set(),
topic,
})
}
}
#[must_use]
pub fn register_pipeline_entity(&self) -> EntityKey {
self.controller_context
.telemetry_registry_handle
.register_entity(self.pipeline_attribute_set())
}
#[must_use]
pub fn extension_context(&self) -> ExtensionContext {
ExtensionContext::new(
self.controller_context.clone(),
ExtensionScopeAttributeSet::pipeline(self.pipeline_attribute_set()),
)
}
#[must_use]
pub fn register_node_entity(&self) -> EntityKey {
if self.node_telemetry_attrs.is_empty() {
self.controller_context
.telemetry_registry_handle
.register_entity(self.node_attribute_set())
} else {
self.controller_context
.telemetry_registry_handle
.register_entity(self.node_with_custom_attribute_set())
}
}
fn engine_attribute_set(&self) -> EngineAttributeSet {
EngineAttributeSet {
core_id: self.pipeline_context_params.core_id,
numa_node_id: self.pipeline_context_params.numa_node_id,
}
}
#[must_use]
pub fn pipeline_attribute_set(&self) -> PipelineAttributeSet {
PipelineAttributeSet {
engine_attrs: self.engine_attribute_set(),
pipeline_id: self.pipeline_context_params.pipeline_id.clone(),
pipeline_group_id: self.pipeline_context_params.pipeline_group_id.clone(),
deployment_generation: self.deployment_generation,
}
}
#[must_use]
pub fn node_attribute_set(&self) -> NodeAttributeSet {
NodeAttributeSet {
pipeline_attrs: self.pipeline_attribute_set(),
node_id: self.node_id.clone(),
node_urn: self.node_urn.clone().into(),
node_type: self.node_kind.into(),
}
}
#[must_use]
pub fn has_custom_node_attributes(&self) -> bool {
!self.node_telemetry_attrs.is_empty()
}
#[must_use]
pub fn node_with_custom_attribute_set(&self) -> NodeWithCustomAttributeSet {
NodeWithCustomAttributeSet {
node_attrs: self.node_attribute_set(),
custom_attrs: CustomAttributeSet::new(config_map_to_telemetry(
&self.node_telemetry_attrs,
)),
}
}
#[must_use]
pub fn node_channel_attribute_set(
&self,
channel_id: Cow<'static, str>,
node_port: Cow<'static, str>,
channel_kind: ChannelKind,
channel_mode: ChannelMode,
channel_type: ChannelType,
channel_impl: ChannelImplementation,
) -> NodeChannelAttributeSet {
NodeChannelAttributeSet {
node_attrs: self.node_attribute_set(),
node_port,
channel_id,
channel_kind,
channel_mode,
channel_type,
channel_impl,
}
}
#[must_use]
pub fn register_node_channel_entity(
&self,
channel_id: Cow<'static, str>,
node_port: Cow<'static, str>,
channel_kind: ChannelKind,
channel_mode: ChannelMode,
channel_type: ChannelType,
channel_impl: ChannelImplementation,
) -> EntityKey {
let attrs = self.node_channel_attribute_set(
channel_id,
node_port,
channel_kind,
channel_mode,
channel_type,
channel_impl,
);
let registry = &self.controller_context.telemetry_registry_handle;
if self.node_telemetry_attrs.is_empty() {
registry.register_entity(attrs)
} else {
registry.register_entity(NodeWithCustomChannelAttributeSet {
channel_attrs: attrs,
custom_attrs: CustomAttributeSet::new(config_map_to_telemetry(
&self.node_telemetry_attrs,
)),
})
}
}
#[must_use]
pub fn metrics_registry(&self) -> TelemetryRegistryHandle {
self.controller_context.telemetry_registry_handle.clone()
}
#[must_use]
pub const fn node_interests(&self) -> Interests {
self.node_interests
}
pub(crate) fn set_node_interests(&mut self, interests: Interests) {
self.node_interests = interests;
}
#[must_use]
pub const fn node_duration_distribution(&self) -> DistributionTier {
self.node_duration_distribution
}
pub(crate) fn set_node_duration_distribution(&mut self, tier: DistributionTier) {
self.node_duration_distribution = tier;
}
#[must_use]
pub fn with_node_context(
&self,
node_id: ConfigNodeId,
node_urn: NodeUrn,
node_kind: NodeKind,
node_telemetry_attrs: HashMap<String, TelemetryAttribute>,
) -> Self {
Self {
controller_context: self.controller_context.clone(),
pipeline_context_params: self.pipeline_context_params.clone(),
deployment_generation: self.deployment_generation,
pipeline_telemetry_attrs: self.pipeline_telemetry_attrs.clone(),
node_id,
node_urn,
node_kind,
node_interests: Interests::empty(),
node_duration_distribution: DistributionTier::Normal,
node_telemetry_attrs,
admission: crate::admission::AdmissionBinder::none(),
internal_telemetry: None,
node_names: self.node_names.clone(),
topic_set: self.topic_set.clone(),
listener_group_snapshot: Arc::clone(&self.listener_group_snapshot),
compiled_context_bindings: self.compiled_context_bindings.clone(),
}
}
}
impl MetricSetRegistrar for EntityMetricSetRegistrar<'_> {
fn register_metric_set<M: MetricSetHandler + Default + Debug + Send + Sync>(
&self,
) -> MetricSet<M> {
self.pipeline_context
.register_metric_set_for_entity(self.entity_key)
}
fn register_registration_metric_set<M: RegistrationMetricSetHandler + Debug + Send + Sync>(
&self,
registration_attrs: &M::RegistrationAttributes,
) -> MetricSet<M> {
let metrics = self
.pipeline_context
.controller_context
.telemetry_registry_handle
.register_metric_set_with_registration_attributes_for_entity::<M>(
self.entity_key,
registration_attrs,
);
if let Some(telemetry) = current_node_telemetry_handle() {
telemetry.track_metric_set(metrics.metric_set_key());
}
metrics
}
fn register_measurement_metric_set<M: MeasurementMetricSetHandler + Debug + Send + Sync>(
&self,
) -> MeasurementMetricSet<M> {
self.pipeline_context
.register_measurement_metric_set_for_entity(self.entity_key)
}
fn register_registration_and_measurement_metric_set<
M: RegistrationMetricSetHandler + MeasurementMetricSetHandler + Debug + Send + Sync,
>(
&self,
registration_attrs: &M::RegistrationAttributes,
) -> MeasurementMetricSet<M> {
let metrics = self
.pipeline_context
.controller_context
.telemetry_registry_handle
.register_metric_set_with_registration_and_measurement_attributes_for_entity::<M>(
self.entity_key,
registration_attrs,
);
if let Some(telemetry) = current_node_telemetry_handle() {
telemetry.track_metric_set(metrics.metric_set_key());
}
metrics
}
}
impl MetricSetRegistrar for PipelineContext {
fn register_metric_set<M: MetricSetHandler + Default + Debug + Send + Sync>(
&self,
) -> MetricSet<M> {
self.register_scoped_metrics(
|handle, entity_key| handle.register_metric_set_for_entity::<M>(entity_key),
MetricSet::metric_set_key,
|ctx, handle| {
if ctx.node_telemetry_attrs.is_empty() {
handle.register_metric_set::<M>(ctx.node_attribute_set())
} else {
handle.register_metric_set::<M>(ctx.node_with_custom_attribute_set())
}
},
)
}
fn register_registration_metric_set<M: RegistrationMetricSetHandler + Debug + Send + Sync>(
&self,
registration_attrs: &M::RegistrationAttributes,
) -> MetricSet<M> {
self.register_scoped_metrics(
|handle, entity_key| {
handle.register_metric_set_with_registration_attributes_for_entity::<M>(
entity_key,
registration_attrs,
)
},
MetricSet::metric_set_key,
|ctx, handle| {
if ctx.node_telemetry_attrs.is_empty() {
handle.register_metric_set_with_registration_attributes::<M>(
ctx.node_attribute_set(),
registration_attrs,
)
} else {
handle.register_metric_set_with_registration_attributes::<M>(
ctx.node_with_custom_attribute_set(),
registration_attrs,
)
}
},
)
}
fn register_measurement_metric_set<M: MeasurementMetricSetHandler + Debug + Send + Sync>(
&self,
) -> MeasurementMetricSet<M> {
self.register_scoped_metrics(
|handle, entity_key| {
handle.register_metric_set_with_measurement_attributes_for_entity::<M>(entity_key)
},
MeasurementMetricSet::metric_set_key,
|ctx, handle| {
if ctx.node_telemetry_attrs.is_empty() {
handle.register_metric_set_with_measurement_attributes::<M>(
ctx.node_attribute_set(),
)
} else {
handle.register_metric_set_with_measurement_attributes::<M>(
ctx.node_with_custom_attribute_set(),
)
}
},
)
}
fn register_registration_and_measurement_metric_set<
M: RegistrationMetricSetHandler + MeasurementMetricSetHandler + Debug + Send + Sync,
>(
&self,
registration_attrs: &M::RegistrationAttributes,
) -> MeasurementMetricSet<M> {
self.register_scoped_metrics(
|handle, entity_key| {
handle
.register_metric_set_with_registration_and_measurement_attributes_for_entity::<M>(
entity_key,
registration_attrs,
)
},
MeasurementMetricSet::metric_set_key,
|ctx, handle| {
if ctx.node_telemetry_attrs.is_empty() {
handle.register_metric_set_with_registration_and_measurement_attributes::<M>(
ctx.node_attribute_set(),
registration_attrs,
)
} else {
handle.register_metric_set_with_registration_and_measurement_attributes::<M>(
ctx.node_with_custom_attribute_set(),
registration_attrs,
)
}
},
)
}
}
#[derive(Clone, Debug)]
pub struct ExtensionContext {
controller_context: ControllerContext,
extension_scope: ExtensionScopeAttributeSet,
}
impl ExtensionContext {
#[must_use]
pub fn new(
controller_context: ControllerContext,
extension_scope: ExtensionScopeAttributeSet,
) -> Self {
debug_assert!(
!extension_scope.kind.is_empty(),
"ExtensionContext requires a non-empty scope kind"
);
Self {
controller_context,
extension_scope,
}
}
#[must_use]
pub fn metrics_registry(&self) -> TelemetryRegistryHandle {
self.controller_context.telemetry_registry_handle.clone()
}
#[must_use]
pub fn extension_attribute_set(
&self,
extension_id: Cow<'static, str>,
variant: crate::extension::wrapper::ExtensionVariant,
) -> ExtensionAttributeSet {
ExtensionAttributeSet {
extension_id,
extension_variant: Cow::Borrowed(variant.as_str()),
extension_scope: self.extension_scope.clone(),
}
}
#[must_use]
pub fn register_extension_entity(
&self,
extension_id: Cow<'static, str>,
variant: crate::extension::wrapper::ExtensionVariant,
) -> EntityKey {
self.controller_context
.telemetry_registry_handle
.register_entity(self.extension_attribute_set(extension_id, variant))
}
#[must_use]
pub fn extension_channel_attribute_set(
&self,
extension_id: Cow<'static, str>,
variant: crate::extension::wrapper::ExtensionVariant,
channel_id: Cow<'static, str>,
channel_mode: ChannelMode,
channel_impl: ChannelImplementation,
) -> ExtensionChannelAttributeSet {
ExtensionChannelAttributeSet {
extension_attrs: self.extension_attribute_set(extension_id, variant),
channel_id,
channel_mode,
channel_impl,
}
}
#[must_use]
pub fn register_extension_channel_entity(
&self,
extension_id: Cow<'static, str>,
variant: crate::extension::wrapper::ExtensionVariant,
channel_id: Cow<'static, str>,
channel_mode: ChannelMode,
channel_impl: ChannelImplementation,
) -> EntityKey {
let attrs = self.extension_channel_attribute_set(
extension_id,
variant,
channel_id,
channel_mode,
channel_impl,
);
self.controller_context
.telemetry_registry_handle
.register_entity(attrs)
}
#[must_use]
pub fn register_metric_set_for_entity<T: MetricSetHandler + Default + Debug + Send + Sync>(
&self,
entity_key: EntityKey,
) -> MetricSet<T> {
self.controller_context
.telemetry_registry_handle
.register_metric_set_for_entity::<T>(entity_key)
}
}
#[cfg(test)]
mod tests {
use super::*;
use otel_arrow_dfe_config::pipeline::telemetry::AttributeValue;
use otel_arrow_dfe_telemetry::registry::TelemetryRegistryHandle;
use std::collections::HashMap;
#[test]
fn state_directory_is_optional_across_generations() {
let controller = ControllerContext::new(TelemetryRegistryHandle::new());
assert!(controller.state_directory().is_none());
for generation in [0, 1, 2] {
let context = controller.pipeline_context_with_generation(
"g".into(),
"p".into(),
generation as usize,
3,
generation as usize,
generation,
);
assert!(context.state_directory().is_none());
}
}
#[test]
fn resource_attributes_maps_semconv_keys() {
let ctx = ControllerContext::new_with_identity(
TelemetryRegistryHandle::new(),
"proc-123",
"machine-abc",
"container-xyz",
);
assert_eq!(
ctx.resource_attributes(),
vec![
("host.id".to_string(), "machine-abc".to_string()),
("container.id".to_string(), "container-xyz".to_string()),
("service.instance.id".to_string(), "proc-123".to_string()),
]
);
}
#[test]
fn resource_attributes_omits_empty_values() {
let ctx = ControllerContext::new_with_identity(
TelemetryRegistryHandle::new(),
"proc-123",
"",
"",
);
assert_eq!(
ctx.resource_attributes(),
vec![("service.instance.id".to_string(), "proc-123".to_string())]
);
}
#[test]
fn pipeline_context_uses_resolved_numa_node_for_engine_attrs() {
let ctx = ControllerContext::new(TelemetryRegistryHandle::new())
.pipeline_context_with_placement(
PipelineGroupId::from("g"),
PipelineId::from("p"),
4,
1,
0,
0,
2,
);
let attrs = ctx.pipeline_attribute_set();
assert_eq!(attrs.engine_attrs.core_id, 4);
assert_eq!(attrs.engine_attrs.numa_node_id, 2);
}
#[test]
fn pipeline_context_preserves_listener_group_snapshot_across_node_context() {
let addr = "127.0.0.1:4317".parse().unwrap();
let mut ctx = ControllerContext::new(TelemetryRegistryHandle::new()).pipeline_context_with(
Default::default(),
Default::default(),
0,
1,
0,
);
ctx.set_listener_group_snapshot(ListenerGroupSnapshot::new(
3,
vec![crate::listener_group::ListenerGroupPlan {
key: crate::listener_group::ListenerGroupKey::new(
"group".into(),
"pipeline".into(),
"receiver".into(),
addr,
crate::listener_group::ListenerProtocol::Tcp,
),
expected_members: Vec::new(),
}],
));
let node_ctx = ctx.with_node_context(
"receiver".into(),
"urn:otel:receiver:otlp".into(),
NodeKind::Receiver,
HashMap::new(),
);
assert_eq!(node_ctx.listener_group_snapshot().generation, 3);
assert!(
node_ctx
.listener_group_snapshot()
.plan_for(
"receiver",
addr,
crate::listener_group::ListenerProtocol::Tcp
)
.is_some()
);
}
#[test]
fn pipeline_context_preserves_compiled_policy_across_node_context() {
let resolved = otel_arrow_dfe_config::engine::ResolvedOtelDataflowSpec {
engine: Default::default(),
pipelines: Vec::new(),
};
let factory = crate::PipelineFactory::<()>::new(&[], &[], &[], &[]);
let bindings = factory
.compile_initial_context(&resolved)
.expect("context bindings")
.bindings;
let controller = ControllerContext::new(TelemetryRegistryHandle::new());
let mut pipeline =
controller.pipeline_context_with("group".into(), "pipeline".into(), 0, 1, 0);
pipeline.set_compiled_context_bindings(Arc::clone(&bindings));
let node = pipeline.with_node_context(
"node".into(),
"urn:otel:processor:test".into(),
NodeKind::Processor,
HashMap::new(),
);
assert!(Arc::ptr_eq(node.compiled_context_bindings(), &bindings));
}
fn pipeline_ctx_with_custom_attrs(
registry: TelemetryRegistryHandle,
custom: HashMap<String, TelemetryAttribute>,
) -> PipelineContext {
let controller_ctx = ControllerContext::new(registry);
let pipeline_params = PipelineContextParams {
pipeline_group_id: Cow::Borrowed("group1"),
pipeline_id: Cow::Borrowed("pipe1"),
core_id: 0,
num_cores: 1,
thread_id: 0,
numa_node_id: 0,
};
PipelineContext::new(controller_ctx, pipeline_params).with_node_context(
Cow::Borrowed("test-node"),
NodeUrn::parse("urn:otel:receiver:test").unwrap(),
NodeKind::Receiver,
custom,
)
}
fn register_channel(ctx: &PipelineContext) -> EntityKey {
ctx.register_node_channel_entity(
Cow::Borrowed("channel-1"),
Cow::Borrowed("out"),
ChannelKind::Pdata,
ChannelMode::Local,
ChannelType::Mpsc,
ChannelImplementation::Internal,
)
}
#[test]
fn register_node_channel_entity_includes_custom_attributes() {
let registry = TelemetryRegistryHandle::new();
let mut custom = HashMap::new();
let _ = custom.insert(
"custom.identity.foo".to_string(),
TelemetryAttribute::new(AttributeValue::String("bar".to_string())),
);
let ctx = pipeline_ctx_with_custom_attrs(registry.clone(), custom);
let key = register_channel(&ctx);
let (schema, rendered) = registry
.visit_entity(key, |a| (a.schema_name(), a.attributes_to_string()))
.expect("channel entity registered");
assert_eq!(schema, "node.channel.custom.attrs");
assert!(
rendered.contains("custom={custom.identity.foo=bar}"),
"custom identity attributes missing from channel entity: {rendered}"
);
assert!(
rendered.contains("channel.id=channel-1") && rendered.contains("node.id=test-node"),
"base channel attributes must be preserved: {rendered}"
);
}
#[test]
fn register_node_channel_entity_omits_empty_custom_attributes() {
let registry = TelemetryRegistryHandle::new();
let ctx = pipeline_ctx_with_custom_attrs(registry.clone(), HashMap::new());
let key = register_channel(&ctx);
let (schema, rendered) = registry
.visit_entity(key, |a| (a.schema_name(), a.attributes_to_string()))
.expect("channel entity registered");
assert_eq!(schema, "node.channel.attrs");
assert!(
!rendered.contains("custom="),
"nodes without custom attributes must not emit a custom attribute: {rendered}"
);
}
fn register_test_topic(ctx: &PipelineContext) -> EntityKey {
ctx.register_topic_entity(Cow::Borrowed("test-topic"))
}
#[test]
fn register_topic_entity_includes_custom_attributes() {
let registry = TelemetryRegistryHandle::new();
let mut custom = HashMap::new();
let _ = custom.insert(
"custom.identity.foo".to_string(),
TelemetryAttribute::new(AttributeValue::String("bar".to_string())),
);
let ctx = pipeline_ctx_with_custom_attrs(registry.clone(), custom);
let key = register_test_topic(&ctx);
let (schema, rendered) = registry
.visit_entity(key, |a| (a.schema_name(), a.attributes_to_string()))
.expect("topic entity registered");
assert_eq!(schema, "node.custom.topic.attrs");
assert!(
rendered.contains("custom={custom.identity.foo=bar}"),
"custom identity attributes missing from topic entity: {rendered}"
);
assert!(
rendered.contains("topic=test-topic") && rendered.contains("node.id=test-node"),
"base topic attributes must be preserved: {rendered}"
);
}
#[test]
fn register_topic_entity_omits_empty_custom_attributes() {
let registry = TelemetryRegistryHandle::new();
let ctx = pipeline_ctx_with_custom_attrs(registry.clone(), HashMap::new());
let key = register_test_topic(&ctx);
let (schema, rendered) = registry
.visit_entity(key, |a| (a.schema_name(), a.attributes_to_string()))
.expect("topic entity registered");
assert_eq!(schema, "node.topic.attrs");
assert!(
!rendered.contains("custom="),
"nodes without custom attributes must not emit a custom attribute: {rendered}"
);
}
#[test]
fn topic_metric_registrar_links_entity_with_and_without_custom_attrs() {
use crate::flow_metrics::FlowInputMessageMetrics;
let registry = TelemetryRegistryHandle::new();
let ctx = pipeline_ctx_with_custom_attrs(registry.clone(), HashMap::new());
let registrar = ctx.metric_set_registrar_with_topic(Cow::Borrowed("test-topic"));
let metrics = FlowInputMessageMetrics::register(®istrar);
let key = metrics.entity_key();
let (schema, rendered) = registry
.visit_entity(key, |a| (a.schema_name(), a.attributes_to_string()))
.expect("measurement set entity registered without custom attrs");
assert_eq!(schema, "node.topic.attrs");
assert!(
rendered.contains("topic=test-topic") && rendered.contains("node.id=test-node"),
"base topic attributes must be preserved: {rendered}"
);
assert!(
!rendered.contains("custom="),
"nodes without custom attributes must not emit a custom attribute: {rendered}"
);
let registry = TelemetryRegistryHandle::new();
let mut custom = HashMap::new();
let _ = custom.insert(
"custom.identity.foo".to_string(),
TelemetryAttribute::new(AttributeValue::String("bar".to_string())),
);
let ctx = pipeline_ctx_with_custom_attrs(registry.clone(), custom);
let registrar = ctx.metric_set_registrar_with_topic(Cow::Borrowed("test-topic"));
let metrics = FlowInputMessageMetrics::register(®istrar);
let key = metrics.entity_key();
let (schema, rendered) = registry
.visit_entity(key, |a| (a.schema_name(), a.attributes_to_string()))
.expect("measurement set entity registered with custom attrs");
assert_eq!(schema, "node.custom.topic.attrs");
assert!(
rendered.contains("custom={custom.identity.foo=bar}"),
"custom identity attributes missing from topic entity: {rendered}"
);
assert!(
rendered.contains("topic=test-topic") && rendered.contains("node.id=test-node"),
"base topic attributes must be preserved: {rendered}"
);
}
#[test]
fn generic_entity_registration_tracks_node_cleanup() {
use crate::entity_context::{
NodeTelemetryGuard, NodeTelemetryHandle, with_node_telemetry_handle,
};
use crate::flow_metrics::FlowInputMessageMetrics;
let registry = TelemetryRegistryHandle::new();
let ctx = pipeline_ctx_with_custom_attrs(registry.clone(), HashMap::new());
let node_entity = ctx.register_node_entity();
let handle = NodeTelemetryHandle::new(registry.clone(), node_entity);
let guard = NodeTelemetryGuard::new(handle.clone());
with_node_telemetry_handle(handle, || {
let child_entity = ctx.register_entity(NodeWithTopicAttributeSet {
node_attrs: ctx.node_attribute_set(),
topic: Cow::Borrowed("child"),
});
let registrar = ctx.metric_set_registrar_for_entity(child_entity);
let _metrics = FlowInputMessageMetrics::register(®istrar);
});
assert_eq!(registry.entity_count(), 2);
assert_eq!(registry.metric_set_count(), 1);
drop(guard);
assert_eq!(registry.entity_count(), 0);
assert_eq!(registry.metric_set_count(), 0);
}
}