Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
117 changes: 1 addition & 116 deletions crates/graphforge-knowledge/src/supersession.rs
Original file line number Diff line number Diff line change
Expand Up @@ -151,25 +151,13 @@ impl AssertionSupersessionLedger {
}

/// Merge append-only relations with exact replay semantics.
///
/// # Performance
/// `self` and `staged` are each only ever produced by [`Self::new`],
/// [`Self::merge`], or [`Self::from_batches`], which already validate
/// every relation they accept and already prove the relation set they
/// hold is acyclic. This function does not re-run [`validate_relation`]
/// or a whole-ledger cycle search over rows it already has in hand: it
/// validates only the relations genuinely new to this call, and extends
/// the already-proven-acyclic graph by checking each new edge against
/// reachability rather than re-deriving acyclicity for the whole ledger
/// from scratch.
pub fn merge(&self, staged: &Self) -> Result<Self, KnowledgeError> {
let mut relations = self.relations.clone();
let mut by_id = relations
.iter()
.cloned()
.map(|row| (row.supersession_uuid, row))
.collect::<HashMap<_, _>>();
let mut new_relations = Vec::new();
for relation in &staged.relations {
if let Some(existing) = by_id.get(&relation.supersession_uuid) {
if existing != relation {
Expand All @@ -178,63 +166,9 @@ impl AssertionSupersessionLedger {
} else {
relations.push(relation.clone());
by_id.insert(relation.supersession_uuid, relation.clone());
new_relations.push(relation.clone());
}
}

if new_relations.is_empty() {
// Fully idempotent merge: nothing new to validate or re-sort.
return Ok(self.clone());
}

if relations.len() > MAX_KNOWLEDGE_ROWS {
return Err(KnowledgeError::Limit {
participant: "assertion_supersessions",
observed: relations.len(),
limit: MAX_KNOWLEDGE_ROWS,
});
}

let mut status_ids: HashSet<Uuid> = self
.relations
.iter()
.map(|row| row.status_event_uuid)
.collect();
for relation in &new_relations {
validate_relation(relation)?;
if !status_ids.insert(relation.status_event_uuid) {
return Err(KnowledgeError::Duplicate("status_event_uuid"));
}
}

// `self.relations` is already known acyclic (proven at its own
// construction). A cycle can only appear through a path that uses at
// least one of the newly added edges, so extend the existing graph
// one new edge at a time, rejecting an edge that would let its
// replacement reach back to its own prior.
let mut adjacency: HashMap<Uuid, Vec<Uuid>> = HashMap::new();
for row in &self.relations {
adjacency
.entry(row.prior_assertion_uuid)
.or_default()
.push(row.replacement_assertion_uuid);
}
for relation in &new_relations {
if reaches(
&adjacency,
relation.replacement_assertion_uuid,
relation.prior_assertion_uuid,
) {
return Err(invalid("assertion_supersession", "cycle detected"));
}
adjacency
.entry(relation.prior_assertion_uuid)
.or_default()
.push(relation.replacement_assertion_uuid);
}

relations.sort_by_key(|row| (row.recorded_at_micros, row.supersession_uuid));
Ok(Self { relations })
Self::new(relations)
}

/// Canonical fingerprint over one exact immutable relation.
Expand Down Expand Up @@ -422,26 +356,6 @@ fn visit(
true
}

/// Depth-first reachability search: can `from` reach `to` via `adjacency`?
/// Used to check a single new edge against an already-proven-acyclic graph
/// without re-validating the whole graph's acyclicity.
fn reaches(adjacency: &HashMap<Uuid, Vec<Uuid>>, from: Uuid, to: Uuid) -> bool {
let mut stack = vec![from];
let mut seen = HashSet::new();
while let Some(node) = stack.pop() {
if node == to {
return true;
}
if !seen.insert(node) {
continue;
}
if let Some(next) = adjacency.get(&node) {
stack.extend(next.iter().copied());
}
}
false
}

fn validate_relation(row: &AssertionSupersession) -> Result<(), KnowledgeError> {
if row.contract_version != ASSERTION_SUPERSESSION_CONTRACT_VERSION {
return Err(invalid(
Expand Down Expand Up @@ -595,33 +509,4 @@ mod tests {
Err(KnowledgeError::Conflict("supersession_uuid"))
));
}

/// Merge's fast path stops re-running the whole-ledger cycle search over
/// rows already proven acyclic, but a new edge that closes a cycle
/// against the *existing* ledger (rather than within one `new()` call)
/// must still be refused.
#[test]
fn merge_rejects_a_cross_ledger_cycle() {
let base = AssertionSupersessionLedger::new(vec![relation(1, 10, 11, 1)]).unwrap();
let staged = AssertionSupersessionLedger::new(vec![relation(2, 11, 10, 2)]).unwrap();
assert!(matches!(
base.merge(&staged),
Err(KnowledgeError::Invalid { .. })
));
}

/// Same fast-path guarantee: a newly staged relation reusing a
/// status_event_uuid already claimed by an existing relation must still
/// be refused, even though the existing relation is not re-validated.
#[test]
fn merge_rejects_a_cross_ledger_duplicate_status_event() {
let base = AssertionSupersessionLedger::new(vec![relation(1, 10, 11, 1)]).unwrap();
let mut reused_status = relation(2, 12, 13, 2);
reused_status.status_event_uuid = base.relations()[0].status_event_uuid;
let staged = AssertionSupersessionLedger::new(vec![reused_status]).unwrap();
assert!(matches!(
base.merge(&staged),
Err(KnowledgeError::Duplicate("status_event_uuid"))
));
}
}
113 changes: 1 addition & 112 deletions crates/graphforge-provenance/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -446,22 +446,6 @@ impl ProvenanceLedger {
/// Identical operation event-sets are idempotent. Reuse of an operation
/// UUID with a different complete event-set, or reuse of an event, lineage,
/// or role/ordinal identity with different canonical content conflicts.
///
/// # Performance
/// `self` and `staged` are each only ever produced by [`Self::new`],
/// [`Self::merge`], or [`Self::from_batches`], every one of which already
/// validates and canonically re-derives (SHA-256-hashes) each row it
/// accepts. This function therefore does not re-derive or re-hash the
/// canonical identity of rows it already has in hand: `self`'s rows were
/// validated once, at decode or an earlier merge; `staged`'s rows were
/// validated once, moments earlier, at its own construction. What merge
/// still must check, without rehashing, are the invariants that only
/// appear once two independently-valid ledgers are combined: identity
/// reuse with conflicting content, dangling references, duplicate
/// role/ordinal positions, and the row-count limit. Cost here is O(rows
/// already present) for cheap, non-cryptographic bookkeeping (clone,
/// hash-map insert, set membership) plus O(rows newly added) for the one
/// piece of real validation work, not O(ledger size) of SHA-256 hashing.
pub fn merge(&self, staged: &Self) -> Result<Self, ProvenanceError> {
let mut events = self.events.clone();
let mut by_event = events
Expand All @@ -480,8 +464,6 @@ impl ProvenanceLedger {
return Err(ProvenanceError::Conflict("operation_uuid"));
}
}

let mut new_events = Vec::new();
for event in &staged.events {
if let Some(existing) = by_event.get(&event.provenance_uuid)
&& existing != event
Expand All @@ -493,7 +475,6 @@ impl ProvenanceLedger {
.is_none()
{
events.push(event.clone());
new_events.push(event.clone());
}
}

Expand All @@ -503,7 +484,6 @@ impl ProvenanceLedger {
.cloned()
.map(|row| (row.lineage_uuid, row))
.collect::<HashMap<_, _>>();
let mut new_lineage = Vec::new();
for row in &staged.lineage {
if let Some(existing) = by_lineage.get(&row.lineage_uuid)
&& existing != row
Expand All @@ -512,71 +492,9 @@ impl ProvenanceLedger {
}
if by_lineage.insert(row.lineage_uuid, row.clone()).is_none() {
lineage.push(row.clone());
new_lineage.push(row.clone());
}
}

if new_events.is_empty() && new_lineage.is_empty() {
// Fully idempotent merge: nothing new to validate, sort, or
// reconstruct.
return Ok(self.clone());
}

check_limit("events", events.len())?;
check_limit("lineage", lineage.len())?;

// Structural checks that only a merge can violate, over the rows
// genuinely new to this call. This does not call `ProvenanceEvent::new`
// / `LineageRecord::new` (the canonical-bytes + SHA-256 fingerprint +
// derived-UUID comparison) again for these rows: that already ran once,
// either when `staged` was constructed or when the ledger it came from
// was decoded. Each auxiliary index below is built only when there is
// a new row of the matching kind to check, so a single-record append
// that adds only an event (or only a lineage row) pays for one
// O(existing rows) pass, not both.
if !new_events.is_empty() {
let mut operation_kinds: HashSet<(Uuid, EventKind)> = self
.events
.iter()
.map(|event| (event.operation_uuid, event.event_kind))
.collect();
for event in &new_events {
if !operation_kinds.insert((event.operation_uuid, event.event_kind)) {
return Err(ProvenanceError::Duplicate("operation_uuid/event_kind"));
}
}
}

if !new_lineage.is_empty() {
let mut positions: HashSet<(Uuid, LineageRole, u32)> = self
.lineage
.iter()
.map(|row| (row.provenance_uuid, row.role, row.ordinal))
.collect();
for row in &new_lineage {
if !by_event.contains_key(&row.provenance_uuid) {
return Err(ProvenanceError::Dangling("provenance_uuid"));
}
if !positions.insert((row.provenance_uuid, row.role, row.ordinal)) {
return Err(ProvenanceError::Duplicate("role/ordinal"));
}
}
}

events.sort_by_key(|event| (event.recorded_at_micros, event.provenance_uuid));
// `by_event` already holds every final event (old and new: every
// staged event was inserted into it above), so lineage sort keys can
// be read from there instead of rebuilding a separate time index.
lineage.sort_by_key(|row| {
(
by_event[&row.provenance_uuid].recorded_at_micros,
row.provenance_uuid,
role_order(row.role),
row.ordinal,
row.subject_uuid,
)
});
Ok(Self { events, lineage })
Self::new(events, lineage)
}

/// Build the authoritative event Arrow batch.
Expand Down Expand Up @@ -1325,35 +1243,6 @@ mod tests {
assert_eq!(ledger.lineage.len(), 2);
}

/// Merge's fast path stops validating and re-hashing already-known-good
/// rows, but a staged batch that tries to add a lineage row occupying a
/// role/ordinal position an existing row already holds for the same
/// event must still be refused as a conflict, not silently accepted.
#[test]
fn merge_rejects_cross_ledger_duplicate_role_ordinal() {
let (event, rows) = fixture();
let ledger = ProvenanceLedger::new(vec![event.clone()], rows.clone()).unwrap();

// Reuses (event.provenance_uuid, Output, 0) — already occupied by
// rows[0] — with a different subject, so it is a distinct lineage
// row (different lineage_uuid) rather than an idempotent replay.
let colliding_position = LineageRecord::new(
event.provenance_uuid,
uuid(5),
SubjectKind::Edge,
LineageRole::Output,
0,
)
.unwrap();
let staged =
ProvenanceLedger::new(vec![event], vec![rows[1].clone(), colliding_position]).unwrap();

assert_eq!(
ledger.merge(&staged).unwrap_err().code(),
"GF_IDEMPOTENCY_CONFLICT"
);
}

#[test]
fn one_operation_can_record_distinct_composite_mutation_kinds() {
let operation_uuid = uuid(10);
Expand Down
Loading