Skip to main content

codecs/decoding/format/
vrl.rs

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/// Config used to build a `VrlDeserializer`.
21#[configurable_component]
22#[derive(Debug, Clone, Default)]
23pub struct VrlDeserializerConfig {
24    /// VRL-specific decoding options.
25    pub vrl: VrlDeserializerOptions,
26}
27
28/// VRL-specific decoding options.
29#[configurable_component]
30#[derive(Debug, Clone, PartialEq, Eq, Default)]
31pub struct VrlDeserializerOptions {
32    /// The [Vector Remap Language][vrl] (VRL) program to execute for each event.
33    /// The final contents of the `.` target are used as the decoding result.
34    /// Compilation errors or use of `abort` in the program result in a decoding error.
35    ///
36    ///
37    /// [vrl]: https://vector.dev/docs/reference/vrl
38    pub source: String,
39
40    /// The name of the timezone to apply to timestamp conversions that do not contain an explicit
41    /// time zone. The time zone name may be any name in the [TZ database][tz_database], or `local`
42    /// to indicate system local time.
43    ///
44    /// If not set, `local` is used.
45    ///
46    /// [tz_database]: https://en.wikipedia.org/wiki/List_of_tz_database_time_zones
47    #[serde(default)]
48    pub timezone: Option<TimeZone>,
49}
50
51impl VrlDeserializerConfig {
52    /// Build the `VrlDeserializer` from this configuration.
53    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    /// Return the type of event build by this deserializer.
77    pub fn output_type(&self) -> DataType {
78        DataType::Log
79    }
80
81    /// The schema produced by the deserializer.
82    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/// Deserializer that builds `Event`s from a byte frame containing logs compatible with VRL.
95#[derive(Debug, Clone)]
96pub struct VrlDeserializer {
97    program: Program,
98    timezone: TimeZone,
99    /// Per-call metadata template set by the source before decoding. When
100    /// present, every `%`-prefixed path in the template is accessible from
101    /// within the VRL program (e.g. `%splunk_hec.host`, `%vector.secrets.*`).
102    metadata_template: Option<EventMetadata>,
103}
104
105impl VrlDeserializer {
106    /// Attach a metadata template that will be pre-populated on each synthetic
107    /// event before the VRL program runs.
108    ///
109    /// Sources call this once per decode call with the per-request context they
110    /// have assembled (envelope fields, auth tokens, etc.). VRL can then read
111    /// those values via `%`-prefixed paths.
112    #[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            // Pre-populate the synthetic event with the source-assembled metadata so
135            // every `%`-prefixed path is in scope when VRL executes. This lets
136            // user programs read `%splunk_hec.host`, `%vector.secrets.*`, etc.
137            *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        // Syslog input
260        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        // CEF input
285        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    /// A VRL program that uses `get_secret!()` can read a secret injected via
355    /// `with_metadata_template`.
356    #[test]
357    fn test_with_metadata_template_vrl_can_read_secret() {
358        // VRL program copies the injected secret into an event field so we can
359        // assert on its value. The input bytes become `.message` (Legacy namespace)
360        // and we add `.secret_value` alongside it.
361        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    /// Secrets explicitly set by the VRL program win over the template because
377    /// `set_secret!` runs after the template is pre-populated.
378    #[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}