saluki_components/forwarders/cluster_agent/
mod.rs

1//! Cluster Agent forwarder.
2
3use agent_data_plane_config::shared::SharedConfiguration;
4use async_trait::async_trait;
5use http::{
6    header::AUTHORIZATION,
7    uri::{Authority, Scheme},
8    HeaderName, HeaderValue, Request, Uri,
9};
10use saluki_common::buf::FrozenChunkedBytesBuffer;
11use saluki_config::GenericConfiguration;
12use saluki_core::accounting::{MemoryBounds, MemoryBoundsBuilder, UsageExpr};
13use saluki_core::{
14    components::{forwarders::*, ComponentContext},
15    data_model::payload::PayloadType,
16    observability::ComponentMetricsExt as _,
17};
18use saluki_error::{generic_error, GenericError};
19use saluki_metrics::MetricsBuilder;
20use stringtheory::MetaString;
21use tokio::select;
22use tracing::debug;
23
24use crate::common::datadog::{
25    config::ForwarderConfiguration,
26    endpoints::{ResolvedEndpoint, SingleDestination},
27    io::{EndpointRequestMapper, EndpointRequestMapperFactory, TransactionForwarder},
28    telemetry::ComponentTelemetry,
29    transaction::{Metadata, Transaction, TransactionBody},
30    DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT,
31};
32
33const CLUSTER_AGENT_SERIES_PATH: &str = "/series";
34
35static DD_API_KEY_HEADER: HeaderName = HeaderName::from_static("dd-api-key");
36
37/// Cluster Agent forwarder configuration.
38///
39/// This forwarder sends encoded metric series payloads to the Cluster Agent `/series` endpoint using bearer-token
40/// authorization. It intentionally does not use Datadog intake endpoint configuration or `DD-Api-Key` auth.
41pub struct ClusterAgentForwarderConfiguration {
42    forwarder_config: ForwarderConfiguration,
43    auth_header_value: HeaderValue,
44}
45
46impl ClusterAgentForwarderConfiguration {
47    /// Creates a new `ClusterAgentForwarderConfiguration` from the given Cluster Agent endpoint and bearer token.
48    ///
49    /// The Cluster Agent is the only destination this forwarder sends to: the configured Datadog
50    /// intake endpoints, the alternate metrics intakes, and dual shipping do not apply to it, and it
51    /// accepts only V2 series payloads. Because the destination is part of construction, no later
52    /// step can overwrite it.
53    ///
54    /// # Errors
55    ///
56    /// Returns an error if the bearer token cannot be represented in an HTTP header.
57    pub fn from_configuration(
58        shared: &SharedConfiguration, config: &GenericConfiguration, endpoint_url: String, auth_token: String,
59    ) -> Result<Self, GenericError> {
60        let auth_header_value = bearer_auth_header_value(&auth_token)?;
61        let destination = SingleDestination {
62            url: endpoint_url,
63            api_key: auth_token,
64            api_key_refresh_config_path: None,
65            accepts_v3_series: false,
66        };
67        let forwarder_config = ForwarderConfiguration::for_single_destination(shared, config, &destination)
68            .with_allow_arbitrary_tags(false);
69
70        Ok(Self {
71            forwarder_config,
72            auth_header_value,
73        })
74    }
75}
76
77#[async_trait]
78impl ForwarderBuilder for ClusterAgentForwarderConfiguration {
79    fn input_payload_type(&self) -> PayloadType {
80        PayloadType::Http
81    }
82
83    async fn build(&self, context: ComponentContext) -> Result<Box<dyn Forwarder + Send>, GenericError> {
84        let metrics_builder = MetricsBuilder::from_component_context(&context);
85        let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
86        let endpoint_request_mapper_factory = cluster_agent_request_mapper_factory(self.auth_header_value.clone());
87        let forwarder = TransactionForwarder::from_config_with_endpoint_request_mapper(
88            context,
89            self.forwarder_config.clone(),
90            None,
91            get_cluster_agent_endpoint_name,
92            telemetry.clone(),
93            metrics_builder,
94            endpoint_request_mapper_factory,
95        )?;
96
97        Ok(Box::new(ClusterAgentForwarder { forwarder }))
98    }
99}
100
101impl MemoryBounds for ClusterAgentForwarderConfiguration {
102    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
103        builder
104            .minimum()
105            .with_single_value::<ClusterAgentForwarder>("component struct")
106            .with_array::<Transaction<FrozenChunkedBytesBuffer>>("requests channel", 8);
107
108        builder.firm().with_expr(UsageExpr::sum(
109            "in-flight requests",
110            UsageExpr::config(
111                "forwarder_retry_queue_payloads_max_size",
112                self.forwarder_config.retry().queue_max_size_bytes() as usize,
113            ),
114            UsageExpr::product(
115                "high priority queue",
116                UsageExpr::config(
117                    "forwarder_high_prio_buffer_size",
118                    self.forwarder_config.endpoint_buffer_size(),
119                ),
120                UsageExpr::constant("maximum compressed payload size", DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT),
121            ),
122        ));
123    }
124}
125
126/// Cluster Agent forwarder.
127pub struct ClusterAgentForwarder {
128    forwarder: TransactionForwarder<FrozenChunkedBytesBuffer>,
129}
130
131#[async_trait]
132impl Forwarder for ClusterAgentForwarder {
133    async fn run(mut self: Box<Self>, mut context: ForwarderContext) -> Result<(), GenericError> {
134        let Self { forwarder } = *self;
135
136        let mut health = context.take_health_handle();
137        let forwarder = forwarder.spawn().await;
138
139        health.mark_ready();
140        debug!("Cluster Agent forwarder started.");
141
142        loop {
143            select! {
144                _ = health.live() => continue,
145                maybe_payload = context.payloads().next() => match maybe_payload {
146                    Some(payload) => if let Some(http_payload) = payload.try_into_http_payload() {
147                        let (payload_meta, request) = http_payload.into_parts();
148                        let transaction_meta = Metadata::from_event_and_data_point_count(
149                            payload_meta.event_count(),
150                            payload_meta.data_point_count(),
151                        );
152                        let transaction = Transaction::from_original(transaction_meta, request);
153
154                        forwarder.send_transaction(transaction).await?;
155                    }
156                    None => break,
157                },
158            }
159        }
160
161        forwarder.shutdown().await;
162
163        debug!("Cluster Agent forwarder stopped.");
164
165        Ok(())
166    }
167}
168
169fn bearer_auth_header_value(auth_token: &str) -> Result<HeaderValue, GenericError> {
170    let raw_value = format!("Bearer {auth_token}");
171    HeaderValue::from_str(&raw_value)
172        .map_err(|_| generic_error!("cluster_agent.auth_token contains characters that are invalid in HTTP headers."))
173}
174
175fn cluster_agent_request_mapper_factory<B>(auth_header_value: HeaderValue) -> EndpointRequestMapperFactory<B>
176where
177    B: 'static,
178{
179    std::sync::Arc::new(move |endpoint| cluster_agent_request_mapper(endpoint, auth_header_value.clone()))
180}
181
182fn cluster_agent_request_mapper<B>(
183    endpoint: ResolvedEndpoint, auth_header_value: HeaderValue,
184) -> EndpointRequestMapper<B> {
185    let new_uri_authority = Authority::try_from(endpoint.endpoint().authority())
186        .expect("should not fail to construct new endpoint authority");
187    let new_uri_scheme =
188        Scheme::try_from(endpoint.endpoint().scheme()).expect("should not fail to construct new endpoint scheme");
189
190    Box::new(move |mut request: Request<TransactionBody<B>>| {
191        let new_uri = Uri::builder()
192            .scheme(new_uri_scheme.clone())
193            .authority(new_uri_authority.clone())
194            .path_and_query(CLUSTER_AGENT_SERIES_PATH)
195            .build()
196            .expect("should not fail to construct Cluster Agent URI");
197        *request.uri_mut() = new_uri;
198        request.headers_mut().remove(&DD_API_KEY_HEADER);
199        request.headers_mut().insert(AUTHORIZATION, auth_header_value.clone());
200
201        request
202    })
203}
204
205fn get_cluster_agent_endpoint_name(_uri: &Uri) -> Option<MetaString> {
206    Some(MetaString::from_static("cluster_agent_series"))
207}
208
209#[cfg(test)]
210mod tests {
211    use std::collections::HashMap;
212
213    use agent_data_plane_config::{
214        shared::{AltMetricsIntake, V3SeriesMode},
215        ConfigValue,
216    };
217    use http::Method;
218    use saluki_config::ConfigurationLoader;
219
220    use super::*;
221    use crate::common::datadog::{endpoints::EndpointRoute, test_util::shared_configuration};
222
223    #[test]
224    fn request_mapper_preserves_cluster_agent_series_identity_and_sets_bearer_auth() {
225        let auth_header_value = bearer_auth_header_value("secret-token").expect("auth header should be valid");
226        let endpoint = ResolvedEndpoint::from_raw_endpoint("https://cluster-agent.example.com:5005", "secret-token")
227            .expect("endpoint should resolve");
228        let mut mapper = cluster_agent_request_mapper::<()>(endpoint, auth_header_value);
229        let request = Request::builder()
230            .method(Method::POST)
231            .uri("/api/v2/series")
232            .header("dd-api-key", "primary-api-key")
233            .body(TransactionBody::<()>::Rehydrated(None))
234            .expect("request should build");
235
236        let input_endpoint_name = get_cluster_agent_endpoint_name(request.uri());
237        let request = mapper(request);
238        let mapped_endpoint_name = get_cluster_agent_endpoint_name(request.uri());
239
240        assert_eq!(input_endpoint_name.as_deref(), Some("cluster_agent_series"));
241        assert_eq!(mapped_endpoint_name.as_deref(), Some("cluster_agent_series"));
242        assert_eq!(
243            request.uri().to_string(),
244            "https://cluster-agent.example.com:5005/series"
245        );
246        assert_eq!(request.headers().get(AUTHORIZATION).unwrap(), "Bearer secret-token");
247        assert!(request.headers().get("dd-api-key").is_none());
248    }
249
250    #[test]
251    fn bearer_auth_header_rejects_invalid_token() {
252        assert!(bearer_auth_header_value("bad\ntoken").is_err());
253    }
254
255    #[tokio::test]
256    async fn configuration_uses_only_cluster_agent_endpoint() {
257        // Every configured Datadog intake setting here conflicts with the Cluster Agent destination:
258        // a different API key, an explicit `dd_url`, a `site`, an additional endpoint, an alternate
259        // metrics intake, and V3 series routing. None of them may reach the built forwarder.
260        let (raw_config, _) = ConfigurationLoader::for_tests(None, None, false).await;
261        let mut shared = shared_configuration();
262        shared.endpoints.api_key = "primary-api-key".to_string();
263        shared.endpoints.site = ConfigValue::explicit("datadoghq.eu".to_string());
264        shared.endpoints.dd_url = ConfigValue::explicit("https://app.datadoghq.com".to_string());
265        shared.endpoints.additional_endpoints = HashMap::from([(
266            "https://additional.example.com".to_string(),
267            vec!["additional-api-key".to_string()],
268        )]);
269        shared.endpoints.opw_intake = AltMetricsIntake {
270            enabled: true,
271            url: "https://opw.example.com".to_string(),
272            use_v3_series: true,
273        };
274        shared.metrics_encoding.v3_series_mode = V3SeriesMode::Enabled;
275        shared.metrics_encoding.v3_series_endpoint_modes =
276            HashMap::from([("https://app.datadoghq.com".to_string(), V3SeriesMode::Enabled)]);
277
278        // The same configuration drives a plain Datadog forwarder, which does honor all of it.
279        let datadog_forwarder = ForwarderConfiguration::from_configuration(&shared, &raw_config);
280        assert_eq!(V3SeriesMode::Enabled, datadog_forwarder.use_v3_api_series().enabled);
281
282        let config = ClusterAgentForwarderConfiguration::from_configuration(
283            &shared,
284            &raw_config,
285            "https://cluster-agent.example.com".to_string(),
286            "secret-token".to_string(),
287        )
288        .expect("Cluster Agent forwarder configuration should parse");
289        let endpoints = config
290            .forwarder_config
291            .build_routable_endpoints(None)
292            .expect("endpoint should resolve");
293
294        assert_eq!(endpoints.len(), 1);
295        assert_eq!(endpoints[0].route(), EndpointRoute::Primary);
296        assert_eq!(
297            endpoints[0].endpoint().endpoint().as_str(),
298            "https://cluster-agent.example.com/"
299        );
300        assert_eq!(endpoints[0].endpoint().cached_api_key(), "secret-token");
301        assert_eq!(
302            V3SeriesMode::Disabled,
303            config.forwarder_config.use_v3_api_series().enabled
304        );
305        assert!(config.forwarder_config.use_v3_api_series().endpoints.is_empty());
306        assert!(config.forwarder_config.v3_api().series.endpoints.is_empty());
307        assert!(!config.forwarder_config.allow_arbitrary_tags());
308    }
309}