Skip to main content

omni_dev/datadog/
client.rs

1//! Datadog REST API client.
2//!
3//! Thin `reqwest` wrapper that injects the `DD-API-KEY` and
4//! `DD-APPLICATION-KEY` headers on every request and retries 429 responses via
5//! the shared [`retry_429`](crate::utils::http::retry_429) driver (which honours
6//! `Retry-After` / `X-RateLimit-Reset`).
7
8use anyhow::{Context, Result};
9use reqwest::Client;
10use url::Url;
11
12use crate::datadog::auth::{base_url_for_site, DatadogCredentials};
13use crate::datadog::error::DatadogError;
14use crate::request_log;
15use crate::utils::http::{connect_timeout, read_timeout, retry_429};
16use crate::utils::secret::Secret;
17
18/// HTTP client for Datadog REST APIs.
19#[derive(Debug)]
20pub struct DatadogClient {
21    client: Client,
22    base_url: String,
23    api_key: Secret,
24    app_key: Secret,
25}
26
27impl DatadogClient {
28    /// Creates a new Datadog API client.
29    ///
30    /// `base_url` should be the full API host, e.g. `https://api.datadoghq.com`.
31    /// For production use, construct via [`Self::from_credentials`]; tests
32    /// pass a wiremock URL directly.
33    pub fn new(base_url: &str, api_key: &str, app_key: &str) -> Result<Self> {
34        let client = Client::builder()
35            .connect_timeout(connect_timeout())
36            .read_timeout(read_timeout())
37            .build()
38            .context("Failed to build HTTP client")?;
39
40        Ok(Self {
41            client,
42            base_url: base_url.trim_end_matches('/').to_string(),
43            api_key: api_key.into(),
44            app_key: app_key.into(),
45        })
46    }
47
48    /// Creates a client from stored credentials.
49    ///
50    /// Respects `DATADOG_API_URL` as an optional override: when set in the
51    /// process environment it replaces the site-derived base URL. Used for
52    /// tests (wiremock) and on-prem Datadog installs.
53    pub fn from_credentials(creds: &DatadogCredentials) -> Result<Self> {
54        Self::from_credentials_with(&crate::utils::env::SystemEnv, creds)
55    }
56
57    /// [`from_credentials`](Self::from_credentials) over an injected
58    /// [`EnvSource`](crate::utils::env::EnvSource).
59    ///
60    /// Tests pass a pure `MapEnv` to exercise the `DATADOG_API_URL` override
61    /// without mutating the process environment (issue #1030).
62    pub(crate) fn from_credentials_with(
63        env: &impl crate::utils::env::EnvSource,
64        creds: &DatadogCredentials,
65    ) -> Result<Self> {
66        let base_url = env
67            .var(crate::datadog::auth::DATADOG_API_URL)
68            .filter(|s| !s.is_empty())
69            .unwrap_or_else(|| base_url_for_site(&creds.site));
70        Self::new(
71            &base_url,
72            creds.api_key.expose_secret(),
73            creds.app_key.expose_secret(),
74        )
75    }
76
77    /// Returns the API base URL (without trailing slash).
78    #[must_use]
79    pub fn base_url(&self) -> &str {
80        &self.base_url
81    }
82
83    /// Builds an absolute API URL by joining `path` onto `base_url`.
84    ///
85    /// `path` is the full path portion including the leading `/api/…` segment
86    /// (the version varies: `/api/v1/…`, `/api/v2/…`). Centralises the
87    /// `Url::parse(…).context("Invalid Datadog base URL")` spelling repeated
88    /// across the `*_api.rs` modules. Takes `base_url` (rather than `&self`) so
89    /// the free `build_*_url` functions — and their unit tests, which pass
90    /// literal base URLs — can call it unchanged.
91    pub(crate) fn api_url(base_url: &str, path: &str) -> Result<Url> {
92        Url::parse(&format!("{base_url}{path}")).context("Invalid Datadog base URL")
93    }
94
95    /// Checks `response` for success and deserialises its JSON body into `T`.
96    ///
97    /// Non-success responses become a [`DatadogError`] via
98    /// [`Self::response_to_error`] (preserving the 429 rate-limit summary); on
99    /// success the body is parsed with `context` attached on failure. Used by
100    /// the paginated and POST call sites that already hold a response;
101    /// single-shot GETs use [`Self::get_parsed`].
102    pub(crate) async fn parse_response<T: serde::de::DeserializeOwned>(
103        &self,
104        response: reqwest::Response,
105        context: &'static str,
106    ) -> Result<T> {
107        if !response.status().is_success() {
108            return Err(Self::response_to_error(response).await.into());
109        }
110        response.json().await.context(context)
111    }
112
113    /// Sends an authenticated GET and deserialises the JSON body into `T`.
114    ///
115    /// Convenience wrapper over [`Self::get_json`] + [`Self::parse_response`]
116    /// for the common single-request GET-then-parse pattern.
117    pub(crate) async fn get_parsed<T: serde::de::DeserializeOwned>(
118        &self,
119        url: &str,
120        context: &'static str,
121    ) -> Result<T> {
122        let response = self.get_json(url).await?;
123        self.parse_response(response, context).await
124    }
125
126    /// Sends an authenticated GET request and returns the raw response.
127    pub async fn get_json(&self, url: &str) -> Result<reqwest::Response> {
128        retry_429(
129            || {
130                self.client
131                    .get(url)
132                    .header("DD-API-KEY", self.api_key.expose_secret())
133                    .header("DD-APPLICATION-KEY", self.app_key.expose_secret())
134                    .header("Accept", "application/json")
135            },
136            |started, result| {
137                request_log::record_http_result("datadog", "GET", url, started, result);
138            },
139        )
140        .await
141        .context("Failed to send GET request to Datadog API")
142    }
143
144    /// Sends an authenticated POST request with a JSON body and returns the raw response.
145    pub async fn post_json<T: serde::Serialize + Sync + ?Sized>(
146        &self,
147        url: &str,
148        body: &T,
149    ) -> Result<reqwest::Response> {
150        retry_429(
151            || {
152                self.client
153                    .post(url)
154                    .header("DD-API-KEY", self.api_key.expose_secret())
155                    .header("DD-APPLICATION-KEY", self.app_key.expose_secret())
156                    .header("Content-Type", "application/json")
157                    .header("Accept", "application/json")
158                    .json(body)
159            },
160            |started, result| {
161                request_log::record_http_result("datadog", "POST", url, started, result);
162            },
163        )
164        .await
165        .context("Failed to send POST request to Datadog API")
166    }
167
168    /// Consumes a non-success response and turns it into a [`DatadogError`].
169    ///
170    /// For 429 responses, appends a human-readable rate-limit summary
171    /// (extracted from `X-RateLimit-*` headers) to the body, so the caller
172    /// sees why the retry loop gave up.
173    pub async fn response_to_error(response: reqwest::Response) -> DatadogError {
174        let status = response.status().as_u16();
175        let headers = response.headers().clone();
176        let body = response.text().await.unwrap_or_default();
177        let body = if status == 429 {
178            match format_rate_limit(&headers) {
179                Some(suffix) => format!("{body} {suffix}").trim().to_string(),
180                None => body,
181            }
182        } else {
183            body
184        };
185        DatadogError::ApiRequestFailed { status, body }
186    }
187}
188
189fn format_rate_limit(headers: &reqwest::header::HeaderMap) -> Option<String> {
190    let remaining = headers
191        .get("X-RateLimit-Remaining")
192        .and_then(|v| v.to_str().ok());
193    let reset = headers
194        .get("X-RateLimit-Reset")
195        .and_then(|v| v.to_str().ok());
196    let limit = headers
197        .get("X-RateLimit-Limit")
198        .and_then(|v| v.to_str().ok());
199
200    if remaining.is_none() && reset.is_none() && limit.is_none() {
201        return None;
202    }
203
204    let mut parts = Vec::new();
205    if let Some(v) = remaining {
206        parts.push(format!("remaining={v}"));
207    }
208    if let Some(v) = limit {
209        parts.push(format!("limit={v}"));
210    }
211    if let Some(v) = reset {
212        parts.push(format!("reset_in={v}s"));
213    }
214    Some(format!("[rate-limit: {}]", parts.join(", ")))
215}
216
217#[cfg(test)]
218#[allow(clippy::unwrap_used, clippy::expect_used)]
219mod tests {
220    use super::*;
221
222    #[test]
223    fn new_client_strips_trailing_slash() {
224        let client = DatadogClient::new("https://api.datadoghq.com/", "api", "app").unwrap();
225        assert_eq!(client.base_url(), "https://api.datadoghq.com");
226    }
227
228    #[test]
229    fn new_client_preserves_clean_url() {
230        let client = DatadogClient::new("https://api.datadoghq.com", "api", "app").unwrap();
231        assert_eq!(client.base_url(), "https://api.datadoghq.com");
232    }
233
234    #[test]
235    fn client_debug_redacts_keys() {
236        let client = DatadogClient::new(
237            "https://api.datadoghq.com",
238            "sekret-api-key",
239            "sekret-app-key",
240        )
241        .unwrap();
242        // Debug must never print the key values (#1131).
243        let debug = format!("{client:?}");
244        assert!(!debug.contains("sekret-api-key"), "leaked api_key: {debug}");
245        assert!(!debug.contains("sekret-app-key"), "leaked app_key: {debug}");
246        assert!(debug.contains("api_key: <redacted>"));
247        assert!(debug.contains("app_key: <redacted>"));
248    }
249
250    #[test]
251    fn from_credentials_builds_base_url_from_site() {
252        let env = crate::test_support::env::MapEnv::new();
253        let creds = DatadogCredentials {
254            api_key: "api".into(),
255            app_key: "app".into(),
256            site: "us5.datadoghq.com".to_string(),
257        };
258        let client = DatadogClient::from_credentials_with(&env, &creds).unwrap();
259        assert_eq!(client.base_url(), "https://api.us5.datadoghq.com");
260    }
261
262    #[test]
263    fn from_credentials_honours_api_url_override() {
264        let env = crate::test_support::env::MapEnv::new().with(
265            crate::datadog::auth::DATADOG_API_URL,
266            "http://proxy.example:8080",
267        );
268        let creds = DatadogCredentials {
269            api_key: "api".into(),
270            app_key: "app".into(),
271            site: "us5.datadoghq.com".to_string(),
272        };
273        let client = DatadogClient::from_credentials_with(&env, &creds).unwrap();
274        assert_eq!(client.base_url(), "http://proxy.example:8080");
275    }
276
277    #[test]
278    fn from_credentials_ignores_empty_api_url_override() {
279        let env =
280            crate::test_support::env::MapEnv::new().with(crate::datadog::auth::DATADOG_API_URL, "");
281        let creds = DatadogCredentials {
282            api_key: "api".into(),
283            app_key: "app".into(),
284            site: "datadoghq.com".to_string(),
285        };
286        let client = DatadogClient::from_credentials_with(&env, &creds).unwrap();
287        assert_eq!(client.base_url(), "https://api.datadoghq.com");
288    }
289
290    #[tokio::test]
291    async fn get_json_sends_auth_headers() {
292        let server = wiremock::MockServer::start().await;
293        wiremock::Mock::given(wiremock::matchers::method("GET"))
294            .and(wiremock::matchers::path("/test"))
295            .and(wiremock::matchers::header("DD-API-KEY", "my-api"))
296            .and(wiremock::matchers::header("DD-APPLICATION-KEY", "my-app"))
297            .and(wiremock::matchers::header("Accept", "application/json"))
298            .respond_with(
299                wiremock::ResponseTemplate::new(200).set_body_json(serde_json::json!({"ok": true})),
300            )
301            .expect(1)
302            .mount(&server)
303            .await;
304
305        let client = DatadogClient::new(&server.uri(), "my-api", "my-app").unwrap();
306        let resp = client
307            .get_json(&format!("{}/test", server.uri()))
308            .await
309            .unwrap();
310        assert!(resp.status().is_success());
311    }
312
313    #[tokio::test]
314    async fn post_json_sends_body_and_auth() {
315        let server = wiremock::MockServer::start().await;
316        wiremock::Mock::given(wiremock::matchers::method("POST"))
317            .and(wiremock::matchers::path("/test"))
318            .and(wiremock::matchers::header("DD-API-KEY", "my-api"))
319            .and(wiremock::matchers::header("DD-APPLICATION-KEY", "my-app"))
320            .and(wiremock::matchers::header(
321                "Content-Type",
322                "application/json",
323            ))
324            .and(wiremock::matchers::body_json(serde_json::json!({
325                "query": "hello"
326            })))
327            .respond_with(
328                wiremock::ResponseTemplate::new(200).set_body_json(serde_json::json!({"id": "1"})),
329            )
330            .expect(1)
331            .mount(&server)
332            .await;
333
334        let client = DatadogClient::new(&server.uri(), "my-api", "my-app").unwrap();
335        let body = serde_json::json!({"query": "hello"});
336        let resp = client
337            .post_json(&format!("{}/test", server.uri()), &body)
338            .await
339            .unwrap();
340        assert!(resp.status().is_success());
341    }
342
343    #[tokio::test]
344    async fn get_json_retries_on_429_via_retry_after() {
345        let server = wiremock::MockServer::start().await;
346        wiremock::Mock::given(wiremock::matchers::method("GET"))
347            .and(wiremock::matchers::path("/test"))
348            .respond_with(wiremock::ResponseTemplate::new(429).append_header("Retry-After", "0"))
349            .up_to_n_times(1)
350            .mount(&server)
351            .await;
352        wiremock::Mock::given(wiremock::matchers::method("GET"))
353            .and(wiremock::matchers::path("/test"))
354            .respond_with(
355                wiremock::ResponseTemplate::new(200).set_body_json(serde_json::json!({"ok": true})),
356            )
357            .up_to_n_times(1)
358            .mount(&server)
359            .await;
360
361        let client = DatadogClient::new(&server.uri(), "api", "app").unwrap();
362        let resp = client
363            .get_json(&format!("{}/test", server.uri()))
364            .await
365            .unwrap();
366        assert!(resp.status().is_success());
367    }
368
369    #[tokio::test]
370    async fn get_json_retries_on_429_via_x_ratelimit_reset() {
371        let server = wiremock::MockServer::start().await;
372        wiremock::Mock::given(wiremock::matchers::method("GET"))
373            .and(wiremock::matchers::path("/test"))
374            .respond_with(
375                wiremock::ResponseTemplate::new(429).append_header("X-RateLimit-Reset", "0"),
376            )
377            .up_to_n_times(1)
378            .mount(&server)
379            .await;
380        wiremock::Mock::given(wiremock::matchers::method("GET"))
381            .and(wiremock::matchers::path("/test"))
382            .respond_with(
383                wiremock::ResponseTemplate::new(200).set_body_json(serde_json::json!({"ok": true})),
384            )
385            .up_to_n_times(1)
386            .mount(&server)
387            .await;
388
389        let client = DatadogClient::new(&server.uri(), "api", "app").unwrap();
390        let resp = client
391            .get_json(&format!("{}/test", server.uri()))
392            .await
393            .unwrap();
394        assert!(resp.status().is_success());
395    }
396
397    #[tokio::test]
398    async fn post_json_retries_on_429() {
399        let server = wiremock::MockServer::start().await;
400        wiremock::Mock::given(wiremock::matchers::method("POST"))
401            .and(wiremock::matchers::path("/test"))
402            .respond_with(wiremock::ResponseTemplate::new(429).append_header("Retry-After", "0"))
403            .up_to_n_times(1)
404            .mount(&server)
405            .await;
406        wiremock::Mock::given(wiremock::matchers::method("POST"))
407            .and(wiremock::matchers::path("/test"))
408            .respond_with(wiremock::ResponseTemplate::new(201))
409            .up_to_n_times(1)
410            .mount(&server)
411            .await;
412
413        let client = DatadogClient::new(&server.uri(), "api", "app").unwrap();
414        let resp = client
415            .post_json(
416                &format!("{}/test", server.uri()),
417                &serde_json::json!({"k": "v"}),
418            )
419            .await
420            .unwrap();
421        assert_eq!(resp.status().as_u16(), 201);
422    }
423
424    #[tokio::test]
425    async fn get_json_returns_429_after_max_retries() {
426        let server = wiremock::MockServer::start().await;
427        wiremock::Mock::given(wiremock::matchers::method("GET"))
428            .and(wiremock::matchers::path("/test"))
429            .respond_with(wiremock::ResponseTemplate::new(429).append_header("Retry-After", "0"))
430            .mount(&server)
431            .await;
432
433        let client = DatadogClient::new(&server.uri(), "api", "app").unwrap();
434        let resp = client
435            .get_json(&format!("{}/test", server.uri()))
436            .await
437            .unwrap();
438        assert_eq!(resp.status().as_u16(), 429);
439    }
440
441    #[tokio::test]
442    async fn response_to_error_surfaces_rate_limit_headers_on_429() {
443        let server = wiremock::MockServer::start().await;
444        wiremock::Mock::given(wiremock::matchers::method("GET"))
445            .and(wiremock::matchers::path("/test"))
446            .respond_with(
447                wiremock::ResponseTemplate::new(429)
448                    .append_header("Retry-After", "0")
449                    .append_header("X-RateLimit-Remaining", "0")
450                    .append_header("X-RateLimit-Reset", "42")
451                    .append_header("X-RateLimit-Limit", "100")
452                    .set_body_string("too many"),
453            )
454            .mount(&server)
455            .await;
456
457        let client = DatadogClient::new(&server.uri(), "api", "app").unwrap();
458        let resp = client
459            .get_json(&format!("{}/test", server.uri()))
460            .await
461            .unwrap();
462        let err = DatadogClient::response_to_error(resp).await;
463        let msg = err.to_string();
464        assert!(msg.contains("429"));
465        assert!(msg.contains("too many"));
466        assert!(msg.contains("remaining=0"));
467        assert!(msg.contains("limit=100"));
468        assert!(msg.contains("reset_in=42s"));
469    }
470
471    #[tokio::test]
472    async fn response_to_error_does_not_add_rate_limit_suffix_on_non_429() {
473        let server = wiremock::MockServer::start().await;
474        wiremock::Mock::given(wiremock::matchers::method("GET"))
475            .and(wiremock::matchers::path("/test"))
476            .respond_with(wiremock::ResponseTemplate::new(401).set_body_string("Unauthorized"))
477            .mount(&server)
478            .await;
479
480        let client = DatadogClient::new(&server.uri(), "api", "app").unwrap();
481        let resp = client
482            .get_json(&format!("{}/test", server.uri()))
483            .await
484            .unwrap();
485        let err = DatadogClient::response_to_error(resp).await;
486        let msg = err.to_string();
487        assert!(msg.contains("401"));
488        assert!(msg.contains("Unauthorized"));
489        assert!(!msg.contains("rate-limit"));
490    }
491
492    #[tokio::test]
493    async fn response_to_error_omits_suffix_when_no_rate_limit_headers() {
494        let server = wiremock::MockServer::start().await;
495        wiremock::Mock::given(wiremock::matchers::method("GET"))
496            .and(wiremock::matchers::path("/test"))
497            .respond_with(
498                wiremock::ResponseTemplate::new(429)
499                    .append_header("Retry-After", "0")
500                    .set_body_string("slow down"),
501            )
502            .mount(&server)
503            .await;
504
505        let client = DatadogClient::new(&server.uri(), "api", "app").unwrap();
506        let resp = client
507            .get_json(&format!("{}/test", server.uri()))
508            .await
509            .unwrap();
510        let err = DatadogClient::response_to_error(resp).await;
511        let msg = err.to_string();
512        assert!(msg.contains("slow down"));
513        assert!(!msg.contains("rate-limit"));
514    }
515}