-
Notifications
You must be signed in to change notification settings - Fork 129
Add experimental resource validator processor #1956
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Merged
jmacd
merged 23 commits into
open-telemetry:main
from
lalitb:resource-validator-processor
Feb 9, 2026
Merged
Changes from 16 commits
Commits
Show all changes
23 commits
Select commit
Hold shift + click to select a range
5657140
initial commit
lalitb 56cd438
lint
lalitb e153a70
fix
lalitb d236880
Merge branch 'main' into resource-validator-processor
lalitb 54190d2
add arrow validation
lalitb 422d2ed
review
lalitb 5c644ba
nit opt
lalitb bd13a97
Merge branch 'main' into resource-validator-processor
lalitb cc1e8d6
fix doc
lalitb 62ea9a7
fix
lalitb 2669243
Merge branch 'main' into resource-validator-processor
lalitb 32a180b
review comment
lalitb 4bc0b30
add comment for duplicates
lalitb 0d9556f
Merge branch 'main' into resource-validator-processor
lalitb 215e1f4
Merge branch 'main' into resource-validator-processor
lalitb ede9f5e
Merge branch 'main' into resource-validator-processor
cijothomas afa32f9
Merge branch 'main' into resource-validator-processor
lalitb 3a44358
fix
lalitb fa6fce7
fix
lalitb 71f25c1
Review comments
lalitb b3ef9d6
review comemnt
lalitb 18f811b
lint
lalitb 188915c
fmt
lalitb File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
125 changes: 125 additions & 0 deletions
125
...ap-dataflow/crates/otap/src/experimental/resource_validator_processor/README.md
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,125 @@ | ||
| # Resource Validator Processor | ||
|
|
||
| The Resource Validator Processor validates that telemetry data contains a | ||
| required resource attribute with a value from an allowed list. Requests that | ||
| fail validation are permanently NACKed, enabling clients to detect | ||
| misconfiguration immediately rather than having data silently dropped. | ||
|
|
||
| ## Example Use Case | ||
|
|
||
| In multi-tenant cloud environments, telemetry includes a resource attribute | ||
| (e.g., `cloud.resource_id`) containing an identifier for the resource. | ||
| The collector must: | ||
|
|
||
| 1. Validate the attribute exists on the Resource | ||
| 2. Check value is in an allowed list | ||
| 3. Reject with error (HTTP 400 / gRPC INVALID_ARGUMENT) on failure | ||
|
|
||
| This enables clients to detect misconfiguration immediately rather than having | ||
| data silently dropped. | ||
|
|
||
| ## Why Existing Processors Don't Work | ||
|
|
||
| | Processor | Behavior | Gap | | ||
| | -------------------- | -------------------------------- | ------------------ | | ||
| | Filter Processor | Silently drops non-matching data | No error to client | | ||
| | Transform Processor | Modifies/transforms data | Cannot reject | | ||
| | Attributes Processor | Adds/updates/deletes attributes | Cannot reject | | ||
|
|
||
| None validate attribute values against an allowlist with error propagation | ||
| to client. | ||
|
|
||
| ## Configuration | ||
|
|
||
| ```yaml | ||
| processors: | ||
| resource_validator: | ||
| # The resource attribute key that must be present (required, no default) | ||
| required_attribute: "cloud.resource_id" | ||
|
|
||
| # List of allowed values (required - empty list rejects all values) | ||
| allowed_values: | ||
| - "/subscriptions/xxx/resourceGroups/yyy/..." | ||
| - "/subscriptions/aaa/resourceGroups/bbb/..." | ||
|
|
||
| # Case-sensitive comparison (default: true) | ||
| case_sensitive: false | ||
| ``` | ||
|
|
||
| ## Behavior | ||
|
|
||
| | Condition | Result | | ||
| | -------------------------------------------- | ------------------------- | | ||
| | Attribute present with value in allowed list | Pass through | | ||
| | Attribute present, empty allowed_values list | Permanent NACK - HTTP 400 | | ||
| | Attribute missing | Permanent NACK - HTTP 400 | | ||
| | Attribute wrong type (not string) | Permanent NACK - HTTP 400 | | ||
| | Attribute value not in allowed list | Permanent NACK - HTTP 400 | | ||
|
lalitb marked this conversation as resolved.
Outdated
|
||
|
|
||
| ## Metrics | ||
|
|
||
| - `resource_validator_batches_accepted` - Batches that passed validation | ||
| - `resource_validator_batches_rejected_missing` - Rejected: missing attribute | ||
| - `resource_validator_batches_rejected_not_allowed` - Rejected: invalid value | ||
|
|
||
| ## Feature Flag | ||
|
|
||
| This processor is experimental and requires the `resource-validator-processor` | ||
| feature flag: | ||
|
|
||
| ```toml | ||
| [dependencies] | ||
| otap-df-otap = { version = "...", features = ["resource-validator-processor"] } | ||
| ``` | ||
|
|
||
| ## Extensibility for Dynamic Auth Context | ||
|
|
||
| The processor is designed to be extensible for future dynamic validation via | ||
| auth context: | ||
|
|
||
| ### Current Implementation (Static) | ||
|
|
||
| Uses `allowed_values` from configuration: | ||
|
|
||
| ```yaml | ||
| processors: | ||
| resource_validator: | ||
| allowed_values: | ||
| - "/subscriptions/xxx/..." # Static list from config | ||
| ``` | ||
|
|
||
| ### Future Implementation (Dynamic) | ||
|
|
||
| When SAT auth extension is ready, allowed values can come from request auth | ||
| context: | ||
|
|
||
| ```text | ||
| +-----------+ +------------------+ +--------------------+ +----------+ | ||
| | Client |-->| OTLP Receiver |-->| Resource Validator |-->| Exporter | | ||
| | + token | | + Auth Extension | | Processor | | | | ||
| +-----------+ +------------------+ +--------------------+ +----------+ | ||
| | | | ||
| v v | ||
| Auth Extension: Processor reads from | ||
| 1. Validates token context instead of config | ||
| 2. Extracts claims via get_allowed_values() | ||
| 3. Sets ctx.auth | ||
| ``` | ||
|
|
||
| ### Extension Points | ||
|
|
||
| The processor provides these extension points for dynamic auth: | ||
|
|
||
| 1. **`AllowedValuesSource` enum**: Supports `Static` (current) and `Dynamic` | ||
| (future) modes | ||
| 2. **`get_allowed_values()` method**: Returns allowed values for a request; | ||
| can be extended to check `pdata.context().auth()` first | ||
| 3. **Fallback support**: Dynamic mode includes config fallback for requests | ||
| without auth context | ||
|
|
||
| When the SAT auth extension is implemented, the `get_allowed_values()` method | ||
| can be updated to: | ||
|
|
||
| 1. Check if the request context contains auth information | ||
| 2. Extract allowed resource IDs from auth claims | ||
| 3. Fall back to static config if auth context is unavailable | ||
150 changes: 150 additions & 0 deletions
150
rust/otap-dataflow/crates/otap/src/experimental/resource_validator_processor/config.rs
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,150 @@ | ||
| // Copyright The OpenTelemetry Authors | ||
| // SPDX-License-Identifier: Apache-2.0 | ||
|
|
||
| //! Configuration for the Resource Validator Processor | ||
|
|
||
| use otap_df_config::error::Error as ConfigError; | ||
| use serde::{Deserialize, Serialize}; | ||
| use std::collections::HashSet; | ||
|
|
||
| /// Configuration for the Resource Validator Processor | ||
| /// | ||
| /// This processor validates that a required resource attribute exists and its value | ||
| /// is in an allowed list. Requests that fail validation are NACKed with a permanent | ||
| /// error, enabling clients to detect misconfiguration immediately. | ||
| /// | ||
| /// # Example Configuration | ||
| /// ```yaml | ||
| /// processors: | ||
| /// resource_validator: | ||
| /// required_attribute: "cloud.resource_id" | ||
| /// allowed_values: | ||
| /// - "/subscriptions/xxx/resourceGroups/yyy/..." | ||
| /// case_sensitive: false # optional, defaults to true | ||
| /// ``` | ||
| #[derive(Debug, Clone, Serialize, Deserialize)] | ||
| pub struct Config { | ||
| /// The resource attribute key that must be present on all resources. | ||
| /// This is a required field with no default. | ||
| pub required_attribute: String, | ||
|
lalitb marked this conversation as resolved.
Outdated
|
||
|
|
||
| /// List of allowed values for the required attribute. | ||
| /// Empty list rejects all values. | ||
|
lalitb marked this conversation as resolved.
|
||
| #[serde(default)] | ||
| pub allowed_values: Vec<String>, | ||
|
|
||
| /// Whether to perform case-sensitive comparison of attribute values. | ||
| /// Default: true | ||
| #[serde(default = "default_case_sensitive")] | ||
| pub case_sensitive: bool, | ||
|
lalitb marked this conversation as resolved.
|
||
| } | ||
|
|
||
| fn default_case_sensitive() -> bool { | ||
|
lalitb marked this conversation as resolved.
Outdated
|
||
| true | ||
| } | ||
|
|
||
| impl Config { | ||
| /// Validates the configuration. | ||
| pub fn validate(&self) -> Result<(), ConfigError> { | ||
| if self.required_attribute.is_empty() { | ||
|
lalitb marked this conversation as resolved.
Outdated
|
||
| return Err(ConfigError::InvalidUserConfig { | ||
| error: "required_attribute cannot be empty".to_string(), | ||
| }); | ||
| } | ||
| if self.allowed_values.is_empty() { | ||
| return Err(ConfigError::InvalidUserConfig { | ||
| error: "allowed_values cannot be empty (would reject all data)".to_string(), | ||
| }); | ||
| } | ||
| Ok(()) | ||
| } | ||
|
|
||
| /// Returns a pre-processed set of allowed values for efficient lookup. | ||
| /// If case_sensitive is false, all values are lowercased. | ||
| #[must_use] | ||
| pub fn allowed_values_set(&self) -> HashSet<String> { | ||
| if self.case_sensitive { | ||
| self.allowed_values.iter().cloned().collect() | ||
| } else { | ||
| self.allowed_values | ||
| .iter() | ||
| .map(|v| v.to_lowercase()) | ||
| .collect() | ||
| } | ||
| } | ||
| } | ||
|
|
||
| #[cfg(test)] | ||
| mod tests { | ||
| use super::*; | ||
|
|
||
| /// Helper to create a config for testing | ||
| fn test_config(required_attribute: &str, allowed_values: Vec<&str>) -> Config { | ||
| Config { | ||
| required_attribute: required_attribute.to_string(), | ||
| allowed_values: allowed_values.into_iter().map(String::from).collect(), | ||
| case_sensitive: true, | ||
| } | ||
| } | ||
|
|
||
| #[test] | ||
| fn test_validate_empty_attribute() { | ||
| let config = Config { | ||
| required_attribute: "".to_string(), | ||
| allowed_values: vec!["value".to_string()], | ||
| case_sensitive: true, | ||
| }; | ||
| assert!(config.validate().is_err()); | ||
| } | ||
|
|
||
| #[test] | ||
| fn test_validate_empty_allowed_values() { | ||
| let config = test_config("cloud.resource_id", vec![]); | ||
| assert!(config.validate().is_err()); | ||
| } | ||
|
|
||
| #[test] | ||
| fn test_validate_valid_config() { | ||
| let config = test_config("cloud.resource_id", vec!["/subscriptions/123"]); | ||
| assert!(config.validate().is_ok()); | ||
| } | ||
|
|
||
| #[test] | ||
| fn test_allowed_values_set_case_sensitive() { | ||
| let config = Config { | ||
| required_attribute: "cloud.resource_id".to_string(), | ||
| allowed_values: vec!["Value1".to_string(), "Value2".to_string()], | ||
| case_sensitive: true, | ||
| }; | ||
| let set = config.allowed_values_set(); | ||
| assert!(set.contains("Value1")); | ||
| assert!(set.contains("Value2")); | ||
| assert!(!set.contains("value1")); | ||
| } | ||
|
|
||
| #[test] | ||
| fn test_allowed_values_set_case_insensitive() { | ||
| let config = Config { | ||
| required_attribute: "cloud.resource_id".to_string(), | ||
| allowed_values: vec!["Value1".to_string(), "VALUE2".to_string()], | ||
| case_sensitive: false, | ||
| }; | ||
| let set = config.allowed_values_set(); | ||
| assert!(set.contains("value1")); | ||
| assert!(set.contains("value2")); | ||
| assert!(!set.contains("Value1")); | ||
| } | ||
|
|
||
| #[test] | ||
| fn test_deserialize_config() { | ||
| let json = r#"{ | ||
| "required_attribute": "my.attribute", | ||
| "allowed_values": ["val1", "val2"], | ||
| "case_sensitive": false | ||
| }"#; | ||
| let config: Config = serde_json::from_str(json).unwrap(); | ||
| assert_eq!(config.required_attribute, "my.attribute"); | ||
| assert_eq!(config.allowed_values, vec!["val1", "val2"]); | ||
| assert!(!config.case_sensitive); | ||
| } | ||
| } | ||
28 changes: 28 additions & 0 deletions
28
rust/otap-dataflow/crates/otap/src/experimental/resource_validator_processor/metrics.rs
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,28 @@ | ||
| // Copyright The OpenTelemetry Authors | ||
| // SPDX-License-Identifier: Apache-2.0 | ||
|
|
||
| //! Metrics for the Resource Validator Processor | ||
|
|
||
| use otap_df_telemetry::instrument::Counter; | ||
| use otap_df_telemetry_macros::metric_set; | ||
|
|
||
| /// Metrics collected by the Resource Validator Processor. | ||
| /// | ||
| /// Note: Metrics are counted at the batch level, not per telemetry item. | ||
| /// This is intentional because validation is pass/fail for the entire batch - | ||
| /// if any resource in a batch fails validation, the whole batch is NACKed. | ||
|
lalitb marked this conversation as resolved.
Outdated
|
||
| #[metric_set(name = "resource_validator.processor.metrics")] | ||
| #[derive(Debug, Default, Clone)] | ||
| pub struct ResourceValidatorMetrics { | ||
| /// Number of batches that passed validation | ||
|
lalitb marked this conversation as resolved.
|
||
| #[metric(unit = "{batch}")] | ||
| pub batches_accepted: Counter<u64>, | ||
|
|
||
| /// Number of batches rejected due to missing required attribute | ||
| #[metric(unit = "{batch}")] | ||
| pub batches_rejected_missing: Counter<u64>, | ||
|
|
||
| /// Number of batches rejected due to value not in allowed list | ||
| #[metric(unit = "{batch}")] | ||
| pub batches_rejected_not_allowed: Counter<u64>, | ||
| } | ||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.