datadog_agent_commons/ipc/client/
mod.rs1use std::time::Duration;
4
5use backon::{ConstantBuilder, Retryable as _};
6use datadog_protos::agent::v1::{
7 RefreshRemoteAgentRequest, RegisterRemoteAgentRequest, RegisterRemoteAgentResponse, ReportRemoteAgentEventRequest,
8 ReportRemoteAgentEventResponse,
9};
10use datadog_protos::agent::{
11 AgentClient, AgentSecureClient, AutodiscoveryStreamResponse, ConfigEvent, ConfigStreamRequest, EntityId,
12 FetchEntityRequest, HostTagReply, HostTagRequest, HostnameRequest, RemoteAgentClient as RemoteAgentServiceClient,
13 StreamTagsRequest, StreamTagsResponse, TagCardinality, WorkloadmetaEventType, WorkloadmetaFilter, WorkloadmetaKind,
14 WorkloadmetaSource, WorkloadmetaStreamRequest, WorkloadmetaStreamResponse,
15};
16use datadog_protos::remote_config::{ClientGetConfigsRequest, ClientGetConfigsResponse};
17use saluki_error::{generic_error, ErrorContext as _, GenericError};
18use saluki_io::net::client::http::HttpsCapableConnectorBuilder;
19use tonic::{
20 service::interceptor::InterceptedService,
21 transport::{Channel, Endpoint},
22 Code, Request, Response, Status,
23};
24use tracing::warn;
25
26use crate::ipc::{config::RemoteAgentClientConfiguration, session::SessionId, tls::build_ipc_client_ipc_tls_config};
27
28mod bearer_auth;
29use self::bearer_auth::BearerAuthInterceptor;
30
31mod streaming;
32pub use self::streaming::StreamingResponse;
33
34const CONNECT_RETRY_ATTEMPTS: usize = 10;
35const CONNECT_RETRY_BACKOFF: Duration = Duration::from_secs(2);
36
37#[derive(Clone)]
39pub struct RemoteAgentClient {
40 client: AgentClient<InterceptedService<Channel, BearerAuthInterceptor>>,
41 secure_client: AgentSecureClient<InterceptedService<Channel, BearerAuthInterceptor>>,
42 remote_agent_client: RemoteAgentServiceClient<InterceptedService<Channel, BearerAuthInterceptor>>,
43}
44
45impl RemoteAgentClient {
46 pub async fn connect(config: &RemoteAgentClientConfiguration) -> Result<Self, GenericError> {
53 let service_builder = || async {
64 let auth_interceptor = BearerAuthInterceptor::from_file(config.auth.auth_token_file_path()).await?;
65 let ipc_cert_file_path = config.auth.ipc_cert_file_path();
66 let client_tls_config = build_ipc_client_ipc_tls_config(ipc_cert_file_path).await?;
67 let connector_builder = HttpsCapableConnectorBuilder::default();
68 #[cfg(target_os = "linux")]
69 let connector_builder = if let Some(addr) = config.vsock_addr() {
70 connector_builder.with_vsock_addr(addr)
71 } else {
72 connector_builder
73 };
74 let https_connector = connector_builder.build(client_tls_config)?;
75 let endpoint = config.endpoint();
76 let channel = Endpoint::from(endpoint.clone())
77 .connect_timeout(Duration::from_secs(2))
78 .connect_with_connector(https_connector)
79 .await
80 .with_error_context(|| format!("Failed to connect to Datadog Agent API at '{}'.", endpoint))?;
81
82 let service = InterceptedService::new(channel, auth_interceptor);
83
84 let mut secure_client =
87 AgentSecureClient::new(service.clone()).max_decoding_message_size(config.grpc_max_message_size);
88 try_query_agent_api(&mut secure_client).await?;
89
90 Ok::<_, GenericError>(service)
91 };
92
93 let service = service_builder
94 .retry(
95 ConstantBuilder::default()
96 .with_delay(CONNECT_RETRY_BACKOFF)
97 .with_max_times(CONNECT_RETRY_ATTEMPTS),
98 )
99 .notify(|e, delay| {
100 warn!(error = %e, "Failed to create Datadog Agent API client. Retrying in {:?}...", delay);
101 })
102 .await
103 .error_context("Failed to create Datadog Agent API client.")?;
104
105 let client = AgentClient::new(service.clone()).max_decoding_message_size(config.grpc_max_message_size);
106 let secure_client =
107 AgentSecureClient::new(service.clone()).max_decoding_message_size(config.grpc_max_message_size);
108 let remote_agent_client =
109 RemoteAgentServiceClient::new(service).max_decoding_message_size(config.grpc_max_message_size);
110
111 Ok(Self {
112 client,
113 secure_client,
114 remote_agent_client,
115 })
116 }
117
118 pub async fn get_hostname(&mut self) -> Result<String, GenericError> {
124 let response = self
125 .client
126 .get_hostname(HostnameRequest {})
127 .await
128 .map(|r| r.into_inner())?;
129
130 Ok(response.hostname)
131 }
132
133 pub fn get_tagger_stream(
142 &mut self, cardinality: TagCardinality, prefixes: Option<Vec<String>>,
143 ) -> StreamingResponse<StreamTagsResponse> {
144 let mut client = self.secure_client.clone();
145 StreamingResponse::from_response_future(async move {
146 client
147 .tagger_stream_entities(StreamTagsRequest {
148 cardinality: cardinality.into(),
149 prefixes: prefixes.unwrap_or_default(),
150 ..Default::default()
151 })
152 .await
153 })
154 }
155
156 pub fn get_workloadmeta_stream(&mut self) -> StreamingResponse<WorkloadmetaStreamResponse> {
161 let mut client = self.secure_client.clone();
162 StreamingResponse::from_response_future(async move {
163 client
164 .workloadmeta_stream_entities(WorkloadmetaStreamRequest {
165 filter: Some(WorkloadmetaFilter {
166 kinds: vec![
167 WorkloadmetaKind::Container.into(),
168 WorkloadmetaKind::KubernetesPod.into(),
169 WorkloadmetaKind::EcsTask.into(),
170 ],
171 source: WorkloadmetaSource::All.into(),
172 event_type: WorkloadmetaEventType::EventTypeAll.into(),
173 }),
174 })
175 .await
176 })
177 }
178
179 pub async fn register_remote_agent(
185 &mut self, pid: u32, display_name: &str, flavor: &str, api_endpoint: &str, services: Vec<String>,
186 ) -> Result<Response<RegisterRemoteAgentResponse>, GenericError> {
187 let mut client = self.remote_agent_client.clone();
188 let response = client
189 .register_remote_agent(RegisterRemoteAgentRequest {
190 pid: pid.to_string(),
191 flavor: flavor.to_string(),
192 display_name: display_name.to_string(),
193 api_endpoint_uri: api_endpoint.to_string(),
194 services,
195 })
196 .await?;
197 Ok(response)
198 }
199
200 pub async fn refresh_remote_agent(&mut self, session_id: &SessionId) -> Result<Response<()>, GenericError> {
206 let mut client = self.remote_agent_client.clone();
207 let response = client
208 .refresh_remote_agent(RefreshRemoteAgentRequest {
209 session_id: session_id.to_string(),
210 })
211 .await?
212 .map(|_| ());
213 Ok(response)
214 }
215
216 pub async fn report_remote_agent_event(
222 &mut self, request: ReportRemoteAgentEventRequest,
223 ) -> Result<Response<ReportRemoteAgentEventResponse>, GenericError> {
224 let mut client = self.remote_agent_client.clone();
225 let response = client.report_remote_agent_event(request).await?;
226 Ok(response)
227 }
228
229 pub async fn get_host_tags(&self) -> Result<Response<HostTagReply>, GenericError> {
235 let mut client = self.secure_client.clone();
236 let response = client.get_host_tags(HostTagRequest {}).await?;
237 Ok(response)
238 }
239
240 pub async fn client_get_configs(
247 &self, request: ClientGetConfigsRequest,
248 ) -> Result<ClientGetConfigsResponse, Status> {
249 let mut client = self.secure_client.clone();
250 let response = client.client_get_configs(request).await?;
251 Ok(response.into_inner())
252 }
253
254 pub fn get_autodiscovery_stream(&mut self) -> StreamingResponse<AutodiscoveryStreamResponse> {
259 let mut client = self.secure_client.clone();
260 StreamingResponse::from_response_future(async move { client.autodiscovery_stream_config(()).await })
261 }
262
263 pub fn stream_config_events(&mut self, session_id: &SessionId) -> StreamingResponse<ConfigEvent> {
268 let mut client = self.secure_client.clone();
269 let app_details = saluki_metadata::get_app_details();
270 let formatted_full_name = app_details
271 .full_name()
272 .replace(" ", "-")
273 .replace("_", "-")
274 .to_lowercase();
275
276 let mut request = Request::new(ConfigStreamRequest {
277 name: formatted_full_name,
278 });
279
280 request
281 .metadata_mut()
282 .insert("session_id", session_id.to_grpc_header_value());
283
284 StreamingResponse::from_response_future(async move { client.stream_config_events(request).await })
285 }
286}
287
288async fn try_query_agent_api(
289 client: &mut AgentSecureClient<InterceptedService<Channel, BearerAuthInterceptor>>,
290) -> Result<(), GenericError> {
291 let noop_fetch_request = FetchEntityRequest {
292 id: Some(EntityId {
293 prefix: "container_id".to_string(),
294 uid: "nonexistent".to_string(),
295 }),
296 cardinality: TagCardinality::High.into(),
297 };
298 match client.tagger_fetch_entity(noop_fetch_request).await {
299 Ok(_) => Ok(()),
300 Err(e) => match e.code() {
301 Code::Unauthenticated => Err(generic_error!(
302 "Failed to authenticate to Datadog Agent API. Check that the configured authentication token is correct."
303 )),
304 _ => Err(e.into()),
305 },
306 }
307}