1#![forbid(unsafe_code)]
8
9use std::collections::{HashMap, HashSet};
10use std::sync::{Arc, LazyLock};
11
12use arrow::array::{
13 Array, FixedSizeBinaryArray, FixedSizeBinaryBuilder, StringArray, TimestampMicrosecondArray,
14 UInt32Array,
15};
16use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
17use arrow::record_batch::RecordBatch;
18use graphforge_core::canonical::{
19 CANONICAL_CONTRACT_VERSION, CanonicalDomain, CanonicalError, CanonicalWriter, fingerprint,
20 uuid_v8,
21};
22use uuid::Uuid;
23
24pub const PROVENANCE_CAPABILITY_VERSION: u32 = 1;
26pub const PROVENANCE_EVENT_CONTRACT_VERSION: u32 = 1;
28pub const LINEAGE_CONTRACT_VERSION: u32 = 1;
30pub const EVENT_KIND_REGISTRY_VERSION: u32 = 5;
32pub const SUBJECT_KIND_REGISTRY_VERSION: u32 = 1;
34pub const LINEAGE_ROLE_REGISTRY_VERSION: u32 = 1;
36pub const MAX_PROVENANCE_ROWS: usize = 1_000_000;
38
39pub static PROVENANCE_EVENT_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
41 Arc::new(Schema::new(vec![
42 uuid_field("provenance_uuid", false),
43 uuid_field("operation_uuid", false),
44 Field::new("event_kind", DataType::Utf8, false),
45 uuid_field("actor_uuid", true),
46 Field::new(
47 "recorded_at",
48 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
49 false,
50 ),
51 Field::new("contract_version", DataType::UInt32, false),
52 ]))
53});
54
55pub static LINEAGE_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
57 Arc::new(Schema::new(vec![
58 uuid_field("lineage_uuid", false),
59 uuid_field("provenance_uuid", false),
60 uuid_field("subject_uuid", false),
61 Field::new("subject_kind", DataType::Utf8, false),
62 Field::new("role", DataType::Utf8, false),
63 Field::new("ordinal", DataType::UInt32, false),
64 Field::new("contract_version", DataType::UInt32, false),
65 ]))
66});
67
68static PROVENANCE_EVENT_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
69 fingerprint(
70 CanonicalDomain::Schema,
71 CANONICAL_CONTRACT_VERSION,
72 b"provenance_event/1|provenance_uuid:fixed[16]:required|operation_uuid:fixed[16]:required|event_kind:utf8:required|actor_uuid:fixed[16]:nullable|recorded_at:timestamp_us_utc:required|contract_version:u32:required",
73 )
74 .expect("registered provenance event schema is within canonical bounds")
75});
76
77static LINEAGE_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
78 fingerprint(
79 CanonicalDomain::Schema,
80 CANONICAL_CONTRACT_VERSION,
81 b"lineage/1|lineage_uuid:fixed[16]:required|provenance_uuid:fixed[16]:required|subject_uuid:fixed[16]:required|subject_kind:utf8:required|role:utf8:required|ordinal:u32:required|contract_version:u32:required",
82 )
83 .expect("registered lineage schema is within canonical bounds")
84});
85
86fn uuid_field(name: &str, nullable: bool) -> Field {
87 Field::new(name, DataType::FixedSizeBinary(16), nullable)
88}
89
90#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
92#[serde(rename_all = "snake_case")]
93pub enum EventKind {
94 CreateNode,
96 CreateEdge,
98 MergeCreate,
100 MergeMatchedNoop,
102 SetProperty,
104 RemoveProperty,
106 AddLabel,
108 RemoveLabel,
110 Delete,
112 DetachDelete,
114 OntologyInference,
116 CreateAssertion,
118 AssessConfidence,
120 RecordEvidence,
122 RecordAlgorithmRun,
124 RecordBeliefProjectionAttachment,
126}
127
128impl EventKind {
129 #[must_use]
131 pub const fn as_str(self) -> &'static str {
132 match self {
133 Self::CreateNode => "create_node",
134 Self::CreateEdge => "create_edge",
135 Self::MergeCreate => "merge_create",
136 Self::MergeMatchedNoop => "merge_matched_noop",
137 Self::SetProperty => "set_property",
138 Self::RemoveProperty => "remove_property",
139 Self::AddLabel => "add_label",
140 Self::RemoveLabel => "remove_label",
141 Self::Delete => "delete",
142 Self::DetachDelete => "detach_delete",
143 Self::OntologyInference => "ontology_inference",
144 Self::CreateAssertion => "create_assertion",
145 Self::AssessConfidence => "assess_confidence",
146 Self::RecordEvidence => "record_evidence",
147 Self::RecordAlgorithmRun => "record_algorithm_run",
148 Self::RecordBeliefProjectionAttachment => "record_belief_projection_attachment",
149 }
150 }
151
152 fn parse(value: &str) -> Result<Self, ProvenanceError> {
153 match value {
154 "create_node" => Ok(Self::CreateNode),
155 "create_edge" => Ok(Self::CreateEdge),
156 "merge_create" => Ok(Self::MergeCreate),
157 "merge_matched_noop" => Ok(Self::MergeMatchedNoop),
158 "set_property" => Ok(Self::SetProperty),
159 "remove_property" => Ok(Self::RemoveProperty),
160 "add_label" => Ok(Self::AddLabel),
161 "remove_label" => Ok(Self::RemoveLabel),
162 "delete" => Ok(Self::Delete),
163 "detach_delete" => Ok(Self::DetachDelete),
164 "ontology_inference" => Ok(Self::OntologyInference),
165 "create_assertion" => Ok(Self::CreateAssertion),
166 "assess_confidence" => Ok(Self::AssessConfidence),
167 "record_evidence" => Ok(Self::RecordEvidence),
168 "record_algorithm_run" => Ok(Self::RecordAlgorithmRun),
169 "record_belief_projection_attachment" => Ok(Self::RecordBeliefProjectionAttachment),
170 _ => Err(invalid("event_kind", "unknown closed value")),
171 }
172 }
173}
174
175#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
177#[serde(rename_all = "snake_case")]
178pub enum SubjectKind {
179 Node,
181 Edge,
183 Assertion,
185 EvidenceLink,
187 ConfidenceAssessment,
189 AlgorithmRun,
191 BeliefProjectionAttachment,
193}
194
195impl SubjectKind {
196 #[must_use]
198 pub const fn as_str(self) -> &'static str {
199 match self {
200 Self::Node => "node",
201 Self::Edge => "edge",
202 Self::Assertion => "assertion",
203 Self::EvidenceLink => "evidence_link",
204 Self::ConfidenceAssessment => "confidence_assessment",
205 Self::AlgorithmRun => "algorithm_run",
206 Self::BeliefProjectionAttachment => "belief_projection_attachment",
207 }
208 }
209
210 fn parse(value: &str) -> Result<Self, ProvenanceError> {
211 match value {
212 "node" => Ok(Self::Node),
213 "edge" => Ok(Self::Edge),
214 "assertion" => Ok(Self::Assertion),
215 "evidence_link" => Ok(Self::EvidenceLink),
216 "confidence_assessment" => Ok(Self::ConfidenceAssessment),
217 "algorithm_run" => Ok(Self::AlgorithmRun),
218 "belief_projection_attachment" => Ok(Self::BeliefProjectionAttachment),
219 _ => Err(invalid("subject_kind", "unknown closed value")),
220 }
221 }
222}
223
224#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
226#[serde(rename_all = "snake_case")]
227pub enum LineageRole {
228 Input,
230 Output,
232}
233
234impl LineageRole {
235 #[must_use]
237 pub const fn as_str(self) -> &'static str {
238 match self {
239 Self::Input => "input",
240 Self::Output => "output",
241 }
242 }
243
244 fn parse(value: &str) -> Result<Self, ProvenanceError> {
245 match value {
246 "input" => Ok(Self::Input),
247 "output" => Ok(Self::Output),
248 _ => Err(invalid("role", "unknown closed value")),
249 }
250 }
251}
252
253#[derive(Clone, Debug, PartialEq, Eq)]
255pub struct ProvenanceEvent {
256 pub provenance_uuid: Uuid,
258 pub operation_uuid: Uuid,
260 pub event_kind: EventKind,
262 pub actor_uuid: Option<Uuid>,
264 pub recorded_at_micros: i64,
266 pub contract_version: u32,
268}
269
270impl ProvenanceEvent {
271 pub fn new(
276 operation_uuid: Uuid,
277 event_kind: EventKind,
278 actor_uuid: Option<Uuid>,
279 recorded_at_micros: i64,
280 ) -> Result<Self, ProvenanceError> {
281 require_uuid(operation_uuid, "operation_uuid")?;
282 if let Some(actor_uuid) = actor_uuid {
283 require_uuid(actor_uuid, "actor_uuid")?;
284 }
285 let canonical =
286 event_canonical_bytes(operation_uuid, event_kind, actor_uuid, recorded_at_micros)?;
287 let provenance_uuid = uuid_v8(fingerprint(
288 CanonicalDomain::ProvenanceEvent,
289 CANONICAL_CONTRACT_VERSION,
290 &canonical,
291 )?);
292 Ok(Self {
293 provenance_uuid,
294 operation_uuid,
295 event_kind,
296 actor_uuid,
297 recorded_at_micros,
298 contract_version: PROVENANCE_EVENT_CONTRACT_VERSION,
299 })
300 }
301
302 pub fn canonical_bytes(&self) -> Result<Vec<u8>, ProvenanceError> {
307 Ok(event_canonical_bytes(
308 self.operation_uuid,
309 self.event_kind,
310 self.actor_uuid,
311 self.recorded_at_micros,
312 )?)
313 }
314
315 pub fn fingerprint(&self) -> Result<[u8; 32], ProvenanceError> {
320 Ok(fingerprint(
321 CanonicalDomain::ProvenanceEvent,
322 CANONICAL_CONTRACT_VERSION,
323 &self.canonical_bytes()?,
324 )?)
325 }
326}
327
328#[derive(Clone, Debug, PartialEq, Eq)]
330pub struct LineageRecord {
331 pub lineage_uuid: Uuid,
333 pub provenance_uuid: Uuid,
335 pub subject_uuid: Uuid,
337 pub subject_kind: SubjectKind,
339 pub role: LineageRole,
341 pub ordinal: u32,
343 pub contract_version: u32,
345}
346
347impl LineageRecord {
348 pub fn new(
353 provenance_uuid: Uuid,
354 subject_uuid: Uuid,
355 subject_kind: SubjectKind,
356 role: LineageRole,
357 ordinal: u32,
358 ) -> Result<Self, ProvenanceError> {
359 require_uuid(provenance_uuid, "provenance_uuid")?;
360 require_uuid(subject_uuid, "subject_uuid")?;
361 let canonical =
362 lineage_canonical_bytes(provenance_uuid, subject_uuid, subject_kind, role, ordinal)?;
363 let lineage_uuid = uuid_v8(fingerprint(
364 CanonicalDomain::Lineage,
365 CANONICAL_CONTRACT_VERSION,
366 &canonical,
367 )?);
368 Ok(Self {
369 lineage_uuid,
370 provenance_uuid,
371 subject_uuid,
372 subject_kind,
373 role,
374 ordinal,
375 contract_version: LINEAGE_CONTRACT_VERSION,
376 })
377 }
378
379 pub fn canonical_bytes(&self) -> Result<Vec<u8>, ProvenanceError> {
384 Ok(lineage_canonical_bytes(
385 self.provenance_uuid,
386 self.subject_uuid,
387 self.subject_kind,
388 self.role,
389 self.ordinal,
390 )?)
391 }
392
393 pub fn fingerprint(&self) -> Result<[u8; 32], ProvenanceError> {
398 Ok(fingerprint(
399 CanonicalDomain::Lineage,
400 CANONICAL_CONTRACT_VERSION,
401 &self.canonical_bytes()?,
402 )?)
403 }
404}
405
406#[derive(Clone, Debug, Default, PartialEq, Eq)]
408pub struct ProvenanceLedger {
409 pub events: Vec<ProvenanceEvent>,
411 pub lineage: Vec<LineageRecord>,
413}
414
415impl ProvenanceLedger {
416 pub fn new(
422 mut events: Vec<ProvenanceEvent>,
423 mut lineage: Vec<LineageRecord>,
424 ) -> Result<Self, ProvenanceError> {
425 validate_rows(&events, &lineage)?;
426 let times = events
427 .iter()
428 .map(|event| (event.provenance_uuid, event.recorded_at_micros))
429 .collect::<HashMap<_, _>>();
430 events.sort_by_key(|event| (event.recorded_at_micros, event.provenance_uuid));
431 lineage.sort_by_key(|row| {
432 (
433 times[&row.provenance_uuid],
434 row.provenance_uuid,
435 role_order(row.role),
436 row.ordinal,
437 row.subject_uuid,
438 )
439 });
440 Ok(Self { events, lineage })
441 }
442
443 pub fn merge(&self, staged: &Self) -> Result<Self, ProvenanceError> {
450 let mut events = self.events.clone();
451 let mut by_event = events
452 .iter()
453 .cloned()
454 .map(|event| (event.provenance_uuid, event))
455 .collect::<HashMap<_, _>>();
456 let existing_operations = events_by_operation(&events);
457 let staged_operations = events_by_operation(&staged.events);
458 for (operation_uuid, staged_events) in &staged_operations {
459 if let Some(existing_events) = existing_operations.get(operation_uuid)
460 && (existing_events != staged_events
461 || operation_lineage(&self.lineage, existing_events)
462 != operation_lineage(&staged.lineage, staged_events))
463 {
464 return Err(ProvenanceError::Conflict("operation_uuid"));
465 }
466 }
467 for event in &staged.events {
468 if let Some(existing) = by_event.get(&event.provenance_uuid)
469 && existing != event
470 {
471 return Err(ProvenanceError::Conflict("provenance_uuid"));
472 }
473 if by_event
474 .insert(event.provenance_uuid, event.clone())
475 .is_none()
476 {
477 events.push(event.clone());
478 }
479 }
480
481 let mut lineage = self.lineage.clone();
482 let mut by_lineage = lineage
483 .iter()
484 .cloned()
485 .map(|row| (row.lineage_uuid, row))
486 .collect::<HashMap<_, _>>();
487 for row in &staged.lineage {
488 if let Some(existing) = by_lineage.get(&row.lineage_uuid)
489 && existing != row
490 {
491 return Err(ProvenanceError::Conflict("lineage_uuid"));
492 }
493 if by_lineage.insert(row.lineage_uuid, row.clone()).is_none() {
494 lineage.push(row.clone());
495 }
496 }
497 Self::new(events, lineage)
498 }
499
500 pub fn event_batch(&self) -> Result<RecordBatch, ProvenanceError> {
505 event_batch(&self.events)
506 }
507
508 pub fn lineage_batch(&self) -> Result<RecordBatch, ProvenanceError> {
513 lineage_batch(&self.lineage)
514 }
515
516 pub fn from_batches(
522 event_batches: &[RecordBatch],
523 lineage_batches: &[RecordBatch],
524 ) -> Result<Self, ProvenanceError> {
525 let mut events = Vec::new();
526 for batch in event_batches {
527 require_schema(batch, &PROVENANCE_EVENT_SCHEMA, "event.schema")?;
528 let provenance = fixed_column(batch, "provenance_uuid")?;
529 let operations = fixed_column(batch, "operation_uuid")?;
530 let kinds = string_column(batch, "event_kind")?;
531 let actors = fixed_column(batch, "actor_uuid")?;
532 let recorded = timestamp_column(batch, "recorded_at")?;
533 let versions = u32_column(batch, "contract_version")?;
534 for row in 0..batch.num_rows() {
535 let actor_uuid = if actors.is_null(row) {
536 None
537 } else {
538 Some(uuid_at(actors, row, "actor_uuid")?)
539 };
540 events.push(ProvenanceEvent {
541 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
542 operation_uuid: uuid_at(operations, row, "operation_uuid")?,
543 event_kind: EventKind::parse(required_text(kinds, row, "event_kind")?)?,
544 actor_uuid,
545 recorded_at_micros: required_i64(recorded, row, "recorded_at")?,
546 contract_version: required_u32(versions, row, "contract_version")?,
547 });
548 }
549 }
550
551 let mut lineage = Vec::new();
552 for batch in lineage_batches {
553 require_schema(batch, &LINEAGE_SCHEMA, "lineage.schema")?;
554 let lineage_ids = fixed_column(batch, "lineage_uuid")?;
555 let provenance = fixed_column(batch, "provenance_uuid")?;
556 let subjects = fixed_column(batch, "subject_uuid")?;
557 let subject_kinds = string_column(batch, "subject_kind")?;
558 let roles = string_column(batch, "role")?;
559 let ordinals = u32_column(batch, "ordinal")?;
560 let versions = u32_column(batch, "contract_version")?;
561 for row in 0..batch.num_rows() {
562 lineage.push(LineageRecord {
563 lineage_uuid: uuid_at(lineage_ids, row, "lineage_uuid")?,
564 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
565 subject_uuid: uuid_at(subjects, row, "subject_uuid")?,
566 subject_kind: SubjectKind::parse(required_text(
567 subject_kinds,
568 row,
569 "subject_kind",
570 )?)?,
571 role: LineageRole::parse(required_text(roles, row, "role")?)?,
572 ordinal: required_u32(ordinals, row, "ordinal")?,
573 contract_version: required_u32(versions, row, "contract_version")?,
574 });
575 }
576 }
577 Self::new(events, lineage)
578 }
579}
580
581#[derive(Clone, Debug)]
583pub struct SchemaRegistryEntry {
584 pub capability_id: &'static str,
586 pub capability_version: u32,
588 pub record_family: &'static str,
590 pub record_version: u32,
592 pub schema: SchemaRef,
594 pub schema_fingerprint: [u8; 32],
596 pub enum_registry_versions: &'static [(&'static str, u32)],
598 pub sort_key: &'static [&'static str],
600 pub diff_identity_fields: &'static [&'static str],
602 pub diff_record_uuid_field: Option<&'static str>,
604 pub fingerprint_domain: CanonicalDomain,
606 pub owner: &'static str,
608 pub implementation_issue: u64,
610 pub max_rows: usize,
612}
613
614impl SchemaRegistryEntry {
615 #[must_use]
617 pub const fn diff_identity_fingerprint_domain(&self) -> CanonicalDomain {
618 CanonicalDomain::ArrowResult
619 }
620
621 #[must_use]
623 pub const fn diff_record_fingerprint_domain(&self) -> CanonicalDomain {
624 self.fingerprint_domain
625 }
626}
627
628#[must_use]
630pub fn schema_registry() -> Vec<SchemaRegistryEntry> {
631 vec![
632 SchemaRegistryEntry {
633 capability_id: "provenance",
634 capability_version: PROVENANCE_CAPABILITY_VERSION,
635 record_family: "events",
636 record_version: PROVENANCE_EVENT_CONTRACT_VERSION,
637 schema: Arc::clone(&PROVENANCE_EVENT_SCHEMA),
638 schema_fingerprint: *PROVENANCE_EVENT_SCHEMA_FINGERPRINT,
639 enum_registry_versions: &[("event_kind", EVENT_KIND_REGISTRY_VERSION)],
640 sort_key: &["recorded_at", "provenance_uuid"],
641 diff_identity_fields: &["provenance_uuid"],
642 diff_record_uuid_field: Some("provenance_uuid"),
643 fingerprint_domain: CanonicalDomain::ProvenanceEvent,
644 owner: "graphforge-provenance",
645 implementation_issue: 773,
646 max_rows: MAX_PROVENANCE_ROWS,
647 },
648 SchemaRegistryEntry {
649 capability_id: "provenance",
650 capability_version: PROVENANCE_CAPABILITY_VERSION,
651 record_family: "lineage",
652 record_version: LINEAGE_CONTRACT_VERSION,
653 schema: Arc::clone(&LINEAGE_SCHEMA),
654 schema_fingerprint: *LINEAGE_SCHEMA_FINGERPRINT,
655 enum_registry_versions: &[
656 ("subject_kind", SUBJECT_KIND_REGISTRY_VERSION),
657 ("role", LINEAGE_ROLE_REGISTRY_VERSION),
658 ],
659 sort_key: &[
660 "recorded_at",
661 "provenance_uuid",
662 "role",
663 "ordinal",
664 "subject_uuid",
665 ],
666 diff_identity_fields: &["lineage_uuid"],
667 diff_record_uuid_field: Some("lineage_uuid"),
668 fingerprint_domain: CanonicalDomain::Lineage,
669 owner: "graphforge-provenance",
670 implementation_issue: 773,
671 max_rows: MAX_PROVENANCE_ROWS,
672 },
673 ]
674}
675
676#[derive(thiserror::Error, Debug)]
678pub enum ProvenanceError {
679 #[error("invalid provenance {field}: {message}")]
681 Invalid {
682 field: &'static str,
684 message: &'static str,
686 },
687 #[error("provenance {participant} row limit exceeded: observed {observed}, limit {limit}")]
689 Limit {
690 participant: &'static str,
692 observed: usize,
694 limit: usize,
696 },
697 #[error("duplicate provenance identity: {0}")]
699 Duplicate(&'static str),
700 #[error("dangling provenance reference: {0}")]
702 Dangling(&'static str),
703 #[error("provenance idempotency conflict: {0}")]
705 Conflict(&'static str),
706 #[error(transparent)]
708 Canonical(#[from] CanonicalError),
709 #[error("provenance Arrow failure: {0}")]
711 Arrow(#[from] arrow::error::ArrowError),
712}
713
714impl ProvenanceError {
715 #[must_use]
717 pub const fn code(&self) -> &'static str {
718 match self {
719 Self::Invalid { .. } => "GF_PROVENANCE_INVALID",
720 Self::Limit { .. } => "GF_RESOURCE_LIMIT",
721 Self::Duplicate(_) => "GF_PROVENANCE_DUPLICATE",
722 Self::Dangling(_) => "GF_PROVENANCE_DANGLING",
723 Self::Conflict(_) => "GF_IDEMPOTENCY_CONFLICT",
724 Self::Canonical(error) => error.code(),
725 Self::Arrow(_) => "GF_SCHEMA_MISMATCH",
726 }
727 }
728}
729
730fn event_canonical_bytes(
731 operation_uuid: Uuid,
732 event_kind: EventKind,
733 actor_uuid: Option<Uuid>,
734 recorded_at_micros: i64,
735) -> Result<Vec<u8>, CanonicalError> {
736 let mut writer = CanonicalWriter::new();
737 writer.raw(b"GFPE")?;
738 writer.u32(PROVENANCE_EVENT_CONTRACT_VERSION)?;
739 writer.raw(operation_uuid.as_bytes())?;
740 writer.text(event_kind.as_str())?;
741 match actor_uuid {
742 Some(actor_uuid) => {
743 writer.u8(1)?;
744 writer.raw(actor_uuid.as_bytes())?;
745 }
746 None => writer.u8(0)?,
747 }
748 writer.i64(recorded_at_micros)?;
749 Ok(writer.finish())
750}
751
752fn lineage_canonical_bytes(
753 provenance_uuid: Uuid,
754 subject_uuid: Uuid,
755 subject_kind: SubjectKind,
756 role: LineageRole,
757 ordinal: u32,
758) -> Result<Vec<u8>, CanonicalError> {
759 let mut writer = CanonicalWriter::new();
760 writer.raw(b"GFPL")?;
761 writer.u32(LINEAGE_CONTRACT_VERSION)?;
762 writer.raw(provenance_uuid.as_bytes())?;
763 writer.raw(subject_uuid.as_bytes())?;
764 writer.text(subject_kind.as_str())?;
765 writer.text(role.as_str())?;
766 writer.u32(ordinal)?;
767 Ok(writer.finish())
768}
769
770fn validate_rows(
771 events: &[ProvenanceEvent],
772 lineage: &[LineageRecord],
773) -> Result<(), ProvenanceError> {
774 check_limit("events", events.len())?;
775 check_limit("lineage", lineage.len())?;
776 let mut event_ids = HashSet::with_capacity(events.len());
777 let mut operation_kinds = HashSet::with_capacity(events.len());
778 for event in events {
779 if event.contract_version != PROVENANCE_EVENT_CONTRACT_VERSION {
780 return Err(ProvenanceError::Invalid {
781 field: "event.contract_version",
782 message: "unsupported version",
783 });
784 }
785 require_uuid(event.provenance_uuid, "provenance_uuid")?;
786 require_uuid(event.operation_uuid, "operation_uuid")?;
787 if let Some(actor_uuid) = event.actor_uuid {
788 require_uuid(actor_uuid, "actor_uuid")?;
789 }
790 if !event_ids.insert(event.provenance_uuid) {
791 return Err(ProvenanceError::Duplicate("provenance_uuid"));
792 }
793 if !operation_kinds.insert((event.operation_uuid, event.event_kind)) {
794 return Err(ProvenanceError::Duplicate("operation_uuid/event_kind"));
795 }
796 let expected = ProvenanceEvent::new(
797 event.operation_uuid,
798 event.event_kind,
799 event.actor_uuid,
800 event.recorded_at_micros,
801 )?;
802 if expected.provenance_uuid != event.provenance_uuid {
803 return Err(ProvenanceError::Invalid {
804 field: "provenance_uuid",
805 message: "does not match canonical event",
806 });
807 }
808 }
809
810 let mut lineage_ids = HashSet::with_capacity(lineage.len());
811 let mut positions = HashSet::with_capacity(lineage.len());
812 for row in lineage {
813 if row.contract_version != LINEAGE_CONTRACT_VERSION {
814 return Err(ProvenanceError::Invalid {
815 field: "lineage.contract_version",
816 message: "unsupported version",
817 });
818 }
819 require_uuid(row.lineage_uuid, "lineage_uuid")?;
820 require_uuid(row.subject_uuid, "subject_uuid")?;
821 if !event_ids.contains(&row.provenance_uuid) {
822 return Err(ProvenanceError::Dangling("provenance_uuid"));
823 }
824 if !lineage_ids.insert(row.lineage_uuid) {
825 return Err(ProvenanceError::Duplicate("lineage_uuid"));
826 }
827 if !positions.insert((row.provenance_uuid, row.role, row.ordinal)) {
828 return Err(ProvenanceError::Duplicate("role/ordinal"));
829 }
830 let expected = LineageRecord::new(
831 row.provenance_uuid,
832 row.subject_uuid,
833 row.subject_kind,
834 row.role,
835 row.ordinal,
836 )?;
837 if expected.lineage_uuid != row.lineage_uuid {
838 return Err(ProvenanceError::Invalid {
839 field: "lineage_uuid",
840 message: "does not match canonical lineage",
841 });
842 }
843 }
844 Ok(())
845}
846
847type OperationEventIdentity = (Uuid, EventKind, Option<Uuid>, i64, u32);
848
849fn events_by_operation(events: &[ProvenanceEvent]) -> HashMap<Uuid, Vec<OperationEventIdentity>> {
850 let mut grouped = HashMap::<_, Vec<_>>::new();
851 for event in events {
852 grouped.entry(event.operation_uuid).or_default().push((
853 event.provenance_uuid,
854 event.event_kind,
855 event.actor_uuid,
856 event.recorded_at_micros,
857 event.contract_version,
858 ));
859 }
860 for operation_events in grouped.values_mut() {
861 operation_events.sort_by_key(|event| event.0);
862 }
863 grouped
864}
865
866fn operation_lineage(
867 lineage: &[LineageRecord],
868 events: &[OperationEventIdentity],
869) -> Vec<(Uuid, SubjectKind, LineageRole, u32, Uuid, u32)> {
870 let event_ids = events.iter().map(|event| event.0).collect::<HashSet<_>>();
871 let mut rows = lineage
872 .iter()
873 .filter(|row| event_ids.contains(&row.provenance_uuid))
874 .map(|row| {
875 (
876 row.lineage_uuid,
877 row.subject_kind,
878 row.role,
879 row.ordinal,
880 row.subject_uuid,
881 row.contract_version,
882 )
883 })
884 .collect::<Vec<_>>();
885 rows.sort_unstable_by_key(|row| (row.0, row.3, row.4));
886 rows
887}
888
889fn check_limit(participant: &'static str, observed: usize) -> Result<(), ProvenanceError> {
890 if observed > MAX_PROVENANCE_ROWS {
891 Err(ProvenanceError::Limit {
892 participant,
893 observed,
894 limit: MAX_PROVENANCE_ROWS,
895 })
896 } else {
897 Ok(())
898 }
899}
900
901fn require_uuid(value: Uuid, field: &'static str) -> Result<(), ProvenanceError> {
902 if value.is_nil() {
903 Err(ProvenanceError::Invalid {
904 field,
905 message: "nil UUID is forbidden",
906 })
907 } else {
908 Ok(())
909 }
910}
911
912const fn role_order(role: LineageRole) -> u8 {
913 match role {
914 LineageRole::Input => 0,
915 LineageRole::Output => 1,
916 }
917}
918
919fn event_batch(events: &[ProvenanceEvent]) -> Result<RecordBatch, ProvenanceError> {
920 let provenance = fixed_uuid(events.iter().map(|event| event.provenance_uuid))?;
921 let operation = fixed_uuid(events.iter().map(|event| event.operation_uuid))?;
922 let kinds = StringArray::from_iter_values(events.iter().map(|event| event.event_kind.as_str()));
923 let mut actor = FixedSizeBinaryBuilder::with_capacity(events.len(), 16);
924 for event in events {
925 match event.actor_uuid {
926 Some(value) => actor.append_value(value.as_bytes())?,
927 None => actor.append_null(),
928 }
929 }
930 let recorded_at = TimestampMicrosecondArray::from(
931 events
932 .iter()
933 .map(|event| event.recorded_at_micros)
934 .collect::<Vec<_>>(),
935 )
936 .with_timezone("UTC");
937 let versions = UInt32Array::from_iter_values(events.iter().map(|event| event.contract_version));
938 Ok(RecordBatch::try_new(
939 Arc::clone(&PROVENANCE_EVENT_SCHEMA),
940 vec![
941 Arc::new(provenance),
942 Arc::new(operation),
943 Arc::new(kinds),
944 Arc::new(actor.finish()),
945 Arc::new(recorded_at),
946 Arc::new(versions),
947 ],
948 )?)
949}
950
951fn lineage_batch(rows: &[LineageRecord]) -> Result<RecordBatch, ProvenanceError> {
952 let lineage = fixed_uuid(rows.iter().map(|row| row.lineage_uuid))?;
953 let provenance = fixed_uuid(rows.iter().map(|row| row.provenance_uuid))?;
954 let subjects = fixed_uuid(rows.iter().map(|row| row.subject_uuid))?;
955 let subject_kinds =
956 StringArray::from_iter_values(rows.iter().map(|row| row.subject_kind.as_str()));
957 let roles = StringArray::from_iter_values(rows.iter().map(|row| row.role.as_str()));
958 let ordinals = UInt32Array::from_iter_values(rows.iter().map(|row| row.ordinal));
959 let versions = UInt32Array::from_iter_values(rows.iter().map(|row| row.contract_version));
960 Ok(RecordBatch::try_new(
961 Arc::clone(&LINEAGE_SCHEMA),
962 vec![
963 Arc::new(lineage),
964 Arc::new(provenance),
965 Arc::new(subjects),
966 Arc::new(subject_kinds),
967 Arc::new(roles),
968 Arc::new(ordinals),
969 Arc::new(versions),
970 ],
971 )?)
972}
973
974fn fixed_uuid(
975 values: impl IntoIterator<Item = Uuid>,
976) -> Result<FixedSizeBinaryArray, arrow::error::ArrowError> {
977 let values = values.into_iter();
978 let (lower, _) = values.size_hint();
979 let mut builder = FixedSizeBinaryBuilder::with_capacity(lower, 16);
980 for value in values {
981 builder.append_value(value.as_bytes())?;
982 }
983 Ok(builder.finish())
984}
985
986fn require_schema(
987 batch: &RecordBatch,
988 expected: &SchemaRef,
989 field: &'static str,
990) -> Result<(), ProvenanceError> {
991 if batch.schema().as_ref() == expected.as_ref() {
992 Ok(())
993 } else {
994 Err(invalid(field, "schema mismatch"))
995 }
996}
997
998fn fixed_column<'a>(
999 batch: &'a RecordBatch,
1000 name: &'static str,
1001) -> Result<&'a FixedSizeBinaryArray, ProvenanceError> {
1002 batch
1003 .column_by_name(name)
1004 .and_then(|array| array.as_any().downcast_ref::<FixedSizeBinaryArray>())
1005 .ok_or_else(|| invalid(name, "column type mismatch"))
1006}
1007
1008fn string_column<'a>(
1009 batch: &'a RecordBatch,
1010 name: &'static str,
1011) -> Result<&'a StringArray, ProvenanceError> {
1012 batch
1013 .column_by_name(name)
1014 .and_then(|array| array.as_any().downcast_ref::<StringArray>())
1015 .ok_or_else(|| invalid(name, "column type mismatch"))
1016}
1017
1018fn timestamp_column<'a>(
1019 batch: &'a RecordBatch,
1020 name: &'static str,
1021) -> Result<&'a TimestampMicrosecondArray, ProvenanceError> {
1022 batch
1023 .column_by_name(name)
1024 .and_then(|array| array.as_any().downcast_ref::<TimestampMicrosecondArray>())
1025 .ok_or_else(|| invalid(name, "column type mismatch"))
1026}
1027
1028fn u32_column<'a>(
1029 batch: &'a RecordBatch,
1030 name: &'static str,
1031) -> Result<&'a UInt32Array, ProvenanceError> {
1032 batch
1033 .column_by_name(name)
1034 .and_then(|array| array.as_any().downcast_ref::<UInt32Array>())
1035 .ok_or_else(|| invalid(name, "column type mismatch"))
1036}
1037
1038fn uuid_at(
1039 array: &FixedSizeBinaryArray,
1040 row: usize,
1041 field: &'static str,
1042) -> Result<Uuid, ProvenanceError> {
1043 if array.is_null(row) {
1044 return Err(invalid(field, "null is forbidden"));
1045 }
1046 Uuid::from_slice(array.value(row)).map_err(|_| invalid(field, "invalid UUID bytes"))
1047}
1048
1049fn required_text<'a>(
1050 array: &'a StringArray,
1051 row: usize,
1052 field: &'static str,
1053) -> Result<&'a str, ProvenanceError> {
1054 if array.is_null(row) {
1055 Err(invalid(field, "null is forbidden"))
1056 } else {
1057 Ok(array.value(row))
1058 }
1059}
1060
1061fn required_i64(
1062 array: &TimestampMicrosecondArray,
1063 row: usize,
1064 field: &'static str,
1065) -> Result<i64, ProvenanceError> {
1066 if array.is_null(row) {
1067 Err(invalid(field, "null is forbidden"))
1068 } else {
1069 Ok(array.value(row))
1070 }
1071}
1072
1073fn required_u32(
1074 array: &UInt32Array,
1075 row: usize,
1076 field: &'static str,
1077) -> Result<u32, ProvenanceError> {
1078 if array.is_null(row) {
1079 Err(invalid(field, "null is forbidden"))
1080 } else {
1081 Ok(array.value(row))
1082 }
1083}
1084
1085const fn invalid(field: &'static str, message: &'static str) -> ProvenanceError {
1086 ProvenanceError::Invalid { field, message }
1087}
1088
1089#[cfg(test)]
1090mod tests {
1091 use arrow::array::{Array, FixedSizeBinaryArray, StringArray};
1092
1093 use super::*;
1094
1095 fn uuid(value: u128) -> Uuid {
1096 Uuid::from_u128(value)
1097 }
1098
1099 fn fixture() -> (ProvenanceEvent, Vec<LineageRecord>) {
1100 let event =
1101 ProvenanceEvent::new(uuid(1), EventKind::CreateEdge, Some(uuid(2)), 123).unwrap();
1102 let rows = vec![
1103 LineageRecord::new(
1104 event.provenance_uuid,
1105 uuid(4),
1106 SubjectKind::Edge,
1107 LineageRole::Output,
1108 0,
1109 )
1110 .unwrap(),
1111 LineageRecord::new(
1112 event.provenance_uuid,
1113 uuid(3),
1114 SubjectKind::Node,
1115 LineageRole::Input,
1116 0,
1117 )
1118 .unwrap(),
1119 ];
1120 (event, rows)
1121 }
1122
1123 #[test]
1124 fn canonical_ids_and_bytes_are_stable() {
1125 let (event, rows) = fixture();
1126 assert_eq!(
1127 event.provenance_uuid.to_string(),
1128 "1255afb8-f9f5-806e-8086-f79f9ab73376"
1129 );
1130 assert_eq!(
1131 rows[0].lineage_uuid.to_string(),
1132 "b2fbd21f-798c-8105-babf-0a1af2d64ea9"
1133 );
1134 assert_eq!(event, event.clone());
1135 assert_eq!(event.fingerprint().unwrap(), event.fingerprint().unwrap());
1136 }
1137
1138 #[test]
1139 fn ledger_orders_history_and_shapes_authoritative_arrow() {
1140 let (event, mut rows) = fixture();
1141 rows.reverse();
1142 let ledger = ProvenanceLedger::new(vec![event.clone()], rows).unwrap();
1143 assert_eq!(ledger.lineage[0].role, LineageRole::Input);
1144 assert_eq!(ledger.lineage[1].role, LineageRole::Output);
1145
1146 let events = ledger.event_batch().unwrap();
1147 assert_eq!(events.schema(), *PROVENANCE_EVENT_SCHEMA);
1148 assert_eq!(events.num_rows(), 1);
1149 assert_eq!(
1150 events
1151 .column_by_name("event_kind")
1152 .unwrap()
1153 .as_any()
1154 .downcast_ref::<StringArray>()
1155 .unwrap()
1156 .value(0),
1157 "create_edge"
1158 );
1159 assert!(events.column_by_name("actor_uuid").unwrap().is_valid(0));
1160
1161 let lineage = ledger.lineage_batch().unwrap();
1162 assert_eq!(lineage.schema(), *LINEAGE_SCHEMA);
1163 assert_eq!(lineage.num_rows(), 2);
1164 assert_eq!(
1165 lineage
1166 .column_by_name("subject_uuid")
1167 .unwrap()
1168 .as_any()
1169 .downcast_ref::<FixedSizeBinaryArray>()
1170 .unwrap()
1171 .value(0),
1172 uuid(3).as_bytes()
1173 );
1174 assert_eq!(
1175 ProvenanceLedger::from_batches(&[events], &[lineage]).unwrap(),
1176 ledger
1177 );
1178 }
1179
1180 #[test]
1181 fn validation_rejects_nil_tampered_duplicate_and_dangling_rows() {
1182 assert_eq!(
1183 ProvenanceEvent::new(Uuid::nil(), EventKind::CreateNode, None, 0)
1184 .unwrap_err()
1185 .code(),
1186 "GF_PROVENANCE_INVALID"
1187 );
1188
1189 let (event, rows) = fixture();
1190 let mut tampered = event.clone();
1191 tampered.provenance_uuid = uuid(99);
1192 assert_eq!(
1193 ProvenanceLedger::new(vec![tampered], vec![])
1194 .unwrap_err()
1195 .code(),
1196 "GF_PROVENANCE_INVALID"
1197 );
1198 assert_eq!(
1199 ProvenanceLedger::new(vec![event.clone(), event.clone()], vec![])
1200 .unwrap_err()
1201 .code(),
1202 "GF_PROVENANCE_DUPLICATE"
1203 );
1204 assert_eq!(
1205 ProvenanceLedger::new(vec![], rows).unwrap_err().code(),
1206 "GF_PROVENANCE_DANGLING"
1207 );
1208 }
1209
1210 #[test]
1211 fn identical_merge_is_idempotent_and_conflicts_are_atomic() {
1212 let (event, rows) = fixture();
1213 let ledger = ProvenanceLedger::new(vec![event.clone()], rows.clone()).unwrap();
1214 assert_eq!(ledger.merge(&ledger).unwrap(), ledger);
1215
1216 let conflicting =
1217 ProvenanceEvent::new(event.operation_uuid, EventKind::Delete, None, 124).unwrap();
1218 let staged = ProvenanceLedger::new(vec![conflicting], vec![]).unwrap();
1219 assert_eq!(
1220 ledger.merge(&staged).unwrap_err().code(),
1221 "GF_IDEMPOTENCY_CONFLICT"
1222 );
1223 let conflicting_lineage = ProvenanceLedger::new(
1224 vec![event.clone()],
1225 vec![
1226 rows[1].clone(),
1227 LineageRecord::new(
1228 event.provenance_uuid,
1229 uuid(99),
1230 SubjectKind::Edge,
1231 LineageRole::Output,
1232 0,
1233 )
1234 .unwrap(),
1235 ],
1236 )
1237 .unwrap();
1238 assert_eq!(
1239 ledger.merge(&conflicting_lineage).unwrap_err().code(),
1240 "GF_IDEMPOTENCY_CONFLICT"
1241 );
1242 assert_eq!(ledger.events.len(), 1);
1243 assert_eq!(ledger.lineage.len(), 2);
1244 }
1245
1246 #[test]
1247 fn one_operation_can_record_distinct_composite_mutation_kinds() {
1248 let operation_uuid = uuid(10);
1249 let node_event =
1250 ProvenanceEvent::new(operation_uuid, EventKind::CreateNode, None, 456).unwrap();
1251 let edge_event =
1252 ProvenanceEvent::new(operation_uuid, EventKind::CreateEdge, None, 456).unwrap();
1253 let ledger =
1254 ProvenanceLedger::new(vec![edge_event.clone(), node_event.clone()], vec![]).unwrap();
1255
1256 assert_eq!(ledger.events.len(), 2);
1257 assert_eq!(ledger.merge(&ledger).unwrap(), ledger);
1258
1259 let partial_retry = ProvenanceLedger::new(vec![node_event], vec![]).unwrap();
1260 assert_eq!(
1261 ledger.merge(&partial_retry).unwrap_err().code(),
1262 "GF_IDEMPOTENCY_CONFLICT"
1263 );
1264
1265 let duplicate_kind =
1266 ProvenanceEvent::new(operation_uuid, EventKind::CreateNode, Some(uuid(11)), 456)
1267 .unwrap();
1268 assert_eq!(
1269 ProvenanceLedger::new(
1270 vec![edge_event, duplicate_kind.clone(), duplicate_kind],
1271 vec![]
1272 )
1273 .unwrap_err()
1274 .code(),
1275 "GF_PROVENANCE_DUPLICATE"
1276 );
1277 }
1278
1279 #[test]
1280 fn registry_has_one_owner_and_exact_frozen_schemas() {
1281 let registry = schema_registry();
1282 assert_eq!(registry.len(), 2);
1283 assert!(registry.iter().all(|entry| {
1284 entry.owner == "graphforge-provenance"
1285 && entry.capability_id == "provenance"
1286 && entry.capability_version == 1
1287 && entry.implementation_issue == 773
1288 }));
1289 assert_eq!(registry[0].diff_identity_fields, &["provenance_uuid"]);
1290 assert_eq!(registry[0].diff_record_uuid_field, Some("provenance_uuid"));
1291 assert_eq!(registry[1].diff_identity_fields, &["lineage_uuid"]);
1292 assert_eq!(registry[1].diff_record_uuid_field, Some("lineage_uuid"));
1293 assert_eq!(
1294 registry[0]
1295 .schema
1296 .fields()
1297 .iter()
1298 .map(|field| field.name().as_str())
1299 .collect::<Vec<_>>(),
1300 [
1301 "provenance_uuid",
1302 "operation_uuid",
1303 "event_kind",
1304 "actor_uuid",
1305 "recorded_at",
1306 "contract_version"
1307 ]
1308 );
1309 assert_eq!(
1310 registry[1]
1311 .schema
1312 .fields()
1313 .iter()
1314 .map(|field| field.name().as_str())
1315 .collect::<Vec<_>>(),
1316 [
1317 "lineage_uuid",
1318 "provenance_uuid",
1319 "subject_uuid",
1320 "subject_kind",
1321 "role",
1322 "ordinal",
1323 "contract_version"
1324 ]
1325 );
1326 }
1327
1328 #[test]
1329 fn role_ordinals_are_unique_within_each_event() {
1330 let (event, mut rows) = fixture();
1331 rows.push(
1332 LineageRecord::new(
1333 event.provenance_uuid,
1334 uuid(5),
1335 SubjectKind::Node,
1336 LineageRole::Input,
1337 0,
1338 )
1339 .unwrap(),
1340 );
1341 assert_eq!(
1342 ProvenanceLedger::new(vec![event], rows).unwrap_err().code(),
1343 "GF_PROVENANCE_DUPLICATE"
1344 );
1345 }
1346}