datadog_agent_commons/ipc/client/
mod.rs

1//! Helpers for interacting with the Datadog Agent.
2
3use 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/// A client for interacting with the Datadog Agent's internal gRPC-based API.
38#[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    /// Connects to the Core Agent's gRPC IPC endpoint using the given client configuration.
47    ///
48    /// # Errors
49    ///
50    /// If the connection can't be established, the authentication token or IPC certificate can't be read, or the
51    /// initial health-check RPC fails, an error will be returned. Each of these is retried first.
52    pub async fn connect(config: &RemoteAgentClientConfiguration) -> Result<Self, GenericError> {
53        // TODO: We need to write a Tower middleware service that allows applying a backoff between failed calls,
54        // specifically so that we can throttle reconnection attempts.
55        //
56        // When the remote Agent endpoint is not available -- Agent isn't running, etc -- the gRPC client will
57        // essentially freewheel, trying to reconnect as quickly as possible, which spams the logs, wastes resources, so
58        // on and so forth. We would want to essentially apply a backoff like any other client would for the RPC calls
59        // themselves, but use it with the _connector_ instead.
60        //
61        // We could potentially just use a retry middleware, but Tonic does have its own reconnection logic, so we'd
62        // have to test it out to make sure it behaves sensibly.
63        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            // Health check inside the retried region. A transient partition can let the TCP connect succeed but break
85            // this first RPC stream, so retrying here keeps a boot-time blip from failing client construction.
86            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    /// Gets the detected hostname from the Agent.
119    ///
120    /// # Errors
121    ///
122    /// If there is an error querying the Agent API, an error will be returned.
123    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    /// Gets a stream of tagger entities at the given cardinality, optionally limited to a set of entity ID prefixes.
134    ///
135    /// When `prefixes` is given, the server only sends events for entities whose ID carries one of those prefixes,
136    /// which avoids paying the serialization and bandwidth cost of entities that the caller has no use for. Passing
137    /// `None` applies no filtering, and the server sends events for every prefix it knows about.
138    ///
139    /// If there is an error with the initial request, or an error occurs while streaming, the next message in the
140    /// stream will be `Some(Err(status))`, where the status indicates the underlying error.
141    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    /// Gets a stream of all workloadmeta entities.
157    ///
158    /// If there is an error with the initial request, or an error occurs while streaming, the next message in the
159    /// stream will be `Some(Err(status))`, where the status indicates the underlying error.
160    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    /// Registers a Remote Agent with the Agent.
180    ///
181    /// # Errors
182    ///
183    /// If there is an error sending the request to the Agent API, an error will be returned.
184    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    /// Refreshes the given remote agent session with the Agent.
201    ///
202    /// # Errors
203    ///
204    /// If there is an error sending the request to the Agent API, an error will be returned.
205    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    /// Reports one or more operational events for a remote agent session to the Agent.
217    ///
218    /// # Errors
219    ///
220    /// If there is an error sending the request to the Agent API, an error will be returned.
221    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    /// Gets the host tags from the Agent.
230    ///
231    /// # Errors
232    ///
233    /// If there is an error querying the Agent API, an error will be returned.
234    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    /// Polls the Agent for the Remote Configuration assigned to a client.
241    ///
242    /// # Errors
243    ///
244    /// Returns the gRPC status unchanged, so that callers can distinguish `Unimplemented`, which the Agent returns when
245    /// Remote Configuration is disabled, from other failures.
246    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    /// Gets a stream of autodiscovery config updates.
255    ///
256    /// If there is an error with the initial request, or an error occurs while streaming, the next message in the
257    /// stream will be `Some(Err(status))`, where the status indicates the underlying error.
258    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    /// Gets a stream of config events.
264    ///
265    /// If there is an error with the initial request, or an error occurs while streaming, the next message in the
266    /// stream will be `Some(Err(status))`, where the status indicates the underlying error.
267    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}