vector/transforms/lua/v1/
mod.rs1#![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#[configurable_component]
30#[derive(Clone, Debug)]
31#[serde(deny_unknown_fields)]
32pub struct LuaConfig {
33 source: String,
35
36 #[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 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
79const 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#[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 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}