1use std::{convert::Infallible, net::SocketAddr, time::Duration};
2
3use bytes::Bytes;
4use futures_util::FutureExt;
5use http::StatusCode;
6use hyper::{Server, service::make_service_fn};
7use prost::Message;
8use snafu::Snafu;
9use tokio::net::TcpStream;
10use tower::ServiceBuilder;
11use tracing::Span;
12use vector_lib::{
13 EstimatedJsonEncodedSizeOf,
14 codecs::decoding::{OtlpDeserializer, format::Deserializer},
15 config::LogNamespace,
16 event::{BatchNotifier, BatchStatus},
17 internal_event::{
18 ByteSize, BytesReceived, CountByteSize, InternalEventHandle as _, Registered,
19 },
20 opentelemetry::proto::collector::{
21 logs::v1::{ExportLogsServiceRequest, ExportLogsServiceResponse},
22 metrics::v1::{ExportMetricsServiceRequest, ExportMetricsServiceResponse},
23 trace::v1::{ExportTraceServiceRequest, ExportTraceServiceResponse},
24 },
25 tls::MaybeTlsIncomingStream,
26};
27use warp::{
28 Filter, Reply, filters::BoxedFilter, http::HeaderMap, reject::Rejection, reply::Response,
29};
30
31use super::{reply::protobuf, status::Status};
32use crate::{
33 SourceSender,
34 common::http::ErrorMessage,
35 event::Event,
36 http::{KeepaliveConfig, MaxConnectionAgeLayer, build_http_trace_layer},
37 internal_events::{EventsReceived, HttpBadRequest, StreamClosedError},
38 shutdown::ShutdownSignal,
39 sources::{
40 http_server::HttpConfigParamKind,
41 opentelemetry::config::{LOGS, METRICS, OpentelemetryConfig, TRACES},
42 util::{add_headers, decompress_body, http::capped_body},
43 },
44 tls::MaybeTlsSettings,
45};
46
47#[derive(Clone, Copy, Debug, Snafu)]
48pub(crate) enum ApiError {
49 ServerShutdown,
50}
51
52impl warp::reject::Reject for ApiError {}
53
54pub(crate) async fn run_http_server(
55 address: SocketAddr,
56 tls_settings: MaybeTlsSettings,
57 filters: BoxedFilter<(Response,)>,
58 shutdown: ShutdownSignal,
59 keepalive_settings: KeepaliveConfig,
60) -> crate::Result<()> {
61 let listener = tls_settings.bind(&address).await?;
62 let routes = filters.recover(handle_rejection);
63
64 info!(message = "Building HTTP server.", address = %address);
65
66 let span = Span::current();
67 let make_svc = make_service_fn(move |conn: &MaybeTlsIncomingStream<TcpStream>| {
68 let svc = ServiceBuilder::new()
69 .layer(build_http_trace_layer(span.clone()))
70 .option_layer(keepalive_settings.max_connection_age_secs.map(|secs| {
71 MaxConnectionAgeLayer::new(
72 Duration::from_secs(secs),
73 keepalive_settings.max_connection_age_jitter_factor,
74 conn.peer_addr(),
75 )
76 }))
77 .service(warp::service(routes.clone()));
78 futures_util::future::ok::<_, Infallible>(svc)
79 });
80
81 Server::builder(hyper::server::accept::from_stream(listener.accept_stream()))
82 .serve(make_svc)
83 .with_graceful_shutdown(shutdown.map(|_| ()))
84 .await?;
85
86 Ok(())
87}
88
89#[allow(clippy::too_many_arguments)] pub(crate) fn build_warp_filter(
91 acknowledgements: bool,
92 log_namespace: LogNamespace,
93 out: SourceSender,
94 bytes_received: Registered<BytesReceived>,
95 events_received: Registered<EventsReceived>,
96 headers: Vec<HttpConfigParamKind>,
97 logs_deserializer: Option<OtlpDeserializer>,
98 metrics_deserializer: Option<OtlpDeserializer>,
99 traces_deserializer: Option<OtlpDeserializer>,
100) -> BoxedFilter<(Response,)> {
101 let log_filters = build_warp_log_filter(
102 acknowledgements,
103 log_namespace,
104 out.clone(),
105 bytes_received.clone(),
106 events_received.clone(),
107 headers.clone(),
108 logs_deserializer,
109 );
110 let metrics_filters = build_warp_metrics_filter(
111 acknowledgements,
112 log_namespace,
113 out.clone(),
114 bytes_received.clone(),
115 events_received.clone(),
116 headers.clone(),
117 metrics_deserializer,
118 );
119 let trace_filters = build_warp_trace_filter(
120 acknowledgements,
121 out.clone(),
122 bytes_received,
123 events_received,
124 headers.clone(),
125 traces_deserializer,
126 );
127 log_filters
128 .or(trace_filters)
129 .unify()
130 .or(metrics_filters)
131 .unify()
132 .boxed()
133}
134
135fn enrich_events(
136 events: &mut [Event],
137 headers_config: &[HttpConfigParamKind],
138 headers: &HeaderMap,
139 log_namespace: LogNamespace,
140) {
141 add_headers(
142 events,
143 headers_config,
144 headers,
145 log_namespace,
146 OpentelemetryConfig::NAME,
147 );
148}
149
150fn emit_decode_error(error: impl std::fmt::Display) -> ErrorMessage {
151 let message = format!("Could not decode request: {error}");
152 emit!(HttpBadRequest::new(
153 StatusCode::BAD_REQUEST.as_u16(),
154 &message
155 ));
156 ErrorMessage::new(StatusCode::BAD_REQUEST, message)
157}
158
159fn parse_with_deserializer(
160 deserializer: &OtlpDeserializer,
161 body: Bytes,
162 log_namespace: LogNamespace,
163 events_received: &Registered<EventsReceived>,
164) -> Result<Vec<Event>, ErrorMessage> {
165 let events = deserializer
166 .parse(body, log_namespace)
167 .map(|r| r.into_vec())
168 .map_err(emit_decode_error)?;
169
170 let count = super::count_otlp_items(&events);
172 events_received.emit(CountByteSize(
173 count,
174 events.estimated_json_encoded_size_of(),
175 ));
176
177 Ok(events)
178}
179
180fn build_ingest_filter<Resp, F>(
181 telemetry_type: &'static str,
182 acknowledgements: bool,
183 out: SourceSender,
184 make_events: F,
185) -> BoxedFilter<(Response,)>
186where
187 Resp: prost::Message + Default + Send + 'static,
188 F: Clone
189 + Send
190 + Sync
191 + 'static
192 + Fn(Option<String>, HeaderMap, Bytes) -> Result<Vec<Event>, ErrorMessage>,
193{
194 let body_filter = capped_body();
195
196 warp::post()
197 .and(warp::path("v1"))
198 .and(warp::path(telemetry_type))
199 .and(warp::path::end())
200 .and(warp::header::exact_ignore_case(
201 "content-type",
202 "application/x-protobuf",
203 ))
204 .and(warp::header::optional::<String>("content-encoding"))
205 .and(warp::header::headers_cloned())
206 .and(body_filter)
207 .and_then(
208 move |encoding_header: Option<String>, headers: HeaderMap, body: Bytes| {
209 let events = make_events(encoding_header, headers, body);
210 handle_request(
211 events,
212 acknowledgements,
213 out.clone(),
214 telemetry_type,
215 Resp::default(),
216 )
217 },
218 )
219 .boxed()
220}
221
222fn build_warp_log_filter(
223 acknowledgements: bool,
224 log_namespace: LogNamespace,
225 source_sender: SourceSender,
226 bytes_received: Registered<BytesReceived>,
227 events_received: Registered<EventsReceived>,
228 headers_cfg: Vec<HttpConfigParamKind>,
229 deserializer: Option<OtlpDeserializer>,
230) -> BoxedFilter<(Response,)> {
231 let make_events = move |encoding_header: Option<String>, headers: HeaderMap, body: Bytes| {
232 decompress_body(encoding_header.as_deref(), body)
233 .inspect_err(|err| {
234 if err.status_code() == StatusCode::UNSUPPORTED_MEDIA_TYPE {
236 emit!(HttpBadRequest::new(
237 err.status_code().as_u16(),
238 err.message()
239 ));
240 }
241 })
242 .and_then(|decoded_body| {
243 bytes_received.emit(ByteSize(decoded_body.len()));
244 if let Some(d) = deserializer.as_ref() {
245 parse_with_deserializer(d, decoded_body, log_namespace, &events_received)
246 } else {
247 decode_log_body(decoded_body, log_namespace, &events_received)
248 }
249 .map(|mut events| {
250 enrich_events(&mut events, &headers_cfg, &headers, log_namespace);
251 events
252 })
253 })
254 };
255
256 build_ingest_filter::<ExportLogsServiceResponse, _>(
257 LOGS,
258 acknowledgements,
259 source_sender,
260 make_events,
261 )
262}
263fn build_warp_metrics_filter(
264 acknowledgements: bool,
265 log_namespace: LogNamespace,
266 source_sender: SourceSender,
267 bytes_received: Registered<BytesReceived>,
268 events_received: Registered<EventsReceived>,
269 headers_cfg: Vec<HttpConfigParamKind>,
270 deserializer: Option<OtlpDeserializer>,
271) -> BoxedFilter<(Response,)> {
272 let make_events = move |encoding_header: Option<String>, headers: HeaderMap, body: Bytes| {
273 decompress_body(encoding_header.as_deref(), body)
274 .inspect_err(|err| {
275 if err.status_code() == StatusCode::UNSUPPORTED_MEDIA_TYPE {
277 emit!(HttpBadRequest::new(
278 err.status_code().as_u16(),
279 err.message()
280 ));
281 }
282 })
283 .and_then(|decoded_body| {
284 bytes_received.emit(ByteSize(decoded_body.len()));
285 if let Some(d) = deserializer.as_ref() {
286 parse_with_deserializer(d, decoded_body, log_namespace, &events_received)
287 } else {
288 decode_metrics_body(decoded_body, &events_received)
289 }
290 .map(|mut events| {
291 enrich_events(&mut events, &headers_cfg, &headers, log_namespace);
292 events
293 })
294 })
295 };
296
297 build_ingest_filter::<ExportMetricsServiceResponse, _>(
298 METRICS,
299 acknowledgements,
300 source_sender,
301 make_events,
302 )
303}
304
305fn build_warp_trace_filter(
306 acknowledgements: bool,
307 source_sender: SourceSender,
308 bytes_received: Registered<BytesReceived>,
309 events_received: Registered<EventsReceived>,
310 headers_cfg: Vec<HttpConfigParamKind>,
311 deserializer: Option<OtlpDeserializer>,
312) -> BoxedFilter<(Response,)> {
313 let make_events = move |encoding_header: Option<String>, headers: HeaderMap, body: Bytes| {
314 decompress_body(encoding_header.as_deref(), body)
315 .inspect_err(|err| {
316 if err.status_code() == StatusCode::UNSUPPORTED_MEDIA_TYPE {
318 emit!(HttpBadRequest::new(
319 err.status_code().as_u16(),
320 err.message()
321 ));
322 }
323 })
324 .and_then(|decoded_body| {
325 bytes_received.emit(ByteSize(decoded_body.len()));
326 if let Some(d) = deserializer.as_ref() {
327 parse_with_deserializer(
328 d,
329 decoded_body,
330 LogNamespace::default(),
331 &events_received,
332 )
333 } else {
334 decode_trace_body(decoded_body, &events_received)
335 }
336 .map(|mut events| {
337 enrich_events(&mut events, &headers_cfg, &headers, LogNamespace::default());
338 events
339 })
340 })
341 };
342
343 build_ingest_filter::<ExportTraceServiceResponse, _>(
344 TRACES,
345 acknowledgements,
346 source_sender,
347 make_events,
348 )
349}
350
351fn decode_trace_body(
352 body: Bytes,
353 events_received: &Registered<EventsReceived>,
354) -> Result<Vec<Event>, ErrorMessage> {
355 let request = ExportTraceServiceRequest::decode(body).map_err(emit_decode_error)?;
356
357 let events: Vec<Event> = request
358 .resource_spans
359 .into_iter()
360 .flat_map(|v| v.into_event_iter())
361 .collect();
362
363 events_received.emit(CountByteSize(
364 events.len(),
365 events.estimated_json_encoded_size_of(),
366 ));
367
368 Ok(events)
369}
370
371fn decode_log_body(
372 body: Bytes,
373 log_namespace: LogNamespace,
374 events_received: &Registered<EventsReceived>,
375) -> Result<Vec<Event>, ErrorMessage> {
376 let request = ExportLogsServiceRequest::decode(body).map_err(emit_decode_error)?;
377
378 let events: Vec<Event> = request
379 .resource_logs
380 .into_iter()
381 .flat_map(|v| v.into_event_iter(log_namespace))
382 .collect();
383
384 events_received.emit(CountByteSize(
385 events.len(),
386 events.estimated_json_encoded_size_of(),
387 ));
388
389 Ok(events)
390}
391
392fn decode_metrics_body(
393 body: Bytes,
394 events_received: &Registered<EventsReceived>,
395) -> Result<Vec<Event>, ErrorMessage> {
396 let request = ExportMetricsServiceRequest::decode(body).map_err(emit_decode_error)?;
397
398 let events: Vec<Event> = request
399 .resource_metrics
400 .into_iter()
401 .flat_map(|v| v.into_event_iter())
402 .collect();
403
404 events_received.emit(CountByteSize(
405 events.len(),
406 events.estimated_json_encoded_size_of(),
407 ));
408
409 Ok(events)
410}
411
412async fn handle_request(
413 events: Result<Vec<Event>, ErrorMessage>,
414 acknowledgements: bool,
415 mut out: SourceSender,
416 output: &str,
417 resp: impl Message,
418) -> Result<Response, Rejection> {
419 match events {
420 Ok(mut events) => {
421 let receiver = BatchNotifier::maybe_apply_to(acknowledgements, &mut events);
422 let count = events.len();
423
424 out.send_batch_named(output, events).await.map_err(|_| {
425 emit!(StreamClosedError { count });
426 warp::reject::custom(ApiError::ServerShutdown)
427 })?;
428
429 match receiver {
430 None => Ok(protobuf(resp).into_response()),
431 Some(receiver) => match receiver.await {
432 BatchStatus::Delivered => Ok(protobuf(resp).into_response()),
433 BatchStatus::Errored => Err(warp::reject::custom(Status {
434 code: 2, message: "Error delivering contents to sink".into(),
436 ..Default::default()
437 })),
438 BatchStatus::Rejected => Err(warp::reject::custom(Status {
439 code: 2, message: "Contents failed to deliver to sink".into(),
441 ..Default::default()
442 })),
443 },
444 }
445 }
446 Err(err) => Err(warp::reject::custom(err)),
447 }
448}
449
450async fn handle_rejection(err: Rejection) -> Result<impl Reply, std::convert::Infallible> {
451 if let Some(err_msg) = err.find::<ErrorMessage>() {
452 let reply = protobuf(Status {
453 code: 2, message: err_msg.message().into(),
455 ..Default::default()
456 });
457
458 Ok(warp::reply::with_status(reply, err_msg.status_code()))
459 } else {
460 let reply = protobuf(Status {
461 code: 2, message: format!("{err:?}"),
463 ..Default::default()
464 });
465
466 Ok(warp::reply::with_status(
467 reply,
468 StatusCode::INTERNAL_SERVER_ERROR,
469 ))
470 }
471}