saluki_components/forwarders/cluster_agent/
mod.rs1use 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
37pub struct ClusterAgentForwarderConfiguration {
42 forwarder_config: ForwarderConfiguration,
43 auth_header_value: HeaderValue,
44}
45
46impl ClusterAgentForwarderConfiguration {
47 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 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
131pub 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}