Skip to main content

vector/sinks/opentelemetry/
mod.rs

1use indoc::indoc;
2use vector_config::component::GenerateConfig;
3use vector_lib::{
4    codecs::{
5        JsonSerializerConfig,
6        encoding::{FramingConfig, SerializerConfig},
7    },
8    configurable::configurable_component,
9};
10
11use crate::{
12    codecs::{EncodingConfigWithFraming, Transformer},
13    config::{AcknowledgementsConfig, Input, SinkConfig, SinkContext},
14    sinks::{
15        Healthcheck, VectorSink,
16        http::config::{HttpMethod, HttpSinkConfig},
17    },
18};
19
20/// Configuration for the `OpenTelemetry` sink.
21#[configurable_component(sink("opentelemetry", "Deliver OTLP data over HTTP."))]
22#[derive(Clone, Debug, Default)]
23pub struct OpenTelemetryConfig {
24    /// Protocol configuration
25    #[configurable(derived)]
26    protocol: Protocol,
27}
28
29/// The protocol used to send data to OpenTelemetry.
30/// Currently only HTTP is supported, but we plan to support gRPC.
31/// The proto definitions are defined [here](https://github.com/vectordotdev/vector/blob/master/lib/opentelemetry-proto/src/proto/opentelemetry-proto/opentelemetry/proto/README.md).
32#[configurable_component]
33#[derive(Clone, Debug)]
34#[serde(rename_all = "snake_case", tag = "type")]
35#[configurable(metadata(docs::enum_tag_description = "The communication protocol."))]
36pub enum Protocol {
37    /// Send data over HTTP.
38    Http(HttpSinkConfig),
39}
40
41impl Default for Protocol {
42    fn default() -> Self {
43        Protocol::Http(HttpSinkConfig {
44            encoding: EncodingConfigWithFraming::new(
45                Some(FramingConfig::NewlineDelimited),
46                SerializerConfig::Json(JsonSerializerConfig::default()),
47                Transformer::default(),
48            ),
49            uri: Default::default(),
50            method: HttpMethod::Post,
51            auth: Default::default(),
52            compression: Default::default(),
53            payload_prefix: Default::default(),
54            payload_suffix: Default::default(),
55            batch: Default::default(),
56            request: Default::default(),
57            tls: Default::default(),
58            acknowledgements: Default::default(),
59            retry_strategy: Default::default(),
60            confinement: Default::default(),
61        })
62    }
63}
64
65impl GenerateConfig for OpenTelemetryConfig {
66    fn generate_config() -> toml::Value {
67        toml::from_str(indoc! {r#"
68            [protocol]
69            type = "http"
70            uri = "http://localhost:5318/v1/logs"
71            encoding.codec = "json"
72        "#})
73        .unwrap()
74    }
75}
76
77#[async_trait::async_trait]
78#[typetag::serde(name = "opentelemetry")]
79impl SinkConfig for OpenTelemetryConfig {
80    async fn build(&self, cx: SinkContext) -> crate::Result<(VectorSink, Healthcheck)> {
81        match &self.protocol {
82            Protocol::Http(config) => {
83                warn_on_invalid_otlp_batching(config);
84                let result = config
85                    .build_without_confinement_gauge(cx, Self::NAME)
86                    .await?;
87                config.confinement.set_confinement_gauge("sink", Self::NAME);
88                Ok(result)
89            }
90        }
91    }
92
93    fn input(&self) -> Input {
94        match &self.protocol {
95            Protocol::Http(config) => config.input(),
96        }
97    }
98
99    fn acknowledgements(&self) -> &AcknowledgementsConfig {
100        match self.protocol {
101            Protocol::Http(ref config) => config.acknowledgements(),
102        }
103    }
104}
105
106fn warn_on_invalid_otlp_batching(config: &HttpSinkConfig) {
107    let (_, serializer) = config.encoding.config();
108    let is_json = matches!(serializer, SerializerConfig::Json(_));
109    let batches_more_than_one = !matches!(config.batch.max_events, Some(1));
110    if is_json && batches_more_than_one {
111        tracing::warn!(
112            message = "`opentelemetry` sink is configured with `encoding.codec = json` and \
113                       `batch.max_events` greater than 1. This produces invalid OTLP request \
114                       bodies that receivers reject with HTTP 400. Use `encoding.codec = otlp` \
115                       (recommended) or set `batch.max_events = 1`. See \
116                       https://github.com/vectordotdev/vector/issues/22054.",
117        );
118    }
119}
120
121#[cfg(test)]
122mod test {
123    #[test]
124    fn generate_config() {
125        crate::test_util::test_generate_config::<super::OpenTelemetryConfig>();
126    }
127}