MCPcopy Create free account
hub / github.com/atomicdotdev/atomic / canonicalize_provenance_journal

Function canonicalize_provenance_journal

atomic-agent/src/event.rs:737–819  ·  view source on GitHub ↗

Deduplicate and causally order a complete turn journal. Ready events are ordered by generation, timestamp, then stable event ID, so process arrival order cannot affect replay. Causal parents always precede children even when clocks are tied or skewed.

(
    envelopes: Vec<ProvenanceJournalEnvelope>,
)

Source from the content-addressed store, hash-verified

735/// process arrival order cannot affect replay. Causal parents always precede
736/// children even when clocks are tied or skewed.
737pub fn canonicalize_provenance_journal(
738 envelopes: Vec<ProvenanceJournalEnvelope>,
739) -> Result<Vec<ProvenanceJournalEnvelope>, ProvenanceJournalError> {
740 use std::collections::{BTreeMap, BTreeSet};
741
742 let mut by_id = BTreeMap::<String, ProvenanceJournalEnvelope>::new();
743 for mut envelope in envelopes {
744 envelope.validate()?;
745 envelope.causal_parent_ids.sort();
746 envelope.causal_parent_ids.dedup();
747 match by_id.get(&envelope.event_id) {
748 Some(existing) if existing == &envelope => continue,
749 Some(_) => {
750 return Err(ProvenanceJournalError::ConflictingEventId(
751 envelope.event_id,
752 ))
753 }
754 None => {
755 by_id.insert(envelope.event_id.clone(), envelope);
756 }
757 }
758 }
759 let Some(first) = by_id.values().next() else {
760 return Ok(Vec::new());
761 };
762 if by_id
763 .values()
764 .any(|event| event.session_id != first.session_id || event.turn_number != first.turn_number)
765 {
766 return Err(ProvenanceJournalError::MixedTurnIdentity);
767 }
768
769 let mut indegree = BTreeMap::<String, usize>::new();
770 let mut children = BTreeMap::<String, Vec<String>>::new();
771 for envelope in by_id.values() {
772 indegree.insert(envelope.event_id.clone(), envelope.causal_parent_ids.len());
773 for parent in &envelope.causal_parent_ids {
774 if !by_id.contains_key(parent) {
775 return Err(ProvenanceJournalError::MissingCausalParent {
776 event_id: envelope.event_id.clone(),
777 parent_id: parent.clone(),
778 });
779 }
780 children
781 .entry(parent.clone())
782 .or_default()
783 .push(envelope.event_id.clone());
784 }
785 }
786
787 let mut ready = BTreeSet::<(u64, i64, String)>::new();
788 for envelope in by_id.values() {
789 if indegree[&envelope.event_id] == 0 {
790 ready.insert((
791 envelope.generation,
792 envelope.timestamp_ms,
793 envelope.event_id.clone(),
794 ));

Callers 1

Calls 10

sortMethod · 0.80
contains_keyMethod · 0.80
get_mutMethod · 0.80
getMethod · 0.65
validateMethod · 0.45
insertMethod · 0.45
cloneMethod · 0.45
nextMethod · 0.45
lenMethod · 0.45
pushMethod · 0.45

Tested by

no test coverage detected