use dataflow_rs::{ExecutionObserver, TaskEvent};
pub struct MetricsObserver;
impl ExecutionObserver for MetricsObserver {
fn task_finished(&self, event: &TaskEvent<'_>) {
crate::metrics::record_task_duration(
event.workflow_id,
event.task_id,
interned_function_name(event.function),
event.duration.as_secs_f64(),
);
}
}
fn interned_function_name(name: &str) -> &'static str {
super::known_functions()
.find(|known| *known == name)
.unwrap_or("other")
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn every_known_function_interns_to_itself() {
for name in super::super::known_functions() {
assert_eq!(interned_function_name(name), name);
}
}
#[test]
fn an_unknown_function_collapses_to_one_label() {
assert_eq!(interned_function_name("enrich"), "other");
assert_eq!(interned_function_name("__nope__"), "other");
}
#[test]
fn the_validation_alias_the_executor_reports_is_known() {
assert_eq!(interned_function_name("validate"), "validate");
}
#[tokio::test]
async fn a_sync_builtin_is_reported_to_the_observer() {
use std::sync::Mutex;
#[derive(Default)]
struct Seen(Mutex<Vec<(String, String)>>);
impl ExecutionObserver for Seen {
fn task_finished(&self, e: &TaskEvent<'_>) {
self.0
.lock()
.expect("test")
.push((e.task_id.to_string(), e.function.to_string()));
}
}
let seen = std::sync::Arc::new(Seen::default());
let workflow = dataflow_rs::Workflow::from_json(
r#"{"id":"w","name":"w","priority":0,"condition":true,
"tasks":[{"id":"m","name":"m","function":{"name":"map","input":
{"mappings":[{"path":"data.ok","logic":true}]}}}]}"#,
)
.expect("workflow parses");
let engine = dataflow_rs::Engine::new(vec![workflow], std::collections::HashMap::new())
.expect("engine builds")
.with_observer(seen.clone());
let mut message = dataflow_rs::Message::builder().build();
engine
.process_message(&mut message)
.await
.expect("run succeeds");
assert_eq!(
*seen.0.lock().expect("test"),
vec![("m".to_string(), "map".to_string())],
"a built-in dispatched inside the executor must still be timed"
);
}
}