Skip to content
Closed
Changes from all commits
Commits
Show all changes
19 commits
Select commit Hold shift + click to select a range
6316d90
Prototype of OC-DECLARE activity projection
aarkue Apr 11, 2026
9d74580
Merge branch 'main' into feat/oc-declare-act-projection
aarkue Apr 21, 2026
19fa92a
Merge branch 'main' into feat/oc-declare-act-projection
aarkue Apr 22, 2026
e50097e
Merge branch 'main' into feat/oc-declare-act-projection
aarkue Apr 22, 2026
7a8e7c2
Initial version of OCEL Appendable/Readable trait for streaming-like …
aarkue May 3, 2026
d2d0822
More exposed binding functions + SlimLinkedOCEL Tweaks
aarkue May 14, 2026
b299163
Add qualifier-optional relation getters
aarkue May 15, 2026
ec893fa
Fix SQL-based export (timezone handling + float precision)
aarkue May 22, 2026
7effd68
Additional r4pm bindings for KPIs/OC-Performance Analysis
aarkue May 23, 2026
0b132a6
Update bindings (non-aggregated performance analysis)
aarkue Jun 9, 2026
35b8845
Fix double serialization for bindings (instead stick with just JSON Vec)
aarkue Jun 10, 2026
51e6695
Update oc-perf functions to return top-k durations
aarkue Jun 10, 2026
6db786f
Merge branch 'feat/ocel-import-export-traits-improved-memory' into fe…
aarkue Jun 11, 2026
15fa19e
Add OCED-paths functionality
aarkue Jun 17, 2026
8a4c366
Merge branch 'main' into feat/oced-paths
aarkue Jun 17, 2026
f232f36
Merge remote-tracking branch 'origin/main' into feat/oced-paths
aarkue Jun 17, 2026
db98e40
Fix clippy (+MSRV 1.88 merge)
aarkue Jun 17, 2026
d34bbb4
Merge branch 'main' into feat/oced-paths
aarkue Jul 7, 2026
399389c
Merge branch 'main' into feat/oced-paths
aarkue Aug 16, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
300 changes: 300 additions & 0 deletions process_mining/src/bindings/slim_ocel_bindings.rs
Original file line number Diff line number Diff line change
@@ -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::{
Expand Down Expand Up @@ -244,6 +247,91 @@ fn get_obj_activity_trace(ocel: &SlimLinkedOCEL, ob: ObjectIndex) -> Vec<String>
.collect()
}

/// Merge `b` into `a` by summing counts for matching keys. Used as the rayon reduce step.
fn merge_count_maps<K: std::hash::Hash + Eq>(
mut a: HashMap<K, usize>,
b: HashMap<K, usize>,
) -> HashMap<K, usize> {
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<String>, usize)> {
let obs: Vec<ObjectIndex> = ocel.get_obs_of_type(&ob_type).copied().collect();
let counts: HashMap<Vec<usize>, usize> = obs
.into_par_iter()
.fold(HashMap::new, |mut acc, ob| {
let trace: Vec<usize> = 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> =
<SlimLinkedOCEL as LinkedOCELAccess>::get_ev_types(ocel).collect();
let result: Vec<(Vec<String>, usize)> = counts
.into_iter()
.map(|(trace_idx, count)| {
let trace: Vec<String> = 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<ObjectIndex> = 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> =
<SlimLinkedOCEL as LinkedOCELAccess>::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)> {
Expand Down Expand Up @@ -368,3 +456,215 @@ fn get_event_timestamp_of_id(ocel: &SlimLinkedOCEL, ev_id: &String) -> Option<St
ocel.get_ev_by_id(ev_id)
.map(|ev| ev.get_time(ocel).to_string())
}

// ── Analytics ─────────────────────────────────────────────────────────

/// Count E2O relationships per `(event_type, object_type)` pair.
///
/// Each entry is `(event_type, object_type, count)`. A single `(event, object)` pair connected
/// by multiple qualifiers contributes once per qualifier. Pairs with zero relations are omitted.
/// Result row order is unspecified.
#[register_binding]
fn locel_event_object_type_counts(ocel: &SlimLinkedOCEL) -> 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> = <SlimLinkedOCEL as LinkedOCELAccess>::get_ev_types(ocel).collect();
let ob_types: Vec<&str> = <SlimLinkedOCEL as LinkedOCELAccess>::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<ObjectIndex> = 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<Vec<EventIndex>> {
(0..ocel.get_num_obs() as u32)
.into_par_iter()
.map(|i| {
let mut evs: Vec<EventIndex> =
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<EventIndex>],
ocel: &SlimLinkedOCEL,
e: EventIndex,
o: ObjectIndex,
) -> Option<EventIndex> {
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<usize>,
) -> 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<usize>,
) -> 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<K: std::hash::Hash + Eq>(
mut a: HashMap<K, i64>,
b: HashMap<K, i64>,
) -> HashMap<K, i64> {
for (k, v) in b {
*a.entry(k).or_insert(0) += v;
}
a
}
Loading