Skip to main content

vector/transforms/lua/v1/
mod.rs

1// Derivative's Debug impl generates `let _ = field.fmt(f)` which triggers this lint.
2#![allow(clippy::let_underscore_must_use)]
3
4use std::{future::ready, pin::Pin};
5
6use futures::{Stream, StreamExt, stream};
7use mlua::{ExternalError, FromLua};
8use ordered_float::NotNan;
9use snafu::{ResultExt, Snafu};
10use vector_lib::configurable::configurable_component;
11use vrl::path::parse_target_path;
12
13use crate::{
14    config::{DataType, Input, OutputId, TransformOutput},
15    event::{Event, Value},
16    internal_events::{LuaGcTriggered, LuaScriptError},
17    schema,
18    schema::Definition,
19    transforms::{TaskTransform, Transform},
20};
21
22#[derive(Debug, Snafu)]
23enum BuildError {
24    #[snafu(display("Lua error: {}", source))]
25    InvalidLua { source: mlua::Error },
26}
27
28/// Configuration for version one of the `lua` transform.
29#[configurable_component]
30#[derive(Clone, Debug)]
31#[serde(deny_unknown_fields)]
32pub struct LuaConfig {
33    /// The Lua program to execute for each event.
34    source: String,
35
36    /// A list of directories to search when loading a Lua file via the `require` function.
37    ///
38    /// If not specified, the modules are looked up in the configuration directories.
39    #[serde(default)]
40    search_dirs: Vec<String>,
41}
42
43impl LuaConfig {
44    pub fn build(&self) -> crate::Result<Transform> {
45        warn!(
46            "DEPRECATED The `lua` transform API version 1 is deprecated. Please convert your script to version 2."
47        );
48        Lua::new(self.source.clone(), self.search_dirs.clone()).map(Transform::event_task)
49    }
50
51    pub fn input(&self) -> Input {
52        Input::log()
53    }
54
55    pub fn outputs(
56        &self,
57        input_definitions: &[(OutputId, schema::Definition)],
58    ) -> Vec<TransformOutput> {
59        // Lua causes the type definition to be reset
60        let namespaces = input_definitions
61            .iter()
62            .flat_map(|(_output, definition)| definition.log_namespaces().clone())
63            .collect();
64
65        let definition = input_definitions
66            .iter()
67            .map(|(output, _definition)| {
68                (
69                    output.clone(),
70                    Definition::default_for_namespace(&namespaces),
71                )
72            })
73            .collect();
74
75        vec![TransformOutput::new(DataType::Log, definition)]
76    }
77}
78
79// Lua's garbage collector sometimes seems to be not executed automatically on high event rates,
80// which leads to leak-like RAM consumption pattern. This constant sets the number of invocations of
81// the Lua transform after which GC would be called, thus ensuring that the RAM usage is not too high.
82//
83// This constant is larger than 1 because calling GC is an expensive operation, so doing it
84// after each transform would have significant footprint on the performance.
85const GC_INTERVAL: usize = 16;
86
87#[derive(Derivative)]
88#[derivative(Debug)]
89pub struct Lua {
90    #[derivative(Debug = "ignore")]
91    source: String,
92    #[derivative(Debug = "ignore")]
93    search_dirs: Vec<String>,
94    #[derivative(Debug = "ignore")]
95    lua: mlua::Lua,
96    vector_func: mlua::RegistryKey,
97    invocations_after_gc: usize,
98}
99
100impl Clone for Lua {
101    fn clone(&self) -> Self {
102        Lua::new(self.source.clone(), self.search_dirs.clone())
103            .expect("Tried to clone existing valid lua transform. This is an invariant.")
104    }
105}
106
107// This wrapping structure is added in order to make it possible to have independent implementations
108// of `mlua::UserData` trait for event in version 1 and version 2 of the transform.
109#[derive(Clone, FromLua)]
110struct LuaEvent {
111    inner: Event,
112}
113
114impl Lua {
115    pub fn new(source: String, search_dirs: Vec<String>) -> crate::Result<Self> {
116        // In order to support loading C modules in Lua, we need to create unsafe instance
117        // without debug library.
118        let lua = unsafe {
119            mlua::Lua::unsafe_new_with(mlua::StdLib::ALL_SAFE, mlua::LuaOptions::default())
120        };
121
122        let additional_paths = search_dirs
123            .iter()
124            .map(|d| format!("{d}/?.lua"))
125            .collect::<Vec<_>>()
126            .join(";");
127
128        if !additional_paths.is_empty() {
129            let package = lua
130                .globals()
131                .get::<mlua::Table>("package")
132                .context(InvalidLuaSnafu)?;
133            let current_paths = package
134                .get::<String>("path")
135                .unwrap_or_else(|_| ";".to_string());
136            let paths = format!("{additional_paths};{current_paths}");
137            package.set("path", paths).context(InvalidLuaSnafu)?;
138        }
139
140        let func = lua.load(&source).into_function().context(InvalidLuaSnafu)?;
141        let vector_func = lua.create_registry_value(func).context(InvalidLuaSnafu)?;
142
143        Ok(Self {
144            source,
145            search_dirs,
146            lua,
147            vector_func,
148            invocations_after_gc: 0,
149        })
150    }
151
152    fn process(&mut self, event: Event) -> Result<Option<Event>, mlua::Error> {
153        let source_id = event.source_id().cloned();
154        let lua = &self.lua;
155        let globals = lua.globals();
156
157        globals.raw_set("event", LuaEvent { inner: event })?;
158
159        let func = lua.registry_value::<mlua::Function>(&self.vector_func)?;
160        func.call::<()>(())?;
161
162        let result = globals.raw_get::<Option<LuaEvent>>("event").map(|option| {
163            option.map(|lua_event| {
164                let mut event = lua_event.inner;
165                if let Some(source_id) = source_id {
166                    event.set_source_id(source_id);
167                }
168                event
169            })
170        });
171
172        self.invocations_after_gc += 1;
173        if self.invocations_after_gc.is_multiple_of(GC_INTERVAL) {
174            emit!(LuaGcTriggered {
175                used_memory: self.lua.used_memory()
176            });
177            self.lua.gc_collect()?;
178            self.invocations_after_gc = 0;
179        }
180
181        result
182    }
183
184    pub fn transform_one(&mut self, event: Event) -> Option<Event> {
185        match self.process(event) {
186            Ok(event) => event,
187            Err(error) => {
188                emit!(LuaScriptError { error });
189                None
190            }
191        }
192    }
193}
194
195impl TaskTransform<Event> for Lua {
196    fn transform(
197        self: Box<Self>,
198        task: Pin<Box<dyn Stream<Item = Event> + Send>>,
199    ) -> Pin<Box<dyn Stream<Item = Event> + Send>>
200    where
201        Self: 'static,
202    {
203        let mut inner = self;
204        Box::pin(
205            task.filter_map(move |event| {
206                let mut output = Vec::with_capacity(1);
207                ready(match inner.process(event) {
208                    Ok(event) => {
209                        output.extend(event);
210                        Some(stream::iter(output))
211                    }
212                    Err(error) => {
213                        emit!(LuaScriptError { error });
214                        None
215                    }
216                })
217            })
218            .flatten(),
219        )
220    }
221}
222
223impl mlua::UserData for LuaEvent {
224    fn add_methods<M: mlua::UserDataMethods<Self>>(methods: &mut M) {
225        methods.add_meta_method_mut(
226            mlua::MetaMethod::NewIndex,
227            |_lua, this, (key, value): (String, Option<mlua::Value>)| {
228                let key_path = parse_target_path(key.as_str()).map_err(|e| e.into_lua_err())?;
229                match value {
230                    Some(mlua::Value::String(string)) => {
231                        this.inner.as_mut_log().insert(
232                            &key_path,
233                            Value::from(string.to_str().expect("Expected UTF-8.").to_owned()),
234                        );
235                    }
236                    Some(mlua::Value::Integer(integer)) => {
237                        this.inner
238                            .as_mut_log()
239                            .insert(&key_path, Value::Integer(integer));
240                    }
241                    Some(mlua::Value::Number(number)) if !number.is_nan() => {
242                        this.inner
243                            .as_mut_log()
244                            .insert(&key_path, Value::Float(NotNan::new(number).unwrap()));
245                    }
246                    Some(mlua::Value::Boolean(boolean)) => {
247                        this.inner
248                            .as_mut_log()
249                            .insert(&key_path, Value::Boolean(boolean));
250                    }
251                    Some(mlua::Value::Nil) | None => {
252                        this.inner.as_mut_log().remove(&key_path);
253                    }
254                    _ => {
255                        info!(
256                            message =
257                                "Could not set field to Lua value of invalid type, dropping field.",
258                            field = key.as_str()
259                        );
260                        this.inner.as_mut_log().remove(&key_path);
261                    }
262                }
263
264                Ok(())
265            },
266        );
267
268        methods.add_meta_method(mlua::MetaMethod::Index, |lua, this, key: String| {
269            if let Some(value) = this
270                .inner
271                .as_log()
272                .parse_path_and_get_value(key.as_str())
273                .ok()
274                .flatten()
275            {
276                let string = lua.create_string(value.coerce_to_bytes())?;
277                Ok(Some(string))
278            } else {
279                Ok(None)
280            }
281        });
282
283        methods.add_meta_function(mlua::MetaMethod::Pairs, |lua, event: LuaEvent| {
284            let state = lua.create_table()?;
285            {
286                if let Some(keys) = event.inner.as_log().keys() {
287                    let keys = lua.create_table_from(keys.map(|k| (k, true)))?;
288                    state.raw_set("keys", keys)?;
289                }
290                state.raw_set("event", event)?;
291            }
292            let function =
293                lua.create_function(|lua, (state, prev): (mlua::Table, Option<String>)| {
294                    let event: LuaEvent = state.raw_get("event")?;
295                    let keys: mlua::Table = state.raw_get("keys")?;
296                    let next: mlua::Function = lua.globals().raw_get("next")?;
297                    let key: Option<String> = next.call((keys, prev))?;
298                    let value = key.clone().and_then(|k| {
299                        event
300                            .inner
301                            .as_log()
302                            .parse_path_and_get_value(k.as_str())
303                            .ok()
304                            .flatten()
305                    });
306                    match value {
307                        Some(value) => Ok((key, Some(lua.create_string(value.coerce_to_bytes())?))),
308                        None => Ok((None, None)),
309                    }
310                })?;
311            Ok((function, state))
312        });
313    }
314}
315
316pub fn format_error(error: &mlua::Error) -> String {
317    match error {
318        mlua::Error::CallbackError { traceback, cause } => format_error(cause) + "\n" + traceback,
319        err => err.to_string(),
320    }
321}
322
323#[cfg(test)]
324mod tests {
325    use std::sync::Arc;
326
327    use vrl::event_path;
328
329    use super::*;
330    use crate::{
331        config::ComponentKey,
332        event::{Event, LogEvent, Value},
333        test_util,
334    };
335
336    #[test]
337    fn lua_add_field() {
338        let event = transform_one(
339            r#"
340              event["hello"] = "goodbye"
341            "#,
342            LogEvent::from("program me"),
343        )
344        .unwrap();
345
346        assert_eq!(event.as_log()["hello"], "goodbye".into());
347    }
348
349    #[test]
350    fn lua_read_field() {
351        let event = transform_one(
352            r#"
353              _, _, name = string.find(event["message"], "Hello, my name is (%a+).")
354              event["name"] = name
355            "#,
356            LogEvent::from("Hello, my name is Bob."),
357        )
358        .unwrap();
359
360        assert_eq!(event.as_log()["name"], "Bob".into());
361    }
362
363    #[test]
364    fn lua_remove_field() {
365        let mut log = LogEvent::default();
366        log.insert(event_path!("name"), "Bob");
367        let event = transform_one(
368            r#"
369              event["name"] = nil
370            "#,
371            log,
372        )
373        .unwrap();
374
375        assert!(event.as_log().get(event_path!("name")).is_none());
376    }
377
378    #[test]
379    fn lua_drop_event() {
380        let mut log = LogEvent::default();
381        log.insert(event_path!("name"), "Bob");
382        let event = transform_one(
383            r"
384              event = nil
385            ",
386            log,
387        );
388
389        assert!(event.is_none());
390    }
391
392    #[test]
393    fn lua_read_empty_field() {
394        let event = transform_one(
395            r#"
396              if event["non-existent"] == nil then
397                event["result"] = "empty"
398              else
399                event["result"] = "found"
400              end
401            "#,
402            LogEvent::default(),
403        )
404        .unwrap();
405
406        assert_eq!(event.as_log()["result"], "empty".into());
407    }
408
409    #[test]
410    fn lua_integer_value() {
411        let event = transform_one(
412            r#"
413              event["number"] = 3
414            "#,
415            LogEvent::default(),
416        )
417        .unwrap();
418        assert_eq!(event.as_log()["number"], Value::Integer(3));
419    }
420
421    #[test]
422    fn lua_numeric_value() {
423        let event = transform_one(
424            r#"
425              event["number"] = 3.14159
426            "#,
427            LogEvent::default(),
428        )
429        .unwrap();
430        assert_eq!(event.as_log()["number"], Value::from(3.14159));
431    }
432
433    #[test]
434    fn lua_boolean_value() {
435        let event = transform_one(
436            r#"
437              event["bool"] = true
438            "#,
439            LogEvent::default(),
440        )
441        .unwrap();
442        assert_eq!(event.as_log()["bool"], Value::Boolean(true));
443    }
444
445    #[test]
446    fn lua_non_coercible_value() {
447        let event = transform_one(
448            r#"
449              event["junk"] = {"asdf"}
450            "#,
451            LogEvent::default(),
452        )
453        .unwrap();
454        assert_eq!(event.as_log().get(event_path!("junk")), None);
455    }
456
457    #[test]
458    fn lua_non_string_key_write() {
459        crate::test_util::trace_init();
460        let mut transform = Lua::new(
461            r#"
462              event[false] = "hello"
463            "#
464            .to_string(),
465            vec![],
466        )
467        .unwrap();
468
469        let err = transform.process(LogEvent::default().into()).unwrap_err();
470        let err = format_error(&err);
471        assert!(
472            err.contains("error converting Lua boolean to String"),
473            "{}",
474            err
475        );
476    }
477
478    #[test]
479    fn lua_non_string_key_read() {
480        crate::test_util::trace_init();
481        let mut transform = Lua::new(
482            r"
483              print(event[false])
484            "
485            .to_string(),
486            vec![],
487        )
488        .unwrap();
489
490        let err = transform.process(LogEvent::default().into()).unwrap_err();
491        let err = format_error(&err);
492        assert!(
493            err.contains("error converting Lua boolean to String"),
494            "{}",
495            err
496        );
497    }
498
499    #[test]
500    fn lua_script_error() {
501        crate::test_util::trace_init();
502        let mut transform = Lua::new(
503            r#"
504              error("this is an error")
505            "#
506            .to_string(),
507            vec![],
508        )
509        .unwrap();
510
511        let err = transform.process(LogEvent::default().into()).unwrap_err();
512        let err = format_error(&err);
513        assert!(err.contains("this is an error"), "{}", err);
514    }
515
516    #[test]
517    fn lua_syntax_error() {
518        crate::test_util::trace_init();
519        let err = Lua::new(
520            r"
521              1234 = sadf <>&*!#@
522            "
523            .to_string(),
524            vec![],
525        )
526        .map(|_| ())
527        .unwrap_err()
528        .to_string();
529
530        assert!(err.contains("syntax error:"), "{}", err);
531    }
532
533    #[test]
534    fn lua_load_file() {
535        use std::{fs::File, io::Write};
536        crate::test_util::trace_init();
537
538        let dir = tempfile::tempdir().unwrap();
539
540        let mut file = File::create(dir.path().join("script2.lua")).unwrap();
541        write!(
542            &mut file,
543            r#"
544              local M = {{}}
545
546              local function modify(event2)
547                event2["\"new field\""] = "new value"
548              end
549              M.modify = modify
550
551              return M
552            "#
553        )
554        .unwrap();
555
556        let source = r#"
557          local script2 = require("script2")
558          script2.modify(event)
559        "#
560        .to_string();
561
562        let mut transform =
563            Lua::new(source, vec![dir.path().to_string_lossy().into_owned()]).unwrap();
564        let event = transform.transform_one(LogEvent::default().into()).unwrap();
565        assert_eq!(event.as_log()["\"new field\""], "new value".into());
566    }
567
568    #[test]
569    fn lua_pairs() {
570        let mut event = LogEvent::default();
571        event.insert(event_path!("name"), "Bob");
572        event.insert(event_path!("friend"), "Alice");
573
574        let event = transform_one(
575            r"
576              for k,v in pairs(event) do
577                event[k] = k .. v
578              end
579            ",
580            event,
581        )
582        .unwrap();
583
584        assert_eq!(event.as_log()["name"], "nameBob".into());
585        assert_eq!(event.as_log()["friend"], "friendAlice".into());
586    }
587
588    fn transform_one(transform: &str, event: impl Into<Event>) -> Option<Event> {
589        crate::test_util::trace_init();
590
591        let source = source_id();
592        let mut event = event.into();
593        event.set_source_id(Arc::clone(&source));
594
595        let mut transform = Lua::new(transform.to_string(), vec![]).unwrap();
596        let event = transform.transform_one(event);
597
598        if let Some(event) = &event {
599            assert_eq!(event.source_id(), Some(&source));
600        }
601
602        event
603    }
604
605    fn source_id() -> Arc<ComponentKey> {
606        Arc::new(ComponentKey::from(test_util::random_string(16)))
607    }
608}