1#![forbid(unsafe_code)]
7
8mod belief_projection;
9mod hypothesis;
10mod reasoning;
11mod status;
12mod supersession;
13mod valid_time;
14
15pub use hypothesis::{
16 HYPOTHESIS_GROUP_CONTRACT_VERSION, HYPOTHESIS_GROUP_SCHEMA, HYPOTHESIS_KEY_POLICY_VERSION,
17 HYPOTHESIS_MEMBERSHIP_CONTRACT_VERSION, HYPOTHESIS_MEMBERSHIP_SCHEMA,
18 HYPOTHESIS_SELECTION_CONTRACT_VERSION, HYPOTHESIS_SELECTION_SCHEMA,
19 HYPOTHESIS_STATE_POLICY_VERSION, HypothesisGroup, HypothesisLedger, HypothesisMembershipAction,
20 HypothesisMembershipEvent, HypothesisSelectionEvent, MAX_HYPOTHESIS_QUESTION_KEY_BYTES,
21};
22pub use reasoning::{
23 EPISTEMIC_CAPABILITY_VERSION, MAX_REASONING_CONTENT_BYTES,
24 REASONING_CONTENT_FORMAT_REGISTRY_VERSION, REASONING_CONTRACT_VERSION,
25 REASONING_KIND_REGISTRY_VERSION, REASONING_SCHEMA, ReasoningContentFormat, ReasoningKind,
26 ReasoningLedger, ReasoningRecord,
27};
28pub use status::{
29 ASSERTION_STATUS_CONTRACT_VERSION, ASSERTION_STATUS_REGISTRY_VERSION, ASSERTION_STATUS_SCHEMA,
30 AssertionStatus, AssertionStatusEvent, AssertionStatusLedger,
31};
32pub use supersession::{
33 ASSERTION_SUPERSESSION_CONTRACT_VERSION, ASSERTION_SUPERSESSION_POLICY_VERSION,
34 ASSERTION_SUPERSESSION_SCHEMA, AssertionSupersession, AssertionSupersessionLedger,
35};
36pub use valid_time::{
37 ASSERTION_VALIDITY_CONTRACT_VERSION, ASSERTION_VALIDITY_POLICY_VERSION,
38 ASSERTION_VALIDITY_SCHEMA, AssertionValidityEvent, AssertionValidityLedger,
39 VALID_TIME_CAPABILITY_VERSION,
40};
41
42use std::collections::{HashMap, HashSet};
43use std::sync::{Arc, LazyLock};
44
45use arrow::array::{
46 Array, BinaryArray, FixedSizeBinaryArray, FixedSizeBinaryBuilder, Float64Array, StringArray,
47 TimestampMicrosecondArray, UInt32Array,
48};
49use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
50use arrow::record_batch::RecordBatch;
51use graphforge_core::canonical::{
52 CANONICAL_CONTRACT_VERSION, CanonicalDomain, CanonicalError, CanonicalWriter, fingerprint,
53};
54use uuid::{Uuid, Version};
55
56pub const KNOWLEDGE_CAPABILITY_VERSION: u32 = 1;
58pub const ASSERTION_CONTRACT_VERSION: u32 = 1;
60pub const ASSERTION_GRAPH_REF_CONTRACT_VERSION: u32 = 1;
62pub const CONFIDENCE_ASSESSMENT_CONTRACT_VERSION: u32 = 1;
64pub const CONFIDENCE_INPUT_CONTRACT_VERSION: u32 = 1;
66pub const EVIDENCE_LINK_CONTRACT_VERSION: u32 = 1;
68pub const ALGORITHM_RUN_CONTRACT_VERSION: u32 = 1;
70pub const ALGORITHM_RUN_EVENT_CONTRACT_VERSION: u32 = 1;
72pub const CONFIDENCE_POLICY_REGISTRY_VERSION: u32 = 1;
74pub const GRAPH_OBJECT_KIND_REGISTRY_VERSION: u32 = 1;
76pub const ASSERTION_GRAPH_ROLE_REGISTRY_VERSION: u32 = 1;
78pub const EVIDENCE_SOURCE_KIND_REGISTRY_VERSION: u32 = 1;
80pub const EVIDENCE_ROLE_REGISTRY_VERSION: u32 = 1;
82pub const ALGORITHM_RUN_STATE_REGISTRY_VERSION: u32 = 1;
84pub const MAX_KNOWLEDGE_ROWS: usize = 1_000_000;
86
87pub static ASSERTION_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
89 Arc::new(Schema::new(vec![
90 uuid_field("assertion_uuid", false),
91 Field::new("claim", DataType::Utf8, false),
92 uuid_field("provenance_uuid", false),
93 Field::new(
94 "recorded_at",
95 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
96 false,
97 ),
98 Field::new("contract_version", DataType::UInt32, false),
99 ]))
100});
101
102pub static ASSERTION_GRAPH_REF_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
104 Arc::new(Schema::new(vec![
105 uuid_field("assertion_uuid", false),
106 uuid_field("graph_uuid", false),
107 Field::new("graph_kind", DataType::Utf8, false),
108 Field::new("role", DataType::Utf8, false),
109 Field::new("ordinal", DataType::UInt32, false),
110 Field::new("contract_version", DataType::UInt32, false),
111 ]))
112});
113
114pub static CONFIDENCE_ASSESSMENT_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
116 Arc::new(Schema::new(vec![
117 uuid_field("confidence_uuid", false),
118 uuid_field("assertion_uuid", false),
119 Field::new("policy", DataType::Utf8, false),
120 Field::new("policy_version", DataType::UInt32, false),
121 Field::new("value", DataType::Float64, true),
122 uuid_field("provenance_uuid", false),
123 Field::new(
124 "recorded_at",
125 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
126 false,
127 ),
128 Field::new("contract_version", DataType::UInt32, false),
129 ]))
130});
131
132pub static CONFIDENCE_INPUT_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
134 Arc::new(Schema::new(vec![
135 uuid_field("confidence_uuid", false),
136 uuid_field("input_confidence_uuid", false),
137 Field::new("input_value", DataType::Float64, true),
138 Field::new("ordinal", DataType::UInt32, false),
139 Field::new("contract_version", DataType::UInt32, false),
140 ]))
141});
142
143pub static EVIDENCE_LINK_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
145 Arc::new(Schema::new(vec![
146 uuid_field("evidence_uuid", false),
147 uuid_field("assertion_uuid", false),
148 uuid_field("source_uuid", false),
149 Field::new("source_kind", DataType::Utf8, false),
150 Field::new("role", DataType::Utf8, false),
151 Field::new("weight", DataType::Float64, true),
152 uuid_field("provenance_uuid", false),
153 Field::new(
154 "recorded_at",
155 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
156 false,
157 ),
158 Field::new("contract_version", DataType::UInt32, false),
159 ]))
160});
161
162pub static ALGORITHM_RUN_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
164 Arc::new(Schema::new(vec![
165 uuid_field("run_uuid", false),
166 Field::new("algorithm", DataType::Utf8, false),
167 Field::new("algorithm_version", DataType::UInt32, false),
168 Field::new("descriptor_version", DataType::UInt32, false),
169 Field::new("descriptor", DataType::Binary, false),
170 Field::new(
171 "projection_fingerprint",
172 DataType::FixedSizeBinary(32),
173 false,
174 ),
175 uuid_field("provenance_uuid", false),
176 Field::new(
177 "started_at",
178 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
179 false,
180 ),
181 Field::new("contract_version", DataType::UInt32, false),
182 ]))
183});
184
185pub static ALGORITHM_RUN_EVENT_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
187 Arc::new(Schema::new(vec![
188 uuid_field("event_uuid", false),
189 uuid_field("run_uuid", false),
190 Field::new("state", DataType::Utf8, false),
191 Field::new("result_fingerprint", DataType::FixedSizeBinary(32), true),
192 Field::new("error_code", DataType::Utf8, true),
193 Field::new(
194 "recorded_at",
195 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
196 false,
197 ),
198 uuid_field("provenance_uuid", false),
199 Field::new("contract_version", DataType::UInt32, false),
200 ]))
201});
202
203static ASSERTION_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
204 fingerprint(
205 CanonicalDomain::Schema,
206 CANONICAL_CONTRACT_VERSION,
207 b"assertion/1|assertion_uuid:fixed[16]:required|claim:utf8:required|provenance_uuid:fixed[16]:required|recorded_at:timestamp_us_utc:required|contract_version:u32:required",
208 )
209 .expect("registered assertion schema is within canonical bounds")
210});
211
212static ASSERTION_GRAPH_REF_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
213 fingerprint(
214 CanonicalDomain::Schema,
215 CANONICAL_CONTRACT_VERSION,
216 b"assertion_graph_ref/1|assertion_uuid:fixed[16]:required|graph_uuid:fixed[16]:required|graph_kind:utf8:required|role:utf8:required|ordinal:u32:required|contract_version:u32:required",
217 )
218 .expect("registered assertion graph-ref schema is within canonical bounds")
219});
220
221static CONFIDENCE_ASSESSMENT_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
222 fingerprint(
223 CanonicalDomain::Schema,
224 CANONICAL_CONTRACT_VERSION,
225 b"confidence_assessment/1|confidence_uuid:fixed[16]:required|assertion_uuid:fixed[16]:required|policy:utf8:required|policy_version:u32:required|value:f64:nullable|provenance_uuid:fixed[16]:required|recorded_at:timestamp_us_utc:required|contract_version:u32:required",
226 )
227 .expect("registered confidence-assessment schema is within canonical bounds")
228});
229
230static CONFIDENCE_INPUT_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
231 fingerprint(
232 CanonicalDomain::Schema,
233 CANONICAL_CONTRACT_VERSION,
234 b"confidence_input/1|confidence_uuid:fixed[16]:required|input_confidence_uuid:fixed[16]:required|input_value:f64:nullable|ordinal:u32:required|contract_version:u32:required",
235 )
236 .expect("registered confidence-input schema is within canonical bounds")
237});
238
239static EVIDENCE_LINK_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
240 fingerprint(
241 CanonicalDomain::Schema,
242 CANONICAL_CONTRACT_VERSION,
243 b"evidence_link/1|evidence_uuid:fixed[16]:required|assertion_uuid:fixed[16]:required|source_uuid:fixed[16]:required|source_kind:utf8:required|role:utf8:required|weight:f64:nullable|provenance_uuid:fixed[16]:required|recorded_at:timestamp_us_utc:required|contract_version:u32:required",
244 )
245 .expect("registered evidence-link schema is within canonical bounds")
246});
247
248static ALGORITHM_RUN_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
249 fingerprint(
250 CanonicalDomain::Schema,
251 CANONICAL_CONTRACT_VERSION,
252 b"algorithm_run/1|run_uuid:fixed[16]:required|algorithm:utf8:required|algorithm_version:u32:required|descriptor_version:u32:required|descriptor:binary:required|projection_fingerprint:fixed[32]:required|provenance_uuid:fixed[16]:required|started_at:timestamp_us_utc:required|contract_version:u32:required",
253 )
254 .expect("registered algorithm-run schema is within canonical bounds")
255});
256
257static ALGORITHM_RUN_EVENT_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
258 fingerprint(
259 CanonicalDomain::Schema,
260 CANONICAL_CONTRACT_VERSION,
261 b"algorithm_run_event/1|event_uuid:fixed[16]:required|run_uuid:fixed[16]:required|state:utf8:required|result_fingerprint:fixed[32]:nullable|error_code:utf8:nullable|recorded_at:timestamp_us_utc:required|provenance_uuid:fixed[16]:required|contract_version:u32:required",
262 )
263 .expect("registered algorithm-run-event schema is within canonical bounds")
264});
265
266fn uuid_field(name: &str, nullable: bool) -> Field {
267 Field::new(name, DataType::FixedSizeBinary(16), nullable)
268}
269
270#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
272#[serde(rename_all = "snake_case")]
273pub enum GraphObjectKind {
274 Node,
276 Edge,
278}
279
280impl GraphObjectKind {
281 #[must_use]
283 pub const fn as_str(self) -> &'static str {
284 match self {
285 Self::Node => "node",
286 Self::Edge => "edge",
287 }
288 }
289
290 fn parse(value: &str) -> Result<Self, KnowledgeError> {
291 match value {
292 "node" => Ok(Self::Node),
293 "edge" => Ok(Self::Edge),
294 _ => Err(invalid("graph_kind", "unknown closed value")),
295 }
296 }
297}
298
299#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
301#[serde(rename_all = "snake_case")]
302pub enum AssertionGraphRole {
303 Subject,
305 Object,
307 Context,
309}
310
311impl AssertionGraphRole {
312 #[must_use]
314 pub const fn as_str(self) -> &'static str {
315 match self {
316 Self::Subject => "subject",
317 Self::Object => "object",
318 Self::Context => "context",
319 }
320 }
321
322 fn parse(value: &str) -> Result<Self, KnowledgeError> {
323 match value {
324 "subject" => Ok(Self::Subject),
325 "object" => Ok(Self::Object),
326 "context" => Ok(Self::Context),
327 _ => Err(invalid("role", "unknown closed value")),
328 }
329 }
330}
331
332#[derive(Clone, Debug, PartialEq, Eq)]
334pub struct Assertion {
335 pub assertion_uuid: Uuid,
337 pub claim: String,
339 pub provenance_uuid: Uuid,
341 pub recorded_at_micros: i64,
343 pub contract_version: u32,
345}
346
347impl Assertion {
348 pub fn new(
350 assertion_uuid: Uuid,
351 claim: String,
352 provenance_uuid: Uuid,
353 recorded_at_micros: i64,
354 ) -> Result<Self, KnowledgeError> {
355 require_v7(assertion_uuid, "assertion_uuid")?;
356 require_uuid(provenance_uuid, "provenance_uuid")?;
357 validate_claim(&claim)?;
358 Ok(Self {
359 assertion_uuid,
360 claim,
361 provenance_uuid,
362 recorded_at_micros,
363 contract_version: ASSERTION_CONTRACT_VERSION,
364 })
365 }
366}
367
368#[derive(Clone, Debug, PartialEq, Eq)]
370pub struct AssertionGraphRef {
371 pub assertion_uuid: Uuid,
373 pub graph_uuid: Uuid,
375 pub graph_kind: GraphObjectKind,
377 pub role: AssertionGraphRole,
379 pub ordinal: u32,
381 pub contract_version: u32,
383}
384
385impl AssertionGraphRef {
386 pub fn new(
388 assertion_uuid: Uuid,
389 graph_uuid: Uuid,
390 graph_kind: GraphObjectKind,
391 role: AssertionGraphRole,
392 ordinal: u32,
393 ) -> Result<Self, KnowledgeError> {
394 require_v7(assertion_uuid, "assertion_uuid")?;
395 require_uuid(graph_uuid, "graph_uuid")?;
396 Ok(Self {
397 assertion_uuid,
398 graph_uuid,
399 graph_kind,
400 role,
401 ordinal,
402 contract_version: ASSERTION_GRAPH_REF_CONTRACT_VERSION,
403 })
404 }
405}
406
407#[derive(Clone, Debug, Default, PartialEq, Eq)]
409pub struct AssertionLedger {
410 pub assertions: Vec<Assertion>,
412 pub graph_refs: Vec<AssertionGraphRef>,
414}
415
416impl AssertionLedger {
417 pub fn new(
419 mut assertions: Vec<Assertion>,
420 mut graph_refs: Vec<AssertionGraphRef>,
421 ) -> Result<Self, KnowledgeError> {
422 validate_rows(&assertions, &graph_refs)?;
423 let times = assertions
424 .iter()
425 .map(|row| (row.assertion_uuid, row.recorded_at_micros))
426 .collect::<HashMap<_, _>>();
427 assertions.sort_by_key(|row| (row.recorded_at_micros, row.assertion_uuid));
428 graph_refs.sort_by_key(|row| {
429 (
430 times[&row.assertion_uuid],
431 row.assertion_uuid,
432 role_order(row.role),
433 row.ordinal,
434 kind_order(row.graph_kind),
435 row.graph_uuid,
436 )
437 });
438 Ok(Self {
439 assertions,
440 graph_refs,
441 })
442 }
443
444 pub fn merge(&self, staged: &Self) -> Result<Self, KnowledgeError> {
446 let mut assertions = self.assertions.clone();
447 let mut refs = self.graph_refs.clone();
448 let mut by_id = assertions
449 .iter()
450 .cloned()
451 .map(|row| (row.assertion_uuid, row))
452 .collect::<HashMap<_, _>>();
453 for row in &staged.assertions {
454 if let Some(existing) = by_id.get(&row.assertion_uuid)
455 && (existing != row
456 || refs_for(&refs, row.assertion_uuid)
457 != refs_for(&staged.graph_refs, row.assertion_uuid))
458 {
459 return Err(KnowledgeError::Conflict("assertion_uuid"));
460 }
461 if by_id.insert(row.assertion_uuid, row.clone()).is_none() {
462 assertions.push(row.clone());
463 refs.extend(
464 staged
465 .graph_refs
466 .iter()
467 .filter(|reference| reference.assertion_uuid == row.assertion_uuid)
468 .cloned(),
469 );
470 }
471 }
472 Self::new(assertions, refs)
473 }
474
475 pub fn assertion_fingerprint(&self, assertion_uuid: Uuid) -> Result<[u8; 32], KnowledgeError> {
477 let assertion = self
478 .assertions
479 .iter()
480 .find(|row| row.assertion_uuid == assertion_uuid)
481 .ok_or(KnowledgeError::Dangling("assertion_uuid"))?;
482 let refs = refs_for(&self.graph_refs, assertion_uuid);
483 let mut writer = CanonicalWriter::new();
484 writer.raw(b"GFAS")?;
485 writer.u32(ASSERTION_CONTRACT_VERSION)?;
486 writer.text(&assertion.claim)?;
487 writer.u64(
488 u64::try_from(refs.len()).map_err(|_| KnowledgeError::Limit {
489 participant: "assertion_graph_refs",
490 observed: refs.len(),
491 limit: MAX_KNOWLEDGE_ROWS,
492 })?,
493 )?;
494 for reference in refs {
495 writer.text(reference.role.as_str())?;
496 writer.u32(reference.ordinal)?;
497 writer.text(reference.graph_kind.as_str())?;
498 writer.raw(reference.graph_uuid.as_bytes())?;
499 }
500 Ok(fingerprint(
501 CanonicalDomain::Assertion,
502 CANONICAL_CONTRACT_VERSION,
503 &writer.finish(),
504 )?)
505 }
506
507 pub fn assertion_batch(&self) -> Result<RecordBatch, KnowledgeError> {
509 assertion_batch(&self.assertions)
510 }
511
512 pub fn graph_ref_batch(&self) -> Result<RecordBatch, KnowledgeError> {
514 graph_ref_batch(&self.graph_refs)
515 }
516
517 pub fn from_batches(
519 assertion_batches: &[RecordBatch],
520 graph_ref_batches: &[RecordBatch],
521 ) -> Result<Self, KnowledgeError> {
522 let mut assertions = Vec::new();
523 for batch in assertion_batches {
524 require_schema(batch, &ASSERTION_SCHEMA, "assertion.schema")?;
525 let ids = fixed_column(batch, "assertion_uuid")?;
526 let claims = string_column(batch, "claim")?;
527 let provenance = fixed_column(batch, "provenance_uuid")?;
528 let recorded = timestamp_column(batch, "recorded_at")?;
529 let versions = u32_column(batch, "contract_version")?;
530 for row in 0..batch.num_rows() {
531 assertions.push(Assertion {
532 assertion_uuid: uuid_at(ids, row, "assertion_uuid")?,
533 claim: required_text(claims, row, "claim")?.to_owned(),
534 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
535 recorded_at_micros: required_i64(recorded, row, "recorded_at")?,
536 contract_version: required_u32(versions, row, "contract_version")?,
537 });
538 }
539 }
540 let mut refs = Vec::new();
541 for batch in graph_ref_batches {
542 require_schema(
543 batch,
544 &ASSERTION_GRAPH_REF_SCHEMA,
545 "assertion_graph_ref.schema",
546 )?;
547 let assertions_col = fixed_column(batch, "assertion_uuid")?;
548 let graph_ids = fixed_column(batch, "graph_uuid")?;
549 let kinds = string_column(batch, "graph_kind")?;
550 let roles = string_column(batch, "role")?;
551 let ordinals = u32_column(batch, "ordinal")?;
552 let versions = u32_column(batch, "contract_version")?;
553 for row in 0..batch.num_rows() {
554 refs.push(AssertionGraphRef {
555 assertion_uuid: uuid_at(assertions_col, row, "assertion_uuid")?,
556 graph_uuid: uuid_at(graph_ids, row, "graph_uuid")?,
557 graph_kind: GraphObjectKind::parse(required_text(kinds, row, "graph_kind")?)?,
558 role: AssertionGraphRole::parse(required_text(roles, row, "role")?)?,
559 ordinal: required_u32(ordinals, row, "ordinal")?,
560 contract_version: required_u32(versions, row, "contract_version")?,
561 });
562 }
563 }
564 Self::new(assertions, refs)
565 }
566}
567
568#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
570#[serde(rename_all = "snake_case")]
571pub enum EvidenceSourceKind {
572 Document,
574 Observation,
576 GraphNode,
578 GraphEdge,
580}
581
582impl EvidenceSourceKind {
583 #[must_use]
585 pub const fn as_str(self) -> &'static str {
586 match self {
587 Self::Document => "document",
588 Self::Observation => "observation",
589 Self::GraphNode => "graph_node",
590 Self::GraphEdge => "graph_edge",
591 }
592 }
593
594 fn parse(value: &str) -> Result<Self, KnowledgeError> {
595 match value {
596 "document" => Ok(Self::Document),
597 "observation" => Ok(Self::Observation),
598 "graph_node" => Ok(Self::GraphNode),
599 "graph_edge" => Ok(Self::GraphEdge),
600 _ => Err(invalid("source_kind", "unknown closed value")),
601 }
602 }
603}
604
605#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
607#[serde(rename_all = "snake_case")]
608pub enum EvidenceRole {
609 Supports,
611 Contradicts,
613 Context,
615}
616
617impl EvidenceRole {
618 #[must_use]
620 pub const fn as_str(self) -> &'static str {
621 match self {
622 Self::Supports => "supports",
623 Self::Contradicts => "contradicts",
624 Self::Context => "context",
625 }
626 }
627
628 fn parse(value: &str) -> Result<Self, KnowledgeError> {
629 match value {
630 "supports" => Ok(Self::Supports),
631 "contradicts" => Ok(Self::Contradicts),
632 "context" => Ok(Self::Context),
633 _ => Err(invalid("role", "unknown closed value")),
634 }
635 }
636}
637
638#[derive(Clone, Debug, PartialEq)]
640pub struct EvidenceLink {
641 pub evidence_uuid: Uuid,
643 pub assertion_uuid: Uuid,
645 pub source_uuid: Uuid,
647 pub source_kind: EvidenceSourceKind,
649 pub role: EvidenceRole,
651 pub weight: Option<f64>,
653 pub provenance_uuid: Uuid,
655 pub recorded_at_micros: i64,
657 pub contract_version: u32,
659}
660
661impl EvidenceLink {
662 #[allow(clippy::too_many_arguments)]
664 pub fn new(
665 evidence_uuid: Uuid,
666 assertion_uuid: Uuid,
667 source_uuid: Uuid,
668 source_kind: EvidenceSourceKind,
669 role: EvidenceRole,
670 weight: Option<f64>,
671 provenance_uuid: Uuid,
672 recorded_at_micros: i64,
673 ) -> Result<Self, KnowledgeError> {
674 require_v7(evidence_uuid, "evidence_uuid")?;
675 require_v7(assertion_uuid, "assertion_uuid")?;
676 require_uuid(source_uuid, "source_uuid")?;
677 require_uuid(provenance_uuid, "provenance_uuid")?;
678 validate_confidence(weight, "weight")?;
679 let weight = weight.map(normalize_zero);
680 Ok(Self {
681 evidence_uuid,
682 assertion_uuid,
683 source_uuid,
684 source_kind,
685 role,
686 weight,
687 provenance_uuid,
688 recorded_at_micros,
689 contract_version: EVIDENCE_LINK_CONTRACT_VERSION,
690 })
691 }
692}
693
694#[derive(Clone, Debug, Default, PartialEq)]
696pub struct EvidenceLedger {
697 pub links: Vec<EvidenceLink>,
699}
700
701impl EvidenceLedger {
702 pub fn new(mut links: Vec<EvidenceLink>) -> Result<Self, KnowledgeError> {
704 if links.len() > MAX_KNOWLEDGE_ROWS {
705 return Err(KnowledgeError::Limit {
706 participant: "evidence",
707 observed: links.len(),
708 limit: MAX_KNOWLEDGE_ROWS,
709 });
710 }
711 let mut ids = HashSet::new();
712 for link in &links {
713 require_v7(link.evidence_uuid, "evidence_uuid")?;
714 require_v7(link.assertion_uuid, "assertion_uuid")?;
715 require_uuid(link.source_uuid, "source_uuid")?;
716 require_uuid(link.provenance_uuid, "provenance_uuid")?;
717 if link.contract_version != EVIDENCE_LINK_CONTRACT_VERSION {
718 return Err(invalid("contract_version", "unsupported evidence version"));
719 }
720 if !ids.insert(link.evidence_uuid) {
721 return Err(KnowledgeError::Duplicate("evidence_uuid"));
722 }
723 if let Some(weight) = link.weight {
724 validate_confidence(Some(weight), "weight")?;
725 }
726 }
727 links.sort_by_key(|row| (row.recorded_at_micros, row.evidence_uuid));
728 Ok(Self { links })
729 }
730
731 pub fn merge(&self, staged: &Self) -> Result<Self, KnowledgeError> {
733 let mut links = self.links.clone();
734 for row in &staged.links {
735 if let Some(existing) = links
736 .iter()
737 .find(|existing| existing.evidence_uuid == row.evidence_uuid)
738 {
739 if existing != row {
740 return Err(KnowledgeError::Conflict("evidence_uuid"));
741 }
742 } else {
743 links.push(row.clone());
744 }
745 }
746 Self::new(links)
747 }
748
749 pub fn evidence_fingerprint(&self, evidence_uuid: Uuid) -> Result<[u8; 32], KnowledgeError> {
751 let row = self
752 .links
753 .iter()
754 .find(|row| row.evidence_uuid == evidence_uuid)
755 .ok_or(KnowledgeError::Dangling("evidence_uuid"))?;
756 let mut writer = CanonicalWriter::new();
757 writer.raw(b"GFEV")?;
758 writer.u32(EVIDENCE_LINK_CONTRACT_VERSION)?;
759 writer.raw(row.assertion_uuid.as_bytes())?;
760 writer.raw(row.source_uuid.as_bytes())?;
761 writer.text(row.source_kind.as_str())?;
762 writer.text(row.role.as_str())?;
763 canonical_optional_f64(&mut writer, row.weight)?;
764 Ok(fingerprint(
765 CanonicalDomain::EvidenceLink,
766 CANONICAL_CONTRACT_VERSION,
767 &writer.finish(),
768 )?)
769 }
770
771 pub fn batch(&self) -> Result<RecordBatch, KnowledgeError> {
773 evidence_batch(&self.links)
774 }
775
776 pub fn from_batches(batches: &[RecordBatch]) -> Result<Self, KnowledgeError> {
778 let mut links = Vec::new();
779 for batch in batches {
780 require_schema(batch, &EVIDENCE_LINK_SCHEMA, "evidence.schema")?;
781 let ids = fixed_column(batch, "evidence_uuid")?;
782 let assertions = fixed_column(batch, "assertion_uuid")?;
783 let sources = fixed_column(batch, "source_uuid")?;
784 let kinds = string_column(batch, "source_kind")?;
785 let roles = string_column(batch, "role")?;
786 let weights = f64_column(batch, "weight")?;
787 let provenance = fixed_column(batch, "provenance_uuid")?;
788 let recorded = timestamp_column(batch, "recorded_at")?;
789 let versions = u32_column(batch, "contract_version")?;
790 for row in 0..batch.num_rows() {
791 links.push(EvidenceLink {
792 evidence_uuid: uuid_at(ids, row, "evidence_uuid")?,
793 assertion_uuid: uuid_at(assertions, row, "assertion_uuid")?,
794 source_uuid: uuid_at(sources, row, "source_uuid")?,
795 source_kind: EvidenceSourceKind::parse(required_text(
796 kinds,
797 row,
798 "source_kind",
799 )?)?,
800 role: EvidenceRole::parse(required_text(roles, row, "role")?)?,
801 weight: optional_f64(weights, row),
802 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
803 recorded_at_micros: required_i64(recorded, row, "recorded_at")?,
804 contract_version: required_u32(versions, row, "contract_version")?,
805 });
806 }
807 }
808 Self::new(links)
809 }
810}
811
812#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
814#[serde(rename_all = "snake_case")]
815pub enum AlgorithmRunState {
816 Started,
818 Completed,
820 Failed,
822 Cancelled,
824 Interrupted,
826}
827
828impl AlgorithmRunState {
829 #[must_use]
831 pub const fn as_str(self) -> &'static str {
832 match self {
833 Self::Started => "started",
834 Self::Completed => "completed",
835 Self::Failed => "failed",
836 Self::Cancelled => "cancelled",
837 Self::Interrupted => "interrupted",
838 }
839 }
840
841 fn parse(value: &str) -> Result<Self, KnowledgeError> {
842 match value {
843 "started" => Ok(Self::Started),
844 "completed" => Ok(Self::Completed),
845 "failed" => Ok(Self::Failed),
846 "cancelled" => Ok(Self::Cancelled),
847 "interrupted" => Ok(Self::Interrupted),
848 _ => Err(invalid("state", "unknown closed value")),
849 }
850 }
851
852 #[must_use]
854 pub const fn is_terminal(self) -> bool {
855 !matches!(self, Self::Started)
856 }
857}
858
859#[derive(Clone, Debug, PartialEq, Eq)]
861pub struct AlgorithmRun {
862 pub run_uuid: Uuid,
864 pub algorithm: String,
866 pub algorithm_version: u32,
868 pub descriptor_version: u32,
870 pub descriptor: Vec<u8>,
872 pub projection_fingerprint: [u8; 32],
874 pub provenance_uuid: Uuid,
876 pub started_at_micros: i64,
878 pub contract_version: u32,
880}
881
882impl AlgorithmRun {
883 #[allow(clippy::too_many_arguments)]
885 pub fn new(
886 run_uuid: Uuid,
887 algorithm: String,
888 algorithm_version: u32,
889 descriptor_version: u32,
890 descriptor: Vec<u8>,
891 projection_fingerprint: [u8; 32],
892 provenance_uuid: Uuid,
893 started_at_micros: i64,
894 ) -> Result<Self, KnowledgeError> {
895 let row = Self {
896 run_uuid,
897 algorithm,
898 algorithm_version,
899 descriptor_version,
900 descriptor,
901 projection_fingerprint,
902 provenance_uuid,
903 started_at_micros,
904 contract_version: ALGORITHM_RUN_CONTRACT_VERSION,
905 };
906 validate_algorithm_run(&row)?;
907 Ok(row)
908 }
909}
910
911#[derive(Clone, Debug, PartialEq, Eq)]
913pub struct AlgorithmRunEvent {
914 pub event_uuid: Uuid,
916 pub run_uuid: Uuid,
918 pub state: AlgorithmRunState,
920 pub result_fingerprint: Option<[u8; 32]>,
922 pub error_code: Option<String>,
924 pub recorded_at_micros: i64,
926 pub provenance_uuid: Uuid,
928 pub contract_version: u32,
930}
931
932impl AlgorithmRunEvent {
933 #[allow(clippy::too_many_arguments)]
935 pub fn new(
936 event_uuid: Uuid,
937 run_uuid: Uuid,
938 state: AlgorithmRunState,
939 result_fingerprint: Option<[u8; 32]>,
940 error_code: Option<String>,
941 recorded_at_micros: i64,
942 provenance_uuid: Uuid,
943 ) -> Result<Self, KnowledgeError> {
944 let row = Self {
945 event_uuid,
946 run_uuid,
947 state,
948 result_fingerprint,
949 error_code,
950 recorded_at_micros,
951 provenance_uuid,
952 contract_version: ALGORITHM_RUN_EVENT_CONTRACT_VERSION,
953 };
954 validate_algorithm_run_event(&row)?;
955 Ok(row)
956 }
957}
958
959#[derive(Clone, Debug, Default, PartialEq, Eq)]
961pub struct AlgorithmRunLedger {
962 pub runs: Vec<AlgorithmRun>,
964 pub events: Vec<AlgorithmRunEvent>,
966}
967
968impl AlgorithmRunLedger {
969 pub fn new(
971 mut runs: Vec<AlgorithmRun>,
972 mut events: Vec<AlgorithmRunEvent>,
973 ) -> Result<Self, KnowledgeError> {
974 runs.sort_by_key(|row| (row.started_at_micros, row.run_uuid));
975 events.sort_by_key(|row| (row.recorded_at_micros, row.event_uuid));
976 validate_algorithm_run_rows(&runs, &events)?;
977 Ok(Self { runs, events })
978 }
979
980 pub fn merge(&self, staged: &Self) -> Result<Self, KnowledgeError> {
982 let mut runs = self.runs.clone();
983 for row in &staged.runs {
984 match runs.iter().find(|current| current.run_uuid == row.run_uuid) {
985 Some(current) if current == row => {}
986 Some(_) => return Err(KnowledgeError::Conflict("run_uuid")),
987 None => runs.push(row.clone()),
988 }
989 }
990 let mut events = self.events.clone();
991 for row in &staged.events {
992 match events
993 .iter()
994 .find(|current| current.event_uuid == row.event_uuid)
995 {
996 Some(current) if current == row => {}
997 Some(_) => return Err(KnowledgeError::Conflict("event_uuid")),
998 None => events.push(row.clone()),
999 }
1000 }
1001 Self::new(runs, events)
1002 }
1003
1004 #[must_use]
1006 pub fn run(&self, run_uuid: Uuid) -> Option<&AlgorithmRun> {
1007 self.runs.iter().find(|row| row.run_uuid == run_uuid)
1008 }
1009
1010 #[must_use]
1012 pub fn events_for(&self, run_uuid: Uuid) -> Vec<AlgorithmRunEvent> {
1013 self.events
1014 .iter()
1015 .filter(|row| row.run_uuid == run_uuid)
1016 .cloned()
1017 .collect()
1018 }
1019
1020 #[must_use]
1022 pub fn terminal_event(&self, run_uuid: Uuid) -> Option<&AlgorithmRunEvent> {
1023 self.events
1024 .iter()
1025 .find(|row| row.run_uuid == run_uuid && row.state.is_terminal())
1026 }
1027
1028 pub fn run_batch(&self) -> Result<RecordBatch, KnowledgeError> {
1030 algorithm_run_batch(&self.runs)
1031 }
1032
1033 pub fn event_batch(&self) -> Result<RecordBatch, KnowledgeError> {
1035 algorithm_run_event_batch(&self.events)
1036 }
1037
1038 pub fn from_batches(
1040 run_batches: &[RecordBatch],
1041 event_batches: &[RecordBatch],
1042 ) -> Result<Self, KnowledgeError> {
1043 let mut runs = Vec::new();
1044 for batch in run_batches {
1045 require_schema(batch, &ALGORITHM_RUN_SCHEMA, "algorithm_runs")?;
1046 let ids = fixed_column(batch, "run_uuid")?;
1047 let algorithms = string_column(batch, "algorithm")?;
1048 let algorithm_versions = u32_column(batch, "algorithm_version")?;
1049 let descriptor_versions = u32_column(batch, "descriptor_version")?;
1050 let descriptors = binary_column(batch, "descriptor")?;
1051 let projections = fixed_column(batch, "projection_fingerprint")?;
1052 let provenance = fixed_column(batch, "provenance_uuid")?;
1053 let started = timestamp_column(batch, "started_at")?;
1054 let contracts = u32_column(batch, "contract_version")?;
1055 for row in 0..batch.num_rows() {
1056 runs.push(AlgorithmRun {
1057 run_uuid: uuid_at(ids, row, "run_uuid")?,
1058 algorithm: required_text(algorithms, row, "algorithm")?.to_owned(),
1059 algorithm_version: required_u32(algorithm_versions, row, "algorithm_version")?,
1060 descriptor_version: required_u32(
1061 descriptor_versions,
1062 row,
1063 "descriptor_version",
1064 )?,
1065 descriptor: required_binary(descriptors, row, "descriptor")?.to_vec(),
1066 projection_fingerprint: fixed_32_at(
1067 projections,
1068 row,
1069 "projection_fingerprint",
1070 )?,
1071 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
1072 started_at_micros: required_i64(started, row, "started_at")?,
1073 contract_version: required_u32(contracts, row, "contract_version")?,
1074 });
1075 }
1076 }
1077 let mut events = Vec::new();
1078 for batch in event_batches {
1079 require_schema(batch, &ALGORITHM_RUN_EVENT_SCHEMA, "algorithm_run_events")?;
1080 let ids = fixed_column(batch, "event_uuid")?;
1081 let runs_column = fixed_column(batch, "run_uuid")?;
1082 let states = string_column(batch, "state")?;
1083 let results = fixed_column(batch, "result_fingerprint")?;
1084 let errors = string_column(batch, "error_code")?;
1085 let recorded = timestamp_column(batch, "recorded_at")?;
1086 let provenance = fixed_column(batch, "provenance_uuid")?;
1087 let contracts = u32_column(batch, "contract_version")?;
1088 for row in 0..batch.num_rows() {
1089 events.push(AlgorithmRunEvent {
1090 event_uuid: uuid_at(ids, row, "event_uuid")?,
1091 run_uuid: uuid_at(runs_column, row, "run_uuid")?,
1092 state: AlgorithmRunState::parse(required_text(states, row, "state")?)?,
1093 result_fingerprint: optional_fixed_32(results, row, "result_fingerprint")?,
1094 error_code: optional_text(errors, row),
1095 recorded_at_micros: required_i64(recorded, row, "recorded_at")?,
1096 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
1097 contract_version: required_u32(contracts, row, "contract_version")?,
1098 });
1099 }
1100 }
1101 Self::new(runs, events)
1102 }
1103}
1104
1105#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
1107#[serde(rename_all = "snake_case")]
1108pub enum ConfidencePolicy {
1109 Explicit,
1111 ConservativeMin,
1113}
1114
1115impl ConfidencePolicy {
1116 #[must_use]
1118 pub const fn as_str(self) -> &'static str {
1119 match self {
1120 Self::Explicit => "explicit",
1121 Self::ConservativeMin => "conservative_min",
1122 }
1123 }
1124
1125 fn parse(value: &str) -> Result<Self, KnowledgeError> {
1126 match value {
1127 "explicit" => Ok(Self::Explicit),
1128 "conservative_min" => Ok(Self::ConservativeMin),
1129 _ => Err(invalid("policy", "unknown closed value")),
1130 }
1131 }
1132}
1133
1134#[derive(Clone, Debug, PartialEq)]
1136pub struct ConfidenceAssessment {
1137 pub confidence_uuid: Uuid,
1139 pub assertion_uuid: Uuid,
1141 pub policy: ConfidencePolicy,
1143 pub policy_version: u32,
1145 pub value: Option<f64>,
1147 pub provenance_uuid: Uuid,
1149 pub recorded_at_micros: i64,
1151 pub contract_version: u32,
1153}
1154
1155impl ConfidenceAssessment {
1156 pub fn new(
1158 confidence_uuid: Uuid,
1159 assertion_uuid: Uuid,
1160 policy: ConfidencePolicy,
1161 value: Option<f64>,
1162 provenance_uuid: Uuid,
1163 recorded_at_micros: i64,
1164 ) -> Result<Self, KnowledgeError> {
1165 require_v7(confidence_uuid, "confidence_uuid")?;
1166 require_v7(assertion_uuid, "assertion_uuid")?;
1167 require_uuid(provenance_uuid, "provenance_uuid")?;
1168 validate_confidence(value, "value")?;
1169 Ok(Self {
1170 confidence_uuid,
1171 assertion_uuid,
1172 policy,
1173 policy_version: 1,
1174 value: value.map(normalize_zero),
1175 provenance_uuid,
1176 recorded_at_micros,
1177 contract_version: CONFIDENCE_ASSESSMENT_CONTRACT_VERSION,
1178 })
1179 }
1180}
1181
1182#[derive(Clone, Debug, PartialEq)]
1184pub struct ConfidenceInput {
1185 pub confidence_uuid: Uuid,
1187 pub input_confidence_uuid: Uuid,
1189 pub input_value: Option<f64>,
1191 pub ordinal: u32,
1193 pub contract_version: u32,
1195}
1196
1197impl ConfidenceInput {
1198 pub fn new(
1200 confidence_uuid: Uuid,
1201 input_confidence_uuid: Uuid,
1202 input_value: Option<f64>,
1203 ordinal: u32,
1204 ) -> Result<Self, KnowledgeError> {
1205 require_v7(confidence_uuid, "confidence_uuid")?;
1206 require_v7(input_confidence_uuid, "input_confidence_uuid")?;
1207 validate_confidence(input_value, "input_value")?;
1208 Ok(Self {
1209 confidence_uuid,
1210 input_confidence_uuid,
1211 input_value: input_value.map(normalize_zero),
1212 ordinal,
1213 contract_version: CONFIDENCE_INPUT_CONTRACT_VERSION,
1214 })
1215 }
1216}
1217
1218#[derive(Clone, Debug, Default, PartialEq)]
1220pub struct ConfidenceLedger {
1221 pub assessments: Vec<ConfidenceAssessment>,
1223 pub inputs: Vec<ConfidenceInput>,
1225}
1226
1227impl ConfidenceLedger {
1228 pub fn new(
1230 mut assessments: Vec<ConfidenceAssessment>,
1231 mut inputs: Vec<ConfidenceInput>,
1232 ) -> Result<Self, KnowledgeError> {
1233 inputs.sort_by_key(|row| (row.confidence_uuid, row.ordinal, row.input_confidence_uuid));
1234 validate_confidence_rows(&assessments, &inputs)?;
1235 let times = assessments
1236 .iter()
1237 .map(|row| (row.confidence_uuid, row.recorded_at_micros))
1238 .collect::<HashMap<_, _>>();
1239 assessments.sort_by_key(|row| (row.recorded_at_micros, row.confidence_uuid));
1240 inputs.sort_by_key(|row| {
1241 (
1242 times[&row.confidence_uuid],
1243 row.confidence_uuid,
1244 row.ordinal,
1245 row.input_confidence_uuid,
1246 )
1247 });
1248 Ok(Self {
1249 assessments,
1250 inputs,
1251 })
1252 }
1253
1254 pub fn explicit(
1256 confidence_uuid: Uuid,
1257 assertion_uuid: Uuid,
1258 value: f64,
1259 provenance_uuid: Uuid,
1260 recorded_at_micros: i64,
1261 ) -> Result<Self, KnowledgeError> {
1262 Self::new(
1263 vec![ConfidenceAssessment::new(
1264 confidence_uuid,
1265 assertion_uuid,
1266 ConfidencePolicy::Explicit,
1267 Some(value),
1268 provenance_uuid,
1269 recorded_at_micros,
1270 )?],
1271 vec![],
1272 )
1273 }
1274
1275 pub fn conservative_min(
1277 &self,
1278 confidence_uuid: Uuid,
1279 assertion_uuid: Uuid,
1280 mut requested: Vec<Uuid>,
1281 provenance_uuid: Uuid,
1282 recorded_at_micros: i64,
1283 ) -> Result<Self, KnowledgeError> {
1284 requested.sort_unstable();
1285 if requested.windows(2).any(|pair| pair[0] == pair[1]) {
1286 return Err(KnowledgeError::Duplicate("input_confidence_uuid"));
1287 }
1288 let values = self
1289 .assessments
1290 .iter()
1291 .map(|row| (row.confidence_uuid, row.value))
1292 .collect::<HashMap<_, _>>();
1293 let mut minimum = None;
1294 let mut complete = !requested.is_empty();
1295 let mut inputs = Vec::with_capacity(requested.len());
1296 for (ordinal, input_uuid) in requested.into_iter().enumerate() {
1297 require_v7(input_uuid, "input_confidence_uuid")?;
1298 let observed = values.get(&input_uuid).copied().flatten();
1299 if let Some(value) = observed {
1300 minimum = Some(minimum.map_or(value, |current: f64| current.min(value)));
1301 } else {
1302 complete = false;
1303 }
1304 inputs.push(ConfidenceInput::new(
1305 confidence_uuid,
1306 input_uuid,
1307 observed,
1308 u32::try_from(ordinal).map_err(|_| KnowledgeError::Limit {
1309 participant: "confidence_inputs",
1310 observed: ordinal,
1311 limit: u32::MAX as usize,
1312 })?,
1313 )?);
1314 }
1315 Self::new(
1316 vec![ConfidenceAssessment::new(
1317 confidence_uuid,
1318 assertion_uuid,
1319 ConfidencePolicy::ConservativeMin,
1320 complete.then_some(minimum).flatten(),
1321 provenance_uuid,
1322 recorded_at_micros,
1323 )?],
1324 inputs,
1325 )
1326 }
1327
1328 pub fn merge(&self, staged: &Self) -> Result<Self, KnowledgeError> {
1330 let mut assessments = self.assessments.clone();
1331 let mut inputs = self.inputs.clone();
1332 for row in &staged.assessments {
1333 if let Some(existing) = assessments
1334 .iter()
1335 .find(|existing| existing.confidence_uuid == row.confidence_uuid)
1336 {
1337 if existing != row
1338 || inputs_for(&inputs, row.confidence_uuid)
1339 != inputs_for(&staged.inputs, row.confidence_uuid)
1340 {
1341 return Err(KnowledgeError::Conflict("confidence_uuid"));
1342 }
1343 } else {
1344 assessments.push(row.clone());
1345 inputs.extend(
1346 staged
1347 .inputs
1348 .iter()
1349 .filter(|input| input.confidence_uuid == row.confidence_uuid)
1350 .cloned(),
1351 );
1352 }
1353 }
1354 Self::new(assessments, inputs)
1355 }
1356
1357 pub fn assessment_fingerprint(
1359 &self,
1360 confidence_uuid: Uuid,
1361 ) -> Result<[u8; 32], KnowledgeError> {
1362 let row = self
1363 .assessments
1364 .iter()
1365 .find(|row| row.confidence_uuid == confidence_uuid)
1366 .ok_or(KnowledgeError::Dangling("confidence_uuid"))?;
1367 let inputs = inputs_for(&self.inputs, confidence_uuid);
1368 let mut writer = CanonicalWriter::new();
1369 writer.raw(b"GFCA")?;
1370 writer.u32(CONFIDENCE_ASSESSMENT_CONTRACT_VERSION)?;
1371 writer.raw(row.assertion_uuid.as_bytes())?;
1372 writer.text(row.policy.as_str())?;
1373 writer.u32(row.policy_version)?;
1374 canonical_optional_f64(&mut writer, row.value)?;
1375 writer.u64(inputs.len() as u64)?;
1376 for input in inputs {
1377 writer.raw(input.input_confidence_uuid.as_bytes())?;
1378 canonical_optional_f64(&mut writer, input.input_value)?;
1379 writer.u32(input.ordinal)?;
1380 }
1381 Ok(fingerprint(
1382 CanonicalDomain::ConfidenceAssessment,
1383 CANONICAL_CONTRACT_VERSION,
1384 &writer.finish(),
1385 )?)
1386 }
1387
1388 pub fn assessment_batch(&self) -> Result<RecordBatch, KnowledgeError> {
1390 confidence_assessment_batch(&self.assessments)
1391 }
1392
1393 pub fn input_batch(&self) -> Result<RecordBatch, KnowledgeError> {
1395 confidence_input_batch(&self.inputs)
1396 }
1397
1398 pub fn from_batches(
1400 assessment_batches: &[RecordBatch],
1401 input_batches: &[RecordBatch],
1402 ) -> Result<Self, KnowledgeError> {
1403 let mut assessments = Vec::new();
1404 for batch in assessment_batches {
1405 require_schema(batch, &CONFIDENCE_ASSESSMENT_SCHEMA, "confidence.schema")?;
1406 let ids = fixed_column(batch, "confidence_uuid")?;
1407 let assertions = fixed_column(batch, "assertion_uuid")?;
1408 let policies = string_column(batch, "policy")?;
1409 let policy_versions = u32_column(batch, "policy_version")?;
1410 let values = f64_column(batch, "value")?;
1411 let provenance = fixed_column(batch, "provenance_uuid")?;
1412 let recorded = timestamp_column(batch, "recorded_at")?;
1413 let versions = u32_column(batch, "contract_version")?;
1414 for row in 0..batch.num_rows() {
1415 assessments.push(ConfidenceAssessment {
1416 confidence_uuid: uuid_at(ids, row, "confidence_uuid")?,
1417 assertion_uuid: uuid_at(assertions, row, "assertion_uuid")?,
1418 policy: ConfidencePolicy::parse(required_text(policies, row, "policy")?)?,
1419 policy_version: required_u32(policy_versions, row, "policy_version")?,
1420 value: optional_f64(values, row),
1421 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
1422 recorded_at_micros: required_i64(recorded, row, "recorded_at")?,
1423 contract_version: required_u32(versions, row, "contract_version")?,
1424 });
1425 }
1426 }
1427 let mut inputs = Vec::new();
1428 for batch in input_batches {
1429 require_schema(batch, &CONFIDENCE_INPUT_SCHEMA, "confidence_input.schema")?;
1430 let owners = fixed_column(batch, "confidence_uuid")?;
1431 let ids = fixed_column(batch, "input_confidence_uuid")?;
1432 let values = f64_column(batch, "input_value")?;
1433 let ordinals = u32_column(batch, "ordinal")?;
1434 let versions = u32_column(batch, "contract_version")?;
1435 for row in 0..batch.num_rows() {
1436 inputs.push(ConfidenceInput {
1437 confidence_uuid: uuid_at(owners, row, "confidence_uuid")?,
1438 input_confidence_uuid: uuid_at(ids, row, "input_confidence_uuid")?,
1439 input_value: optional_f64(values, row),
1440 ordinal: required_u32(ordinals, row, "ordinal")?,
1441 contract_version: required_u32(versions, row, "contract_version")?,
1442 });
1443 }
1444 }
1445 Self::new(assessments, inputs)
1446 }
1447}
1448
1449#[derive(Clone, Debug)]
1451pub struct SchemaRegistryEntry {
1452 pub capability_id: &'static str,
1454 pub capability_version: u32,
1456 pub record_family: &'static str,
1458 pub record_version: u32,
1460 pub schema: SchemaRef,
1462 pub schema_fingerprint: [u8; 32],
1464 pub enum_registry_versions: &'static [(&'static str, u32)],
1466 pub sort_key: &'static [&'static str],
1468 pub diff_identity_fields: &'static [&'static str],
1470 pub diff_record_uuid_field: Option<&'static str>,
1472 pub fingerprint_domain: CanonicalDomain,
1474 pub owner: &'static str,
1476 pub implementation_issue: u64,
1478 pub max_rows: usize,
1480}
1481
1482impl SchemaRegistryEntry {
1483 #[must_use]
1485 pub const fn diff_identity_fingerprint_domain(&self) -> CanonicalDomain {
1486 CanonicalDomain::ArrowResult
1487 }
1488
1489 #[must_use]
1491 pub const fn diff_record_fingerprint_domain(&self) -> CanonicalDomain {
1492 self.fingerprint_domain
1493 }
1494}
1495
1496#[must_use]
1498pub fn schema_registry() -> Vec<SchemaRegistryEntry> {
1499 let mut entries = base_schema_registry_entries();
1500 entries.extend(algorithm_run_schema_entries());
1501 entries.push(reasoning::schema_registry_entry());
1502 entries.push(status::schema_registry_entry());
1503 entries.push(supersession::schema_registry_entry());
1504 entries.extend(hypothesis::schema_registry_entries());
1505 entries.push(valid_time::schema_registry_entry());
1506 entries.push(belief_projection::schema_registry_entry());
1507 entries
1508}
1509
1510fn base_schema_registry_entries() -> Vec<SchemaRegistryEntry> {
1511 vec![
1512 SchemaRegistryEntry {
1513 capability_id: "knowledge",
1514 capability_version: KNOWLEDGE_CAPABILITY_VERSION,
1515 record_family: "assertions",
1516 record_version: ASSERTION_CONTRACT_VERSION,
1517 schema: Arc::clone(&ASSERTION_SCHEMA),
1518 schema_fingerprint: *ASSERTION_SCHEMA_FINGERPRINT,
1519 enum_registry_versions: &[],
1520 sort_key: &["recorded_at", "assertion_uuid"],
1521 diff_identity_fields: &["assertion_uuid"],
1522 diff_record_uuid_field: Some("assertion_uuid"),
1523 fingerprint_domain: CanonicalDomain::Assertion,
1524 owner: "graphforge-knowledge",
1525 implementation_issue: 2411,
1526 max_rows: MAX_KNOWLEDGE_ROWS,
1527 },
1528 SchemaRegistryEntry {
1529 capability_id: "knowledge",
1530 capability_version: KNOWLEDGE_CAPABILITY_VERSION,
1531 record_family: "assertion_graph_refs",
1532 record_version: ASSERTION_GRAPH_REF_CONTRACT_VERSION,
1533 schema: Arc::clone(&ASSERTION_GRAPH_REF_SCHEMA),
1534 schema_fingerprint: *ASSERTION_GRAPH_REF_SCHEMA_FINGERPRINT,
1535 enum_registry_versions: &[
1536 ("graph_kind", GRAPH_OBJECT_KIND_REGISTRY_VERSION),
1537 ("role", ASSERTION_GRAPH_ROLE_REGISTRY_VERSION),
1538 ],
1539 sort_key: &[
1540 "assertion_uuid",
1541 "role",
1542 "ordinal",
1543 "graph_kind",
1544 "graph_uuid",
1545 ],
1546 diff_identity_fields: &["assertion_uuid", "graph_uuid", "role", "ordinal"],
1547 diff_record_uuid_field: None,
1548 fingerprint_domain: CanonicalDomain::Assertion,
1549 owner: "graphforge-knowledge",
1550 implementation_issue: 2411,
1551 max_rows: MAX_KNOWLEDGE_ROWS,
1552 },
1553 SchemaRegistryEntry {
1554 capability_id: "knowledge",
1555 capability_version: KNOWLEDGE_CAPABILITY_VERSION,
1556 record_family: "confidence_assessments",
1557 record_version: CONFIDENCE_ASSESSMENT_CONTRACT_VERSION,
1558 schema: Arc::clone(&CONFIDENCE_ASSESSMENT_SCHEMA),
1559 schema_fingerprint: *CONFIDENCE_ASSESSMENT_SCHEMA_FINGERPRINT,
1560 enum_registry_versions: &[("confidence_policy", CONFIDENCE_POLICY_REGISTRY_VERSION)],
1561 sort_key: &["recorded_at", "confidence_uuid"],
1562 diff_identity_fields: &["confidence_uuid"],
1563 diff_record_uuid_field: Some("confidence_uuid"),
1564 fingerprint_domain: CanonicalDomain::ConfidenceAssessment,
1565 owner: "graphforge-knowledge",
1566 implementation_issue: 774,
1567 max_rows: MAX_KNOWLEDGE_ROWS,
1568 },
1569 SchemaRegistryEntry {
1570 capability_id: "knowledge",
1571 capability_version: KNOWLEDGE_CAPABILITY_VERSION,
1572 record_family: "confidence_inputs",
1573 record_version: CONFIDENCE_INPUT_CONTRACT_VERSION,
1574 schema: Arc::clone(&CONFIDENCE_INPUT_SCHEMA),
1575 schema_fingerprint: *CONFIDENCE_INPUT_SCHEMA_FINGERPRINT,
1576 enum_registry_versions: &[],
1577 sort_key: &["confidence_uuid", "ordinal", "input_confidence_uuid"],
1578 diff_identity_fields: &["confidence_uuid", "input_confidence_uuid"],
1579 diff_record_uuid_field: None,
1580 fingerprint_domain: CanonicalDomain::ConfidenceAssessment,
1581 owner: "graphforge-knowledge",
1582 implementation_issue: 774,
1583 max_rows: MAX_KNOWLEDGE_ROWS,
1584 },
1585 SchemaRegistryEntry {
1586 capability_id: "knowledge",
1587 capability_version: KNOWLEDGE_CAPABILITY_VERSION,
1588 record_family: "evidence",
1589 record_version: EVIDENCE_LINK_CONTRACT_VERSION,
1590 schema: Arc::clone(&EVIDENCE_LINK_SCHEMA),
1591 schema_fingerprint: *EVIDENCE_LINK_SCHEMA_FINGERPRINT,
1592 enum_registry_versions: &[
1593 (
1594 "evidence_source_kind",
1595 EVIDENCE_SOURCE_KIND_REGISTRY_VERSION,
1596 ),
1597 ("evidence_role", EVIDENCE_ROLE_REGISTRY_VERSION),
1598 ],
1599 sort_key: &["recorded_at", "evidence_uuid"],
1600 diff_identity_fields: &["evidence_uuid"],
1601 diff_record_uuid_field: Some("evidence_uuid"),
1602 fingerprint_domain: CanonicalDomain::EvidenceLink,
1603 owner: "graphforge-knowledge",
1604 implementation_issue: 775,
1605 max_rows: MAX_KNOWLEDGE_ROWS,
1606 },
1607 ]
1608}
1609
1610fn algorithm_run_schema_entries() -> [SchemaRegistryEntry; 2] {
1611 [
1612 SchemaRegistryEntry {
1613 capability_id: "knowledge",
1614 capability_version: KNOWLEDGE_CAPABILITY_VERSION,
1615 record_family: "algorithm_runs",
1616 record_version: ALGORITHM_RUN_CONTRACT_VERSION,
1617 schema: Arc::clone(&ALGORITHM_RUN_SCHEMA),
1618 schema_fingerprint: *ALGORITHM_RUN_SCHEMA_FINGERPRINT,
1619 enum_registry_versions: &[],
1620 sort_key: &["started_at", "run_uuid"],
1621 diff_identity_fields: &["run_uuid"],
1622 diff_record_uuid_field: Some("run_uuid"),
1623 fingerprint_domain: CanonicalDomain::InvocationDescriptor,
1624 owner: "graphforge-knowledge",
1625 implementation_issue: 2003,
1626 max_rows: MAX_KNOWLEDGE_ROWS,
1627 },
1628 SchemaRegistryEntry {
1629 capability_id: "knowledge",
1630 capability_version: KNOWLEDGE_CAPABILITY_VERSION,
1631 record_family: "algorithm_run_events",
1632 record_version: ALGORITHM_RUN_EVENT_CONTRACT_VERSION,
1633 schema: Arc::clone(&ALGORITHM_RUN_EVENT_SCHEMA),
1634 schema_fingerprint: *ALGORITHM_RUN_EVENT_SCHEMA_FINGERPRINT,
1635 enum_registry_versions: &[(
1636 "algorithm_run_state",
1637 ALGORITHM_RUN_STATE_REGISTRY_VERSION,
1638 )],
1639 sort_key: &["recorded_at", "event_uuid"],
1640 diff_identity_fields: &["event_uuid"],
1641 diff_record_uuid_field: Some("event_uuid"),
1642 fingerprint_domain: CanonicalDomain::ArrowResult,
1643 owner: "graphforge-knowledge",
1644 implementation_issue: 2003,
1645 max_rows: MAX_KNOWLEDGE_ROWS,
1646 },
1647 ]
1648}
1649
1650#[derive(thiserror::Error, Debug)]
1652pub enum KnowledgeError {
1653 #[error("invalid knowledge {field}: {message}")]
1655 Invalid {
1656 field: &'static str,
1658 message: &'static str,
1660 },
1661 #[error("knowledge {participant} row limit exceeded: observed {observed}, limit {limit}")]
1663 Limit {
1664 participant: &'static str,
1666 observed: usize,
1668 limit: usize,
1670 },
1671 #[error("duplicate knowledge identity: {0}")]
1673 Duplicate(&'static str),
1674 #[error("dangling knowledge reference: {0}")]
1676 Dangling(&'static str),
1677 #[error("knowledge idempotency conflict: {0}")]
1679 Conflict(&'static str),
1680 #[error("knowledge transaction conflict: {0}")]
1682 TransactionConflict(&'static str),
1683 #[error(transparent)]
1685 Canonical(#[from] CanonicalError),
1686 #[error("knowledge Arrow failure: {0}")]
1688 Arrow(#[from] arrow::error::ArrowError),
1689}
1690
1691impl KnowledgeError {
1692 #[must_use]
1694 pub const fn code(&self) -> &'static str {
1695 match self {
1696 Self::Invalid { .. } => "GF_KNOWLEDGE_INVALID",
1697 Self::Limit { .. } => "GF_RESOURCE_LIMIT",
1698 Self::Duplicate(_) => "GF_KNOWLEDGE_DUPLICATE",
1699 Self::Dangling(_) => "GF_KNOWLEDGE_DANGLING",
1700 Self::Conflict(_) => "GF_IDEMPOTENCY_CONFLICT",
1701 Self::TransactionConflict(_) => "GF_TRANSACTION_CONFLICT",
1702 Self::Canonical(error) => error.code(),
1703 Self::Arrow(_) => "GF_SCHEMA_MISMATCH",
1704 }
1705 }
1706}
1707
1708fn validate_rows(
1709 assertions: &[Assertion],
1710 refs: &[AssertionGraphRef],
1711) -> Result<(), KnowledgeError> {
1712 check_limit("assertions", assertions.len())?;
1713 check_limit("assertion_graph_refs", refs.len())?;
1714 let mut assertion_ids = HashSet::with_capacity(assertions.len());
1715 for assertion in assertions {
1716 if assertion.contract_version != ASSERTION_CONTRACT_VERSION {
1717 return Err(invalid("assertion.contract_version", "unsupported version"));
1718 }
1719 require_v7(assertion.assertion_uuid, "assertion_uuid")?;
1720 require_uuid(assertion.provenance_uuid, "provenance_uuid")?;
1721 validate_claim(&assertion.claim)?;
1722 if !assertion_ids.insert(assertion.assertion_uuid) {
1723 return Err(KnowledgeError::Duplicate("assertion_uuid"));
1724 }
1725 }
1726 let mut tuples = HashSet::with_capacity(refs.len());
1727 let mut role_ordinals: HashMap<(Uuid, AssertionGraphRole), Vec<u32>> = HashMap::new();
1728 let mut covered_assertions = HashSet::with_capacity(assertions.len());
1729 for reference in refs {
1730 if reference.contract_version != ASSERTION_GRAPH_REF_CONTRACT_VERSION {
1731 return Err(invalid(
1732 "assertion_graph_ref.contract_version",
1733 "unsupported version",
1734 ));
1735 }
1736 require_v7(reference.assertion_uuid, "assertion_uuid")?;
1737 require_uuid(reference.graph_uuid, "graph_uuid")?;
1738 if !assertion_ids.contains(&reference.assertion_uuid) {
1739 return Err(KnowledgeError::Dangling("assertion_uuid"));
1740 }
1741 if !tuples.insert((
1742 reference.assertion_uuid,
1743 reference.graph_uuid,
1744 reference.role,
1745 reference.ordinal,
1746 )) {
1747 return Err(KnowledgeError::Duplicate(
1748 "assertion_uuid/graph_uuid/role/ordinal",
1749 ));
1750 }
1751 role_ordinals
1752 .entry((reference.assertion_uuid, reference.role))
1753 .or_default()
1754 .push(reference.ordinal);
1755 covered_assertions.insert(reference.assertion_uuid);
1756 }
1757 for assertion_uuid in assertion_ids {
1758 if !covered_assertions.contains(&assertion_uuid) {
1759 return Err(KnowledgeError::Dangling("assertion.graph_refs"));
1760 }
1761 }
1762 for ordinals in role_ordinals.values_mut() {
1763 ordinals.sort_unstable();
1764 if ordinals
1765 .iter()
1766 .enumerate()
1767 .any(|(expected, actual)| usize::try_from(*actual) != Ok(expected))
1768 {
1769 return Err(invalid("ordinal", "must be contiguous from zero per role"));
1770 }
1771 }
1772 Ok(())
1773}
1774
1775fn assertion_batch(rows: &[Assertion]) -> Result<RecordBatch, KnowledgeError> {
1776 let mut ids = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1777 let mut provenance = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1778 for row in rows {
1779 ids.append_value(row.assertion_uuid.as_bytes())?;
1780 provenance.append_value(row.provenance_uuid.as_bytes())?;
1781 }
1782 RecordBatch::try_new(
1783 Arc::clone(&ASSERTION_SCHEMA),
1784 vec![
1785 Arc::new(ids.finish()),
1786 Arc::new(StringArray::from_iter_values(
1787 rows.iter().map(|row| row.claim.as_str()),
1788 )),
1789 Arc::new(provenance.finish()),
1790 Arc::new(
1791 TimestampMicrosecondArray::from_iter_values(
1792 rows.iter().map(|row| row.recorded_at_micros),
1793 )
1794 .with_timezone("UTC"),
1795 ),
1796 Arc::new(UInt32Array::from_iter_values(
1797 rows.iter().map(|row| row.contract_version),
1798 )),
1799 ],
1800 )
1801 .map_err(Into::into)
1802}
1803
1804fn graph_ref_batch(rows: &[AssertionGraphRef]) -> Result<RecordBatch, KnowledgeError> {
1805 let mut assertions = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1806 let mut graph_ids = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1807 for row in rows {
1808 assertions.append_value(row.assertion_uuid.as_bytes())?;
1809 graph_ids.append_value(row.graph_uuid.as_bytes())?;
1810 }
1811 RecordBatch::try_new(
1812 Arc::clone(&ASSERTION_GRAPH_REF_SCHEMA),
1813 vec![
1814 Arc::new(assertions.finish()),
1815 Arc::new(graph_ids.finish()),
1816 Arc::new(StringArray::from_iter_values(
1817 rows.iter().map(|row| row.graph_kind.as_str()),
1818 )),
1819 Arc::new(StringArray::from_iter_values(
1820 rows.iter().map(|row| row.role.as_str()),
1821 )),
1822 Arc::new(UInt32Array::from_iter_values(
1823 rows.iter().map(|row| row.ordinal),
1824 )),
1825 Arc::new(UInt32Array::from_iter_values(
1826 rows.iter().map(|row| row.contract_version),
1827 )),
1828 ],
1829 )
1830 .map_err(Into::into)
1831}
1832
1833fn confidence_assessment_batch(
1834 rows: &[ConfidenceAssessment],
1835) -> Result<RecordBatch, KnowledgeError> {
1836 let mut ids = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1837 let mut assertions = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1838 let mut provenance = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1839 for row in rows {
1840 ids.append_value(row.confidence_uuid.as_bytes())?;
1841 assertions.append_value(row.assertion_uuid.as_bytes())?;
1842 provenance.append_value(row.provenance_uuid.as_bytes())?;
1843 }
1844 RecordBatch::try_new(
1845 Arc::clone(&CONFIDENCE_ASSESSMENT_SCHEMA),
1846 vec![
1847 Arc::new(ids.finish()),
1848 Arc::new(assertions.finish()),
1849 Arc::new(StringArray::from_iter_values(
1850 rows.iter().map(|row| row.policy.as_str()),
1851 )),
1852 Arc::new(UInt32Array::from_iter_values(
1853 rows.iter().map(|row| row.policy_version),
1854 )),
1855 Arc::new(Float64Array::from(
1856 rows.iter().map(|row| row.value).collect::<Vec<_>>(),
1857 )),
1858 Arc::new(provenance.finish()),
1859 Arc::new(
1860 TimestampMicrosecondArray::from_iter_values(
1861 rows.iter().map(|row| row.recorded_at_micros),
1862 )
1863 .with_timezone("UTC"),
1864 ),
1865 Arc::new(UInt32Array::from_iter_values(
1866 rows.iter().map(|row| row.contract_version),
1867 )),
1868 ],
1869 )
1870 .map_err(Into::into)
1871}
1872
1873fn confidence_input_batch(rows: &[ConfidenceInput]) -> Result<RecordBatch, KnowledgeError> {
1874 let mut owners = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1875 let mut ids = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1876 for row in rows {
1877 owners.append_value(row.confidence_uuid.as_bytes())?;
1878 ids.append_value(row.input_confidence_uuid.as_bytes())?;
1879 }
1880 RecordBatch::try_new(
1881 Arc::clone(&CONFIDENCE_INPUT_SCHEMA),
1882 vec![
1883 Arc::new(owners.finish()),
1884 Arc::new(ids.finish()),
1885 Arc::new(Float64Array::from(
1886 rows.iter().map(|row| row.input_value).collect::<Vec<_>>(),
1887 )),
1888 Arc::new(UInt32Array::from_iter_values(
1889 rows.iter().map(|row| row.ordinal),
1890 )),
1891 Arc::new(UInt32Array::from_iter_values(
1892 rows.iter().map(|row| row.contract_version),
1893 )),
1894 ],
1895 )
1896 .map_err(Into::into)
1897}
1898
1899fn evidence_batch(rows: &[EvidenceLink]) -> Result<RecordBatch, KnowledgeError> {
1900 let mut ids = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1901 let mut assertions = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1902 let mut sources = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1903 let mut provenance = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
1904 for row in rows {
1905 ids.append_value(row.evidence_uuid.as_bytes())?;
1906 assertions.append_value(row.assertion_uuid.as_bytes())?;
1907 sources.append_value(row.source_uuid.as_bytes())?;
1908 provenance.append_value(row.provenance_uuid.as_bytes())?;
1909 }
1910 RecordBatch::try_new(
1911 Arc::clone(&EVIDENCE_LINK_SCHEMA),
1912 vec![
1913 Arc::new(ids.finish()),
1914 Arc::new(assertions.finish()),
1915 Arc::new(sources.finish()),
1916 Arc::new(StringArray::from_iter_values(
1917 rows.iter().map(|row| row.source_kind.as_str()),
1918 )),
1919 Arc::new(StringArray::from_iter_values(
1920 rows.iter().map(|row| row.role.as_str()),
1921 )),
1922 Arc::new(rows.iter().map(|row| row.weight).collect::<Float64Array>()),
1923 Arc::new(provenance.finish()),
1924 Arc::new(
1925 TimestampMicrosecondArray::from_iter_values(
1926 rows.iter().map(|row| row.recorded_at_micros),
1927 )
1928 .with_timezone("UTC"),
1929 ),
1930 Arc::new(UInt32Array::from_iter_values(
1931 rows.iter().map(|row| row.contract_version),
1932 )),
1933 ],
1934 )
1935 .map_err(Into::into)
1936}
1937
1938fn validate_algorithm_run(row: &AlgorithmRun) -> Result<(), KnowledgeError> {
1939 require_v7(row.run_uuid, "run_uuid")?;
1940 require_uuid(row.provenance_uuid, "provenance_uuid")?;
1941 if row.algorithm.is_empty()
1942 || row.algorithm.len() as u64 > graphforge_core::canonical::MAX_CANONICAL_TEXT_BYTES
1943 {
1944 return Err(invalid("algorithm", "must be bounded non-empty UTF-8"));
1945 }
1946 if row.algorithm_version != 1 {
1947 return Err(invalid("algorithm_version", "unsupported version"));
1948 }
1949 if row.descriptor_version != 1 {
1950 return Err(invalid("descriptor_version", "unsupported version"));
1951 }
1952 if row.descriptor.is_empty()
1953 || row.descriptor.len() as u64 > graphforge_core::canonical::MAX_CANONICAL_BINARY_BYTES
1954 {
1955 return Err(invalid("descriptor", "must be bounded and non-empty"));
1956 }
1957 if row.contract_version != ALGORITHM_RUN_CONTRACT_VERSION {
1958 return Err(invalid(
1959 "contract_version",
1960 "unsupported algorithm-run version",
1961 ));
1962 }
1963 Ok(())
1964}
1965
1966fn validate_algorithm_run_event(row: &AlgorithmRunEvent) -> Result<(), KnowledgeError> {
1967 require_uuid(row.event_uuid, "event_uuid")?;
1968 require_v7(row.run_uuid, "run_uuid")?;
1969 require_uuid(row.provenance_uuid, "provenance_uuid")?;
1970 if row.contract_version != ALGORITHM_RUN_EVENT_CONTRACT_VERSION {
1971 return Err(invalid(
1972 "contract_version",
1973 "unsupported algorithm-run-event version",
1974 ));
1975 }
1976 if row.error_code.as_ref().is_some_and(|code| {
1977 code.is_empty()
1978 || code.len() as u64 > graphforge_core::canonical::MAX_CANONICAL_TEXT_BYTES
1979 || !code.starts_with("GF_")
1980 }) {
1981 return Err(invalid("error_code", "must be a bounded stable GF_ code"));
1982 }
1983 match row.state {
1984 AlgorithmRunState::Started => {
1985 if row.result_fingerprint.is_some() || row.error_code.is_some() {
1986 return Err(invalid("state", "started has no terminal payload"));
1987 }
1988 }
1989 AlgorithmRunState::Completed => {
1990 if row.result_fingerprint.is_none() || row.error_code.is_some() {
1991 return Err(invalid(
1992 "state",
1993 "completed requires only a result fingerprint",
1994 ));
1995 }
1996 }
1997 AlgorithmRunState::Failed
1998 | AlgorithmRunState::Cancelled
1999 | AlgorithmRunState::Interrupted => {
2000 if row.result_fingerprint.is_some() || row.error_code.is_none() {
2001 return Err(invalid(
2002 "state",
2003 "non-success terminal requires only an error code",
2004 ));
2005 }
2006 }
2007 }
2008 Ok(())
2009}
2010
2011fn validate_algorithm_run_rows(
2012 runs: &[AlgorithmRun],
2013 events: &[AlgorithmRunEvent],
2014) -> Result<(), KnowledgeError> {
2015 check_limit("algorithm_runs", runs.len())?;
2016 check_limit("algorithm_run_events", events.len())?;
2017 let mut run_ids = HashSet::with_capacity(runs.len());
2018 let mut run_index = HashMap::with_capacity(runs.len());
2019 for row in runs {
2020 validate_algorithm_run(row)?;
2021 if !run_ids.insert(row.run_uuid) {
2022 return Err(KnowledgeError::Duplicate("run_uuid"));
2023 }
2024 run_index.insert(row.run_uuid, row);
2025 }
2026 let mut event_ids = HashSet::with_capacity(events.len());
2027 let mut per_run: HashMap<Uuid, (usize, usize)> = HashMap::new();
2028 for row in events {
2029 validate_algorithm_run_event(row)?;
2030 if !event_ids.insert(row.event_uuid) {
2031 return Err(KnowledgeError::Duplicate("event_uuid"));
2032 }
2033 let run = run_index
2034 .get(&row.run_uuid)
2035 .ok_or(KnowledgeError::Dangling("run_uuid"))?;
2036 if row.recorded_at_micros < run.started_at_micros {
2037 return Err(invalid("recorded_at", "event precedes run start"));
2038 }
2039 let counts = per_run.entry(row.run_uuid).or_default();
2040 if row.state == AlgorithmRunState::Started {
2041 counts.0 += 1;
2042 if row.recorded_at_micros != run.started_at_micros
2043 || row.provenance_uuid != run.provenance_uuid
2044 {
2045 return Err(invalid(
2046 "started",
2047 "start event must match immutable run identity",
2048 ));
2049 }
2050 } else {
2051 counts.1 += 1;
2052 }
2053 }
2054 for run in runs {
2055 let (started, terminal) = per_run.get(&run.run_uuid).copied().unwrap_or_default();
2056 if started != 1 {
2057 return Err(invalid("started", "run requires exactly one start event"));
2058 }
2059 if terminal > 1 {
2060 return Err(invalid(
2061 "terminal",
2062 "run permits at most one terminal event",
2063 ));
2064 }
2065 }
2066 Ok(())
2067}
2068
2069fn algorithm_run_batch(rows: &[AlgorithmRun]) -> Result<RecordBatch, KnowledgeError> {
2070 let mut ids = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
2071 let mut projections = FixedSizeBinaryBuilder::with_capacity(rows.len(), 32);
2072 let mut provenance = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
2073 for row in rows {
2074 ids.append_value(row.run_uuid.as_bytes())?;
2075 projections.append_value(row.projection_fingerprint)?;
2076 provenance.append_value(row.provenance_uuid.as_bytes())?;
2077 }
2078 RecordBatch::try_new(
2079 Arc::clone(&ALGORITHM_RUN_SCHEMA),
2080 vec![
2081 Arc::new(ids.finish()),
2082 Arc::new(StringArray::from_iter_values(
2083 rows.iter().map(|row| row.algorithm.as_str()),
2084 )),
2085 Arc::new(UInt32Array::from_iter_values(
2086 rows.iter().map(|row| row.algorithm_version),
2087 )),
2088 Arc::new(UInt32Array::from_iter_values(
2089 rows.iter().map(|row| row.descriptor_version),
2090 )),
2091 Arc::new(BinaryArray::from_iter_values(
2092 rows.iter().map(|row| row.descriptor.as_slice()),
2093 )),
2094 Arc::new(projections.finish()),
2095 Arc::new(provenance.finish()),
2096 Arc::new(
2097 TimestampMicrosecondArray::from_iter_values(
2098 rows.iter().map(|row| row.started_at_micros),
2099 )
2100 .with_timezone("UTC"),
2101 ),
2102 Arc::new(UInt32Array::from_iter_values(
2103 rows.iter().map(|row| row.contract_version),
2104 )),
2105 ],
2106 )
2107 .map_err(Into::into)
2108}
2109
2110fn algorithm_run_event_batch(rows: &[AlgorithmRunEvent]) -> Result<RecordBatch, KnowledgeError> {
2111 let mut ids = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
2112 let mut runs = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
2113 let mut results = FixedSizeBinaryBuilder::with_capacity(rows.len(), 32);
2114 let mut provenance = FixedSizeBinaryBuilder::with_capacity(rows.len(), 16);
2115 for row in rows {
2116 ids.append_value(row.event_uuid.as_bytes())?;
2117 runs.append_value(row.run_uuid.as_bytes())?;
2118 match row.result_fingerprint {
2119 Some(value) => results.append_value(value)?,
2120 None => results.append_null(),
2121 }
2122 provenance.append_value(row.provenance_uuid.as_bytes())?;
2123 }
2124 RecordBatch::try_new(
2125 Arc::clone(&ALGORITHM_RUN_EVENT_SCHEMA),
2126 vec![
2127 Arc::new(ids.finish()),
2128 Arc::new(runs.finish()),
2129 Arc::new(StringArray::from_iter_values(
2130 rows.iter().map(|row| row.state.as_str()),
2131 )),
2132 Arc::new(results.finish()),
2133 Arc::new(StringArray::from(
2134 rows.iter()
2135 .map(|row| row.error_code.as_deref())
2136 .collect::<Vec<_>>(),
2137 )),
2138 Arc::new(
2139 TimestampMicrosecondArray::from_iter_values(
2140 rows.iter().map(|row| row.recorded_at_micros),
2141 )
2142 .with_timezone("UTC"),
2143 ),
2144 Arc::new(provenance.finish()),
2145 Arc::new(UInt32Array::from_iter_values(
2146 rows.iter().map(|row| row.contract_version),
2147 )),
2148 ],
2149 )
2150 .map_err(Into::into)
2151}
2152
2153fn validate_confidence_rows(
2154 assessments: &[ConfidenceAssessment],
2155 inputs: &[ConfidenceInput],
2156) -> Result<(), KnowledgeError> {
2157 check_limit("confidence_assessments", assessments.len())?;
2158 check_limit("confidence_inputs", inputs.len())?;
2159 let mut ids = HashSet::with_capacity(assessments.len());
2160 let mut policies = HashMap::with_capacity(assessments.len());
2161 for row in assessments {
2162 require_v7(row.confidence_uuid, "confidence_uuid")?;
2163 require_v7(row.assertion_uuid, "assertion_uuid")?;
2164 require_uuid(row.provenance_uuid, "provenance_uuid")?;
2165 validate_confidence(row.value, "value")?;
2166 if row.policy_version != 1 {
2167 return Err(invalid("policy_version", "unsupported version"));
2168 }
2169 if row.contract_version != CONFIDENCE_ASSESSMENT_CONTRACT_VERSION {
2170 return Err(invalid(
2171 "confidence.contract_version",
2172 "unsupported version",
2173 ));
2174 }
2175 if !ids.insert(row.confidence_uuid) {
2176 return Err(KnowledgeError::Duplicate("confidence_uuid"));
2177 }
2178 policies.insert(row.confidence_uuid, row.policy);
2179 }
2180 let mut input_ids = HashSet::with_capacity(inputs.len());
2181 let mut normalized_inputs: HashMap<Uuid, Vec<(u32, Uuid, Option<f64>)>> = HashMap::new();
2182 for input in inputs {
2183 require_v7(input.confidence_uuid, "confidence_uuid")?;
2184 require_v7(input.input_confidence_uuid, "input_confidence_uuid")?;
2185 validate_confidence(input.input_value, "input_value")?;
2186 if input.contract_version != CONFIDENCE_INPUT_CONTRACT_VERSION {
2187 return Err(invalid(
2188 "confidence_input.contract_version",
2189 "unsupported version",
2190 ));
2191 }
2192 if !ids.contains(&input.confidence_uuid) {
2193 return Err(KnowledgeError::Dangling("confidence_uuid"));
2194 }
2195 if !input_ids.insert((input.confidence_uuid, input.input_confidence_uuid)) {
2196 return Err(KnowledgeError::Duplicate("input_confidence_uuid"));
2197 }
2198 normalized_inputs
2199 .entry(input.confidence_uuid)
2200 .or_default()
2201 .push((
2202 input.ordinal,
2203 input.input_confidence_uuid,
2204 input.input_value,
2205 ));
2206 }
2207 for (confidence_uuid, policy) in policies {
2208 let assessment = assessments
2209 .iter()
2210 .find(|row| row.confidence_uuid == confidence_uuid)
2211 .expect("validated assessment identity");
2212 let values = normalized_inputs.entry(confidence_uuid).or_default();
2213 values.sort_by_key(|(ordinal, _, _)| *ordinal);
2214 if values
2215 .iter()
2216 .enumerate()
2217 .any(|(expected, (actual, _, _))| usize::try_from(*actual) != Ok(expected))
2218 {
2219 return Err(invalid("ordinal", "must be contiguous from zero"));
2220 }
2221 if values.windows(2).any(|pair| pair[0].1 >= pair[1].1) {
2222 return Err(invalid(
2223 "input_confidence_uuid",
2224 "must be unique UUID-normalized order",
2225 ));
2226 }
2227 validate_policy_snapshot(assessment, policy, values)?;
2228 }
2229 Ok(())
2230}
2231
2232fn validate_policy_snapshot(
2233 assessment: &ConfidenceAssessment,
2234 policy: ConfidencePolicy,
2235 values: &[(u32, Uuid, Option<f64>)],
2236) -> Result<(), KnowledgeError> {
2237 match policy {
2238 ConfidencePolicy::Explicit => {
2239 if !values.is_empty() {
2240 return Err(invalid(
2241 "confidence_inputs",
2242 "explicit policy has no inputs",
2243 ));
2244 }
2245 if assessment.value.is_none() {
2246 return Err(invalid("value", "explicit policy requires a value"));
2247 }
2248 }
2249 ConfidencePolicy::ConservativeMin => {
2250 let expected =
2251 if values.is_empty() || values.iter().any(|(_, _, value)| value.is_none()) {
2252 None
2253 } else {
2254 values
2255 .iter()
2256 .filter_map(|(_, _, value)| *value)
2257 .reduce(f64::min)
2258 };
2259 if assessment.value != expected {
2260 return Err(invalid(
2261 "value",
2262 "does not match conservative_min input snapshot",
2263 ));
2264 }
2265 }
2266 }
2267 Ok(())
2268}
2269
2270fn inputs_for(rows: &[ConfidenceInput], confidence_uuid: Uuid) -> Vec<ConfidenceInput> {
2271 let mut inputs = rows
2272 .iter()
2273 .filter(|row| row.confidence_uuid == confidence_uuid)
2274 .cloned()
2275 .collect::<Vec<_>>();
2276 inputs.sort_by_key(|row| (row.ordinal, row.input_confidence_uuid));
2277 inputs
2278}
2279
2280fn validate_confidence(value: Option<f64>, field: &'static str) -> Result<(), KnowledgeError> {
2281 if value.is_some_and(|value| !value.is_finite() || !(0.0..=1.0).contains(&value)) {
2282 Err(invalid(field, "must be finite and in [0,1]"))
2283 } else {
2284 Ok(())
2285 }
2286}
2287
2288fn normalize_zero(value: f64) -> f64 {
2289 if value == 0.0 { 0.0 } else { value }
2290}
2291
2292fn canonical_optional_f64(
2293 writer: &mut CanonicalWriter,
2294 value: Option<f64>,
2295) -> Result<(), KnowledgeError> {
2296 match value {
2297 None => writer.u8(0)?,
2298 Some(value) => {
2299 writer.u8(1)?;
2300 writer.u64(normalize_zero(value).to_bits())?;
2301 }
2302 }
2303 Ok(())
2304}
2305
2306fn refs_for(rows: &[AssertionGraphRef], assertion_uuid: Uuid) -> Vec<AssertionGraphRef> {
2307 let mut refs = rows
2308 .iter()
2309 .filter(|row| row.assertion_uuid == assertion_uuid)
2310 .cloned()
2311 .collect::<Vec<_>>();
2312 refs.sort_by_key(|row| {
2313 (
2314 role_order(row.role),
2315 row.ordinal,
2316 kind_order(row.graph_kind),
2317 row.graph_uuid,
2318 )
2319 });
2320 refs
2321}
2322
2323const fn role_order(role: AssertionGraphRole) -> u8 {
2324 match role {
2325 AssertionGraphRole::Subject => 0,
2326 AssertionGraphRole::Object => 1,
2327 AssertionGraphRole::Context => 2,
2328 }
2329}
2330
2331const fn kind_order(kind: GraphObjectKind) -> u8 {
2332 match kind {
2333 GraphObjectKind::Node => 0,
2334 GraphObjectKind::Edge => 1,
2335 }
2336}
2337
2338fn validate_claim(claim: &str) -> Result<(), KnowledgeError> {
2339 if claim.is_empty() {
2340 return Err(invalid("claim", "must not be empty"));
2341 }
2342 if claim.len() as u64 > graphforge_core::canonical::MAX_CANONICAL_TEXT_BYTES {
2343 return Err(invalid("claim", "exceeds canonical UTF-8 limit"));
2344 }
2345 Ok(())
2346}
2347
2348fn require_uuid(value: Uuid, field: &'static str) -> Result<(), KnowledgeError> {
2349 if value.is_nil() {
2350 Err(invalid(field, "must not be nil"))
2351 } else {
2352 Ok(())
2353 }
2354}
2355
2356fn require_v7(value: Uuid, field: &'static str) -> Result<(), KnowledgeError> {
2357 require_uuid(value, field)?;
2358 if value.get_version() != Some(Version::SortRand) {
2359 return Err(invalid(field, "must be UUIDv7"));
2360 }
2361 Ok(())
2362}
2363
2364fn check_limit(participant: &'static str, observed: usize) -> Result<(), KnowledgeError> {
2365 if observed > MAX_KNOWLEDGE_ROWS {
2366 Err(KnowledgeError::Limit {
2367 participant,
2368 observed,
2369 limit: MAX_KNOWLEDGE_ROWS,
2370 })
2371 } else {
2372 Ok(())
2373 }
2374}
2375
2376const fn invalid(field: &'static str, message: &'static str) -> KnowledgeError {
2377 KnowledgeError::Invalid { field, message }
2378}
2379
2380fn require_schema(
2381 batch: &RecordBatch,
2382 expected: &SchemaRef,
2383 field: &'static str,
2384) -> Result<(), KnowledgeError> {
2385 if batch.schema().as_ref() == expected.as_ref() {
2386 Ok(())
2387 } else {
2388 Err(invalid(field, "schema mismatch"))
2389 }
2390}
2391
2392fn fixed_column<'a>(
2393 batch: &'a RecordBatch,
2394 name: &'static str,
2395) -> Result<&'a FixedSizeBinaryArray, KnowledgeError> {
2396 batch
2397 .column_by_name(name)
2398 .and_then(|column| column.as_any().downcast_ref())
2399 .ok_or_else(|| invalid(name, "missing or wrong Arrow type"))
2400}
2401
2402fn string_column<'a>(
2403 batch: &'a RecordBatch,
2404 name: &'static str,
2405) -> Result<&'a StringArray, KnowledgeError> {
2406 batch
2407 .column_by_name(name)
2408 .and_then(|column| column.as_any().downcast_ref())
2409 .ok_or_else(|| invalid(name, "missing or wrong Arrow type"))
2410}
2411
2412fn binary_column<'a>(
2413 batch: &'a RecordBatch,
2414 name: &'static str,
2415) -> Result<&'a BinaryArray, KnowledgeError> {
2416 batch
2417 .column_by_name(name)
2418 .and_then(|column| column.as_any().downcast_ref())
2419 .ok_or_else(|| invalid(name, "missing or wrong Arrow type"))
2420}
2421
2422fn timestamp_column<'a>(
2423 batch: &'a RecordBatch,
2424 name: &'static str,
2425) -> Result<&'a TimestampMicrosecondArray, KnowledgeError> {
2426 batch
2427 .column_by_name(name)
2428 .and_then(|column| column.as_any().downcast_ref())
2429 .ok_or_else(|| invalid(name, "missing or wrong Arrow type"))
2430}
2431
2432fn u32_column<'a>(
2433 batch: &'a RecordBatch,
2434 name: &'static str,
2435) -> Result<&'a UInt32Array, KnowledgeError> {
2436 batch
2437 .column_by_name(name)
2438 .and_then(|column| column.as_any().downcast_ref())
2439 .ok_or_else(|| invalid(name, "missing or wrong Arrow type"))
2440}
2441
2442fn f64_column<'a>(
2443 batch: &'a RecordBatch,
2444 name: &'static str,
2445) -> Result<&'a Float64Array, KnowledgeError> {
2446 batch
2447 .column_by_name(name)
2448 .and_then(|column| column.as_any().downcast_ref())
2449 .ok_or_else(|| invalid(name, "missing or wrong Arrow type"))
2450}
2451
2452fn optional_f64(array: &Float64Array, row: usize) -> Option<f64> {
2453 (!array.is_null(row)).then(|| normalize_zero(array.value(row)))
2454}
2455
2456fn uuid_at(
2457 array: &FixedSizeBinaryArray,
2458 row: usize,
2459 field: &'static str,
2460) -> Result<Uuid, KnowledgeError> {
2461 if array.is_null(row) {
2462 return Err(invalid(field, "must not be null"));
2463 }
2464 Uuid::from_slice(array.value(row)).map_err(|_| invalid(field, "malformed UUID"))
2465}
2466
2467fn fixed_32_at(
2468 array: &FixedSizeBinaryArray,
2469 row: usize,
2470 field: &'static str,
2471) -> Result<[u8; 32], KnowledgeError> {
2472 if array.is_null(row) || array.value_length() != 32 {
2473 return Err(invalid(field, "must be a 32-byte value"));
2474 }
2475 Ok(array
2476 .value(row)
2477 .try_into()
2478 .expect("validated fixed-size binary width"))
2479}
2480
2481fn optional_fixed_32(
2482 array: &FixedSizeBinaryArray,
2483 row: usize,
2484 field: &'static str,
2485) -> Result<Option<[u8; 32]>, KnowledgeError> {
2486 if array.is_null(row) {
2487 Ok(None)
2488 } else {
2489 fixed_32_at(array, row, field).map(Some)
2490 }
2491}
2492
2493fn required_text<'a>(
2494 array: &'a StringArray,
2495 row: usize,
2496 field: &'static str,
2497) -> Result<&'a str, KnowledgeError> {
2498 if array.is_null(row) {
2499 Err(invalid(field, "must not be null"))
2500 } else {
2501 Ok(array.value(row))
2502 }
2503}
2504
2505fn optional_text(array: &StringArray, row: usize) -> Option<String> {
2506 (!array.is_null(row)).then(|| array.value(row).to_owned())
2507}
2508
2509fn required_binary<'a>(
2510 array: &'a BinaryArray,
2511 row: usize,
2512 field: &'static str,
2513) -> Result<&'a [u8], KnowledgeError> {
2514 if array.is_null(row) {
2515 Err(invalid(field, "must not be null"))
2516 } else {
2517 Ok(array.value(row))
2518 }
2519}
2520
2521fn required_i64(
2522 array: &TimestampMicrosecondArray,
2523 row: usize,
2524 field: &'static str,
2525) -> Result<i64, KnowledgeError> {
2526 if array.is_null(row) {
2527 Err(invalid(field, "must not be null"))
2528 } else {
2529 Ok(array.value(row))
2530 }
2531}
2532
2533fn required_u32(
2534 array: &UInt32Array,
2535 row: usize,
2536 field: &'static str,
2537) -> Result<u32, KnowledgeError> {
2538 if array.is_null(row) {
2539 Err(invalid(field, "must not be null"))
2540 } else {
2541 Ok(array.value(row))
2542 }
2543}
2544
2545#[cfg(test)]
2546mod tests {
2547 use super::*;
2548
2549 fn uuid7(seed: u8) -> Uuid {
2550 let mut bytes = [seed; 16];
2551 bytes[6] = (bytes[6] & 0x0f) | 0x70;
2552 bytes[8] = (bytes[8] & 0x3f) | 0x80;
2553 Uuid::from_bytes(bytes)
2554 }
2555
2556 fn fixture() -> AssertionLedger {
2557 let assertion_uuid = uuid7(1);
2558 AssertionLedger::new(
2559 vec![Assertion::new(assertion_uuid, "exact claim".into(), uuid7(2), 10).unwrap()],
2560 vec![
2561 AssertionGraphRef::new(
2562 assertion_uuid,
2563 uuid7(3),
2564 GraphObjectKind::Node,
2565 AssertionGraphRole::Subject,
2566 0,
2567 )
2568 .unwrap(),
2569 AssertionGraphRef::new(
2570 assertion_uuid,
2571 uuid7(4),
2572 GraphObjectKind::Edge,
2573 AssertionGraphRole::Context,
2574 0,
2575 )
2576 .unwrap(),
2577 ],
2578 )
2579 .unwrap()
2580 }
2581
2582 #[test]
2583 fn exact_claim_and_sorted_refs_have_stable_fingerprint() {
2584 let first = fixture();
2585 let mut reversed = first.graph_refs.clone();
2586 reversed.reverse();
2587 let second = AssertionLedger::new(first.assertions.clone(), reversed).unwrap();
2588 let fingerprint = first.assertion_fingerprint(uuid7(1)).unwrap();
2589 assert_eq!(fingerprint, second.assertion_fingerprint(uuid7(1)).unwrap());
2590 let encoded = fingerprint
2591 .iter()
2592 .map(|byte| format!("{byte:02x}"))
2593 .collect::<String>();
2594 assert_eq!(
2597 encoded,
2598 "6024d4f2e22b35d5850f6a2f1c9f0cf77d28ab1e7ba83342416a1e7660f7b4bb"
2599 );
2600 }
2601
2602 #[test]
2603 fn arrow_round_trip_preserves_exact_records_across_chunking() {
2604 let ledger = fixture();
2605 let assertions = ledger.assertion_batch().unwrap();
2606 let refs = ledger.graph_ref_batch().unwrap();
2607 assert_eq!(
2608 AssertionLedger::from_batches(
2609 &[assertions.slice(0, 0), assertions.slice(0, 1)],
2610 &[refs.slice(0, 1), refs.slice(1, 1)],
2611 )
2612 .unwrap(),
2613 ledger
2614 );
2615 }
2616
2617 #[test]
2618 fn invalid_claims_refs_and_ordinals_fail_structurally() {
2619 let assertion_uuid = uuid7(1);
2620 let assertion = Assertion::new(assertion_uuid, "claim".into(), uuid7(2), 10).unwrap();
2621 assert!(AssertionLedger::new(vec![assertion.clone()], vec![]).is_err());
2622 assert!(matches!(
2623 Assertion::new(assertion_uuid, String::new(), uuid7(2), 10),
2624 Err(KnowledgeError::Invalid { field: "claim", .. })
2625 ));
2626 let oversized =
2627 "x".repeat(graphforge_core::canonical::MAX_CANONICAL_TEXT_BYTES as usize + 1);
2628 assert!(matches!(
2629 Assertion::new(assertion_uuid, oversized, uuid7(2), 10),
2630 Err(KnowledgeError::Invalid { field: "claim", .. })
2631 ));
2632
2633 let non_contiguous = AssertionGraphRef::new(
2634 assertion_uuid,
2635 uuid7(3),
2636 GraphObjectKind::Node,
2637 AssertionGraphRole::Subject,
2638 1,
2639 )
2640 .unwrap();
2641 assert!(matches!(
2642 AssertionLedger::new(vec![assertion.clone()], vec![non_contiguous]),
2643 Err(KnowledgeError::Invalid {
2644 field: "ordinal",
2645 ..
2646 })
2647 ));
2648
2649 let duplicate = AssertionGraphRef::new(
2650 assertion_uuid,
2651 uuid7(3),
2652 GraphObjectKind::Node,
2653 AssertionGraphRole::Subject,
2654 0,
2655 )
2656 .unwrap();
2657 assert!(matches!(
2658 AssertionLedger::new(vec![assertion], vec![duplicate.clone(), duplicate]),
2659 Err(KnowledgeError::Duplicate(
2660 "assertion_uuid/graph_uuid/role/ordinal"
2661 ))
2662 ));
2663 }
2664
2665 #[test]
2666 fn closed_values_fail_during_arrow_decode() {
2667 let ledger = fixture();
2668 let assertions = ledger.assertion_batch().unwrap();
2669 let refs = ledger.graph_ref_batch().unwrap();
2670 let bad_kind = RecordBatch::try_new(
2671 Arc::clone(&ASSERTION_GRAPH_REF_SCHEMA),
2672 vec![
2673 Arc::clone(refs.column(0)),
2674 Arc::clone(refs.column(1)),
2675 Arc::new(StringArray::from(vec!["vertex", "edge"])),
2676 Arc::clone(refs.column(3)),
2677 Arc::clone(refs.column(4)),
2678 Arc::clone(refs.column(5)),
2679 ],
2680 )
2681 .unwrap();
2682 assert!(matches!(
2683 AssertionLedger::from_batches(&[assertions.clone()], &[bad_kind]),
2684 Err(KnowledgeError::Invalid {
2685 field: "graph_kind",
2686 ..
2687 })
2688 ));
2689
2690 let bad_role = RecordBatch::try_new(
2691 Arc::clone(&ASSERTION_GRAPH_REF_SCHEMA),
2692 vec![
2693 Arc::clone(refs.column(0)),
2694 Arc::clone(refs.column(1)),
2695 Arc::clone(refs.column(2)),
2696 Arc::new(StringArray::from(vec!["target", "context"])),
2697 Arc::clone(refs.column(4)),
2698 Arc::clone(refs.column(5)),
2699 ],
2700 )
2701 .unwrap();
2702 assert!(matches!(
2703 AssertionLedger::from_batches(&[assertions], &[bad_role]),
2704 Err(KnowledgeError::Invalid { field: "role", .. })
2705 ));
2706 }
2707
2708 #[test]
2709 fn m20_schema_registry_excludes_every_m21_field() {
2710 let fields = schema_registry()
2711 .into_iter()
2712 .filter(|entry| entry.capability_id == "knowledge")
2713 .flat_map(|entry| {
2714 entry
2715 .schema
2716 .fields()
2717 .iter()
2718 .map(|field| field.name().clone())
2719 .collect::<Vec<_>>()
2720 })
2721 .collect::<HashSet<_>>();
2722 for deferred in [
2723 "status",
2724 "confidence",
2725 "evidence",
2726 "hypothesis_uuid",
2727 "reasoning",
2728 "supersedes_uuid",
2729 "valid_from",
2730 "valid_to",
2731 ] {
2732 assert!(!fields.contains(deferred));
2733 }
2734 }
2735
2736 #[test]
2737 fn schema_registry_owns_record_diff_identity_contracts() {
2738 let identities = schema_registry()
2739 .into_iter()
2740 .map(|entry| {
2741 for field in entry.diff_identity_fields {
2742 assert!(entry.schema.field_with_name(field).is_ok());
2743 }
2744 if let Some(field) = entry.diff_record_uuid_field {
2745 assert!(entry.diff_identity_fields.contains(&field));
2746 assert_eq!(
2747 entry.schema.field_with_name(field).unwrap().data_type(),
2748 &DataType::FixedSizeBinary(16)
2749 );
2750 }
2751 (
2752 entry.record_family,
2753 entry.diff_identity_fields,
2754 entry.diff_record_uuid_field,
2755 )
2756 })
2757 .collect::<Vec<_>>();
2758
2759 assert!(identities.contains(&(
2760 "assertions",
2761 &["assertion_uuid"][..],
2762 Some("assertion_uuid")
2763 )));
2764 assert!(identities.contains(&(
2765 "assertion_graph_refs",
2766 &["assertion_uuid", "graph_uuid", "role", "ordinal"][..],
2767 None
2768 )));
2769 assert!(identities.contains(&(
2770 "confidence_inputs",
2771 &["confidence_uuid", "input_confidence_uuid"][..],
2772 None
2773 )));
2774 }
2775
2776 #[test]
2777 fn explicit_confidence_round_trips_and_normalizes_negative_zero() {
2778 let ledger = ConfidenceLedger::explicit(uuid7(10), uuid7(1), -0.0, uuid7(2), 20).unwrap();
2779 assert_eq!(
2780 ledger.assessments[0].value.unwrap().to_bits(),
2781 0.0f64.to_bits()
2782 );
2783 assert!(ledger.inputs.is_empty());
2784 let assessments = ledger.assessment_batch().unwrap();
2785 let inputs = ledger.input_batch().unwrap();
2786 assert_eq!(
2787 ConfidenceLedger::from_batches(&[assessments], &[inputs]).unwrap(),
2788 ledger
2789 );
2790 }
2791
2792 #[test]
2793 fn conservative_min_is_uuid_normalized_and_snapshots_missing_values() {
2794 let first = ConfidenceLedger::explicit(uuid7(10), uuid7(1), 0.8, uuid7(2), 10).unwrap();
2795 let second = ConfidenceLedger::explicit(uuid7(11), uuid7(1), 0.3, uuid7(3), 11).unwrap();
2796 let existing = first.merge(&second).unwrap();
2797 let staged = existing
2798 .conservative_min(
2799 uuid7(20),
2800 uuid7(1),
2801 vec![uuid7(12), uuid7(11), uuid7(10)],
2802 uuid7(4),
2803 20,
2804 )
2805 .unwrap();
2806 assert_eq!(staged.assessments[0].value, None);
2807 assert_eq!(
2808 staged
2809 .inputs
2810 .iter()
2811 .map(|row| row.input_confidence_uuid)
2812 .collect::<Vec<_>>(),
2813 vec![uuid7(10), uuid7(11), uuid7(12)]
2814 );
2815 assert_eq!(
2816 staged
2817 .inputs
2818 .iter()
2819 .map(|row| row.input_value)
2820 .collect::<Vec<_>>(),
2821 vec![Some(0.8), Some(0.3), None]
2822 );
2823
2824 let complete = existing
2825 .conservative_min(
2826 uuid7(21),
2827 uuid7(1),
2828 vec![uuid7(11), uuid7(10)],
2829 uuid7(5),
2830 21,
2831 )
2832 .unwrap();
2833 assert_eq!(complete.assessments[0].value, Some(0.3));
2834 assert_eq!(
2835 staged.assessment_fingerprint(uuid7(20)).unwrap(),
2836 existing
2837 .conservative_min(
2838 uuid7(20),
2839 uuid7(1),
2840 vec![uuid7(10), uuid7(12), uuid7(11)],
2841 uuid7(4),
2842 20,
2843 )
2844 .unwrap()
2845 .assessment_fingerprint(uuid7(20))
2846 .unwrap()
2847 );
2848 let encoded = staged
2849 .assessment_fingerprint(uuid7(20))
2850 .unwrap()
2851 .iter()
2852 .map(|byte| format!("{byte:02x}"))
2853 .collect::<String>();
2854 assert_eq!(
2856 encoded,
2857 "052104e7e9573852f29255e5f4d7942bd0f409a00b8cda80aba2d09112e40bbb"
2858 );
2859 }
2860
2861 #[test]
2862 fn confidence_idempotency_and_validation_are_fail_closed() {
2863 for invalid_value in [f64::NAN, f64::INFINITY, -0.01, 1.01] {
2864 assert!(matches!(
2865 ConfidenceLedger::explicit(uuid7(10), uuid7(1), invalid_value, uuid7(2), 20),
2866 Err(KnowledgeError::Invalid { field: "value", .. })
2867 ));
2868 }
2869 let existing = ConfidenceLedger::explicit(uuid7(10), uuid7(1), 0.8, uuid7(2), 20).unwrap();
2870 assert_eq!(
2871 existing.merge(&existing).unwrap().assessments.len(),
2872 existing.assessments.len()
2873 );
2874 let conflict = ConfidenceLedger::explicit(uuid7(10), uuid7(1), 0.7, uuid7(2), 20).unwrap();
2875 assert!(matches!(
2876 existing.merge(&conflict),
2877 Err(KnowledgeError::Conflict("confidence_uuid"))
2878 ));
2879 assert!(matches!(
2880 existing.conservative_min(
2881 uuid7(20),
2882 uuid7(1),
2883 vec![uuid7(10), uuid7(10)],
2884 uuid7(2),
2885 20,
2886 ),
2887 Err(KnowledgeError::Duplicate("input_confidence_uuid"))
2888 ));
2889 }
2890
2891 #[test]
2892 fn confidence_arrow_decode_rejects_schema_and_closed_policy_drift() {
2893 let ledger = ConfidenceLedger::explicit(uuid7(10), uuid7(1), 0.8, uuid7(2), 20).unwrap();
2894 let batch = ledger.assessment_batch().unwrap();
2895 let bad = RecordBatch::try_new(
2896 Arc::clone(&CONFIDENCE_ASSESSMENT_SCHEMA),
2897 vec![
2898 Arc::clone(batch.column(0)),
2899 Arc::clone(batch.column(1)),
2900 Arc::new(StringArray::from(vec!["average"])),
2901 Arc::clone(batch.column(3)),
2902 Arc::clone(batch.column(4)),
2903 Arc::clone(batch.column(5)),
2904 Arc::clone(batch.column(6)),
2905 Arc::clone(batch.column(7)),
2906 ],
2907 )
2908 .unwrap();
2909 assert!(matches!(
2910 ConfidenceLedger::from_batches(&[bad], &[]),
2911 Err(KnowledgeError::Invalid {
2912 field: "policy",
2913 ..
2914 })
2915 ));
2916 }
2917
2918 #[test]
2919 fn evidence_round_trips_orders_fingerprints_and_merges_idempotently() {
2920 let later = EvidenceLink::new(
2921 uuid7(31),
2922 uuid7(1),
2923 uuid7(41),
2924 EvidenceSourceKind::Observation,
2925 EvidenceRole::Contradicts,
2926 Some(0.25),
2927 uuid7(51),
2928 20,
2929 )
2930 .unwrap();
2931 let earlier = EvidenceLink::new(
2932 uuid7(30),
2933 uuid7(1),
2934 uuid7(40),
2935 EvidenceSourceKind::Document,
2936 EvidenceRole::Supports,
2937 Some(-0.0),
2938 uuid7(50),
2939 10,
2940 )
2941 .unwrap();
2942 let ledger = EvidenceLedger::new(vec![later, earlier]).unwrap();
2943 assert_eq!(ledger.links[0].evidence_uuid, uuid7(30));
2944 assert_eq!(ledger.links[0].weight.unwrap().to_bits(), 0.0f64.to_bits());
2945 assert_eq!(
2946 EvidenceLedger::from_batches(&[ledger.batch().unwrap()]).unwrap(),
2947 ledger
2948 );
2949 assert_eq!(ledger.merge(&ledger).unwrap(), ledger);
2950 let encoded = ledger
2951 .evidence_fingerprint(uuid7(30))
2952 .unwrap()
2953 .iter()
2954 .map(|byte| format!("{byte:02x}"))
2955 .collect::<String>();
2956 assert_eq!(
2957 encoded,
2958 "363146c1fc94d57cc02c098a7187d753d85a7b10ebd715ecbae76e85e1d3fd5e"
2959 );
2960 }
2961
2962 #[test]
2963 fn evidence_validation_and_conflicting_identity_are_fail_closed() {
2964 for weight in [f64::NAN, f64::INFINITY, -0.01, 1.01] {
2965 assert!(matches!(
2966 EvidenceLink::new(
2967 uuid7(30),
2968 uuid7(1),
2969 uuid7(40),
2970 EvidenceSourceKind::Document,
2971 EvidenceRole::Supports,
2972 Some(weight),
2973 uuid7(50),
2974 10,
2975 ),
2976 Err(KnowledgeError::Invalid {
2977 field: "weight",
2978 ..
2979 })
2980 ));
2981 }
2982 let first = EvidenceLedger::new(vec![
2983 EvidenceLink::new(
2984 uuid7(30),
2985 uuid7(1),
2986 uuid7(40),
2987 EvidenceSourceKind::Document,
2988 EvidenceRole::Supports,
2989 None,
2990 uuid7(50),
2991 10,
2992 )
2993 .unwrap(),
2994 ])
2995 .unwrap();
2996 let conflict = EvidenceLedger::new(vec![
2997 EvidenceLink::new(
2998 uuid7(30),
2999 uuid7(1),
3000 uuid7(41),
3001 EvidenceSourceKind::Observation,
3002 EvidenceRole::Context,
3003 None,
3004 uuid7(50),
3005 10,
3006 )
3007 .unwrap(),
3008 ])
3009 .unwrap();
3010 assert!(matches!(
3011 first.merge(&conflict),
3012 Err(KnowledgeError::Conflict("evidence_uuid"))
3013 ));
3014 }
3015
3016 #[test]
3017 fn algorithm_run_lifecycle_round_trips_and_rejects_second_terminal() {
3018 assert_eq!(
3019 ALGORITHM_RUN_SCHEMA_FINGERPRINT
3020 .iter()
3021 .map(|byte| format!("{byte:02x}"))
3022 .collect::<String>(),
3023 "ff52080371c5956aa9bc8b0cf9c1022e2c3271d576d43ab57f4b75eb33cf64ed"
3024 );
3025 assert_eq!(
3026 ALGORITHM_RUN_EVENT_SCHEMA_FINGERPRINT
3027 .iter()
3028 .map(|byte| format!("{byte:02x}"))
3029 .collect::<String>(),
3030 "f055828f73440cca1bd2ea746cb4599d52a736229d2a4df8ddec70972ceef0aa"
3031 );
3032 let run = AlgorithmRun::new(
3033 uuid7(40),
3034 "rank.degree".into(),
3035 1,
3036 1,
3037 b"canonical descriptor".to_vec(),
3038 [7; 32],
3039 uuid7(41),
3040 10,
3041 )
3042 .unwrap();
3043 let started = AlgorithmRunEvent::new(
3044 uuid7(42),
3045 run.run_uuid,
3046 AlgorithmRunState::Started,
3047 None,
3048 None,
3049 10,
3050 run.provenance_uuid,
3051 )
3052 .unwrap();
3053 let completed = AlgorithmRunEvent::new(
3054 uuid7(43),
3055 run.run_uuid,
3056 AlgorithmRunState::Completed,
3057 Some([8; 32]),
3058 None,
3059 11,
3060 uuid7(44),
3061 )
3062 .unwrap();
3063 let ledger = AlgorithmRunLedger::new(vec![run.clone()], vec![completed, started]).unwrap();
3064 assert_eq!(
3065 AlgorithmRunLedger::from_batches(
3066 &[ledger.run_batch().unwrap()],
3067 &[ledger.event_batch().unwrap()],
3068 )
3069 .unwrap(),
3070 ledger
3071 );
3072 let failed = AlgorithmRunEvent::new(
3073 uuid7(45),
3074 run.run_uuid,
3075 AlgorithmRunState::Failed,
3076 None,
3077 Some("GF_EXECUTION".into()),
3078 12,
3079 uuid7(46),
3080 )
3081 .unwrap();
3082 let mut events = ledger.events.clone();
3083 events.push(failed);
3084 assert!(matches!(
3085 AlgorithmRunLedger::new(ledger.runs.clone(), events),
3086 Err(KnowledgeError::Invalid {
3087 field: "terminal",
3088 ..
3089 })
3090 ));
3091 }
3092
3093 #[test]
3094 fn closed_domain_vocabularies_round_trip_and_reject_unknown_tokens() {
3095 for value in [GraphObjectKind::Node, GraphObjectKind::Edge] {
3096 assert_eq!(GraphObjectKind::parse(value.as_str()).unwrap(), value);
3097 }
3098 assert!(GraphObjectKind::parse("vertex").is_err());
3099 for value in [
3100 AssertionGraphRole::Subject,
3101 AssertionGraphRole::Object,
3102 AssertionGraphRole::Context,
3103 ] {
3104 assert_eq!(AssertionGraphRole::parse(value.as_str()).unwrap(), value);
3105 }
3106 assert!(AssertionGraphRole::parse("target").is_err());
3107 for value in [
3108 EvidenceSourceKind::Document,
3109 EvidenceSourceKind::Observation,
3110 EvidenceSourceKind::GraphNode,
3111 EvidenceSourceKind::GraphEdge,
3112 ] {
3113 assert_eq!(EvidenceSourceKind::parse(value.as_str()).unwrap(), value);
3114 }
3115 assert!(EvidenceSourceKind::parse("web").is_err());
3116 for value in [
3117 EvidenceRole::Supports,
3118 EvidenceRole::Contradicts,
3119 EvidenceRole::Context,
3120 ] {
3121 assert_eq!(EvidenceRole::parse(value.as_str()).unwrap(), value);
3122 }
3123 assert!(EvidenceRole::parse("proves").is_err());
3124 for value in [
3125 ConfidencePolicy::Explicit,
3126 ConfidencePolicy::ConservativeMin,
3127 ] {
3128 assert_eq!(ConfidencePolicy::parse(value.as_str()).unwrap(), value);
3129 }
3130 assert!(ConfidencePolicy::parse("average").is_err());
3131 for value in [
3132 AlgorithmRunState::Started,
3133 AlgorithmRunState::Completed,
3134 AlgorithmRunState::Failed,
3135 AlgorithmRunState::Cancelled,
3136 AlgorithmRunState::Interrupted,
3137 ] {
3138 assert_eq!(AlgorithmRunState::parse(value.as_str()).unwrap(), value);
3139 }
3140 assert!(!AlgorithmRunState::Started.is_terminal());
3141 assert!(AlgorithmRunState::Completed.is_terminal());
3142 assert!(AlgorithmRunState::parse("running").is_err());
3143 }
3144
3145 #[test]
3146 fn knowledge_error_codes_are_closed_and_exact() {
3147 let cases = [
3148 (invalid("field", "bad"), "GF_KNOWLEDGE_INVALID"),
3149 (
3150 KnowledgeError::Limit {
3151 participant: "assertions",
3152 observed: 2,
3153 limit: 1,
3154 },
3155 "GF_RESOURCE_LIMIT",
3156 ),
3157 (KnowledgeError::Duplicate("id"), "GF_KNOWLEDGE_DUPLICATE"),
3158 (KnowledgeError::Dangling("id"), "GF_KNOWLEDGE_DANGLING"),
3159 (KnowledgeError::Conflict("id"), "GF_IDEMPOTENCY_CONFLICT"),
3160 (
3161 KnowledgeError::TransactionConflict("id"),
3162 "GF_TRANSACTION_CONFLICT",
3163 ),
3164 (
3165 KnowledgeError::Canonical(CanonicalError::Malformed("bad")),
3166 "GF_CANONICAL_INVALID",
3167 ),
3168 (
3169 KnowledgeError::Arrow(arrow::error::ArrowError::SchemaError("bad".into())),
3170 "GF_SCHEMA_MISMATCH",
3171 ),
3172 ];
3173 for (error, code) in cases {
3174 assert_eq!(error.code(), code);
3175 assert!(!error.to_string().is_empty());
3176 }
3177 }
3178
3179 #[test]
3180 fn algorithm_run_validation_rejects_every_malformed_identity_and_transition() {
3181 let run = AlgorithmRun::new(
3182 uuid7(40),
3183 "pagerank".into(),
3184 1,
3185 1,
3186 vec![1],
3187 [2; 32],
3188 uuid7(41),
3189 100,
3190 )
3191 .unwrap();
3192 for invalid_run in [
3193 AlgorithmRun {
3194 algorithm: String::new(),
3195 ..run.clone()
3196 },
3197 AlgorithmRun {
3198 algorithm_version: 2,
3199 ..run.clone()
3200 },
3201 AlgorithmRun {
3202 descriptor_version: 2,
3203 ..run.clone()
3204 },
3205 AlgorithmRun {
3206 descriptor: Vec::new(),
3207 ..run.clone()
3208 },
3209 AlgorithmRun {
3210 contract_version: 2,
3211 ..run.clone()
3212 },
3213 ] {
3214 assert!(matches!(
3215 validate_algorithm_run(&invalid_run),
3216 Err(KnowledgeError::Invalid { .. })
3217 ));
3218 }
3219
3220 let start = AlgorithmRunEvent::new(
3221 uuid7(42),
3222 run.run_uuid,
3223 AlgorithmRunState::Started,
3224 None,
3225 None,
3226 100,
3227 run.provenance_uuid,
3228 )
3229 .unwrap();
3230 let completed = AlgorithmRunEvent::new(
3231 uuid7(43),
3232 run.run_uuid,
3233 AlgorithmRunState::Completed,
3234 Some([3; 32]),
3235 None,
3236 101,
3237 uuid7(44),
3238 )
3239 .unwrap();
3240 assert_eq!(
3241 AlgorithmRunLedger::new(vec![run.clone()], vec![start.clone(), completed.clone()])
3242 .unwrap()
3243 .events_for(run.run_uuid)
3244 .len(),
3245 2
3246 );
3247 assert!(matches!(
3248 validate_algorithm_run_event(&AlgorithmRunEvent {
3249 contract_version: 2,
3250 ..start.clone()
3251 }),
3252 Err(KnowledgeError::Invalid { .. })
3253 ));
3254 for event in [
3255 AlgorithmRunEvent {
3256 result_fingerprint: Some([1; 32]),
3257 ..start.clone()
3258 },
3259 AlgorithmRunEvent {
3260 result_fingerprint: None,
3261 ..completed.clone()
3262 },
3263 AlgorithmRunEvent {
3264 state: AlgorithmRunState::Failed,
3265 result_fingerprint: None,
3266 error_code: None,
3267 ..completed.clone()
3268 },
3269 AlgorithmRunEvent {
3270 state: AlgorithmRunState::Failed,
3271 result_fingerprint: None,
3272 error_code: Some("bad".into()),
3273 ..completed.clone()
3274 },
3275 ] {
3276 assert!(matches!(
3277 validate_algorithm_run_event(&event),
3278 Err(KnowledgeError::Invalid { .. })
3279 ));
3280 }
3281
3282 assert!(matches!(
3283 AlgorithmRunLedger::new(vec![run.clone(), run.clone()], vec![start.clone()]),
3284 Err(KnowledgeError::Duplicate("run_uuid"))
3285 ));
3286 assert!(matches!(
3287 AlgorithmRunLedger::new(vec![run.clone()], vec![start.clone(), start.clone()]),
3288 Err(KnowledgeError::Duplicate("event_uuid"))
3289 ));
3290 assert!(matches!(
3291 AlgorithmRunLedger::new(
3292 vec![run.clone()],
3293 vec![AlgorithmRunEvent {
3294 run_uuid: uuid7(50),
3295 ..start.clone()
3296 }]
3297 ),
3298 Err(KnowledgeError::Dangling("run_uuid"))
3299 ));
3300 assert!(matches!(
3301 AlgorithmRunLedger::new(
3302 vec![run.clone()],
3303 vec![AlgorithmRunEvent {
3304 recorded_at_micros: 99,
3305 ..start.clone()
3306 }]
3307 ),
3308 Err(KnowledgeError::Invalid {
3309 field: "recorded_at",
3310 ..
3311 })
3312 ));
3313 assert!(matches!(
3314 AlgorithmRunLedger::new(vec![run.clone()], Vec::new()),
3315 Err(KnowledgeError::Invalid {
3316 field: "started",
3317 ..
3318 })
3319 ));
3320 assert!(matches!(
3321 AlgorithmRunLedger::new(
3322 vec![run],
3323 vec![
3324 start,
3325 completed.clone(),
3326 AlgorithmRunEvent {
3327 event_uuid: uuid7(51),
3328 ..completed
3329 },
3330 ]
3331 ),
3332 Err(KnowledgeError::Invalid {
3333 field: "terminal",
3334 ..
3335 })
3336 ));
3337 }
3338
3339 #[test]
3340 fn assertion_and_algorithm_run_merges_cover_idempotent_append_and_conflict_paths() {
3341 let base = fixture();
3342 assert_eq!(base.merge(&base).unwrap(), base);
3343 let second_id = uuid7(20);
3344 let second = AssertionLedger::new(
3345 vec![Assertion::new(second_id, "second".into(), uuid7(21), 20).unwrap()],
3346 vec![
3347 AssertionGraphRef::new(
3348 second_id,
3349 uuid7(22),
3350 GraphObjectKind::Node,
3351 AssertionGraphRole::Subject,
3352 0,
3353 )
3354 .unwrap(),
3355 ],
3356 )
3357 .unwrap();
3358 let merged = base.merge(&second).unwrap();
3359 assert_eq!(merged.assertions.len(), 2);
3360 let conflicting = AssertionLedger::new(
3361 vec![Assertion::new(uuid7(1), "different".into(), uuid7(2), 10).unwrap()],
3362 base.graph_refs.clone(),
3363 )
3364 .unwrap();
3365 assert!(matches!(
3366 base.merge(&conflicting),
3367 Err(KnowledgeError::Conflict("assertion_uuid"))
3368 ));
3369
3370 let run = AlgorithmRun::new(
3371 uuid7(40),
3372 "pagerank".into(),
3373 1,
3374 1,
3375 vec![1],
3376 [2; 32],
3377 uuid7(41),
3378 100,
3379 )
3380 .unwrap();
3381 let started = AlgorithmRunEvent::new(
3382 uuid7(42),
3383 run.run_uuid,
3384 AlgorithmRunState::Started,
3385 None,
3386 None,
3387 100,
3388 run.provenance_uuid,
3389 )
3390 .unwrap();
3391 let initial = AlgorithmRunLedger::new(vec![run.clone()], vec![started.clone()]).unwrap();
3392 assert_eq!(initial.run(run.run_uuid), Some(&run));
3393 assert_eq!(initial.events_for(run.run_uuid), vec![started.clone()]);
3394 assert!(initial.terminal_event(run.run_uuid).is_none());
3395 assert_eq!(initial.merge(&initial).unwrap(), initial);
3396
3397 let completed = AlgorithmRunEvent::new(
3398 uuid7(43),
3399 run.run_uuid,
3400 AlgorithmRunState::Completed,
3401 Some([3; 32]),
3402 None,
3403 101,
3404 uuid7(44),
3405 )
3406 .unwrap();
3407 let staged =
3408 AlgorithmRunLedger::new(vec![run.clone()], vec![started, completed.clone()]).unwrap();
3409 let merged = initial.merge(&staged).unwrap();
3410 assert_eq!(merged.terminal_event(run.run_uuid), Some(&completed));
3411
3412 let conflicting_run = AlgorithmRun {
3413 algorithm: "hits".into(),
3414 ..run.clone()
3415 };
3416 let conflict = AlgorithmRunLedger {
3417 runs: vec![conflicting_run],
3418 events: vec![],
3419 };
3420 assert!(matches!(
3421 initial.merge(&conflict),
3422 Err(KnowledgeError::Conflict("run_uuid"))
3423 ));
3424 let conflicting_event = AlgorithmRunEvent {
3425 error_code: Some("GF_EXECUTION".into()),
3426 state: AlgorithmRunState::Failed,
3427 ..completed
3428 };
3429 let conflict = AlgorithmRunLedger {
3430 runs: vec![],
3431 events: vec![conflicting_event],
3432 };
3433 assert!(matches!(
3434 merged.merge(&conflict),
3435 Err(KnowledgeError::Conflict("event_uuid"))
3436 ));
3437 }
3438}
3439pub use belief_projection::{
3440 ALGORITHM_INTERPRETATION_ATTACHMENT_SCHEMA, BELIEF_PROJECTION_ATTACHMENT_CONTRACT_VERSION,
3441 BeliefProjectionAttachment, BeliefProjectionAttachmentLedger,
3442};