dataflow_rs/engine/functions/
log.rs1use crate::engine::error::Result;
2use crate::engine::executor::{ArenaContext, with_arena};
3use crate::engine::message::{Change, Message};
4use crate::engine::task_outcome::TaskOutcome;
5use datalogic_rs::{Engine, Logic};
6use datavalue::DataValue;
7use log::{debug, error, info, trace, warn};
8use serde::Deserialize;
9use serde_json::Value;
10use std::collections::HashMap;
11use std::sync::Arc;
12
13#[derive(Debug, Clone, Default, Deserialize)]
15#[serde(rename_all = "lowercase")]
16pub enum LogLevel {
17 Trace,
18 Debug,
19 #[default]
20 Info,
21 Warn,
22 Error,
23}
24
25impl LogLevel {
26 fn as_log_level(&self) -> log::Level {
28 match self {
29 LogLevel::Trace => log::Level::Trace,
30 LogLevel::Debug => log::Level::Debug,
31 LogLevel::Info => log::Level::Info,
32 LogLevel::Warn => log::Level::Warn,
33 LogLevel::Error => log::Level::Error,
34 }
35 }
36}
37
38#[derive(Debug, Clone, Deserialize)]
42pub struct LogConfig {
43 #[serde(default)]
45 pub level: LogLevel,
46
47 pub message: Value,
49
50 #[serde(default)]
52 pub fields: HashMap<String, Value>,
53
54 #[serde(skip)]
56 pub compiled_message: Option<Arc<Logic>>,
57
58 #[serde(skip)]
62 pub compiled_fields: Vec<(String, Option<Arc<Logic>>)>,
63}
64
65impl LogConfig {
66 pub fn execute(
73 &self,
74 message: &mut Message,
75 engine: &Arc<Engine>,
76 ) -> Result<(TaskOutcome, Vec<Change>)> {
77 with_arena(|arena| {
78 let mut arena_ctx = ArenaContext::from_owned(&message.context, arena);
79 self.execute_in_arena(message, &mut arena_ctx, engine)
80 })
81 }
82
83 pub(crate) fn execute_in_arena(
87 &self,
88 _message: &mut Message,
89 arena_ctx: &mut ArenaContext<'_>,
90 engine: &Arc<Engine>,
91 ) -> Result<(TaskOutcome, Vec<Change>)> {
92 if !log::log_enabled!(target: "dataflow::log", self.level.as_log_level()) {
96 return Ok((TaskOutcome::Success, vec![]));
97 }
98
99 let arena = arena_ctx.arena();
100 let ctx_av = arena_ctx.as_data_value();
101
102 let stringify = |compiled: &Logic| -> String {
106 match engine.evaluate(compiled, ctx_av, arena) {
107 Ok(DataValue::String(s)) => (*s).to_string(),
108 Ok(other) => other.to_string(),
109 Err(e) => {
110 error!("Log: Failed to evaluate expression: {:?}", e);
111 "<eval error>".to_string()
112 }
113 }
114 };
115
116 let log_message = match &self.compiled_message {
117 Some(compiled) => stringify(compiled),
118 None => "<uncompiled message>".to_string(),
119 };
120
121 let mut field_parts = Vec::with_capacity(self.compiled_fields.len());
122 for (key, compiled_opt) in &self.compiled_fields {
123 let val = match compiled_opt {
124 Some(compiled) => stringify(compiled),
125 None => "<uncompiled>".to_string(),
126 };
127 field_parts.push(format!("{}={}", key, val));
128 }
129
130 let full_message = if field_parts.is_empty() {
131 log_message
132 } else {
133 format!("{} [{}]", log_message, field_parts.join(", "))
134 };
135
136 match self.level {
137 LogLevel::Trace => trace!(target: "dataflow::log", "{}", full_message),
138 LogLevel::Debug => debug!(target: "dataflow::log", "{}", full_message),
139 LogLevel::Info => info!(target: "dataflow::log", "{}", full_message),
140 LogLevel::Warn => warn!(target: "dataflow::log", "{}", full_message),
141 LogLevel::Error => error!(target: "dataflow::log", "{}", full_message),
142 }
143
144 Ok((TaskOutcome::Success, vec![]))
146 }
147}