Skip to main content

vector/sources/opentelemetry/
http.rs

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)] // TODO change to a builder struct
90pub(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    // Count individual items within OTLP batches for consistency with other sources
171    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                // Other status codes are already handled by `sources::util::decompress_body` (tech debt).
235                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                // Other status codes are already handled by `sources::util::decompress_body` (tech debt).
276                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                // Other status codes are already handled by `sources::util::decompress_body` (tech debt).
317                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, // UNKNOWN - OTLP doesn't require use of status.code, but we can't encode a None here
435                        message: "Error delivering contents to sink".into(),
436                        ..Default::default()
437                    })),
438                    BatchStatus::Rejected => Err(warp::reject::custom(Status {
439                        code: 2, // UNKNOWN - OTLP doesn't require use of status.code, but we can't encode a None here
440                        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, // UNKNOWN - OTLP doesn't require use of status.code, but we can't encode a None here
454            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, // UNKNOWN - OTLP doesn't require use of status.code, but we can't encode a None here
462            message: format!("{err:?}"),
463            ..Default::default()
464        });
465
466        Ok(warp::reply::with_status(
467            reply,
468            StatusCode::INTERNAL_SERVER_ERROR,
469        ))
470    }
471}