mod layout;
mod propagation;
pub use layout::*;
pub use propagation::CompiledHeaderPropagationPolicy;
use propagation::CompiledHeaderPropagationPolicy as HeaderPropagationPolicy;
use crate::PipelineFactory;
use crate::error::Error as EngineError;
use otel_arrow_dfe_config::authorized_identity_policy::AuthorizedIdentityPolicy;
use otel_arrow_dfe_config::context_policy::{
ContextEntryDeclaration as ConfigContextEntryDeclaration, ContextEntryPart,
};
use otel_arrow_dfe_config::engine::ResolvedOtelDataflowSpec;
use otel_arrow_dfe_config::error::Error;
use otel_arrow_dfe_config::node::{NodeKind, NodeUserConfig};
use otel_arrow_dfe_config::transport_headers_policy::{
CompiledHeaderCapturePolicy, HeaderCapturePolicy, TransportHeadersPolicy,
};
use otel_arrow_dfe_config::{ContextEntryName, NodeId as ConfigNodeId, PipelineKey};
use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::sync::Arc;
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub enum ContextEntryTarget {
Primitive {
domain: ContextDomain,
name: ContextEntryName,
},
CompositeMember {
composite: ContextEntryName,
member: ContextEntryName,
},
Composite {
name: ContextEntryName,
},
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum SelectedSource<'a> {
Field(ContextDomain, &'a ContextEntryName),
Constant(&'a ContextEntryName),
}
impl<'a> SelectedSource<'a> {
fn from_part(part: &'a ContextEntryPart) -> Option<Self> {
match part {
ContextEntryPart::Constant { name, .. } => Some(Self::Constant(name)),
ContextEntryPart::TransportHeader { name, .. } => {
Some(Self::Field(ContextDomain::TransportHeader, name))
}
ContextEntryPart::AuthorizedIdentity { name, .. } => {
Some(Self::Field(ContextDomain::AuthorizedIdentity, name))
}
ContextEntryPart::TransportHeaderMatch { .. } => None,
}
}
}
impl ContextEntryTarget {
fn composite_name(&self) -> Option<&ContextEntryName> {
match self {
Self::Primitive { .. } => None,
Self::CompositeMember { composite, .. } => Some(composite),
Self::Composite { name } => Some(name),
}
}
fn selected_sources<'a>(
&'a self,
composites: &'a [ConfigContextEntryDeclaration],
) -> Result<Vec<SelectedSource<'a>>, Error> {
match self {
Self::Primitive { domain, name } => Ok(vec![SelectedSource::Field(*domain, name)]),
Self::CompositeMember { composite, member } => {
let declaration = composite_declaration(composite, composites)?;
let part = declaration
.definition
.0
.iter()
.find(|part| part.member_name() == Some(member))
.ok_or_else(|| {
invalid_context(format!("unknown context member `{composite}:{member}`"))
})?;
Ok(vec![
SelectedSource::from_part(part)
.expect("selected composite member is value-bearing"),
])
}
Self::Composite { name } => {
let declaration = composite_declaration(name, composites)?;
Ok(declaration
.definition
.0
.iter()
.filter_map(SelectedSource::from_part)
.collect())
}
}
}
}
fn invalid_context(error: impl Into<String>) -> Error {
Error::InvalidUserConfig {
error: error.into(),
}
}
fn composite_declaration<'a>(
name: &ContextEntryName,
declarations: &'a [ConfigContextEntryDeclaration],
) -> Result<&'a ConfigContextEntryDeclaration, Error> {
let mut matching = declarations.iter().filter(|entry| &entry.name == name);
let declaration = matching
.next()
.ok_or_else(|| invalid_context(format!("unknown composite context entry `{name}`")))?;
if matching.next().is_some() {
return Err(invalid_context(format!(
"duplicate composite context entry `{name}`"
)));
}
Ok(declaration)
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct ContextEntrySelector {
pub target: ContextEntryTarget,
pub form: ContextEntrySelectorForm,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub enum ContextEntrySelectorForm {
Value,
StoredKeyValue,
OriginalKeyValue,
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub enum ContextConsumerSelector {
Entries {
entries: Box<[ContextEntrySelector]>,
},
AllStored {
domain: ContextDomain,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub enum ContextDeclaration {
Produces {
domain: ContextDomain,
entry: ContextEntryName,
},
Consumes {
selector: ContextConsumerSelector,
},
HeaderCapture {
policy: HeaderCapturePolicy,
},
HeaderPropagation {
policy: HeaderPropagationPolicy,
},
AuthorizedIdentityCapture {
policy: AuthorizedIdentityPolicy,
},
}
impl ContextDeclaration {
fn is_component_declaration(&self) -> bool {
matches!(self, Self::Produces { .. } | Self::Consumes { .. })
}
fn context_runtime_requirements(
&self,
composites: &[ConfigContextEntryDeclaration],
) -> Result<ContextRuntimeRequirements, Error> {
let mut requirements = ContextRuntimeRequirements::none();
match self {
Self::Consumes {
selector: ContextConsumerSelector::Entries { entries },
} => {
for entry in entries {
let sources = entry.target.selected_sources(composites)?;
if entry.form != ContextEntrySelectorForm::OriginalKeyValue {
continue;
}
for source in sources {
match source {
SelectedSource::Field(ContextDomain::TransportHeader, name) => {
_ = requirements
.original_name_retention
.overrides
.insert(original_name_key(name), true);
}
SelectedSource::Field(domain, name) => {
return Err(invalid_context(format!(
"original wire name requested for {domain:?} context entry `{name}`; only transport headers have original wire names"
)));
}
SelectedSource::Constant(name) => {
return Err(invalid_context(format!(
"original wire name requested for constant context entry `{name}`; constants have no original wire names"
)));
}
}
}
}
}
Self::Consumes {
selector: ContextConsumerSelector::AllStored { .. },
}
| Self::Produces { .. }
| Self::HeaderCapture { .. }
| Self::AuthorizedIdentityCapture { .. } => {}
Self::HeaderPropagation { policy } => {
requirements
.original_name_retention
.default_preserve_original = policy.propagates_original_name_by_default();
policy.visit_original_name_requirement_names(|name| {
let preserve_original = policy.propagates_original_name(name);
if preserve_original
!= requirements
.original_name_retention
.default_preserve_original
{
_ = requirements
.original_name_retention
.overrides
.insert(original_name_key(name), preserve_original);
}
});
}
}
Ok(requirements)
}
}
#[derive(Clone, Copy)]
pub struct ContextDeclarationProvider {
pub declarations: ContextDeclarationFn,
}
pub type ContextDeclarationFn = fn(&serde_json::Value) -> Result<NodeContextDeclarations, Error>;
pub trait ConfigNodeContextDeclaration: serde::de::DeserializeOwned {
fn context_declarations(&self) -> NodeContextDeclarations;
fn validate_context_declarations(
&self,
pipeline_ctx: &crate::context::PipelineContext,
) -> Result<(), Error> {
pipeline_ctx
.compiled_context_bindings()
.validate_node_declarations(
&pipeline_ctx.pipeline_key(),
&pipeline_ctx.node_id(),
&self.context_declarations(),
)
}
}
#[derive(Debug, Default, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct NodeContextDeclarations {
declarations: Box<[ContextDeclaration]>,
}
impl FromIterator<ContextDeclaration> for NodeContextDeclarations {
fn from_iter<T>(iter: T) -> Self
where
T: IntoIterator<Item = ContextDeclaration>,
{
let mut uniq: Vec<_> = iter.into_iter().collect();
uniq.sort();
uniq.dedup();
Self {
declarations: uniq.into_boxed_slice(),
}
}
}
impl IntoIterator for NodeContextDeclarations {
type Item = ContextDeclaration;
type IntoIter = std::vec::IntoIter<ContextDeclaration>;
fn into_iter(self) -> Self::IntoIter {
self.declarations.into_vec().into_iter()
}
}
impl NodeContextDeclarations {
pub fn iter(&self) -> impl Iterator<Item = &ContextDeclaration> {
self.declarations.iter()
}
#[must_use]
pub fn is_empty(&self) -> bool {
self.declarations.is_empty()
}
#[must_use]
pub fn len(&self) -> usize {
self.declarations.len()
}
}
impl ContextDeclarationProvider {
#[must_use]
pub const fn from_typed_config<T>() -> Self
where
T: ConfigNodeContextDeclaration,
{
Self {
declarations: typed_context_declarations::<T>,
}
}
}
fn typed_context_declarations<T>(
config: &serde_json::Value,
) -> Result<NodeContextDeclarations, Error>
where
T: ConfigNodeContextDeclaration,
{
Ok(
otel_arrow_dfe_config::validation::deserialize_typed_config::<T>(config)?
.context_declarations(),
)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CompiledContextBindings {
by_pipeline: HashMap<PipelineKey, HashMap<ConfigNodeId, CompiledNodeBindings>>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct CompiledNodeBindings {
component_declarations: NodeContextDeclarations,
composites: Box<[ConfigContextEntryDeclaration]>,
header_capture: Option<CompiledHeaderCapturePolicy>,
header_propagation: Option<HeaderPropagationPolicy>,
authorized_identity_capture: Option<AuthorizedIdentityPolicy>,
}
type ContextDeclarationsByPipeline =
HashMap<PipelineKey, HashMap<ConfigNodeId, PreparedNodeContextDeclarations>>;
#[derive(Debug)]
struct PreparedNodeContextDeclarations {
declarations: NodeContextDeclarations,
composites: Box<[ConfigContextEntryDeclaration]>,
requirements: ContextRuntimeRequirements,
}
impl PreparedNodeContextDeclarations {
fn new(
declarations: NodeContextDeclarations,
context: &[ConfigContextEntryDeclaration],
) -> Result<Self, Error> {
let composite_names = declarations
.iter()
.filter_map(|declaration| match declaration {
ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries { entries },
} => Some(entries.iter()),
_ => None,
})
.flatten()
.filter_map(|entry| entry.target.composite_name())
.collect::<BTreeSet<_>>();
let composites = composite_names
.into_iter()
.map(|name| {
let mut declaration = composite_declaration(name, context)?.clone();
validate_definition(&declaration)?;
declaration.definition.0.sort_unstable();
Ok(declaration)
})
.collect::<Result<Box<[_]>, Error>>()?;
let requirements = declarations.iter().try_fold(
ContextRuntimeRequirements::none(),
|requirements, declaration| {
Ok::<_, Error>(
requirements.union(declaration.context_runtime_requirements(&composites)?),
)
},
)?;
Ok(Self {
declarations,
composites,
requirements,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ContextRuntimeRequirements {
original_name_retention: OriginalNameRetention,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct OriginalNameRetention {
default_preserve_original: bool,
overrides: BTreeMap<Box<str>, bool>,
}
#[derive(Debug, Clone)]
pub struct PreparedContext {
pub runtime_requirements: ContextRuntimeRequirements,
pub bindings: Arc<CompiledContextBindings>,
}
impl ContextRuntimeRequirements {
fn compile(declarations: &ContextDeclarationsByPipeline) -> Self {
declarations
.values()
.flat_map(HashMap::values)
.fold(Self::none(), |requirements, declaration| {
requirements.union(declaration.requirements.clone())
})
}
fn none() -> Self {
Self {
original_name_retention: OriginalNameRetention {
default_preserve_original: false,
overrides: BTreeMap::new(),
},
}
}
fn union(self, other: Self) -> Self {
Self {
original_name_retention: self
.original_name_retention
.union(other.original_name_retention),
}
}
#[must_use]
pub fn can_satisfy(&self, candidate: &Self) -> bool {
self.original_name_retention
.can_satisfy(&candidate.original_name_retention)
}
#[must_use]
pub fn preserves_original_name(&self, name: &ContextEntryName) -> bool {
self.original_name_retention.preserves_original_name(name)
}
}
impl OriginalNameRetention {
fn union(self, other: Self) -> Self {
let default_preserve_original =
self.default_preserve_original || other.default_preserve_original;
let mut names = self
.overrides
.keys()
.chain(other.overrides.keys())
.cloned()
.collect::<Vec<_>>();
names.sort();
names.dedup();
let overrides = names
.into_iter()
.filter_map(|name| {
let preserve_original =
self.preserves_original_key(&name) || other.preserves_original_key(&name);
(preserve_original != default_preserve_original)
.then_some((name, preserve_original))
})
.collect();
Self {
default_preserve_original,
overrides,
}
}
fn can_satisfy(&self, candidate: &Self) -> bool {
if candidate.default_preserve_original && !self.default_preserve_original {
return false;
}
self.overrides
.keys()
.chain(candidate.overrides.keys())
.all(|name| {
!candidate.preserves_original_key(name) || self.preserves_original_key(name)
})
}
fn preserves_original_name(&self, name: &ContextEntryName) -> bool {
self.preserves_original_key(&original_name_key(name))
}
fn preserves_original_key(&self, name: &str) -> bool {
self.overrides
.get(name)
.copied()
.unwrap_or(self.default_preserve_original)
}
}
fn original_name_key(name: &ContextEntryName) -> Box<str> {
name.as_str().to_ascii_lowercase().into()
}
impl CompiledNodeBindings {
fn compile(
prepared: PreparedNodeContextDeclarations,
requirements: &ContextRuntimeRequirements,
) -> Self {
let mut component_declarations = Vec::new();
let mut header_capture = None;
let mut header_propagation = None;
let mut authorized_identity_capture = None;
for declaration in prepared.declarations {
match declaration {
declaration @ (ContextDeclaration::Produces { .. }
| ContextDeclaration::Consumes { .. }) => {
component_declarations.push(declaration);
}
ContextDeclaration::HeaderCapture { policy } => {
header_capture =
Some(policy.compile(|name| requirements.preserves_original_name(name)));
}
ContextDeclaration::HeaderPropagation { policy } => {
header_propagation = Some(policy);
}
ContextDeclaration::AuthorizedIdentityCapture { policy } => {
authorized_identity_capture = Some(policy);
}
}
}
Self {
component_declarations: component_declarations.into_iter().collect(),
composites: prepared.composites,
header_capture,
header_propagation,
authorized_identity_capture,
}
}
fn is_empty(&self) -> bool {
self.component_declarations.is_empty()
&& self.header_capture.is_none()
&& self.header_propagation.is_none()
&& self.authorized_identity_capture.is_none()
}
}
impl CompiledContextBindings {
#[must_use]
pub fn empty() -> Self {
Self {
by_pipeline: HashMap::new(),
}
}
fn compile(
declarations: ContextDeclarationsByPipeline,
requirements: &ContextRuntimeRequirements,
) -> Self {
let by_pipeline = declarations
.into_iter()
.map(|(pipeline, nodes)| {
let nodes = nodes
.into_iter()
.map(|(node, declarations)| {
(
node,
CompiledNodeBindings::compile(declarations, requirements),
)
})
.collect();
(pipeline, nodes)
})
.collect();
Self { by_pipeline }
}
pub(crate) fn header_capture_policy(
&self,
pipeline: &PipelineKey,
node: &ConfigNodeId,
) -> Option<&CompiledHeaderCapturePolicy> {
self.by_pipeline
.get(pipeline)?
.get(node)?
.header_capture
.as_ref()
}
pub(crate) fn header_propagation_policy(
&self,
pipeline: &PipelineKey,
node: &ConfigNodeId,
) -> Option<&HeaderPropagationPolicy> {
self.by_pipeline
.get(pipeline)?
.get(node)?
.header_propagation
.as_ref()
}
pub(crate) fn authorized_identity_policy(
&self,
pipeline: &PipelineKey,
node: &ConfigNodeId,
) -> Option<&AuthorizedIdentityPolicy> {
self.by_pipeline
.get(pipeline)?
.get(node)?
.authorized_identity_capture
.as_ref()
}
#[must_use]
pub fn pipeline_bindings_match(&self, other: &Self, pipeline: &PipelineKey) -> bool {
let current = self.by_pipeline.get(pipeline);
let candidate = other.by_pipeline.get(pipeline);
let current_binding_count = current
.into_iter()
.flat_map(|nodes| nodes.values())
.filter(|node| !node.is_empty())
.count();
let candidate_binding_count = candidate
.into_iter()
.flat_map(|nodes| nodes.values())
.filter(|node| !node.is_empty())
.count();
current_binding_count == candidate_binding_count
&& current
.into_iter()
.flat_map(|nodes| nodes.iter())
.all(|(node_id, node)| {
node.is_empty()
|| candidate
.and_then(|nodes| nodes.get(node_id))
.is_some_and(|candidate_node| candidate_node == node)
})
}
pub fn validate_node_declarations(
&self,
pipeline: &PipelineKey,
node: &ConfigNodeId,
declarations: &NodeContextDeclarations,
) -> Result<(), Error> {
match self
.by_pipeline
.get(pipeline)
.and_then(|nodes| nodes.get(node))
{
Some(expected) if expected.component_declarations == *declarations => Ok(()),
_ => Err(Error::UnrecognizedContextDeclaration {}),
}
}
}
impl<PData: 'static + Clone + std::fmt::Debug> PipelineFactory<PData> {
pub fn compile_initial_context(
&self,
resolved: &ResolvedOtelDataflowSpec,
) -> Result<PreparedContext, EngineError> {
let declarations = self.context_declarations(resolved)?;
let runtime_requirements = ContextRuntimeRequirements::compile(&declarations);
let bindings = Self::compile_bindings(declarations, &runtime_requirements);
Ok(PreparedContext {
runtime_requirements,
bindings,
})
}
pub fn compile_candidate_context(
&self,
resolved: &ResolvedOtelDataflowSpec,
installed_requirements: &ContextRuntimeRequirements,
) -> Result<PreparedContext, EngineError> {
let declarations = self.context_declarations(resolved)?;
let runtime_requirements = ContextRuntimeRequirements::compile(&declarations);
let bindings = Self::compile_bindings(declarations, installed_requirements);
Ok(PreparedContext {
runtime_requirements,
bindings,
})
}
fn compile_bindings(
declarations: ContextDeclarationsByPipeline,
requirements: &ContextRuntimeRequirements,
) -> Arc<CompiledContextBindings> {
Arc::new(CompiledContextBindings::compile(declarations, requirements))
}
fn context_declarations(
&self,
resolved: &ResolvedOtelDataflowSpec,
) -> Result<ContextDeclarationsByPipeline, EngineError> {
let mut declarations = ContextDeclarationsByPipeline::new();
for pipeline in &resolved.pipelines {
let pipeline_key = PipelineKey::new(
pipeline.pipeline_group_id.clone(),
pipeline.pipeline_id.clone(),
);
let mut declarations_by_node = HashMap::new();
for (node_id, node_config) in pipeline.pipeline.node_iter() {
let component_declarations = self.node_context_declarations(
node_config.kind(),
node_config.r#type.as_ref(),
&node_config.config,
)?;
let wrapper_declarations = Self::wrapper_context_declarations(
node_config,
&pipeline.policies.transport_headers,
&pipeline.policies.authorized_identity,
&pipeline.policies.context,
)?;
let declarations = component_declarations
.into_iter()
.chain(wrapper_declarations)
.collect();
let declarations =
PreparedNodeContextDeclarations::new(declarations, &pipeline.policies.context)
.map_err(|error| EngineError::ConfigError(Box::new(error)))?;
let _ = declarations_by_node.insert(node_id.clone(), declarations);
}
let _ = declarations.insert(pipeline_key, declarations_by_node);
}
Ok(declarations)
}
fn wrapper_context_declarations(
node: &NodeUserConfig,
pipeline_policy: &Option<TransportHeadersPolicy>,
authorized_identity: &Option<AuthorizedIdentityPolicy>,
context: &[ConfigContextEntryDeclaration],
) -> Result<NodeContextDeclarations, EngineError> {
let declarations = match node.kind() {
NodeKind::Receiver => node
.header_capture
.as_ref()
.or_else(|| {
pipeline_policy
.as_ref()
.map(|policy| &policy.header_capture)
})
.cloned()
.map(|policy| ContextDeclaration::HeaderCapture { policy })
.into_iter()
.chain(
authorized_identity
.as_ref()
.filter(|policy| !policy.is_empty())
.cloned()
.map(|policy| ContextDeclaration::AuthorizedIdentityCapture { policy }),
)
.collect(),
NodeKind::Exporter => {
let policy = node.header_propagation.as_ref().or_else(|| {
pipeline_policy
.as_ref()
.map(|policy| &policy.header_propagation)
});
policy
.cloned()
.map(|policy| {
HeaderPropagationPolicy::compile(policy, context)
.map(|policy| ContextDeclaration::HeaderPropagation { policy })
.map_err(|error| {
EngineError::ConfigError(Box::new(Error::InvalidUserConfig {
error,
}))
})
})
.transpose()?
.into_iter()
.collect()
}
NodeKind::Processor => NodeContextDeclarations::default(),
};
Ok(declarations)
}
fn node_context_declarations(
&self,
kind: NodeKind,
urn: &str,
config: &serde_json::Value,
) -> Result<NodeContextDeclarations, EngineError> {
let missing_factory = || {
EngineError::ConfigError(Box::new(Error::InvalidUserConfig {
error: format!("node factory `{urn}` is not registered"),
}))
};
let (validate_config, context_declarations) = match kind {
NodeKind::Receiver => {
let factory = self
.get_receiver_factory_map()
.get(urn)
.ok_or_else(&missing_factory)?;
(factory.validate_config, factory.context_declarations)
}
NodeKind::Processor => {
let factory = self
.get_processor_factory_map()
.get(urn)
.ok_or_else(&missing_factory)?;
(factory.validate_config, factory.context_declarations)
}
NodeKind::Exporter => {
let factory = self
.get_exporter_factory_map()
.get(urn)
.ok_or_else(&missing_factory)?;
(factory.validate_config, factory.context_declarations)
}
};
validate_config(config).map_err(|error| EngineError::ConfigError(Box::new(error)))?;
let declarations = context_declarations
.map(|provider| {
(provider.declarations)(config)
.map_err(|error| EngineError::ConfigError(Box::new(error)))
})
.transpose()?
.unwrap_or_default();
if let Some(declaration) = declarations
.iter()
.find(|declaration| !declaration.is_component_declaration())
{
return Err(EngineError::ConfigError(Box::new(
Error::InvalidUserConfig {
error: format!(
"node factory `{urn}` returned engine-owned context declaration \
`{declaration:?}`"
),
},
)));
}
Ok(declarations)
}
}
#[cfg(test)]
mod tests {
use super::*;
use otel_arrow_dfe_config::transport_headers::{TransportHeader, TransportHeaders};
use otel_arrow_dfe_config::transport_headers_policy::HeaderPropagationPolicy as HeaderPropagationConfig;
use otel_arrow_dfe_config::transport_headers_policy::{CaptureDefaults, CaptureRule};
#[derive(serde::Deserialize)]
struct TestDeclarationConfig {
entry: ContextEntryName,
#[serde(default)]
composite: Option<ContextEntryName>,
#[serde(default)]
original: bool,
}
impl ConfigNodeContextDeclaration for TestDeclarationConfig {
fn context_declarations(&self) -> NodeContextDeclarations {
[ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries {
entries: vec![ContextEntrySelector {
target: match &self.composite {
Some(composite) => ContextEntryTarget::CompositeMember {
composite: composite.clone(),
member: self.entry.clone(),
},
None => ContextEntryTarget::Primitive {
domain: ContextDomain::TransportHeader,
name: self.entry.clone(),
},
},
form: if self.original {
ContextEntrySelectorForm::OriginalKeyValue
} else {
ContextEntrySelectorForm::Value
},
}]
.into_boxed_slice(),
},
}]
.into_iter()
.collect()
}
}
fn pipeline(group: &str, name: &str) -> PipelineKey {
PipelineKey::new(group.to_owned().into(), name.to_owned().into())
}
fn context_name(name: &str) -> ContextEntryName {
name.try_into().expect("valid test context entry name")
}
fn unused_test_receiver(
_: crate::context::PipelineContext,
_: crate::node::NodeId,
_: Arc<NodeUserConfig>,
_: &crate::config::ReceiverConfig,
_: &crate::capability::registry::Capabilities,
) -> Result<crate::receiver::ReceiverWrapper<()>, Error> {
unreachable!("context compilation does not construct test nodes")
}
fn unused_test_exporter(
_: crate::context::PipelineContext,
_: crate::node::NodeId,
_: Arc<NodeUserConfig>,
_: &crate::config::ExporterConfig,
_: &crate::capability::registry::Capabilities,
) -> Result<crate::exporter::ExporterWrapper<()>, Error> {
unreachable!("context compilation does not construct test nodes")
}
fn unused_test_processor(
_: crate::context::PipelineContext,
_: crate::node::NodeId,
_: Arc<NodeUserConfig>,
_: &crate::config::ProcessorConfig,
_: &crate::capability::registry::Capabilities,
) -> Result<crate::processor::ProcessorWrapper<()>, Error> {
unreachable!("context compilation does not construct test nodes")
}
fn accept_test_config(_: &serde_json::Value) -> Result<(), Error> {
Ok(())
}
static TEST_RECEIVERS: [crate::ReceiverFactory<()>; 2] = [
crate::ReceiverFactory {
name: "urn:test:receiver:example",
create: unused_test_receiver,
context_declarations: None,
wiring_contract: crate::wiring_contract::WiringContract::UNRESTRICTED,
validate_config: otel_arrow_dfe_config::validation::no_config,
},
crate::ReceiverFactory {
name: "urn:otel:receiver:internal_telemetry",
create: unused_test_receiver,
context_declarations: None,
wiring_contract: crate::wiring_contract::WiringContract::UNRESTRICTED,
validate_config: accept_test_config,
},
];
static TEST_EXPORTERS: [crate::ExporterFactory<()>; 4] = [
crate::ExporterFactory {
name: "urn:test:exporter:example",
create: unused_test_exporter,
context_declarations: None,
wiring_contract: crate::wiring_contract::WiringContract::UNRESTRICTED,
validate_config: otel_arrow_dfe_config::validation::no_config,
},
crate::ExporterFactory {
name: "urn:otel:exporter:noop",
create: unused_test_exporter,
context_declarations: None,
wiring_contract: crate::wiring_contract::WiringContract::UNRESTRICTED,
validate_config: otel_arrow_dfe_config::validation::no_config,
},
crate::ExporterFactory {
name: "urn:otel:exporter:console",
create: unused_test_exporter,
context_declarations: None,
wiring_contract: crate::wiring_contract::WiringContract::UNRESTRICTED,
validate_config: otel_arrow_dfe_config::validation::no_config,
},
crate::ExporterFactory {
name: "urn:test:exporter:context",
create: unused_test_exporter,
context_declarations: Some(ContextDeclarationProvider::from_typed_config::<
TestDeclarationConfig,
>()),
wiring_contract: crate::wiring_contract::WiringContract::UNRESTRICTED,
validate_config: otel_arrow_dfe_config::validation::validate_typed_config::<
TestDeclarationConfig,
>,
},
];
static TEST_PROCESSORS: [crate::ProcessorFactory<()>; 1] = [crate::ProcessorFactory {
name: "urn:otel:processor:type_router",
create: unused_test_processor,
context_declarations: None,
wiring_contract: crate::wiring_contract::WiringContract::UNRESTRICTED,
validate_config: accept_test_config,
}];
fn test_pipeline_factory() -> PipelineFactory<()> {
PipelineFactory::new(&TEST_RECEIVERS, &TEST_PROCESSORS, &TEST_EXPORTERS, &[])
}
fn conditional_pipeline_yaml(composite: &str, selector: &str) -> String {
format!(
r#"
version: otel_dataflow/v1
policies:
context:
entries:
tenant: {composite}
engine: {{}}
groups:
default:
pipelines:
main:
nodes:
receiver:
type: "urn:test:receiver:example"
config: {{}}
exporter:
type: "urn:test:exporter:example"
header_propagation:
default:
selector:
type: named
named: [{selector}]
name: stored_name
config: {{}}
connections:
- from: receiver
to: exporter
"#
)
}
fn resolve_conditional_pipeline(composite: &str, selector: &str) -> ResolvedOtelDataflowSpec {
otel_arrow_dfe_config::engine::OtelDataflowSpec::from_yaml(&conditional_pipeline_yaml(
composite, selector,
))
.expect("conditional pipeline YAML is valid")
.resolve()
}
fn declarations_by_pipeline(
effective: NodeContextDeclarations,
) -> ContextDeclarationsByPipeline {
HashMap::from([(
pipeline("group", "pipeline"),
HashMap::from([(
ConfigNodeId::from("node"),
PreparedNodeContextDeclarations::new(effective, &[]).expect("valid declarations"),
)]),
)])
}
fn compiled_bindings(effective: NodeContextDeclarations) -> CompiledContextBindings {
let declarations = declarations_by_pipeline(effective);
let requirements = ContextRuntimeRequirements::compile(&declarations);
CompiledContextBindings::compile(declarations, &requirements)
}
fn context_runtime_requirements(
effective: NodeContextDeclarations,
) -> ContextRuntimeRequirements {
ContextRuntimeRequirements::compile(&declarations_by_pipeline(effective))
}
fn primitive_target(domain: ContextDomain, name: &str) -> ContextEntryTarget {
ContextEntryTarget::Primitive {
domain,
name: context_name(name),
}
}
fn member_target(composite: &str, member: &str) -> ContextEntryTarget {
ContextEntryTarget::CompositeMember {
composite: context_name(composite),
member: context_name(member),
}
}
fn consumer(
target: ContextEntryTarget,
form: ContextEntrySelectorForm,
) -> NodeContextDeclarations {
[ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries {
entries: Box::new([ContextEntrySelector { target, form }]),
},
}]
.into_iter()
.collect()
}
fn mixed_composite() -> ConfigContextEntryDeclaration {
use otel_arrow_dfe_config::context_policy::{ContextEntryDefinition, ContextScope};
ConfigContextEntryDeclaration {
scope: ContextScope::Engine,
name: context_name("tenant"),
definition: serde_json::from_value::<ContextEntryDefinition>(serde_json::json!([
{"type": "transport_header", "name": "id", "store_as": "header_id"},
{"type": "authorized_identity", "name": "id", "store_as": "identity_id"},
{"type": "transport_header_match", "name": "environment", "value": "production"}
]))
.expect("valid composite"),
}
}
fn constant_composite() -> ConfigContextEntryDeclaration {
use otel_arrow_dfe_config::context_policy::{ContextEntryDefinition, ContextScope};
ConfigContextEntryDeclaration {
scope: ContextScope::Engine,
name: context_name("route"),
definition: serde_json::from_value::<ContextEntryDefinition>(serde_json::json!([
{"type": "constant", "name": "route_name", "value": "otlp-http-json"},
{"type": "transport_header", "name": "workspace"}
]))
.expect("valid constant composite"),
}
}
#[test]
fn declaration_domains_are_part_of_binding_identity() {
let declarations = |domain| {
[
ContextDeclaration::Produces {
domain,
entry: context_name("id"),
},
ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries {
entries: Box::new([ContextEntrySelector {
target: primitive_target(domain, "id"),
form: ContextEntrySelectorForm::Value,
}]),
},
},
ContextDeclaration::Consumes {
selector: ContextConsumerSelector::AllStored { domain },
},
]
};
let header = declarations(ContextDomain::TransportHeader);
let identity = declarations(ContextDomain::AuthorizedIdentity);
let combined: NodeContextDeclarations =
header.clone().into_iter().chain(identity.clone()).collect();
assert_eq!(combined.len(), 6);
assert!(
!context_runtime_requirements(combined).preserves_original_name(&context_name("id"))
);
for (header, identity) in header.into_iter().zip(identity) {
let header: NodeContextDeclarations = [header].into_iter().collect();
let identity: NodeContextDeclarations = [identity].into_iter().collect();
let installed = compiled_bindings(header.clone());
let candidate = compiled_bindings(identity.clone());
let key = pipeline("group", "pipeline");
let node = ConfigNodeId::from("node");
assert!(
installed
.validate_node_declarations(&key, &node, &header)
.is_ok()
);
assert!(
installed
.validate_node_declarations(&key, &node, &identity)
.is_err()
);
assert!(!installed.pipeline_bindings_match(&candidate, &key));
assert!(!candidate.pipeline_bindings_match(&installed, &key));
}
}
#[test]
fn composite_member_requirements_preserve_source_domain_and_gate() {
let context = [mixed_composite()];
let prepared = PreparedNodeContextDeclarations::new(
consumer(
member_target("tenant", "header_id"),
ContextEntrySelectorForm::OriginalKeyValue,
),
&context,
)
.expect("header member supports original names");
assert!(
prepared
.requirements
.preserves_original_name(&context_name("id"))
);
assert!(
!prepared
.requirements
.preserves_original_name(&context_name("header_id"))
);
assert!(
!prepared
.requirements
.preserves_original_name(&context_name("environment"))
);
assert_eq!(prepared.composites[0].definition.0.len(), 3);
for form in [
ContextEntrySelectorForm::Value,
ContextEntrySelectorForm::StoredKeyValue,
] {
let prepared = PreparedNodeContextDeclarations::new(
consumer(member_target("tenant", "identity_id"), form),
&context,
)
.expect("identity member");
assert!(
!prepared
.requirements
.preserves_original_name(&context_name("id"))
);
}
let whole = ContextEntryTarget::Composite {
name: context_name("tenant"),
};
let source_name = context_name("id");
assert_eq!(
whole
.selected_sources(&context)
.expect("whole composite sources"),
[
SelectedSource::Field(ContextDomain::TransportHeader, &source_name),
SelectedSource::Field(ContextDomain::AuthorizedIdentity, &source_name),
]
);
let prepared = PreparedNodeContextDeclarations::new(
consumer(whole.clone(), ContextEntrySelectorForm::StoredKeyValue),
&context,
)
.expect("mixed composite supports stored values");
assert!(
!prepared
.requirements
.preserves_original_name(&context_name("id"))
);
let mut headers_only = mixed_composite();
_ = headers_only.definition.0.remove(1);
let prepared = PreparedNodeContextDeclarations::new(
consumer(whole, ContextEntrySelectorForm::OriginalKeyValue),
&[headers_only],
)
.expect("header-only composite supports original names");
assert!(
prepared
.requirements
.preserves_original_name(&context_name("id"))
);
assert!(
!prepared
.requirements
.preserves_original_name(&context_name("environment"))
);
}
#[test]
fn constant_members_have_no_external_or_original_name_requirements() {
let context = [constant_composite()];
let whole = ContextEntryTarget::Composite {
name: context_name("route"),
};
let constant_name = context_name("route_name");
let header_name = context_name("workspace");
assert_eq!(
whole
.selected_sources(&context)
.expect("whole composite sources"),
[
SelectedSource::Constant(&constant_name),
SelectedSource::Field(ContextDomain::TransportHeader, &header_name),
]
);
let prepared = PreparedNodeContextDeclarations::new(
consumer(
member_target("route", "route_name"),
ContextEntrySelectorForm::Value,
),
&context,
)
.expect("constant value selection");
assert!(
!prepared
.requirements
.preserves_original_name(&context_name("route_name"))
);
for target in [member_target("route", "route_name"), whole] {
let error = PreparedNodeContextDeclarations::new(
consumer(target, ContextEntrySelectorForm::OriginalKeyValue),
&context,
)
.expect_err("constant has no original wire name");
assert!(
error.to_string().contains(
"original wire name requested for constant context entry `route_name`"
),
"{error}"
);
}
}
#[test]
fn invalid_consumer_targets_and_representations_are_rejected() {
let context = [mixed_composite()];
for (target, form, expected) in [
(
member_target("missing", "header_id"),
ContextEntrySelectorForm::Value,
"unknown composite context entry `missing`",
),
(
member_target("tenant", "missing"),
ContextEntrySelectorForm::Value,
"unknown context member `tenant:missing`",
),
(
member_target("tenant", "environment"),
ContextEntrySelectorForm::Value,
"unknown context member `tenant:environment`",
),
(
primitive_target(ContextDomain::AuthorizedIdentity, "id"),
ContextEntrySelectorForm::OriginalKeyValue,
"original wire name requested for AuthorizedIdentity context entry `id`",
),
(
member_target("tenant", "identity_id"),
ContextEntrySelectorForm::OriginalKeyValue,
"original wire name requested for AuthorizedIdentity context entry `id`",
),
(
ContextEntryTarget::Composite {
name: context_name("tenant"),
},
ContextEntrySelectorForm::OriginalKeyValue,
"original wire name requested for AuthorizedIdentity context entry `id`",
),
] {
let error = PreparedNodeContextDeclarations::new(consumer(target, form), &context)
.expect_err("invalid selection must fail");
assert!(error.to_string().contains(expected), "{error}");
}
}
fn resolve_consumer_pipeline(
composite: &str,
member: &str,
original: bool,
) -> ResolvedOtelDataflowSpec {
let yaml = format!(
r#"
version: otel_dataflow/v1
policies:
context:
entries:
tenant: {composite}
engine: {{}}
groups:
default:
pipelines:
main:
nodes:
receiver:
type: "urn:test:receiver:example"
config: {{}}
exporter:
type: "urn:test:exporter:context"
config:
entry: {member}
composite: tenant
original: {original}
connections:
- from: receiver
to: exporter
"#
);
otel_arrow_dfe_config::engine::OtelDataflowSpec::from_yaml(&yaml)
.expect("consumer pipeline YAML")
.resolve()
}
#[test]
fn full_yaml_preparation_rejects_invalid_consumer_declarations() {
let composite = "[{type: authorized_identity, name: id}]";
let factory = test_pipeline_factory();
let current = factory
.compile_initial_context(&resolve_consumer_pipeline(composite, "id", false))
.expect("valid identity consumer");
for (member, original, expected) in [
("missing", false, "unknown context member `tenant:missing`"),
(
"id",
true,
"original wire name requested for AuthorizedIdentity",
),
] {
let resolved = resolve_consumer_pipeline(composite, member, original);
for result in [
factory.compile_initial_context(&resolved),
factory.compile_candidate_context(&resolved, ¤t.runtime_requirements),
] {
let error = result.expect_err("invalid consumer must fail preparation");
assert!(error.to_string().contains(expected), "{error}");
}
}
}
#[test]
fn full_yaml_bindings_track_selected_composite_definitions() {
let composite = "[{type: transport_header, name: id, store_as: key}, \
{type: authorized_identity, name: account}, \
{type: transport_header_match, name: environment, value: production}]";
let factory = test_pipeline_factory();
let installed = factory
.compile_initial_context(&resolve_consumer_pipeline(composite, "key", false))
.expect("initial consumer");
let key = pipeline("default", "main");
for changed in [
composite.replace("name: id", "name: other"),
composite.replace("type: transport_header,", "type: authorized_identity,"),
composite.replace("name: account", "name: other_account"),
composite.replace("value: production", "value: staging"),
] {
let candidate = factory
.compile_candidate_context(
&resolve_consumer_pipeline(&changed, "key", false),
&installed.runtime_requirements,
)
.expect("changed consumer");
assert!(
!installed
.bindings
.pipeline_bindings_match(&candidate.bindings, &key)
);
assert!(
!candidate
.bindings
.pipeline_bindings_match(&installed.bindings, &key)
);
}
let reordered = "[{type: transport_header_match, name: environment, value: production}, \
{type: authorized_identity, name: account}, \
{type: transport_header, name: id, store_as: key}]";
let candidate = factory
.compile_candidate_context(
&resolve_consumer_pipeline(reordered, "key", false),
&installed.runtime_requirements,
)
.expect("reordered consumer");
assert!(
installed
.bindings
.pipeline_bindings_match(&candidate.bindings, &key)
);
}
#[test]
fn compiled_capture_policy_tracks_each_match_name() {
let capture = HeaderCapturePolicy::new(
CaptureDefaults::default(),
vec![
CaptureRule {
match_names: vec![context_name("x-first"), context_name("x-second")],
store_as: None,
sensitive: false,
value_kind: None,
},
CaptureRule {
match_names: vec![context_name("x-alias-a"), context_name("x-alias-b")],
store_as: Some(context_name("canonical")),
sensitive: false,
value_kind: None,
},
],
);
let bindings = compiled_bindings(
[
ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries {
entries: ["x-first", "canonical"]
.map(|name| ContextEntrySelector {
target: ContextEntryTarget::Primitive {
domain: ContextDomain::TransportHeader,
name: context_name(name),
},
form: ContextEntrySelectorForm::OriginalKeyValue,
})
.into(),
},
},
ContextDeclaration::HeaderCapture { policy: capture },
]
.into_iter()
.collect(),
);
let capture = bindings
.header_capture_policy(&pipeline("group", "pipeline"), &ConfigNodeId::from("node"))
.expect("compiled capture policy");
let mut headers = TransportHeaders::new();
let _ = capture.capture_from_pairs(
[
("X-First", b"first".as_slice()),
("X-Second", b"second".as_slice()),
("X-Alias-A", b"alias-a".as_slice()),
("X-Alias-B", b"alias-b".as_slice()),
]
.into_iter(),
&mut headers,
);
assert_eq!(headers.get(0).expect("first header").wire_name(), "X-First");
assert_eq!(
headers.get(1).expect("second header").wire_name(),
"x-second"
);
assert_eq!(
headers.get(2).expect("first alias").wire_name(),
"X-Alias-A"
);
assert_eq!(
headers.get(3).expect("second alias").wire_name(),
"X-Alias-B"
);
}
#[test]
fn declarations_require_only_the_requested_name_form() {
let declarations: NodeContextDeclarations = [
ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries {
entries: vec![ContextEntrySelector {
target: ContextEntryTarget::Primitive {
domain: ContextDomain::TransportHeader,
name: context_name("original"),
},
form: ContextEntrySelectorForm::OriginalKeyValue,
}]
.into_boxed_slice(),
},
},
ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries {
entries: vec![ContextEntrySelector {
target: ContextEntryTarget::Primitive {
domain: ContextDomain::TransportHeader,
name: context_name("value"),
},
form: ContextEntrySelectorForm::Value,
}]
.into_boxed_slice(),
},
},
]
.into_iter()
.collect();
let requirements = context_runtime_requirements(declarations);
assert!(requirements.preserves_original_name(&context_name("original")));
assert!(!requirements.preserves_original_name(&context_name("value")));
}
#[test]
fn requirements_canonicalize_default_and_overrides() {
let propagation: HeaderPropagationConfig = serde_json::from_value(serde_json::json!({
"default": {
"selector": {"type": "all_captured"},
"name": "preserve"
},
"overrides": [{
"match": {"stored_names": ["Authorization"]},
"name": "stored_name"
}]
}))
.expect("valid propagation policy");
let propagation = HeaderPropagationPolicy::compile(propagation, &[])
.expect("propagation policy compiles");
let requirements = context_runtime_requirements(
[ContextDeclaration::HeaderPropagation {
policy: propagation,
}]
.into_iter()
.collect(),
);
assert!(
requirements
.original_name_retention
.default_preserve_original
);
assert!(!requirements.preserves_original_name(&context_name("AUTHORIZATION")));
assert!(requirements.preserves_original_name(&context_name("X-Tenant")));
assert_eq!(
requirements.original_name_retention.overrides,
BTreeMap::from([(Box::<str>::from("authorization"), false)])
);
}
#[test]
fn installed_requirements_support_only_available_original_names() {
let installed = context_runtime_requirements(
[
ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries {
entries: vec![ContextEntrySelector {
target: ContextEntryTarget::Primitive {
domain: ContextDomain::TransportHeader,
name: context_name("x-tenant"),
},
form: ContextEntrySelectorForm::OriginalKeyValue,
}]
.into_boxed_slice(),
},
},
ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries {
entries: vec![ContextEntrySelector {
target: ContextEntryTarget::Primitive {
domain: ContextDomain::TransportHeader,
name: context_name("authorization"),
},
form: ContextEntrySelectorForm::Value,
}]
.into_boxed_slice(),
},
},
]
.into_iter()
.collect(),
);
let removed = context_runtime_requirements(NodeContextDeclarations::default());
let supported = context_runtime_requirements(
[ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries {
entries: vec![ContextEntrySelector {
target: ContextEntryTarget::Primitive {
domain: ContextDomain::TransportHeader,
name: context_name("X-Tenant"),
},
form: ContextEntrySelectorForm::OriginalKeyValue,
}]
.into_boxed_slice(),
},
}]
.into_iter()
.collect(),
);
let unsupported_name = context_runtime_requirements(
[ContextDeclaration::Consumes {
selector: ContextConsumerSelector::Entries {
entries: vec![ContextEntrySelector {
target: ContextEntryTarget::Primitive {
domain: ContextDomain::TransportHeader,
name: context_name("x-request-id"),
},
form: ContextEntrySelectorForm::OriginalKeyValue,
}]
.into_boxed_slice(),
},
}]
.into_iter()
.collect(),
);
let unsupported_default = context_runtime_requirements(
[ContextDeclaration::HeaderPropagation {
policy: HeaderPropagationPolicy::compile(
serde_json::from_value(serde_json::json!({
"default": {
"selector": {"type": "all_captured"},
"name": "preserve"
}
}))
.expect("valid propagation policy"),
&[],
)
.expect("propagation policy compiles"),
}]
.into_iter()
.collect(),
);
assert!(installed.can_satisfy(&removed));
assert!(installed.can_satisfy(&supported));
assert!(!installed.can_satisfy(&unsupported_name));
assert!(!installed.can_satisfy(&unsupported_default));
}
#[test]
fn header_propagation_policy_is_a_context_declaration() {
let policy: HeaderPropagationConfig = serde_json::from_value(serde_json::json!({
"default": {
"selector": {
"type": "named",
"named": ["preserved"]
},
"name": "preserve"
}
}))
.expect("valid propagation policy");
let policy =
HeaderPropagationPolicy::compile(policy, &[]).expect("propagation policy compiles");
let declarations: NodeContextDeclarations = [ContextDeclaration::HeaderPropagation {
policy: policy.clone(),
}]
.into_iter()
.collect();
let compiled = compiled_bindings(declarations.clone());
let requirements = context_runtime_requirements(declarations.clone());
assert!(requirements.preserves_original_name(&context_name("preserved")));
assert!(!requirements.preserves_original_name(&context_name("other")));
assert_eq!(
compiled.header_propagation_policy(
&pipeline("group", "pipeline"),
&ConfigNodeId::from("node")
),
Some(&policy),
);
}
#[test]
fn wrapper_declarations_resolve_policy_precedence() {
let identity_policy: AuthorizedIdentityPolicy =
serde_json::from_value(serde_json::json!([{"claim": "sub", "store_as": "tenant"}]))
.expect("valid authorized identity policy");
let node_capture = HeaderCapturePolicy::new(
CaptureDefaults::default(),
vec![CaptureRule {
match_names: vec![context_name("node")],
store_as: None,
sensitive: false,
value_kind: None,
}],
);
let pipeline_policy = TransportHeadersPolicy {
header_capture: HeaderCapturePolicy::new(
CaptureDefaults::default(),
vec![CaptureRule {
match_names: vec![context_name("pipeline")],
store_as: None,
sensitive: false,
value_kind: None,
}],
),
..Default::default()
};
let mut receiver = NodeUserConfig::new_receiver_config("urn:test:receiver:example");
receiver.header_capture = Some(node_capture.clone());
assert_eq!(
PipelineFactory::<()>::wrapper_context_declarations(
&receiver,
&Some(pipeline_policy.clone()),
&Some(identity_policy.clone()),
&[],
)
.expect("wrapper declarations"),
[
ContextDeclaration::HeaderCapture {
policy: node_capture,
},
ContextDeclaration::AuthorizedIdentityCapture {
policy: identity_policy.clone(),
},
]
.into_iter()
.collect(),
);
let receiver = NodeUserConfig::new_receiver_config("urn:test:receiver:example");
assert_eq!(
PipelineFactory::<()>::wrapper_context_declarations(
&receiver,
&Some(pipeline_policy.clone()),
&Some(identity_policy.clone()),
&[],
)
.expect("wrapper declarations"),
[
ContextDeclaration::HeaderCapture {
policy: pipeline_policy.header_capture.clone(),
},
ContextDeclaration::AuthorizedIdentityCapture {
policy: identity_policy.clone(),
},
]
.into_iter()
.collect(),
);
let mut exporter = NodeUserConfig::new_exporter_config("urn:test:exporter:example");
let node_propagation = HeaderPropagationConfig::default();
exporter.header_propagation = Some(node_propagation.clone());
assert_eq!(
PipelineFactory::<()>::wrapper_context_declarations(
&exporter,
&Some(pipeline_policy.clone()),
&Some(identity_policy.clone()),
&[],
)
.expect("wrapper declarations"),
[ContextDeclaration::HeaderPropagation {
policy: HeaderPropagationPolicy::compile(node_propagation, &[])
.expect("node propagation policy compiles"),
}]
.into_iter()
.collect(),
);
let exporter = NodeUserConfig::new_exporter_config("urn:test:exporter:example");
assert_eq!(
PipelineFactory::<()>::wrapper_context_declarations(
&exporter,
&Some(pipeline_policy.clone()),
&Some(identity_policy),
&[],
)
.expect("wrapper declarations"),
[ContextDeclaration::HeaderPropagation {
policy: HeaderPropagationPolicy::compile(pipeline_policy.header_propagation, &[])
.expect("pipeline propagation policy compiles"),
}]
.into_iter()
.collect(),
);
}
#[test]
fn wrapper_compiles_conditional_composite_header_propagation() {
let context: otel_arrow_dfe_config::context_policy::ContextPolicy = serde_yaml::from_str(
r#"
entries:
tenant:
- type: transport_header
name: workspace
store_as: workspace_id
- type: transport_header_match
name: environment
value: production
"#,
)
.expect("valid context policy");
let (name, definition) = context.entries.into_iter().next().expect("declaration");
let declaration = ConfigContextEntryDeclaration {
scope: otel_arrow_dfe_config::context_policy::ContextScope::Engine,
name,
definition,
};
let mut exporter = NodeUserConfig::new_exporter_config("urn:test:exporter:example");
exporter.header_propagation = Some(
serde_yaml::from_str(
r#"
default:
selector:
type: named
named: [tenant:workspace_id]
name: stored_name
"#,
)
.expect("valid propagation policy"),
);
let declarations = PipelineFactory::<()>::wrapper_context_declarations(
&exporter,
&None,
&None,
&[declaration],
)
.expect("wrapper declarations");
let ContextDeclaration::HeaderPropagation { policy } =
declarations.iter().next().expect("propagation declaration")
else {
panic!("expected header propagation declaration");
};
let mut headers = TransportHeaders::new();
headers.push(TransportHeader::text(context_name("workspace"), b"acme"));
assert_eq!(policy.propagate(&headers).count(), 0);
headers.push(TransportHeader::text(
context_name("environment"),
b"production",
));
let propagated = policy.propagate(&headers).collect::<Vec<_>>();
assert_eq!(propagated.len(), 1);
assert_eq!(propagated[0].header_name, "workspace_id");
assert_eq!(propagated[0].value, b"acme");
}
#[test]
fn full_yaml_compilation_tracks_conditional_composite_changes() {
let current_composite = "[{type: transport_header, name: workspace, store_as: workspace_id}, \
{type: transport_header, name: account, store_as: account_id}, \
{type: transport_header_match, name: environment, value: production}]";
let factory = test_pipeline_factory();
let current = resolve_conditional_pipeline(current_composite, "tenant:workspace_id");
let installed = factory
.compile_initial_context(¤t)
.expect("initial context compiles");
let pipeline = pipeline("default", "main");
let exporter = ConfigNodeId::from("exporter");
let policy = installed
.bindings
.header_propagation_policy(&pipeline, &exporter)
.expect("compiled exporter propagation policy");
let mut headers = TransportHeaders::new();
headers.push(TransportHeader::text(context_name("workspace"), b"acme"));
headers.push(TransportHeader::text(
context_name("environment"),
b"production",
));
let propagated = policy.propagate(&headers).collect::<Vec<_>>();
assert_eq!(propagated.len(), 1);
assert_eq!(propagated[0].header_name, "workspace_id");
let changed_condition = resolve_conditional_pipeline(
"[{type: transport_header, name: workspace, store_as: workspace_id}, \
{type: transport_header, name: account, store_as: account_id}, \
{type: transport_header_match, name: environment, value: staging}]",
"tenant:workspace_id",
);
let condition_candidate = factory
.compile_candidate_context(&changed_condition, &installed.runtime_requirements)
.expect("condition candidate compiles");
assert!(
!installed
.bindings
.pipeline_bindings_match(&condition_candidate.bindings, &pipeline)
);
let changed_member = resolve_conditional_pipeline(current_composite, "tenant:account_id");
let member_candidate = factory
.compile_candidate_context(&changed_member, &installed.runtime_requirements)
.expect("member candidate compiles");
assert!(
!installed
.bindings
.pipeline_bindings_match(&member_candidate.bindings, &pipeline)
);
}
#[test]
fn full_yaml_compilation_ignores_composite_condition_order() {
let current = resolve_conditional_pipeline(
"[{type: transport_header, name: workspace, store_as: workspace_id}, \
{type: transport_header_match, name: environment, value: production}, \
{type: transport_header_match, name: region, value: us-east}]",
"tenant:workspace_id",
);
let reordered = resolve_conditional_pipeline(
"[{type: transport_header, name: workspace, store_as: workspace_id}, \
{type: transport_header_match, name: region, value: us-east}, \
{type: transport_header_match, name: environment, value: production}]",
"tenant:workspace_id",
);
let factory = test_pipeline_factory();
let installed = factory
.compile_initial_context(¤t)
.expect("initial context compiles");
let candidate = factory
.compile_candidate_context(&reordered, &installed.runtime_requirements)
.expect("reordered context compiles");
assert!(
installed
.bindings
.pipeline_bindings_match(&candidate.bindings, &pipeline("default", "main"))
);
}
#[test]
fn full_yaml_compilation_reports_actionable_composite_selector_errors() {
let cases = [
(
"[{type: transport_header, name: workspace, store_as: workspace_id}]",
"missing:workspace_id",
"unknown composite context entry `missing`",
),
(
"[{type: transport_header, name: workspace, store_as: workspace_id}]",
"tenant:missing",
"context entry reference `tenant:missing` does not select a transport-header member",
),
(
"[{type: authorized_identity, name: customer_id}]",
"tenant:customer_id",
"context entry reference `tenant:customer_id` selects authorized-identity member `customer_id`, which cannot be propagated as a transport header",
),
];
let factory = test_pipeline_factory();
for (composite, selector, expected) in cases {
let resolved = resolve_conditional_pipeline(composite, selector);
let error = factory
.compile_initial_context(&resolved)
.expect_err("invalid selector must fail startup");
let message = error.to_string();
assert!(message.contains(expected), "{message}");
}
}
#[test]
fn empty_authorized_identity_policy_produces_no_binding() {
let receiver = NodeUserConfig::new_receiver_config("urn:test:receiver:example");
for policy in [None, Some(AuthorizedIdentityPolicy::default())] {
let declarations =
PipelineFactory::<()>::wrapper_context_declarations(&receiver, &None, &policy, &[])
.expect("wrapper declarations");
assert!(declarations.is_empty());
let compiled = compiled_bindings(declarations);
let node = compiled
.by_pipeline
.get(&pipeline("group", "pipeline"))
.and_then(|nodes| nodes.get(&ConfigNodeId::from("node")))
.expect("compiled node binding");
assert!(node.is_empty());
assert!(
compiled
.authorized_identity_policy(
&pipeline("group", "pipeline"),
&ConfigNodeId::from("node"),
)
.is_none()
);
}
}
#[test]
fn authorized_identity_policy_is_a_compiled_receiver_binding() {
let policy: AuthorizedIdentityPolicy =
serde_json::from_value(serde_json::json!([{"claim": "sub", "store_as": "tenant"}]))
.expect("valid authorized identity policy");
let declarations: NodeContextDeclarations =
[ContextDeclaration::AuthorizedIdentityCapture {
policy: policy.clone(),
}]
.into_iter()
.collect();
let compiled = compiled_bindings(declarations);
let changed_policy: AuthorizedIdentityPolicy = serde_json::from_value(
serde_json::json!([{"claim": "groups", "store_as": "access_groups"}]),
)
.expect("valid changed authorized identity policy");
let changed = compiled_bindings(
[ContextDeclaration::AuthorizedIdentityCapture {
policy: changed_policy,
}]
.into_iter()
.collect(),
);
let pipeline = pipeline("group", "pipeline");
assert_eq!(
compiled.authorized_identity_policy(&pipeline, &ConfigNodeId::from("node")),
Some(&policy),
);
assert!(!compiled.pipeline_bindings_match(&changed, &pipeline));
assert!(!changed.pipeline_bindings_match(&compiled, &pipeline));
}
#[test]
fn parsed_config_declarations_are_validated_against_compiled_policy() {
let pipeline = pipeline("group", "pipeline");
let node: ConfigNodeId = "node".into();
let matching: TestDeclarationConfig =
serde_json::from_value(serde_json::json!({"entry": "expected"}))
.expect("valid matching config");
let changed: TestDeclarationConfig =
serde_json::from_value(serde_json::json!({"entry": "changed"}))
.expect("valid changed config");
let propagation_declaration = ContextDeclaration::HeaderPropagation {
policy: HeaderPropagationPolicy::default(),
};
let declarations = matching
.context_declarations()
.into_iter()
.chain(std::iter::once(propagation_declaration.clone()))
.collect();
let declarations = HashMap::from([(
pipeline.clone(),
HashMap::from([(
node.clone(),
PreparedNodeContextDeclarations::new(declarations, &[])
.expect("valid declarations"),
)]),
)]);
let requirements = ContextRuntimeRequirements::compile(&declarations);
let bindings = CompiledContextBindings::compile(declarations, &requirements);
assert!(
bindings
.validate_node_declarations(&pipeline, &node, &matching.context_declarations(),)
.is_ok()
);
assert!(
bindings
.validate_node_declarations(&pipeline, &node, &changed.context_declarations())
.is_err()
);
assert!(
bindings
.validate_node_declarations(
&pipeline,
&ConfigNodeId::from("other"),
&matching.context_declarations(),
)
.is_err()
);
assert!(
CompiledContextBindings::empty()
.validate_node_declarations(&pipeline, &node, &matching.context_declarations())
.is_err()
);
let ContextDeclaration::HeaderPropagation {
policy: propagation_policy,
} = propagation_declaration
else {
unreachable!("test declaration is header propagation");
};
assert_eq!(
bindings.header_propagation_policy(&pipeline, &node),
Some(&propagation_policy)
);
}
#[test]
fn empty_declarations_require_a_compiled_node() {
let key = pipeline("group", "pipeline");
let node = ConfigNodeId::from("node");
let declarations = NodeContextDeclarations::default();
let bindings = compiled_bindings(declarations.clone());
assert!(
bindings
.validate_node_declarations(&key, &node, &declarations)
.is_ok()
);
assert!(
bindings
.validate_node_declarations(&pipeline("group", "other"), &node, &declarations)
.is_err()
);
assert!(
CompiledContextBindings::empty()
.validate_node_declarations(&key, &node, &declarations)
.is_err()
);
}
}