use std::collections::{BTreeMap, HashMap};
use chrono::{DateTime, Utc};
use ocel::{
AttrType, AttrValue, AttributeDefinition, Event, EventAttribute, EventType, Object,
ObjectAttribute, ObjectType, Ocel, Relationship, Violation,
};
#[derive(Debug, Clone)]
pub struct StagingEvent {
pub id: String,
pub event_type: String,
pub time: DateTime<Utc>,
pub attributes: Vec<(String, AttrValue)>,
pub relations: Vec<(String, String)>,
}
#[derive(Debug, Clone, Default)]
struct StagingObject {
object_type: String,
attributes: Vec<(String, AttrValue, DateTime<Utc>)>,
relations: Vec<(String, String)>,
}
#[derive(Debug, Default)]
pub struct StagingLog {
events: Vec<StagingEvent>,
object_index: HashMap<String, usize>,
object_ids: Vec<String>,
objects: Vec<StagingObject>,
event_schema: BTreeMap<String, BTreeMap<String, AttrType>>,
object_schema: BTreeMap<String, BTreeMap<String, AttrType>>,
}
fn observe(schema: &mut BTreeMap<String, AttrType>, name: &str, value: &AttrValue) {
let observed = value.attr_type();
match schema.get(name) {
None => {
schema.insert(name.to_owned(), observed);
}
Some(&declared) if declared != observed => {
schema.insert(name.to_owned(), AttrType::String);
}
Some(_) => {}
}
}
fn conform(value: AttrValue, declared: AttrType) -> AttrValue {
if value.attr_type() == declared {
value
} else {
AttrValue::String(value.to_text())
}
}
impl StagingLog {
#[must_use]
pub fn new() -> Self {
Self::default()
}
#[must_use]
pub fn from_ocel(ocel: Ocel) -> Self {
let mut staging = Self::new();
for event_type in &ocel.event_types {
staging
.event_schema
.insert(event_type.name.clone(), seed(&event_type.attributes));
}
for object_type in &ocel.object_types {
staging
.object_schema
.insert(object_type.name.clone(), seed(&object_type.attributes));
}
for object in ocel.objects {
staging.upsert_object(&object.id, &object.object_type);
for attr in object.attributes {
staging.add_object_attribute(&object.id, &attr.name, attr.value, attr.time);
}
for rel in object.relationships {
staging.add_o2o(&object.id, &rel.object_id, &rel.qualifier);
}
}
for event in ocel.events {
staging.add_event(StagingEvent {
id: event.id,
event_type: event.event_type,
time: event.time,
attributes: event
.attributes
.into_iter()
.map(|a| (a.name, a.value))
.collect(),
relations: event
.relationships
.into_iter()
.map(|r| (r.object_id, r.qualifier))
.collect(),
});
}
staging
}
pub fn add_event(&mut self, event: StagingEvent) {
let schema = self
.event_schema
.entry(event.event_type.clone())
.or_default();
for (name, value) in &event.attributes {
observe(schema, name, value);
}
self.events.push(event);
}
pub fn upsert_object(&mut self, id: &str, object_type: &str) {
let index = self.object_slot(id);
if self.objects[index].object_type.is_empty() {
object_type.clone_into(&mut self.objects[index].object_type);
}
self.object_schema
.entry(object_type.to_owned())
.or_default();
}
pub fn add_object_attribute(
&mut self,
id: &str,
name: &str,
value: AttrValue,
time: DateTime<Utc>,
) {
let index = self.object_slot(id);
let object_type = self.objects[index].object_type.clone();
if !object_type.is_empty() {
let schema = self.object_schema.entry(object_type).or_default();
observe(schema, name, &value);
}
self.objects[index]
.attributes
.push((name.to_owned(), value, time));
}
pub fn add_o2o(&mut self, source_id: &str, target_id: &str, qualifier: &str) {
let index = self.object_slot(source_id);
self.objects[index]
.relations
.push((target_id.to_owned(), qualifier.to_owned()));
}
pub fn map_object_ids(&mut self, f: impl Fn(&str) -> String) {
let old_ids = std::mem::take(&mut self.object_ids);
let old_objects = std::mem::take(&mut self.objects);
self.object_index.clear();
for (old_id, mut object) in old_ids.into_iter().zip(old_objects) {
for (target, _) in &mut object.relations {
*target = f(target);
}
let index = self.object_slot(&f(&old_id));
let slot = &mut self.objects[index];
if slot.object_type.is_empty() {
slot.object_type = object.object_type;
}
slot.attributes.append(&mut object.attributes);
slot.relations.append(&mut object.relations);
}
for event in &mut self.events {
for (target, _) in &mut event.relations {
*target = f(target);
}
}
}
pub fn map_events(&mut self, mut f: impl FnMut(&mut StagingEvent)) {
for event in &mut self.events {
f(event);
}
self.rebuild_event_schema();
}
pub fn retain_events(&mut self, mut pred: impl FnMut(&StagingEvent) -> bool) {
self.events.retain(|event| pred(event));
self.rebuild_event_schema();
}
fn rebuild_event_schema(&mut self) {
self.event_schema.clear();
for event in &self.events {
let schema = self
.event_schema
.entry(event.event_type.clone())
.or_default();
for (name, value) in &event.attributes {
observe(schema, name, value);
}
}
}
fn object_slot(&mut self, id: &str) -> usize {
if let Some(&index) = self.object_index.get(id) {
return index;
}
let index = self.objects.len();
self.object_index.insert(id.to_owned(), index);
self.object_ids.push(id.to_owned());
self.objects.push(StagingObject::default());
index
}
pub fn into_ocel(self) -> Result<Ocel, Vec<Violation>> {
let mut builder = Ocel::builder();
for (name, attrs) in &self.event_schema {
builder.add_event_type(EventType {
name: name.clone(),
attributes: attr_defs(attrs),
});
}
for (name, attrs) in &self.object_schema {
builder.add_object_type(ObjectType {
name: name.clone(),
attributes: attr_defs(attrs),
});
}
for event in self.events {
let schema = self.event_schema.get(&event.event_type);
let attributes = event
.attributes
.into_iter()
.map(|(name, value)| {
let declared = schema
.and_then(|s| s.get(&name).copied())
.unwrap_or_else(|| value.attr_type());
EventAttribute {
value: conform(value, declared),
name,
}
})
.collect();
builder.add_event(Event {
id: event.id,
event_type: event.event_type,
time: event.time,
attributes,
relationships: relationships(event.relations),
});
}
for (id, object) in self.object_ids.into_iter().zip(self.objects) {
let schema = self.object_schema.get(&object.object_type);
let mut attributes: Vec<ObjectAttribute> = object
.attributes
.into_iter()
.map(|(name, value, time)| {
let declared = schema
.and_then(|s| s.get(&name).copied())
.unwrap_or_else(|| value.attr_type());
ObjectAttribute {
value: conform(value, declared),
name,
time,
}
})
.collect();
let mut seen: std::collections::HashSet<(String, String, DateTime<Utc>)> =
std::collections::HashSet::new();
attributes.retain(|a| seen.insert((a.name.clone(), a.value.to_text(), a.time)));
attributes.sort_by(|a, b| a.time.cmp(&b.time).then_with(|| a.name.cmp(&b.name)));
let mut relations = object.relations;
relations.sort();
relations.dedup();
builder.add_object(Object {
id,
object_type: object.object_type,
attributes,
relationships: relationships(relations),
});
}
builder.build()
}
}
fn seed(attributes: &[AttributeDefinition]) -> BTreeMap<String, AttrType> {
attributes
.iter()
.map(|a| (a.name.clone(), a.value_type))
.collect()
}
fn attr_defs(attrs: &BTreeMap<String, AttrType>) -> Vec<AttributeDefinition> {
attrs
.iter()
.map(|(name, &value_type)| AttributeDefinition {
name: name.clone(),
value_type,
})
.collect()
}
fn relationships(pairs: Vec<(String, String)>) -> Vec<Relationship> {
pairs
.into_iter()
.map(|(object_id, qualifier)| Relationship {
object_id,
qualifier,
})
.collect()
}