codecs/encoding/format/
text.rs1use 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#[configurable_component]
10#[derive(Debug, Clone, Default)]
11pub struct TextSerializerConfig {
12 #[serde(default, skip_serializing_if = "vector_core::serde::is_default")]
18 pub metric_tag_values: MetricTagValues,
19}
20
21impl TextSerializerConfig {
22 pub const fn new(metric_tag_values: MetricTagValues) -> Self {
24 Self { metric_tag_values }
25 }
26
27 pub const fn build(&self) -> TextSerializer {
29 TextSerializer::new(self.metric_tag_values)
30 }
31
32 pub fn input_type(&self) -> DataType {
34 DataType::Log | DataType::Metric
35 }
36
37 pub fn schema_requirement(&self) -> schema::Requirement {
39 get_serializer_schema_requirement()
40 }
41}
42
43#[derive(Debug, Clone)]
49pub struct TextSerializer {
50 metric_tag_values: MetricTagValues,
51}
52
53impl TextSerializer {
54 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}