saluki_components/forwarders/cluster_agent/
mod.rs

1//! Cluster Agent forwarder.
2
3use agent_data_plane_config::shared::{Endpoints, MetricsEncoding};
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,
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    pub fn from_configuration(
49        config: &GenericConfiguration, endpoint_url: String, auth_token: String,
50    ) -> Result<Self, GenericError> {
51        let auth_header_value = bearer_auth_header_value(&auth_token)?;
52        let mut forwarder_config = ForwarderConfiguration::from_configuration(config)?.with_allow_arbitrary_tags(false);
53
54        let endpoint = forwarder_config.endpoint_mut();
55        endpoint.clear_additional_endpoints();
56        endpoint.set_dd_url(endpoint_url);
57        endpoint.set_api_key(auth_token);
58        forwarder_config.clear_opw_metrics_endpoint();
59        forwarder_config.force_v2_series();
60
61        Ok(Self {
62            forwarder_config,
63            auth_header_value,
64        })
65    }
66
67    /// Creates a Cluster Agent forwarder using authoritative typed metrics-routing configuration.
68    pub fn from_configuration_with_metrics_routing(
69        config: &GenericConfiguration, metrics: &MetricsEncoding, endpoints: &Endpoints, endpoint_url: String,
70        auth_token: String,
71    ) -> Result<Self, GenericError> {
72        let mut config = Self::from_configuration(config, endpoint_url, auth_token)?;
73        config
74            .forwarder_config
75            .apply_typed_metrics_configuration(metrics, endpoints);
76        config.forwarder_config.clear_opw_metrics_endpoint();
77        config.forwarder_config.force_v2_series();
78        Ok(config)
79    }
80}
81
82#[async_trait]
83impl ForwarderBuilder for ClusterAgentForwarderConfiguration {
84    fn input_payload_type(&self) -> PayloadType {
85        PayloadType::Http
86    }
87
88    async fn build(&self, context: ComponentContext) -> Result<Box<dyn Forwarder + Send>, GenericError> {
89        let metrics_builder = MetricsBuilder::from_component_context(&context);
90        let telemetry = ComponentTelemetry::from_builder(&metrics_builder);
91        let endpoint_request_mapper_factory = cluster_agent_request_mapper_factory(self.auth_header_value.clone());
92        let forwarder = TransactionForwarder::from_config_with_endpoint_request_mapper(
93            context,
94            self.forwarder_config.clone(),
95            None,
96            get_cluster_agent_endpoint_name,
97            telemetry.clone(),
98            metrics_builder,
99            endpoint_request_mapper_factory,
100        )?;
101
102        Ok(Box::new(ClusterAgentForwarder { forwarder }))
103    }
104}
105
106impl MemoryBounds for ClusterAgentForwarderConfiguration {
107    fn specify_bounds(&self, builder: &mut MemoryBoundsBuilder) {
108        builder
109            .minimum()
110            .with_single_value::<ClusterAgentForwarder>("component struct")
111            .with_array::<Transaction<FrozenChunkedBytesBuffer>>("requests channel", 8);
112
113        builder.firm().with_expr(UsageExpr::sum(
114            "in-flight requests",
115            UsageExpr::config(
116                "forwarder_retry_queue_payloads_max_size",
117                self.forwarder_config.retry().queue_max_size_bytes() as usize,
118            ),
119            UsageExpr::product(
120                "high priority queue",
121                UsageExpr::config(
122                    "forwarder_high_prio_buffer_size",
123                    self.forwarder_config.endpoint_buffer_size(),
124                ),
125                UsageExpr::constant("maximum compressed payload size", DEFAULT_INTAKE_COMPRESSED_SIZE_LIMIT),
126            ),
127        ));
128    }
129}
130
131/// Cluster Agent forwarder.
132pub struct ClusterAgentForwarder {
133    forwarder: TransactionForwarder<FrozenChunkedBytesBuffer>,
134}
135
136#[async_trait]
137impl Forwarder for ClusterAgentForwarder {
138    async fn run(mut self: Box<Self>, mut context: ForwarderContext) -> Result<(), GenericError> {
139        let Self { forwarder } = *self;
140
141        let mut health = context.take_health_handle();
142        let forwarder = forwarder.spawn().await;
143
144        health.mark_ready();
145        debug!("Cluster Agent forwarder started.");
146
147        loop {
148            select! {
149                _ = health.live() => continue,
150                maybe_payload = context.payloads().next() => match maybe_payload {
151                    Some(payload) => if let Some(http_payload) = payload.try_into_http_payload() {
152                        let (payload_meta, request) = http_payload.into_parts();
153                        let transaction_meta = Metadata::from_event_and_data_point_count(
154                            payload_meta.event_count(),
155                            payload_meta.data_point_count(),
156                        );
157                        let transaction = Transaction::from_original(transaction_meta, request);
158
159                        forwarder.send_transaction(transaction).await?;
160                    }
161                    None => break,
162                },
163            }
164        }
165
166        forwarder.shutdown().await;
167
168        debug!("Cluster Agent forwarder stopped.");
169
170        Ok(())
171    }
172}
173
174fn bearer_auth_header_value(auth_token: &str) -> Result<HeaderValue, GenericError> {
175    let raw_value = format!("Bearer {auth_token}");
176    HeaderValue::from_str(&raw_value)
177        .map_err(|_| generic_error!("cluster_agent.auth_token contains characters that are invalid in HTTP headers."))
178}
179
180fn cluster_agent_request_mapper_factory<B>(auth_header_value: HeaderValue) -> EndpointRequestMapperFactory<B>
181where
182    B: 'static,
183{
184    std::sync::Arc::new(move |endpoint| cluster_agent_request_mapper(endpoint, auth_header_value.clone()))
185}
186
187fn cluster_agent_request_mapper<B>(
188    endpoint: ResolvedEndpoint, auth_header_value: HeaderValue,
189) -> EndpointRequestMapper<B> {
190    let new_uri_authority = Authority::try_from(endpoint.endpoint().authority())
191        .expect("should not fail to construct new endpoint authority");
192    let new_uri_scheme =
193        Scheme::try_from(endpoint.endpoint().scheme()).expect("should not fail to construct new endpoint scheme");
194
195    Box::new(move |mut request: Request<TransactionBody<B>>| {
196        let new_uri = Uri::builder()
197            .scheme(new_uri_scheme.clone())
198            .authority(new_uri_authority.clone())
199            .path_and_query(CLUSTER_AGENT_SERIES_PATH)
200            .build()
201            .expect("should not fail to construct Cluster Agent URI");
202        *request.uri_mut() = new_uri;
203        request.headers_mut().remove(&DD_API_KEY_HEADER);
204        request.headers_mut().insert(AUTHORIZATION, auth_header_value.clone());
205
206        request
207    })
208}
209
210fn get_cluster_agent_endpoint_name(_uri: &Uri) -> Option<MetaString> {
211    Some(MetaString::from_static("cluster_agent_series"))
212}
213
214#[cfg(test)]
215mod tests {
216    use http::Method;
217    use saluki_config::ConfigurationLoader;
218    use serde_json::json;
219
220    use super::*;
221    use crate::common::datadog::endpoints::EndpointRoute;
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        let (config, _) = ConfigurationLoader::for_tests(
258            Some(json!({
259                "api_key": "primary-api-key",
260                "dd_url": "https://app.datadoghq.com",
261                "additional_endpoints": {
262                    "https://additional.example.com": ["additional-api-key"]
263                },
264                "observability_pipelines_worker": {
265                    "metrics": {
266                        "enabled": true,
267                        "url": "https://opw.example.com"
268                    }
269                },
270                "use_v3_api_series_enabled": "true",
271                "serializer_experimental_use_v3_api": {
272                    "series": {
273                        "shadow_sites": ["example.com"]
274                    }
275                }
276            })),
277            None,
278            false,
279        )
280        .await;
281
282        let unmodified_forwarder =
283            ForwarderConfiguration::from_configuration(&config).expect("forwarder configuration should parse");
284        assert_eq!("true", unmodified_forwarder.use_v3_api_series().enabled);
285        assert_eq!(
286            &["example.com".to_string()],
287            unmodified_forwarder.v3_api().series.shadow_sites.as_slice()
288        );
289
290        let mut metrics = MetricsEncoding::default();
291        metrics.v3_series_mode.mode = "true".to_string();
292        metrics.v3_api.series.shadow_sites = vec!["example.com".to_string()];
293        let config = ClusterAgentForwarderConfiguration::from_configuration_with_metrics_routing(
294            &config,
295            &metrics,
296            &Endpoints::default(),
297            "https://cluster-agent.example.com".to_string(),
298            "secret-token".to_string(),
299        )
300        .expect("Cluster Agent forwarder configuration should parse");
301        let endpoints = config
302            .forwarder_config
303            .build_routable_endpoints(None)
304            .expect("endpoint should resolve");
305
306        assert_eq!(endpoints.len(), 1);
307        assert_eq!(endpoints[0].route(), EndpointRoute::Primary);
308        assert_eq!(
309            endpoints[0].endpoint().endpoint().as_str(),
310            "https://cluster-agent.example.com/"
311        );
312        assert_eq!(endpoints[0].endpoint().cached_api_key(), "secret-token");
313        assert_eq!("false", config.forwarder_config.use_v3_api_series().enabled);
314        assert!(config.forwarder_config.use_v3_api_series().endpoints.is_empty());
315        assert!(config.forwarder_config.v3_api().series.endpoints.is_empty());
316        assert!(config.forwarder_config.v3_api().series.shadow_sites.is_empty());
317    }
318}