Skip to main content

graphforge_knowledge/
lib.rs

1//! Immutable UUID-referenced analytical knowledge.
2//!
3//! This crate owns knowledge records, validation, canonical fingerprints,
4//! deterministic ordering, and Arrow schemas. It deliberately has no storage,
5//! graph, execution, or provenance dependency.
6#![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
56/// Knowledge capability contract implemented by this crate.
57pub const KNOWLEDGE_CAPABILITY_VERSION: u32 = 1;
58/// Assertion record contract.
59pub const ASSERTION_CONTRACT_VERSION: u32 = 1;
60/// Assertion-to-graph reference record contract.
61pub const ASSERTION_GRAPH_REF_CONTRACT_VERSION: u32 = 1;
62/// Confidence-assessment record contract.
63pub const CONFIDENCE_ASSESSMENT_CONTRACT_VERSION: u32 = 1;
64/// Confidence-input snapshot record contract.
65pub const CONFIDENCE_INPUT_CONTRACT_VERSION: u32 = 1;
66/// Evidence-link record contract.
67pub const EVIDENCE_LINK_CONTRACT_VERSION: u32 = 1;
68/// Algorithm-run identity record contract.
69pub const ALGORITHM_RUN_CONTRACT_VERSION: u32 = 1;
70/// Algorithm-run lifecycle event contract.
71pub const ALGORITHM_RUN_EVENT_CONTRACT_VERSION: u32 = 1;
72/// Closed confidence-policy registry version.
73pub const CONFIDENCE_POLICY_REGISTRY_VERSION: u32 = 1;
74/// Closed graph-object-kind registry version.
75pub const GRAPH_OBJECT_KIND_REGISTRY_VERSION: u32 = 1;
76/// Closed assertion-role registry version.
77pub const ASSERTION_GRAPH_ROLE_REGISTRY_VERSION: u32 = 1;
78/// Closed evidence source-kind registry version.
79pub const EVIDENCE_SOURCE_KIND_REGISTRY_VERSION: u32 = 1;
80/// Closed evidence role registry version.
81pub const EVIDENCE_ROLE_REGISTRY_VERSION: u32 = 1;
82/// Closed algorithm-run lifecycle registry version.
83pub const ALGORITHM_RUN_STATE_REGISTRY_VERSION: u32 = 1;
84/// Per-participant row bound.
85pub const MAX_KNOWLEDGE_ROWS: usize = 1_000_000;
86
87/// Authoritative `knowledge/assertions.parquet` schema.
88pub 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
102/// Authoritative `knowledge/assertion_graph_refs.parquet` schema.
103pub 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
114/// Authoritative `knowledge/confidence_assessments.parquet` schema.
115pub 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
132/// Authoritative `knowledge/confidence_inputs.parquet` schema.
133pub 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
143/// Authoritative `knowledge/evidence.parquet` schema.
144pub 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
162/// Authoritative `knowledge/algorithm_runs.parquet` schema.
163pub 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
185/// Authoritative `knowledge/algorithm_run_events.parquet` schema.
186pub 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/// Closed graph UUID kind.
271#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
272#[serde(rename_all = "snake_case")]
273pub enum GraphObjectKind {
274    /// Public node UUID.
275    Node,
276    /// Public edge UUID.
277    Edge,
278}
279
280impl GraphObjectKind {
281    /// Canonical persisted spelling.
282    #[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/// Closed assertion-to-graph role.
300#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
301#[serde(rename_all = "snake_case")]
302pub enum AssertionGraphRole {
303    /// Claim subject.
304    Subject,
305    /// Claim object.
306    Object,
307    /// Context needed to interpret the claim.
308    Context,
309}
310
311impl AssertionGraphRole {
312    /// Canonical persisted spelling.
313    #[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/// One immutable analytical assertion.
333#[derive(Clone, Debug, PartialEq, Eq)]
334pub struct Assertion {
335    /// Caller-supplied UUIDv7 identity and idempotency key.
336    pub assertion_uuid: Uuid,
337    /// Exact validated UTF-8 claim bytes.
338    pub claim: String,
339    /// Producing provenance event.
340    pub provenance_uuid: Uuid,
341    /// Transaction time in UTC microseconds.
342    pub recorded_at_micros: i64,
343    /// Assertion record contract.
344    pub contract_version: u32,
345}
346
347impl Assertion {
348    /// Construct one assertion record.
349    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/// One immutable assertion-to-graph UUID reference.
369#[derive(Clone, Debug, PartialEq, Eq)]
370pub struct AssertionGraphRef {
371    /// Owning assertion.
372    pub assertion_uuid: Uuid,
373    /// Referenced public graph UUID.
374    pub graph_uuid: Uuid,
375    /// Node or edge.
376    pub graph_kind: GraphObjectKind,
377    /// Subject/object/context role.
378    pub role: AssertionGraphRole,
379    /// Caller-significant contiguous position within the role.
380    pub ordinal: u32,
381    /// Reference record contract.
382    pub contract_version: u32,
383}
384
385impl AssertionGraphRef {
386    /// Construct one graph reference.
387    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/// Validated immutable assertion participant content.
408#[derive(Clone, Debug, Default, PartialEq, Eq)]
409pub struct AssertionLedger {
410    /// Assertions ordered by `(recorded_at, assertion_uuid)`.
411    pub assertions: Vec<Assertion>,
412    /// References ordered by assertion and the public reference sort key.
413    pub graph_refs: Vec<AssertionGraphRef>,
414}
415
416impl AssertionLedger {
417    /// Validate, sort, and construct assertion content.
418    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    /// Merge a staged assertion set idempotently.
445    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    /// Canonical assertion fingerprint over exact claim bytes and sorted refs.
476    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    /// Build the authoritative assertion Arrow batch.
508    pub fn assertion_batch(&self) -> Result<RecordBatch, KnowledgeError> {
509        assertion_batch(&self.assertions)
510    }
511
512    /// Build the authoritative assertion graph-reference Arrow batch.
513    pub fn graph_ref_batch(&self) -> Result<RecordBatch, KnowledgeError> {
514        graph_ref_batch(&self.graph_refs)
515    }
516
517    /// Decode authoritative Arrow batches and re-run every invariant.
518    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/// Closed evidence source kind.
569#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
570#[serde(rename_all = "snake_case")]
571pub enum EvidenceSourceKind {
572    /// Caller-managed document identity.
573    Document,
574    /// Caller-managed observation identity.
575    Observation,
576    /// Existing graph node identity.
577    GraphNode,
578    /// Existing graph edge identity.
579    GraphEdge,
580}
581
582impl EvidenceSourceKind {
583    /// Canonical persisted spelling.
584    #[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/// Closed relationship between evidence and an assertion.
606#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
607#[serde(rename_all = "snake_case")]
608pub enum EvidenceRole {
609    /// Evidence supports the assertion.
610    Supports,
611    /// Evidence contradicts the assertion.
612    Contradicts,
613    /// Evidence supplies interpretation context.
614    Context,
615}
616
617impl EvidenceRole {
618    /// Canonical persisted spelling.
619    #[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/// One immutable evidence link.
639#[derive(Clone, Debug, PartialEq)]
640pub struct EvidenceLink {
641    /// UUIDv7 identity and idempotency key.
642    pub evidence_uuid: Uuid,
643    /// Existing immutable assertion.
644    pub assertion_uuid: Uuid,
645    /// Caller-managed source identity.
646    pub source_uuid: Uuid,
647    /// Closed source kind.
648    pub source_kind: EvidenceSourceKind,
649    /// Closed relationship to the assertion.
650    pub role: EvidenceRole,
651    /// Optional finite metadata weight in `[0, 1]`.
652    pub weight: Option<f64>,
653    /// Provenance event identity.
654    pub provenance_uuid: Uuid,
655    /// Transaction time in UTC microseconds.
656    pub recorded_at_micros: i64,
657    /// Evidence-link record contract.
658    pub contract_version: u32,
659}
660
661impl EvidenceLink {
662    /// Construct one validated immutable evidence link.
663    #[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/// Validated immutable evidence participant content.
695#[derive(Clone, Debug, Default, PartialEq)]
696pub struct EvidenceLedger {
697    /// Links ordered by `(recorded_at, evidence_uuid)`.
698    pub links: Vec<EvidenceLink>,
699}
700
701impl EvidenceLedger {
702    /// Validate, sort, and construct evidence content.
703    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    /// Merge staged links idempotently.
732    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    /// Canonical fingerprint over normalized immutable content.
750    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    /// Build the authoritative evidence Arrow batch.
772    pub fn batch(&self) -> Result<RecordBatch, KnowledgeError> {
773        evidence_batch(&self.links)
774    }
775
776    /// Decode authoritative Arrow batches and re-run every invariant.
777    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/// Closed append-only algorithm-run lifecycle state.
813#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
814#[serde(rename_all = "snake_case")]
815pub enum AlgorithmRunState {
816    /// Identity was durably published before dispatch.
817    Started,
818    /// Dispatch returned a canonical Arrow result.
819    Completed,
820    /// Dispatch returned a structured failure.
821    Failed,
822    /// Cancellation was observed at a deterministic checkpoint.
823    Cancelled,
824    /// Reopen found a published start without a terminal event.
825    Interrupted,
826}
827
828impl AlgorithmRunState {
829    /// Canonical persisted spelling.
830    #[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    /// Whether this state closes a run.
853    #[must_use]
854    pub const fn is_terminal(self) -> bool {
855        !matches!(self, Self::Started)
856    }
857}
858
859/// Immutable identity for one recorded M18 invocation.
860#[derive(Clone, Debug, PartialEq, Eq)]
861pub struct AlgorithmRun {
862    /// Caller-supplied UUIDv7 run identity.
863    pub run_uuid: Uuid,
864    /// Closed public M18 algorithm name.
865    pub algorithm: String,
866    /// Algorithm contract version.
867    pub algorithm_version: u32,
868    /// Neutral descriptor contract version.
869    pub descriptor_version: u32,
870    /// Exact canonical descriptor bytes.
871    pub descriptor: Vec<u8>,
872    /// Exact resolved graph projection fingerprint.
873    pub projection_fingerprint: [u8; 32],
874    /// Provenance event that published the run identity.
875    pub provenance_uuid: Uuid,
876    /// Durable start transaction time.
877    pub started_at_micros: i64,
878    /// Run-record contract version.
879    pub contract_version: u32,
880}
881
882impl AlgorithmRun {
883    /// Construct one immutable run identity.
884    #[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/// One immutable event in a recorded algorithm lifecycle.
912#[derive(Clone, Debug, PartialEq, Eq)]
913pub struct AlgorithmRunEvent {
914    /// Deterministic event identity.
915    pub event_uuid: Uuid,
916    /// Owning run identity.
917    pub run_uuid: Uuid,
918    /// Closed lifecycle state.
919    pub state: AlgorithmRunState,
920    /// Canonical Arrow fingerprint for a completed result.
921    pub result_fingerprint: Option<[u8; 32]>,
922    /// Sanitized stable error code for non-success terminal states.
923    pub error_code: Option<String>,
924    /// Durable transaction time.
925    pub recorded_at_micros: i64,
926    /// Provenance event for this lifecycle transition.
927    pub provenance_uuid: Uuid,
928    /// Lifecycle-event contract version.
929    pub contract_version: u32,
930}
931
932impl AlgorithmRunEvent {
933    /// Construct one validated lifecycle event.
934    #[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/// Validated immutable run identities and append-only lifecycle events.
960#[derive(Clone, Debug, Default, PartialEq, Eq)]
961pub struct AlgorithmRunLedger {
962    /// Run identities ordered by `(started_at, run_uuid)`.
963    pub runs: Vec<AlgorithmRun>,
964    /// Events ordered by `(recorded_at, event_uuid)`.
965    pub events: Vec<AlgorithmRunEvent>,
966}
967
968impl AlgorithmRunLedger {
969    /// Validate and normalize complete run tables.
970    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    /// Merge immutable identities and events, rejecting conflicting reuse.
981    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    /// Locate one run.
1005    #[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    /// Return lifecycle events for one run in canonical order.
1011    #[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    /// Return the terminal event, when one exists.
1021    #[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    /// Encode the authoritative run table.
1029    pub fn run_batch(&self) -> Result<RecordBatch, KnowledgeError> {
1030        algorithm_run_batch(&self.runs)
1031    }
1032
1033    /// Encode the authoritative event table.
1034    pub fn event_batch(&self) -> Result<RecordBatch, KnowledgeError> {
1035        algorithm_run_event_batch(&self.events)
1036    }
1037
1038    /// Decode, validate, and normalize persisted tables.
1039    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/// Closed confidence policy registry.
1106#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
1107#[serde(rename_all = "snake_case")]
1108pub enum ConfidencePolicy {
1109    /// Caller supplies the assessment value.
1110    Explicit,
1111    /// Minimum of all available, non-null requested inputs; null if any is unavailable.
1112    ConservativeMin,
1113}
1114
1115impl ConfidencePolicy {
1116    /// Canonical persisted spelling.
1117    #[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/// One immutable confidence assessment.
1135#[derive(Clone, Debug, PartialEq)]
1136pub struct ConfidenceAssessment {
1137    /// Caller-supplied UUIDv7 identity and idempotency key.
1138    pub confidence_uuid: Uuid,
1139    /// Assertion being assessed.
1140    pub assertion_uuid: Uuid,
1141    /// Closed policy.
1142    pub policy: ConfidencePolicy,
1143    /// Policy contract version.
1144    pub policy_version: u32,
1145    /// Result in `[0, 1]`, or null when conservative inputs are incomplete.
1146    pub value: Option<f64>,
1147    /// Producing provenance event.
1148    pub provenance_uuid: Uuid,
1149    /// Transaction time in UTC microseconds.
1150    pub recorded_at_micros: i64,
1151    /// Assessment record contract.
1152    pub contract_version: u32,
1153}
1154
1155impl ConfidenceAssessment {
1156    /// Construct one validated assessment.
1157    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/// Immutable snapshot of one requested confidence input.
1183#[derive(Clone, Debug, PartialEq)]
1184pub struct ConfidenceInput {
1185    /// Owning assessment.
1186    pub confidence_uuid: Uuid,
1187    /// Requested immutable assessment identity.
1188    pub input_confidence_uuid: Uuid,
1189    /// Value observed at assessment time; null means absent or null.
1190    pub input_value: Option<f64>,
1191    /// UUID-normalized position.
1192    pub ordinal: u32,
1193    /// Input record contract.
1194    pub contract_version: u32,
1195}
1196
1197impl ConfidenceInput {
1198    /// Construct one validated snapshot input.
1199    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/// Validated append-only confidence participant content.
1219#[derive(Clone, Debug, Default, PartialEq)]
1220pub struct ConfidenceLedger {
1221    /// Assessments ordered by `(recorded_at, confidence_uuid)`.
1222    pub assessments: Vec<ConfidenceAssessment>,
1223    /// Inputs ordered by assessment then `(ordinal, input_confidence_uuid)`.
1224    pub inputs: Vec<ConfidenceInput>,
1225}
1226
1227impl ConfidenceLedger {
1228    /// Validate, sort, and construct confidence content.
1229    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    /// Evaluate and stage an explicit assessment.
1255    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    /// Evaluate `conservative_min@1` and persist the normalized requested-input snapshot.
1276    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    /// Merge staged content idempotently.
1329    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    /// Canonical assessment fingerprint over policy, normalized value, and input snapshot.
1358    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    /// Build the authoritative assessment Arrow batch.
1389    pub fn assessment_batch(&self) -> Result<RecordBatch, KnowledgeError> {
1390        confidence_assessment_batch(&self.assessments)
1391    }
1392
1393    /// Build the authoritative input Arrow batch.
1394    pub fn input_batch(&self) -> Result<RecordBatch, KnowledgeError> {
1395        confidence_input_batch(&self.inputs)
1396    }
1397
1398    /// Decode authoritative Arrow batches and re-run every invariant.
1399    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/// Authoritative registry entry for one knowledge record family.
1450#[derive(Clone, Debug)]
1451pub struct SchemaRegistryEntry {
1452    /// Stable capability ID.
1453    pub capability_id: &'static str,
1454    /// Capability contract version.
1455    pub capability_version: u32,
1456    /// Stable record-family ID.
1457    pub record_family: &'static str,
1458    /// Record contract version.
1459    pub record_version: u32,
1460    /// Exact Arrow schema.
1461    pub schema: SchemaRef,
1462    /// Canonical schema fingerprint.
1463    pub schema_fingerprint: [u8; 32],
1464    /// Closed enum registries used by this family.
1465    pub enum_registry_versions: &'static [(&'static str, u32)],
1466    /// Canonical persisted sort key.
1467    pub sort_key: &'static [&'static str],
1468    /// Logical fields that uniquely identify one record for checkpoint diffs.
1469    pub diff_identity_fields: &'static [&'static str],
1470    /// Logical UUID field surfaced as `record_uuid`, when this family owns one.
1471    pub diff_record_uuid_field: Option<&'static str>,
1472    /// Fingerprint domain.
1473    pub fingerprint_domain: CanonicalDomain,
1474    /// Owning crate.
1475    pub owner: &'static str,
1476    /// Implementation issue.
1477    pub implementation_issue: u64,
1478    /// Maximum accepted rows.
1479    pub max_rows: usize,
1480}
1481
1482impl SchemaRegistryEntry {
1483    /// Domain for owner-declared logical identity projections in checkpoint diffs.
1484    #[must_use]
1485    pub const fn diff_identity_fingerprint_domain(&self) -> CanonicalDomain {
1486        CanonicalDomain::ArrowResult
1487    }
1488
1489    /// Domain for owner-canonical whole-record checkpoint fingerprints.
1490    #[must_use]
1491    pub const fn diff_record_fingerprint_domain(&self) -> CanonicalDomain {
1492        self.fingerprint_domain
1493    }
1494}
1495
1496/// Return the authoritative assertion schema registry.
1497#[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/// Structured knowledge-domain failures.
1651#[derive(thiserror::Error, Debug)]
1652pub enum KnowledgeError {
1653    /// Invalid record value or derived identity.
1654    #[error("invalid knowledge {field}: {message}")]
1655    Invalid {
1656        /// Safe field name.
1657        field: &'static str,
1658        /// Safe failure summary.
1659        message: &'static str,
1660    },
1661    /// Participant row limit exceeded.
1662    #[error("knowledge {participant} row limit exceeded: observed {observed}, limit {limit}")]
1663    Limit {
1664        /// Safe participant name.
1665        participant: &'static str,
1666        /// Observed rows.
1667        observed: usize,
1668        /// Maximum rows.
1669        limit: usize,
1670    },
1671    /// Duplicate identity in one participant.
1672    #[error("duplicate knowledge identity: {0}")]
1673    Duplicate(&'static str),
1674    /// A required assertion or graph UUID is absent.
1675    #[error("dangling knowledge reference: {0}")]
1676    Dangling(&'static str),
1677    /// Idempotency identity was reused for different content.
1678    #[error("knowledge idempotency conflict: {0}")]
1679    Conflict(&'static str),
1680    /// A transaction identity was reused for different immutable content.
1681    #[error("knowledge transaction conflict: {0}")]
1682    TransactionConflict(&'static str),
1683    /// Shared canonicalization failure.
1684    #[error(transparent)]
1685    Canonical(#[from] CanonicalError),
1686    /// Arrow construction failure.
1687    #[error("knowledge Arrow failure: {0}")]
1688    Arrow(#[from] arrow::error::ArrowError),
1689}
1690
1691impl KnowledgeError {
1692    /// Stable public error code.
1693    #[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        // This locks the Assertion domain plus GFAS framing. A digest change is
2595        // a contract break and requires a version bump, never a golden refresh.
2596        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        // Locks `graphforge/confidence-assessment` plus GFCA framing.
2855        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};