1use std::collections::{HashMap, HashSet};
10
11use ocel::Ocel;
12use serde::Serialize;
13
14use crate::trace;
15
16#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
18#[serde(rename_all = "camelCase")]
19pub struct DfgNode {
20 pub activity: String,
21 pub events: usize,
23 pub objects: usize,
25 pub starts: usize,
27 pub ends: usize,
29}
30
31#[derive(Debug, Clone, PartialEq, Serialize)]
33#[serde(rename_all = "camelCase")]
34pub struct DfgEdge {
35 pub from: String,
36 pub to: String,
37 pub frequency: usize,
39 pub objects: usize,
41 pub median_secs: f64,
43 pub mean_secs: f64,
45}
46
47#[derive(Debug, Clone, PartialEq, Serialize)]
49#[serde(rename_all = "camelCase")]
50pub struct Dfg {
51 pub object_type: String,
52 pub objects: usize,
53 pub with_events: usize,
54 pub nodes: Vec<DfgNode>,
56 pub edges: Vec<DfgEdge>,
58}
59
60struct NodeAgg {
61 events: usize,
62 objects: usize,
63 last_slot: usize,
64 starts: usize,
65 ends: usize,
66}
67
68struct EdgeAgg {
69 frequency: usize,
70 objects: usize,
71 last_slot: usize,
72 gaps: Vec<i64>,
73}
74
75#[allow(clippy::cast_precision_loss)]
78fn gap_stats(gaps: &mut [i64]) -> (f64, f64) {
79 gaps.sort_unstable();
80 let n = gaps.len();
81 let median = if n % 2 == 1 {
82 gaps[n / 2] as f64
83 } else {
84 (gaps[n / 2 - 1] + gaps[n / 2]) as f64 / 2.0
85 };
86 let mean = gaps.iter().sum::<i64>() as f64 / n as f64;
87 (median, mean)
88}
89
90#[must_use]
92pub fn dfg(log: &Ocel, object_type: &str) -> Dfg {
93 let traces = trace::build(log, object_type);
94
95 let mut nodes: Vec<NodeAgg> = (0..traces.activity_names.len())
96 .map(|_| NodeAgg {
97 events: 0,
98 objects: 0,
99 last_slot: usize::MAX,
100 starts: 0,
101 ends: 0,
102 })
103 .collect();
104 let mut edges: HashMap<u32, EdgeAgg> = HashMap::new();
105
106 let mut with_events = 0usize;
107 for (slot, steps) in traces.steps.iter().enumerate() {
108 let (Some(&(first, _)), Some(&(last, _))) = (steps.first(), steps.last()) else {
109 continue;
110 };
111 with_events += 1;
112 nodes[first as usize].starts += 1;
113 nodes[last as usize].ends += 1;
114 for &(activity, _) in steps {
115 let node = &mut nodes[activity as usize];
116 node.events += 1;
117 if node.last_slot != slot {
118 node.last_slot = slot;
119 node.objects += 1;
120 }
121 }
122 for pair in steps.windows(2) {
123 let (from, from_time) = pair[0];
124 let (to, to_time) = pair[1];
125 let key = (u32::from(from) << 16) | u32::from(to);
126 let agg = edges.entry(key).or_insert_with(|| EdgeAgg {
127 frequency: 0,
128 objects: 0,
129 last_slot: usize::MAX,
130 gaps: Vec::new(),
131 });
132 agg.frequency += 1;
133 if agg.last_slot != slot {
134 agg.last_slot = slot;
135 agg.objects += 1;
136 }
137 agg.gaps.push((to_time - from_time).num_seconds());
138 }
139 }
140
141 let mut nodes: Vec<DfgNode> = nodes
142 .into_iter()
143 .enumerate()
144 .filter(|(_, agg)| agg.events > 0)
145 .map(|(id, agg)| DfgNode {
146 activity: traces.activity_names[id].to_owned(),
147 events: agg.events,
148 objects: agg.objects,
149 starts: agg.starts,
150 ends: agg.ends,
151 })
152 .collect();
153 nodes.sort_unstable_by(|a, b| {
154 b.events
155 .cmp(&a.events)
156 .then_with(|| a.activity.cmp(&b.activity))
157 });
158
159 let mut edges: Vec<DfgEdge> = edges
160 .into_iter()
161 .map(|(key, mut agg)| {
162 let (median_secs, mean_secs) = gap_stats(&mut agg.gaps);
163 DfgEdge {
164 from: traces.activity_names[(key >> 16) as usize].to_owned(),
165 to: traces.activity_names[(key & 0xffff) as usize].to_owned(),
166 frequency: agg.frequency,
167 objects: agg.objects,
168 median_secs,
169 mean_secs,
170 }
171 })
172 .collect();
173 edges.sort_unstable_by(|a, b| {
174 b.frequency
175 .cmp(&a.frequency)
176 .then_with(|| (a.from.as_str(), a.to.as_str()).cmp(&(b.from.as_str(), b.to.as_str())))
177 });
178
179 Dfg {
180 object_type: object_type.to_owned(),
181 objects: traces.object_ids.len(),
182 with_events,
183 nodes,
184 edges,
185 }
186}
187
188#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
190#[serde(rename_all = "camelCase")]
191pub struct OcTypeCount {
192 pub object_type: String,
193 pub events: usize,
194 pub objects: usize,
195 pub starts: usize,
196 pub ends: usize,
197}
198
199#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
201#[serde(rename_all = "camelCase")]
202pub struct OcActivity {
203 pub activity: String,
204 pub events: usize,
207 pub per_type: Vec<OcTypeCount>,
208}
209
210#[derive(Debug, Clone, PartialEq, Serialize)]
212#[serde(rename_all = "camelCase")]
213pub struct OcDfgEdge {
214 pub object_type: String,
215 #[serde(flatten)]
216 pub edge: DfgEdge,
217}
218
219#[derive(Debug, Clone, PartialEq, Serialize)]
221#[serde(rename_all = "camelCase")]
222pub struct OcDfg {
223 pub object_types: Vec<String>,
224 pub activities: Vec<OcActivity>,
226 pub edges: Vec<OcDfgEdge>,
228}
229
230#[must_use]
232pub fn oc_dfg(log: &Ocel, object_types: &[&str]) -> OcDfg {
233 let included: HashSet<&str> = log
234 .objects
235 .iter()
236 .filter(|o| object_types.contains(&o.object_type.as_str()))
237 .map(|o| o.id.as_str())
238 .collect();
239 let mut event_totals: HashMap<&str, usize> = HashMap::new();
240 for event in &log.events {
241 if event
242 .relationships
243 .iter()
244 .any(|r| included.contains(r.object_id.as_str()))
245 {
246 *event_totals.entry(event.event_type.as_str()).or_insert(0) += 1;
247 }
248 }
249
250 let mut activities: HashMap<String, Vec<OcTypeCount>> = HashMap::new();
251 let mut edges: Vec<OcDfgEdge> = Vec::new();
252 for &object_type in object_types {
253 let graph = dfg(log, object_type);
254 for node in graph.nodes {
255 activities
256 .entry(node.activity)
257 .or_default()
258 .push(OcTypeCount {
259 object_type: object_type.to_owned(),
260 events: node.events,
261 objects: node.objects,
262 starts: node.starts,
263 ends: node.ends,
264 });
265 }
266 edges.extend(graph.edges.into_iter().map(|edge| OcDfgEdge {
267 object_type: object_type.to_owned(),
268 edge,
269 }));
270 }
271
272 let mut activities: Vec<OcActivity> = activities
273 .into_iter()
274 .map(|(activity, per_type)| OcActivity {
275 events: event_totals.get(activity.as_str()).copied().unwrap_or(0),
276 activity,
277 per_type,
278 })
279 .collect();
280 activities.sort_unstable_by(|a, b| {
281 b.events
282 .cmp(&a.events)
283 .then_with(|| a.activity.cmp(&b.activity))
284 });
285 edges.sort_unstable_by_key(|e| std::cmp::Reverse(e.edge.frequency));
286
287 OcDfg {
288 object_types: object_types.iter().map(|&t| t.to_owned()).collect(),
289 activities,
290 edges,
291 }
292}
293
294#[cfg(test)]
295mod tests {
296 use super::*;
297 use chrono::{TimeZone, Utc};
298 use ocel::{Event, EventType, Object, ObjectType, Relationship};
299
300 fn rel(object_id: &str) -> Relationship {
301 Relationship {
302 object_id: object_id.into(),
303 qualifier: "q".into(),
304 }
305 }
306
307 fn event(id: &str, event_type: &str, minute: u32, objects: &[&str]) -> Event {
308 Event {
309 id: id.into(),
310 event_type: event_type.into(),
311 time: Utc.with_ymd_and_hms(2026, 1, 1, 9, minute, 0).unwrap(),
312 attributes: vec![],
313 relationships: objects.iter().map(|o| rel(o)).collect(),
314 }
315 }
316
317 fn log(events: Vec<Event>, objects: &[(&str, &str)]) -> Ocel {
318 let mut builder = Ocel::builder();
319 for name in ["created", "changed", "closed"] {
320 builder.add_event_type(EventType {
321 name: name.into(),
322 attributes: vec![],
323 });
324 }
325 for type_name in ["task", "user"] {
326 builder.add_object_type(ObjectType {
327 name: type_name.into(),
328 attributes: vec![],
329 });
330 }
331 for (id, type_name) in objects {
332 builder.add_object(Object {
333 id: (*id).into(),
334 object_type: (*type_name).into(),
335 attributes: vec![],
336 relationships: vec![],
337 });
338 }
339 for e in events {
340 builder.add_event(e);
341 }
342 builder.build().expect("valid log")
343 }
344
345 #[test]
346 #[allow(clippy::float_cmp)]
348 fn aggregates_edges_with_durations() {
349 let log = log(
350 vec![
351 event("e1", "created", 0, &["a"]),
352 event("e2", "changed", 1, &["a"]),
353 event("e3", "closed", 3, &["a"]),
354 event("e4", "created", 10, &["b"]),
355 event("e5", "changed", 13, &["b"]),
356 event("e6", "closed", 14, &["b"]),
357 event("e7", "created", 20, &["c"]),
358 event("e8", "closed", 21, &["c"]),
359 ],
360 &[("a", "task"), ("b", "task"), ("c", "task")],
361 );
362 let graph = dfg(&log, "task");
363 assert_eq!(graph.objects, 3);
364 assert_eq!(graph.with_events, 3);
365
366 let created = graph
367 .nodes
368 .iter()
369 .find(|n| n.activity == "created")
370 .unwrap();
371 assert_eq!((created.events, created.objects), (3, 3));
372 assert_eq!((created.starts, created.ends), (3, 0));
373 let closed = graph.nodes.iter().find(|n| n.activity == "closed").unwrap();
374 assert_eq!((closed.starts, closed.ends), (0, 3));
375
376 let created_changed = graph
377 .edges
378 .iter()
379 .find(|e| e.from == "created" && e.to == "changed")
380 .unwrap();
381 assert_eq!(created_changed.frequency, 2);
382 assert_eq!(created_changed.objects, 2);
383 assert_eq!(created_changed.median_secs, 120.0);
385 assert_eq!(created_changed.mean_secs, 120.0);
386
387 let created_closed = graph
388 .edges
389 .iter()
390 .find(|e| e.from == "created" && e.to == "closed")
391 .unwrap();
392 assert_eq!(created_closed.frequency, 1);
393 assert_eq!(created_closed.median_secs, 60.0);
394 }
395
396 #[test]
397 fn oc_dfg_counts_shared_events_once() {
398 let log = log(
399 vec![
400 event("e1", "created", 0, &["a", "u"]),
401 event("e2", "closed", 1, &["a", "u"]),
402 ],
403 &[("a", "task"), ("u", "user")],
404 );
405 let graph = oc_dfg(&log, &["task", "user"]);
406 let created = graph
407 .activities
408 .iter()
409 .find(|a| a.activity == "created")
410 .unwrap();
411 assert_eq!(created.events, 1);
413 assert_eq!(created.per_type.len(), 2);
414 assert_eq!(graph.edges.len(), 2);
416 assert!(graph
417 .edges
418 .iter()
419 .all(|e| e.edge.from == "created" && e.edge.to == "closed"));
420 }
421
422 #[test]
423 fn empty_type_yields_empty_graph() {
424 let log = log(vec![], &[]);
425 let graph = dfg(&log, "task");
426 assert_eq!(graph.objects, 0);
427 assert!(graph.nodes.is_empty());
428 assert!(graph.edges.is_empty());
429 }
430}