Skip to main content

codecs/encoding/format/
text.rs

1use bytes::{BufMut, BytesMut};
2use tokio_util::codec::Encoder;
3use vector_config_macros::configurable_component;
4use vector_core::{config::DataType, event::Event, schema};
5
6use crate::{MetricTagValues, encoding::format::common::get_serializer_schema_requirement};
7
8/// Config used to build a `TextSerializer`.
9#[configurable_component]
10#[derive(Debug, Clone, Default)]
11pub struct TextSerializerConfig {
12    /// Controls how metric tag values are encoded.
13    ///
14    /// When set to `single`, only the last non-bare value of tags are displayed with the
15    /// metric.  When set to `full`, all metric tags are exposed as separate assignments.
16    /// When set to `auto`, tag values are encoded using their underlying shape.
17    #[serde(default, skip_serializing_if = "vector_core::serde::is_default")]
18    pub metric_tag_values: MetricTagValues,
19}
20
21impl TextSerializerConfig {
22    /// Creates a new `TextSerializerConfig`.
23    pub const fn new(metric_tag_values: MetricTagValues) -> Self {
24        Self { metric_tag_values }
25    }
26
27    /// Build the `TextSerializer` from this configuration.
28    pub const fn build(&self) -> TextSerializer {
29        TextSerializer::new(self.metric_tag_values)
30    }
31
32    /// The data type of events that are accepted by `TextSerializer`.
33    pub fn input_type(&self) -> DataType {
34        DataType::Log | DataType::Metric
35    }
36
37    /// The schema required by the serializer.
38    pub fn schema_requirement(&self) -> schema::Requirement {
39        get_serializer_schema_requirement()
40    }
41}
42
43/// Serializer that converts a log to bytes by extracting the message key, or converts a metric
44/// to bytes by calling its `Display` implementation.
45///
46/// This serializer exists to emulate the behavior of the `StandardEncoding::Text` for backwards
47/// compatibility, until it is phased out completely.
48#[derive(Debug, Clone)]
49pub struct TextSerializer {
50    metric_tag_values: MetricTagValues,
51}
52
53impl TextSerializer {
54    /// Creates a new `TextSerializer`.
55    pub const fn new(metric_tag_values: MetricTagValues) -> Self {
56        Self { metric_tag_values }
57    }
58}
59
60impl Encoder<Event> for TextSerializer {
61    type Error = vector_common::Error;
62
63    fn encode(&mut self, event: Event, buffer: &mut BytesMut) -> Result<(), Self::Error> {
64        match event {
65            Event::Log(log) => {
66                if let Some(bytes) = log.get_message().map(|value| value.coerce_to_bytes()) {
67                    buffer.put(bytes);
68                }
69            }
70            Event::Metric(mut metric) => {
71                if self.metric_tag_values == MetricTagValues::Single {
72                    metric.reduce_tags_to_single();
73                }
74                let bytes = metric.to_string();
75                buffer.put(bytes.as_ref());
76            }
77            Event::Trace(_) => {}
78        };
79
80        Ok(())
81    }
82}
83
84#[cfg(test)]
85mod tests {
86    use bytes::{Bytes, BytesMut};
87    use vector_core::{
88        event::{LogEvent, Metric, MetricKind, MetricValue},
89        metric_tags,
90    };
91
92    use super::*;
93
94    #[test]
95    fn serialize_log() {
96        let buffer = serialize(
97            TextSerializerConfig::default(),
98            Event::from(LogEvent::from_str_legacy("foo")),
99        );
100        assert_eq!(buffer, Bytes::from("foo"));
101    }
102
103    #[test]
104    fn serialize_metric() {
105        let buffer = serialize(
106            TextSerializerConfig::default(),
107            Event::Metric(Metric::new(
108                "users",
109                MetricKind::Incremental,
110                MetricValue::Set {
111                    values: vec!["bob".into()].into_iter().collect(),
112                },
113            )),
114        );
115        assert_eq!(buffer, Bytes::from("users{} + bob"));
116    }
117
118    #[test]
119    fn serialize_metric_tags_full() {
120        let buffer = serialize(
121            TextSerializerConfig {
122                metric_tag_values: MetricTagValues::Full,
123            },
124            metric2(),
125        );
126        assert_eq!(
127            buffer,
128            Bytes::from(r#"counter{a="first",a,a="second"} + 1"#)
129        );
130    }
131
132    #[test]
133    fn serialize_metric_tags_single() {
134        let buffer = serialize(
135            TextSerializerConfig {
136                metric_tag_values: MetricTagValues::Single,
137            },
138            metric2(),
139        );
140        assert_eq!(buffer, Bytes::from(r#"counter{a="second"} + 1"#));
141    }
142
143    fn metric2() -> Event {
144        Event::Metric(
145            Metric::new(
146                "counter",
147                MetricKind::Incremental,
148                MetricValue::Counter { value: 1.0 },
149            )
150            .with_tags(Some(metric_tags! (
151                "a" => "first",
152                "a" => None,
153                "a" => "second",
154            ))),
155        )
156    }
157
158    fn serialize(config: TextSerializerConfig, input: Event) -> Bytes {
159        let mut buffer = BytesMut::new();
160        config.build().encode(input, &mut buffer).unwrap();
161        buffer.freeze()
162    }
163}