diff --git a/process_mining/src/bindings/slim_ocel_bindings.rs b/process_mining/src/bindings/slim_ocel_bindings.rs index a1fb03a..52ca26e 100644 --- a/process_mining/src/bindings/slim_ocel_bindings.rs +++ b/process_mining/src/bindings/slim_ocel_bindings.rs @@ -1,7 +1,10 @@ //! Binding wrappers for [`SlimLinkedOCEL`] functionality +use std::collections::HashMap; + use chrono::{DateTime, FixedOffset}; use macros_process_mining::register_binding; +use rayon::prelude::*; use crate::core::event_data::object_centric::{ linked_ocel::{ @@ -244,6 +247,91 @@ fn get_obj_activity_trace(ocel: &SlimLinkedOCEL, ob: ObjectIndex) -> Vec .collect() } +/// Merge `b` into `a` by summing counts for matching keys. Used as the rayon reduce step. +fn merge_count_maps( + mut a: HashMap, + b: HashMap, +) -> HashMap { + for (k, v) in b { + *a.entry(k).or_insert(0) += v; + } + a +} + +/// Get all activity-trace variants for objects of the given object type, with their occurrence counts +/// +/// Each entry is a tuple `(activity_trace, count)`, where `activity_trace` is the sequence of event types +/// connected to an object (ordered by event timestamp), and `count` is the number of objects of the +/// requested type that share that exact trace. +#[register_binding] +fn get_variants_of_object_type( + ocel: &SlimLinkedOCEL, + ob_type: String, +) -> Vec<(Vec, usize)> { + let obs: Vec = ocel.get_obs_of_type(&ob_type).copied().collect(); + let counts: HashMap, usize> = obs + .into_par_iter() + .fold(HashMap::new, |mut acc, ob| { + let trace: Vec = ob.get_obj_activity_trace_evtype_indices(ocel).collect(); + *acc.entry(trace).or_insert(0) += 1; + acc + }) + .reduce(HashMap::new, merge_count_maps); + let ev_type_names: Vec<&str> = + ::get_ev_types(ocel).collect(); + let result: Vec<(Vec, usize)> = counts + .into_iter() + .map(|(trace_idx, count)| { + let trace: Vec = trace_idx + .into_iter() + .map(|i| ev_type_names[i].to_string()) + .collect(); + (trace, count) + }) + .collect(); + // result.sort_by(|a, b| b.1.cmp(&a.1).then_with(|| a.0.cmp(&b.0))); + result +} + +/// Get the directly-follows graph (DFG) for objects of the given object type. +/// +/// Each entry is `((from_activity, to_activity), count)`, counting adjacent pairs in each +/// object's timestamp-ordered activity trace. Result order is unspecified. +#[register_binding] +fn get_dfg_of_object_type( + ocel: &SlimLinkedOCEL, + ob_type: String, +) -> Vec<((String, String), usize)> { + let obs: Vec = ocel.get_obs_of_type(&ob_type).copied().collect(); + let counts: HashMap<(usize, usize), usize> = obs + .into_par_iter() + .fold(HashMap::new, |mut acc, ob| { + let mut iter = ob.get_obj_activity_trace_evtype_indices(ocel); + if let Some(mut prev) = iter.next() { + for next in iter { + *acc.entry((prev, next)).or_insert(0) += 1; + prev = next; + } + } + acc + }) + .reduce(HashMap::new, merge_count_maps); + let ev_type_names: Vec<&str> = + ::get_ev_types(ocel).collect(); + counts + .into_iter() + .map(|((from, to), count)| { + ( + ( + ev_type_names[from].to_string(), + ev_type_names[to].to_string(), + ), + count, + ) + }) + .collect() +} + /// Get the outgoing O2O relationships of an object as `(qualifier, object_index)` pairs. #[register_binding] fn locel_get_o2o(ocel: &SlimLinkedOCEL, ob: ObjectIndex) -> Vec<(String, ObjectIndex)> { @@ -368,3 +456,215 @@ fn get_event_timestamp_of_id(ocel: &SlimLinkedOCEL, ev_id: &String) -> Option Vec<(String, String, i64)> { + let num_events = ocel.get_num_evs() as u32; + let counts: HashMap<(usize, usize), i64> = (0..num_events) + .into_par_iter() + .fold(HashMap::new, |mut acc, i| { + let ev = EventIndex::from(i).get_ev(ocel); + for (_q, ob) in &ev.relationships { + let ot = ob.get_ob(ocel).object_type; + *acc.entry((ev.event_type, ot)).or_insert(0) += 1; + } + acc + }) + .reduce(HashMap::new, merge_sum_maps); + let ev_types: Vec<&str> = ::get_ev_types(ocel).collect(); + let ob_types: Vec<&str> = ::get_ob_types(ocel).collect(); + counts + .into_iter() + .map(|((e, o), c)| (ev_types[e].to_string(), ob_types[o].to_string(), c)) + .collect() +} + +/// Conversion rate from `source_type` to `target_type` via O2O, restricted to targets touched by `activity`. +/// +/// Returns the fraction of `source_type` objects that have at least one outgoing O2O edge to a +/// `target_type` object related (via E2O) to some event of the given event type. Returns `0.0` +/// if no `source_type` objects exist. +#[register_binding] +fn locel_conversion_rate( + ocel: &SlimLinkedOCEL, + activity: String, + source_type: String, + target_type: String, +) -> f64 { + let sources: Vec = ocel.get_obs_of_type(&source_type).copied().collect(); + let total = sources.len(); + if total == 0 { + return 0.0; + } + let reached = sources + .par_iter() + .filter(|&&s| { + s.get_o2o(ocel).any(|&t| { + t.get_ob_type(ocel) == &target_type + && t.get_e2o_rev(ocel) + .any(|&e| e.get_ev_type(ocel) == &activity) + }) + }) + .count(); + reached as f64 / total as f64 +} + +/// Each object's reverse-E2O events in `(time, id)` order. +/// +/// Events are sorted per call, so this makes no assumption about global event ordering +fn sorted_events_per_object(ocel: &SlimLinkedOCEL) -> Vec> { + (0..ocel.get_num_obs() as u32) + .into_par_iter() + .map(|i| { + let mut evs: Vec = + ObjectIndex::from(i).get_e2o_rev(ocel).copied().collect(); + evs.sort_by(|a, b| { + a.get_time(ocel) + .cmp(b.get_time(ocel)) + .then_with(|| ocel.get_ev_id(a).cmp(ocel.get_ev_id(b))) + }); + evs + }) + .collect() +} + +/// The `(time, id)`-immediate predecessor of `e` on object `o`, using the +/// per-object sorted lists from [`sorted_events_per_object`]. +/// `None` if `e` is the first event on `o`. +#[inline] +fn df_predecessor( + sorted: &[Vec], + ocel: &SlimLinkedOCEL, + e: EventIndex, + o: ObjectIndex, +) -> Option { + let evs = &sorted[o.into_inner() as usize]; + let key = (e.get_time(ocel), ocel.get_ev_id(&e)); + let pos = evs + .binary_search_by(|x| (x.get_time(ocel), ocel.get_ev_id(x)).cmp(&key)) + .ok()?; + pos.checked_sub(1).map(|p| evs[p]) +} + +/// Per-event synchronization time and the delaying object. +/// +/// For each event with at least one directly-follows predecessor, the synchronization time is +/// `max_predecessor_time - min_predecessor_time` in integer microseconds (the span between its +/// earliest and latest directly-preceding event). The delaying object is the object linking the +/// latest predecessor (ties broken by ascending object id). +/// Returns one row `(event_id, sync_us, delaying_object_id)` per qualifying event. +/// +/// `top_k`: if `Some(k)`, return only the `k` rows with the largest `sync_us`, ties broken by +/// ascending event id, sorted descending. `None` returns every qualifying event. +#[register_binding] +fn locel_oc_perf_sync_per_event( + ocel: &SlimLinkedOCEL, + #[bind(default)] top_k: Option, +) -> Vec<(String, i64, String)> { + let sorted = sorted_events_per_object(ocel); + let mut rows: Vec<(EventIndex, i64, ObjectIndex)> = (0..ocel.get_num_evs() as u32) + .into_par_iter() + .filter_map(|i| { + let e = EventIndex::from(i); + let mut min_us = i64::MAX; + // (latest predecessor time, its object) = the delaying edge. + let mut delaying: Option<(i64, ObjectIndex)> = None; + for &o in e.get_e2o(ocel) { + if let Some(p) = df_predecessor(&sorted, ocel, e, o) { + let t = p.get_time(ocel).timestamp_micros(); + min_us = min_us.min(t); + let keep = match delaying { + Some((bt, bo)) => { + bt > t || (bt == t && ocel.get_ob_id(&bo) <= ocel.get_ob_id(&o)) + } + None => false, + }; + if !keep { + delaying = Some((t, o)); + } + } + } + delaying.map(|(max_us, o)| (e, max_us - min_us, o)) + }) + .collect(); + if let Some(k) = top_k { + let cmp = |a: &(EventIndex, i64, ObjectIndex), b: &(EventIndex, i64, ObjectIndex)| { + b.1.cmp(&a.1) + .then_with(|| ocel.get_ev_id(&a.0).cmp(ocel.get_ev_id(&b.0))) + }; + if k < rows.len() { + rows.select_nth_unstable_by(k, cmp); + rows.truncate(k); + } + rows.sort_unstable_by(cmp); + } + rows.into_iter() + .map(|(e, max_minus_min, o)| { + ( + ocel.get_ev_id(&e).to_string(), + max_minus_min, + ocel.get_ob_id(&o).to_string(), + ) + }) + .collect() +} + +/// Per-event sojourn time. +/// +/// For each event with at least one directly-follows predecessor, the sojourn time is +/// `event_time - latest_predecessor_time` in integer microseconds. Returns one row +/// `(event_id, sojourn_us)` per qualifying event. +/// +/// `top_k`: if `Some(k)`, return only the `k` rows with the largest `sojourn_us`, ties broken by +/// ascending event id, sorted descending. `None` returns every qualifying event. +#[register_binding] +fn locel_oc_perf_sojourn_per_event( + ocel: &SlimLinkedOCEL, + #[bind(default)] top_k: Option, +) -> Vec<(String, i64)> { + let sorted = sorted_events_per_object(ocel); + let mut rows: Vec<(EventIndex, i64)> = (0..ocel.get_num_evs() as u32) + .into_par_iter() + .filter_map(|i| { + let e = EventIndex::from(i); + let latest = e + .get_e2o(ocel) + .filter_map(|&o| df_predecessor(&sorted, ocel, e, o)) + .map(|p| p.get_time(ocel).timestamp_micros()) + .max()?; + Some((e, e.get_time(ocel).timestamp_micros() - latest)) + }) + .collect(); + if let Some(k) = top_k { + let cmp = |a: &(EventIndex, i64), b: &(EventIndex, i64)| { + b.1.cmp(&a.1) + .then_with(|| ocel.get_ev_id(&a.0).cmp(ocel.get_ev_id(&b.0))) + }; + if k < rows.len() { + rows.select_nth_unstable_by(k, cmp); + rows.truncate(k); + } + rows.sort_unstable_by(cmp); + } + rows.into_iter() + .map(|(e, sojourn_us)| (ocel.get_ev_id(&e).to_string(), sojourn_us)) + .collect() +} + +/// Merge `b` into `a` by summing `i64` counts for matching keys. Used as the rayon reduce step. +fn merge_sum_maps( + mut a: HashMap, + b: HashMap, +) -> HashMap { + for (k, v) in b { + *a.entry(k).or_insert(0) += v; + } + a +}