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
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
//! Services created once for one pipeline runtime and injected into nodes.
//!
//! A dependency belongs here when its lifetime is the complete pipeline
//! runtime, multiple node families need the same instance, and central creation
//! or validation protects pipeline-local invariants. Mutable services should
//! also benefit from staying on the pipeline's runtime thread.
//!
//! Node configuration, node-local scheduling and metrics, process-wide state,
//! and data-type-specific control channels should remain with their existing
//! owners. This boundary keeps the service bundle focused and avoids turning it
//! into a container for unrelated runtime objects.
use otel_arrow_dfe_pdata_codec::{CodecService, CodecServiceBuilder, DecodePolicy, RegistryError};
/// Runtime-owned services shared by every effect handler in one pipeline.
///
/// Adding a service here should satisfy the module-level ownership criteria and
/// preserve the single-threaded pipeline runtime's state-sharing semantics.
#[derive(Clone)]
pub struct PipelineRuntimeServices {
codecs: CodecService,
}
impl PipelineRuntimeServices {
/// Validates linked codec extensions and creates lazy pipeline-local state.
pub fn new(decode_policy: DecodePolicy) -> Result<Self, RegistryError> {
Ok(Self {
codecs: CodecServiceBuilder::from_global_registry()?
.with_decode_policy(decode_policy)
.build(),
})
}
/// Pipeline-local codec access.
#[must_use]
pub const fn codecs(&self) -> &CodecService {
&self.codecs
}
/// Returns whether two handles address the same pipeline-owned services.
#[must_use]
pub fn shares_state_with(&self, other: &Self) -> bool {
self.codecs.shares_state_with(&other.codecs)
}
}
/// Common pdata codec access exposed by receiver, processor, and exporter effects.
pub trait CodecEffectHandler {
/// Returns the codec service injected by the pipeline runtime.
fn codec_service(&self) -> &CodecService;
}
#[cfg(test)]
mod tests {
use super::*;
use crate::testing::test_node;
use otel_arrow_dfe_telemetry::reporter::MetricsReporter;
fn requires_codec_effect_handler<T: CodecEffectHandler>() {}
/// Scenario: A pipeline injects its runtime services into all six effect-handler families.
/// Guarantees: Local and shared receivers, processors, and exporters expose one common service contract.
#[test]
fn all_pdata_effect_handlers_use_the_shared_service_contract() {
requires_codec_effect_handler::<crate::local::receiver::EffectHandler<()>>();
requires_codec_effect_handler::<crate::shared::receiver::EffectHandler<()>>();
requires_codec_effect_handler::<crate::local::processor::EffectHandler<()>>();
requires_codec_effect_handler::<crate::shared::processor::EffectHandler<()>>();
requires_codec_effect_handler::<crate::local::exporter::EffectHandler<()>>();
requires_codec_effect_handler::<crate::shared::exporter::EffectHandler<()>>();
let runtime_services = PipelineRuntimeServices::new(DecodePolicy::default()).unwrap();
let (_metrics_rx, metrics_reporter) = MetricsReporter::create_new_and_receiver(1);
let local = crate::local::exporter::EffectHandler::<()>::new(
test_node("local"),
metrics_reporter.clone(),
runtime_services.clone(),
);
let shared = crate::shared::exporter::EffectHandler::<()>::new(
test_node("shared"),
metrics_reporter,
runtime_services.clone(),
);
assert!(
local
.codec_service()
.shares_state_with(shared.codec_service())
);
}
}