saluki_components/forwarders/cluster_agent/
mod.rs1use 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
36pub struct ClusterAgentForwarderConfiguration {
41 forwarder_config: ForwarderConfiguration,
42 auth_header_value: HeaderValue,
43}
44
45impl ClusterAgentForwarderConfiguration {
46 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 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
127pub 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 let forwarder = forwarder.spawn().await;
139
140 health.mark_ready();
141 debug!("Cluster Agent forwarder started.");
142
143 loop {
144 select! {
145 _ = health.live() => continue,
146 maybe_payload = context.payloads().next() => match maybe_payload {
147 Some(payload) => if let Some(http_payload) = payload.try_into_http_payload() {
148 let (payload_meta, request) = http_payload.into_parts();
149 let transaction_meta = Metadata::from_event_and_data_point_count(
150 payload_meta.event_count(),
151 payload_meta.data_point_count(),
152 );
153 let transaction = Transaction::from_original(transaction_meta, request);
154
155 forwarder.send_transaction(transaction).await?;
156 }
157 None => break,
158 },
159 }
160 }
161
162 forwarder.shutdown().await;
163
164 debug!("Cluster Agent forwarder stopped.");
165
166 Ok(())
167 }
168}
169
170fn bearer_auth_header_value(auth_token: &str) -> Result<HeaderValue, GenericError> {
171 let raw_value = format!("Bearer {auth_token}");
172 HeaderValue::from_str(&raw_value)
173 .map_err(|_| generic_error!("cluster_agent.auth_token contains characters that are invalid in HTTP headers."))
174}
175
176fn cluster_agent_request_mapper_factory<B>(auth_header_value: HeaderValue) -> EndpointRequestMapperFactory<B>
177where
178 B: 'static,
179{
180 std::sync::Arc::new(move |endpoint| cluster_agent_request_mapper(endpoint, auth_header_value.clone()))
181}
182
183fn cluster_agent_request_mapper<B>(
184 endpoint: ResolvedEndpoint, auth_header_value: HeaderValue,
185) -> EndpointRequestMapper<B> {
186 let new_uri_authority = Authority::try_from(endpoint.endpoint().authority())
187 .expect("should not fail to construct new endpoint authority");
188 let new_uri_scheme =
189 Scheme::try_from(endpoint.endpoint().scheme()).expect("should not fail to construct new endpoint scheme");
190
191 Box::new(move |mut request: Request<TransactionBody<B>>| {
192 let new_uri = Uri::builder()
193 .scheme(new_uri_scheme.clone())
194 .authority(new_uri_authority.clone())
195 .path_and_query(CLUSTER_AGENT_SERIES_PATH)
196 .build()
197 .expect("should not fail to construct Cluster Agent URI");
198 *request.uri_mut() = new_uri;
199 request.headers_mut().remove(&DD_API_KEY_HEADER);
200 request.headers_mut().insert(AUTHORIZATION, auth_header_value.clone());
201
202 request
203 })
204}
205
206fn get_cluster_agent_endpoint_name(_uri: &Uri) -> Option<MetaString> {
207 Some(MetaString::from_static("cluster_agent_series"))
208}
209
210#[cfg(test)]
211mod tests {
212 use std::collections::HashMap;
213
214 use agent_data_plane_config::{
215 shared::{AltMetricsIntake, V3SeriesMode},
216 ConfigValue,
217 };
218 use http::Method;
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 let mut shared = shared_configuration();
261 shared.endpoints.api_key = "primary-api-key".to_string();
262 shared.endpoints.site = ConfigValue::explicit("datadoghq.eu".to_string());
263 shared.endpoints.dd_url = ConfigValue::explicit("https://app.datadoghq.com".to_string());
264 shared.endpoints.additional_endpoints = HashMap::from([(
265 "https://additional.example.com".to_string(),
266 vec!["additional-api-key".to_string()],
267 )]);
268 shared.endpoints.opw_intake = AltMetricsIntake {
269 enabled: true,
270 url: "https://opw.example.com".to_string(),
271 use_v3_series: true,
272 };
273 shared.metrics_encoding.v3_series_mode = V3SeriesMode::Enabled;
274 shared.metrics_encoding.v3_series_endpoint_modes =
275 HashMap::from([("https://app.datadoghq.com".to_string(), V3SeriesMode::Enabled)]);
276
277 let datadog_forwarder = ForwarderConfiguration::from_configuration(&shared);
279 assert_eq!(V3SeriesMode::Enabled, datadog_forwarder.use_v3_api_series().enabled);
280
281 let config = ClusterAgentForwarderConfiguration::from_configuration(
282 &shared,
283 "https://cluster-agent.example.com".to_string(),
284 "secret-token".to_string(),
285 )
286 .expect("Cluster Agent forwarder configuration should parse");
287 let endpoints = config
288 .forwarder_config
289 .build_routable_endpoints()
290 .expect("endpoint should resolve");
291
292 assert_eq!(endpoints.len(), 1);
293 assert_eq!(endpoints[0].route(), EndpointRoute::Primary);
294 assert_eq!(
295 endpoints[0].endpoint().endpoint().as_str(),
296 "https://cluster-agent.example.com/"
297 );
298 assert_eq!(&*endpoints[0].endpoint().api_key(), "secret-token");
299 assert_eq!(
300 V3SeriesMode::Disabled,
301 config.forwarder_config.use_v3_api_series().enabled
302 );
303 assert!(config.forwarder_config.use_v3_api_series().endpoints.is_empty());
304 assert!(!config.forwarder_config.allow_arbitrary_tags());
305 }
306}