1use bytes::Bytes;
2use smallvec::{SmallVec, smallvec};
3use vector_config_macros::configurable_component;
4use vector_core::{
5 compile_vrl,
6 config::{DataType, LogNamespace},
7 event::{Event, MetricTagMode, TargetEvents, VrlTarget},
8 schema,
9};
10use vrl::{
11 compiler::{CompileConfig, Program, TimeZone, TypeState, runtime::Runtime, state::ExternalEnv},
12 diagnostic::Formatter,
13 value::Kind,
14};
15
16use vector_core::event::EventMetadata;
17
18use crate::decoding::format::Deserializer;
19
20#[configurable_component]
22#[derive(Debug, Clone, Default)]
23pub struct VrlDeserializerConfig {
24 pub vrl: VrlDeserializerOptions,
26}
27
28#[configurable_component]
30#[derive(Debug, Clone, PartialEq, Eq, Default)]
31pub struct VrlDeserializerOptions {
32 pub source: String,
39
40 #[serde(default)]
48 pub timezone: Option<TimeZone>,
49}
50
51impl VrlDeserializerConfig {
52 pub fn build(&self) -> vector_common::Result<VrlDeserializer> {
54 let state = TypeState {
55 local: Default::default(),
56 external: ExternalEnv::default(),
57 };
58
59 match compile_vrl(
60 &self.vrl.source,
61 &vector_vrl_functions::all(),
62 &state,
63 CompileConfig::default(),
64 ) {
65 Ok(result) => Ok(VrlDeserializer {
66 program: result.program,
67 timezone: self.vrl.timezone.unwrap_or(TimeZone::Local),
68 metadata_template: None,
69 }),
70 Err(diagnostics) => Err(Formatter::new(&self.vrl.source, diagnostics)
71 .to_string()
72 .into()),
73 }
74 }
75
76 pub fn output_type(&self) -> DataType {
78 DataType::Log
79 }
80
81 pub fn schema_definition(&self, log_namespace: LogNamespace) -> schema::Definition {
83 match log_namespace {
84 LogNamespace::Legacy => {
85 schema::Definition::empty_legacy_namespace().unknown_fields(Kind::any())
86 }
87 LogNamespace::Vector => {
88 schema::Definition::new_with_default_metadata(Kind::any(), [log_namespace])
89 }
90 }
91 }
92}
93
94#[derive(Debug, Clone)]
96pub struct VrlDeserializer {
97 program: Program,
98 timezone: TimeZone,
99 metadata_template: Option<EventMetadata>,
103}
104
105impl VrlDeserializer {
106 #[must_use]
113 pub fn with_metadata_template(mut self, metadata: EventMetadata) -> Self {
114 self.metadata_template = Some(metadata);
115 self
116 }
117}
118
119fn parse_bytes(bytes: Bytes, log_namespace: LogNamespace) -> Event {
120 use crate::BytesDeserializerConfig;
121 let bytes_deserializer = BytesDeserializerConfig::new().build();
122 let log_event = bytes_deserializer.parse_single(bytes, log_namespace);
123 Event::from(log_event)
124}
125
126impl Deserializer for VrlDeserializer {
127 fn parse(
128 &self,
129 bytes: Bytes,
130 log_namespace: LogNamespace,
131 ) -> vector_common::Result<SmallVec<[Event; 1]>> {
132 let mut event = parse_bytes(bytes, log_namespace);
133 if let Some(template) = &self.metadata_template {
134 *event.metadata_mut() = template.clone();
138 }
139 self.run_vrl(event, log_namespace)
140 }
141}
142
143impl VrlDeserializer {
144 fn run_vrl(
145 &self,
146 event: Event,
147 log_namespace: LogNamespace,
148 ) -> vector_common::Result<SmallVec<[Event; 1]>> {
149 let mut runtime = Runtime::default();
150 let mut target = VrlTarget::new(event, self.program.info(), MetricTagMode::Full);
151 match runtime.resolve(&mut target, &self.program, &self.timezone) {
152 Ok(_) => match target.into_events(log_namespace) {
153 TargetEvents::One(event) => Ok(smallvec![event]),
154 TargetEvents::Logs(events_iter) => Ok(SmallVec::from_iter(events_iter)),
155 TargetEvents::Traces(_) => Err("trace targets are not supported".into()),
156 },
157 Err(e) => Err(e.to_string().into()),
158 }
159 }
160}
161
162#[cfg(test)]
163mod tests {
164 use chrono::{DateTime, Utc};
165 use indoc::indoc;
166 use vrl::{btreemap, event_path, path::OwnedTargetPath, value::Value};
167
168 use super::*;
169
170 fn make_decoder(source: &str) -> VrlDeserializer {
171 VrlDeserializerConfig {
172 vrl: VrlDeserializerOptions {
173 source: source.to_string(),
174 timezone: None,
175 },
176 }
177 .build()
178 .expect("Failed to build VrlDeserializer")
179 }
180
181 #[test]
182 fn test_json_message() {
183 let source = indoc!(
184 r#"
185 %m1 = "metadata"
186 . = string!(.)
187 . = parse_json!(.)
188 "#
189 );
190
191 let decoder = make_decoder(source);
192
193 let log_bytes = Bytes::from(r#"{ "message": "Hello VRL" }"#);
194 let result = decoder.parse(log_bytes, LogNamespace::Vector).unwrap();
195 assert_eq!(result.len(), 1);
196 let event = result.first().unwrap();
197 assert_eq!(
198 *event.as_log().get(&OwnedTargetPath::event_root()).unwrap(),
199 btreemap! { "message" => "Hello VRL" }.into()
200 );
201 assert_eq!(
202 *event
203 .as_log()
204 .get(&OwnedTargetPath::metadata_root())
205 .unwrap(),
206 btreemap! { "m1" => "metadata" }.into()
207 );
208 }
209
210 #[test]
211 fn test_ignored_returned_expression() {
212 let source = indoc!(
213 r#"
214 . = { "a" : 1 }
215 { "b" : 9 }
216 "#
217 );
218
219 let decoder = make_decoder(source);
220
221 let log_bytes = Bytes::from("some bytes");
222 let result = decoder.parse(log_bytes, LogNamespace::Vector).unwrap();
223 assert_eq!(result.len(), 1);
224 let event = result.first().unwrap();
225 assert_eq!(
226 *event.as_log().get(&OwnedTargetPath::event_root()).unwrap(),
227 btreemap! { "a" => 1 }.into()
228 );
229 }
230
231 #[test]
232 fn test_multiple_events() {
233 let source = indoc!(". = [0,1,2]");
234 let decoder = make_decoder(source);
235 let log_bytes = Bytes::from("some bytes");
236 let result = decoder.parse(log_bytes, LogNamespace::Vector).unwrap();
237 assert_eq!(result.len(), 3);
238 for (i, event) in result.iter().enumerate() {
239 assert_eq!(
240 *event.as_log().get(&OwnedTargetPath::event_root()).unwrap(),
241 i.into()
242 );
243 }
244 }
245
246 #[test]
247 fn test_syslog_and_cef_input() {
248 let source = indoc!(
249 r#"
250 if exists(.message) {
251 . = string!(.message)
252 }
253 . = parse_syslog(.) ?? parse_cef(.) ?? null
254 "#
255 );
256
257 let decoder = make_decoder(source);
258
259 let syslog_bytes = Bytes::from(
261 "<34>1 2024-02-06T15:04:05.000Z mymachine.example.com su - ID47 - 'su root' failed for user on /dev/pts/8",
262 );
263 let result = decoder.parse(syslog_bytes, LogNamespace::Vector).unwrap();
264 assert_eq!(result.len(), 1);
265 let syslog_event = result.first().unwrap();
266 assert_eq!(
267 *syslog_event
268 .as_log()
269 .get(&OwnedTargetPath::event_root())
270 .unwrap(),
271 btreemap! {
272 "appname" => "su",
273 "facility" => "auth",
274 "hostname" => "mymachine.example.com",
275 "message" => "'su root' failed for user on /dev/pts/8",
276 "msgid" => "ID47",
277 "severity" => "crit",
278 "timestamp" => "2024-02-06T15:04:05Z".parse::<DateTime<Utc>>().unwrap(),
279 "version" => 1
280 }
281 .into()
282 );
283
284 let cef_bytes = Bytes::from(
286 "CEF:0|Security|Threat Manager|1.0|100|worm successfully stopped|10|src=10.0.0.1 dst=2.1.2.2 spt=1232",
287 );
288 let result = decoder.parse(cef_bytes, LogNamespace::Vector).unwrap();
289 assert_eq!(result.len(), 1);
290 let cef_event = result.first().unwrap();
291 assert_eq!(
292 *cef_event
293 .as_log()
294 .get(&OwnedTargetPath::event_root())
295 .unwrap(),
296 btreemap! {
297 "cefVersion" =>"0",
298 "deviceEventClassId" =>"100",
299 "deviceProduct" =>"Threat Manager",
300 "deviceVendor" =>"Security",
301 "deviceVersion" =>"1.0",
302 "dst" =>"2.1.2.2",
303 "name" =>"worm successfully stopped",
304 "severity" =>"10",
305 "spt" =>"1232",
306 "src" =>"10.0.0.1"
307 }
308 .into()
309 );
310 let random_bytes = Bytes::from("a|- -| x");
311 let result = decoder.parse(random_bytes, LogNamespace::Vector).unwrap();
312 let random_event = result.first().unwrap();
313 assert_eq!(result.len(), 1);
314 assert_eq!(
315 *random_event
316 .as_log()
317 .get(&OwnedTargetPath::event_root())
318 .unwrap(),
319 Value::Null
320 );
321 }
322
323 #[test]
324 fn test_invalid_source() {
325 let error = VrlDeserializerConfig {
326 vrl: VrlDeserializerOptions {
327 source: ". ?".to_string(),
328 timezone: None,
329 },
330 }
331 .build()
332 .unwrap_err()
333 .to_string();
334 assert!(error.contains("error[E203]: syntax error"));
335 }
336
337 #[test]
338 fn test_abort() {
339 let decoder = make_decoder("abort");
340 let log_bytes = Bytes::from(r#"{ "message": "Hello VRL" }"#);
341 let error = decoder
342 .parse(log_bytes, LogNamespace::Vector)
343 .unwrap_err()
344 .to_string();
345 assert!(error.contains("aborted"));
346 }
347
348 fn metadata_with_secret(key: &str, value: &str) -> EventMetadata {
349 let mut metadata = EventMetadata::default();
350 metadata.secrets_mut().insert(key, value);
351 metadata
352 }
353
354 #[test]
357 fn test_with_metadata_template_vrl_can_read_secret() {
358 let decoder = make_decoder(r#".secret_value = get_secret!("my_token")"#)
362 .with_metadata_template(metadata_with_secret("my_token", "super-secret"));
363
364 let bytes = Bytes::from(r#"hello"#);
365 let events = decoder
366 .parse(bytes, LogNamespace::Legacy)
367 .expect("parse should succeed");
368
369 assert_eq!(events.len(), 1);
370 assert_eq!(
371 *events[0].as_log().get(event_path!("secret_value")).unwrap(),
372 Value::from("super-secret")
373 );
374 }
375
376 #[test]
379 fn test_with_metadata_template_codec_wins_on_collision() {
380 let decoder = make_decoder(r#"set_secret!("my_token", "codec-wins")"#)
381 .with_metadata_template(metadata_with_secret("my_token", "template-loses"));
382
383 let bytes = Bytes::from(r#"hello"#);
384 let events = decoder
385 .parse(bytes, LogNamespace::Legacy)
386 .expect("parse should succeed");
387
388 assert_eq!(
389 events[0]
390 .metadata()
391 .secrets()
392 .get("my_token")
393 .unwrap()
394 .as_ref(),
395 "codec-wins"
396 );
397 }
398}