vector/sinks/opentelemetry/
mod.rs1use 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#[configurable_component(sink("opentelemetry", "Deliver OTLP data over HTTP."))]
22#[derive(Clone, Debug, Default)]
23pub struct OpenTelemetryConfig {
24 #[configurable(derived)]
26 protocol: Protocol,
27}
28
29#[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 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}