vector/sinks/datadog/logs/
service.rs1use 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#[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 fn poll_ready(&mut self, _cx: &mut Context) -> Poll<Result<(), Self::Error>> {
146 Poll::Ready(Ok(()))
147 }
148
149 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 headers.insert(name.inner(), value.clone());
172 }
173 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}