Skip to main content

vector/sinks/datadog/logs/
service.rs

1use std::{
2    collections::BTreeMap,
3    sync::Arc,
4    task::{Context, Poll},
5};
6
7use bytes::Bytes;
8use futures::future::BoxFuture;
9use http::{
10    HeaderValue, Request, Uri,
11    header::{CONTENT_ENCODING, CONTENT_LENGTH, CONTENT_TYPE},
12};
13use hyper::{Body, client::connect::Connect};
14use tower::Service;
15use tracing::Instrument;
16use vector_lib::{
17    event::{EventFinalizers, EventStatus, Finalizable},
18    request_metadata::{GroupedCountByteSize, MetaDescriptive, RequestMetadata},
19    stream::DriverResponse,
20};
21
22use crate::{
23    http::{HttpClient, HttpProxyConnector},
24    sinks::{
25        datadog::DatadogApiError,
26        util::{
27            Compression,
28            http::{OrderedHeaderName, validate_headers},
29            retries::RetryLogic,
30        },
31    },
32};
33
34#[derive(Debug, Default, Clone)]
35pub struct LogApiRetry;
36
37impl RetryLogic for LogApiRetry {
38    type Error = DatadogApiError;
39    type Request = LogApiRequest;
40    type Response = LogApiResponse;
41
42    fn is_retriable_error(&self, error: &Self::Error) -> bool {
43        error.is_retriable()
44    }
45}
46
47#[derive(Debug, Clone)]
48pub struct LogApiRequest {
49    pub api_key: Arc<str>,
50    pub compression: Compression,
51    pub body: Bytes,
52    pub finalizers: EventFinalizers,
53    pub uncompressed_size: usize,
54    pub metadata: RequestMetadata,
55}
56
57impl Finalizable for LogApiRequest {
58    fn take_finalizers(&mut self) -> EventFinalizers {
59        std::mem::take(&mut self.finalizers)
60    }
61}
62
63impl MetaDescriptive for LogApiRequest {
64    fn get_metadata(&self) -> &RequestMetadata {
65        &self.metadata
66    }
67
68    fn metadata_mut(&mut self) -> &mut RequestMetadata {
69        &mut self.metadata
70    }
71}
72
73#[derive(Debug)]
74pub struct LogApiResponse {
75    event_status: EventStatus,
76    events_byte_size: GroupedCountByteSize,
77    raw_byte_size: usize,
78}
79
80impl DriverResponse for LogApiResponse {
81    fn event_status(&self) -> EventStatus {
82        self.event_status
83    }
84
85    fn events_sent(&self) -> &GroupedCountByteSize {
86        &self.events_byte_size
87    }
88
89    fn bytes_sent(&self) -> Option<usize> {
90        Some(self.raw_byte_size)
91    }
92}
93
94/// Wrapper for the Datadog API.
95///
96/// Provides a `tower::Service` for the Datadog Logs API, allowing it to be
97/// composed within a Tower "stack", such that we can easily and transparently
98/// provide retries, concurrency limits, rate limits, and more.
99#[derive(Debug, Clone)]
100pub struct LogApiService<C = HttpProxyConnector> {
101    client: HttpClient<Body, C>,
102    uri: Uri,
103    user_provided_headers: BTreeMap<OrderedHeaderName, HeaderValue>,
104    dd_evp_headers: BTreeMap<OrderedHeaderName, HeaderValue>,
105}
106
107impl<C> LogApiService<C>
108where
109    C: Connect + Clone + Send + Sync + 'static,
110{
111    pub fn new(
112        client: HttpClient<Body, C>,
113        uri: Uri,
114        headers: BTreeMap<String, String>,
115        dd_evp_origin: String,
116    ) -> crate::Result<Self> {
117        let user_provided_headers = validate_headers(&headers)?;
118
119        let dd_evp_headers: BTreeMap<String, String> = [
120            ("DD-EVP-ORIGIN".to_string(), dd_evp_origin),
121            ("DD-EVP-ORIGIN-VERSION".to_string(), crate::get_version()),
122        ]
123        .into_iter()
124        .collect();
125        let dd_evp_headers = validate_headers(&dd_evp_headers)?;
126
127        Ok(Self {
128            client,
129            uri,
130            user_provided_headers,
131            dd_evp_headers,
132        })
133    }
134}
135
136impl<C> Service<LogApiRequest> for LogApiService<C>
137where
138    C: Connect + Clone + Send + Sync + 'static,
139{
140    type Response = LogApiResponse;
141    type Error = DatadogApiError;
142    type Future = BoxFuture<'static, Result<Self::Response, Self::Error>>;
143
144    // Emission of Error internal event is handled upstream by the caller
145    fn poll_ready(&mut self, _cx: &mut Context) -> Poll<Result<(), Self::Error>> {
146        Poll::Ready(Ok(()))
147    }
148
149    // Emission of Error internal event is handled upstream by the caller
150    fn call(&mut self, mut request: LogApiRequest) -> Self::Future {
151        let mut client = self.client.clone();
152        let http_request = Request::post(&self.uri)
153            .header(CONTENT_TYPE, "application/json")
154            .header("DD-API-KEY", request.api_key.to_string());
155
156        let http_request = if let Some(ce) = request.compression.content_encoding() {
157            http_request.header(CONTENT_ENCODING, ce)
158        } else {
159            http_request
160        };
161
162        let metadata = std::mem::take(request.metadata_mut());
163        let events_byte_size = metadata.into_events_estimated_json_encoded_byte_size();
164        let raw_byte_size = request.uncompressed_size;
165
166        let mut http_request = http_request.header(CONTENT_LENGTH, request.body.len());
167
168        if let Some(headers) = http_request.headers_mut() {
169            for (name, value) in &self.user_provided_headers {
170                // Replace rather than append to any existing header values
171                headers.insert(name.inner(), value.clone());
172            }
173            // Set DD EVP headers last so that they cannot be overridden.
174            for (name, value) in &self.dd_evp_headers {
175                headers.insert(name.inner(), value.clone());
176            }
177        }
178
179        let http_request = http_request
180            .body(Body::from(request.body))
181            .expect("building HTTP request failed unexpectedly");
182
183        Box::pin(async move {
184            DatadogApiError::from_result(client.call(http_request).in_current_span().await).map(
185                |_| LogApiResponse {
186                    event_status: EventStatus::Delivered,
187                    events_byte_size,
188                    raw_byte_size,
189                },
190            )
191        })
192    }
193}