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