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