use dataflow_rs::{ExecutionObserver, TaskEvent, WorkflowFinished};
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 workflow_finished(&self, event: &WorkflowFinished<'_>) {
crate::metrics::record_workflow_duration(event.workflow_id, 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"
);
}
#[tokio::test]
async fn a_workflow_reports_once_and_only_when_it_runs() {
use std::sync::Mutex;
use std::time::Duration;
#[derive(Default)]
struct Seen {
workflows: Mutex<Vec<(String, Duration)>>,
tasks: Mutex<Duration>,
}
impl ExecutionObserver for Seen {
fn task_finished(&self, e: &TaskEvent<'_>) {
*self.tasks.lock().expect("test") += e.duration;
}
fn workflow_finished(&self, e: &WorkflowFinished<'_>) {
self.workflows
.lock()
.expect("test")
.push((e.workflow_id.to_string(), e.duration));
}
}
fn workflow(id: &str, condition: bool) -> dataflow_rs::Workflow {
dataflow_rs::Workflow::from_json(&format!(
r#"{{"id":"{id}","name":"{id}","priority":0,"condition":{condition},
"tasks":[{{"id":"m","name":"m","function":{{"name":"map","input":
{{"mappings":[{{"path":"data.ok","logic":true}}]}}}}}}]}}"#
))
.expect("workflow parses")
}
let seen = std::sync::Arc::new(Seen::default());
let engine = dataflow_rs::Engine::new(
vec![workflow("runs", true), workflow("skipped", false)],
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");
let workflows = seen.workflows.lock().expect("test");
assert_eq!(
workflows
.iter()
.map(|(id, _)| id.as_str())
.collect::<Vec<_>>(),
["runs"],
"a workflow its condition rejected never starts, so it never finishes"
);
assert!(
workflows[0].1 >= *seen.tasks.lock().expect("test"),
"the workflow's own clock must cover the task bodies inside it, or \
subtracting them for the overhead figure would go negative"
);
}
}