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>,
)
| 735 | /// process arrival order cannot affect replay. Causal parents always precede |
| 736 | /// children even when clocks are tied or skewed. |
| 737 | pub 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 | )); |