Skip to main content

graphforge_knowledge/
belief_projection.rs

1//! Append-only M21 attachments connecting resolved interpretation to completed M20 runs.
2
3use std::collections::{HashMap, HashSet};
4use std::sync::{Arc, LazyLock};
5
6use arrow::array::{
7    Array, BinaryArray, BinaryBuilder, FixedSizeBinaryArray, FixedSizeBinaryBuilder, ListArray,
8    ListBuilder, TimestampMicrosecondArray, TimestampMicrosecondBuilder, UInt32Array,
9    UInt32Builder,
10};
11use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
12use arrow::record_batch::RecordBatch;
13use graphforge_core::canonical::{
14    CANONICAL_CONTRACT_VERSION, CanonicalDomain, CanonicalWriter, MAX_CANONICAL_BINARY_BYTES,
15    fingerprint,
16};
17use uuid::{Uuid, Version};
18
19use crate::{
20    EPISTEMIC_CAPABILITY_VERSION, KnowledgeError, MAX_KNOWLEDGE_ROWS, SchemaRegistryEntry,
21};
22
23/// Frozen `belief_projection_attachment/1` record contract.
24pub const BELIEF_PROJECTION_ATTACHMENT_CONTRACT_VERSION: u32 = 1;
25
26/// Authoritative `algorithm_interpretation_attachments` schema.
27pub static ALGORITHM_INTERPRETATION_ATTACHMENT_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
28    let uuid_list = DataType::List(Arc::new(Field::new(
29        "item",
30        DataType::FixedSizeBinary(16),
31        false,
32    )));
33    Arc::new(Schema::new(vec![
34        uuid_field("attachment_uuid", false),
35        uuid_field("run_uuid", false),
36        uuid_field("source_generation_uuid", false),
37        timestamp_field("transaction_cutoff", false),
38        timestamp_field("valid_time", true),
39        Field::new("policy_version", DataType::UInt32, false),
40        Field::new("policy_bytes", DataType::Binary, false),
41        fingerprint_field("policy_fingerprint", false),
42        fingerprint_field("snapshot_fingerprint", false),
43        fingerprint_field("valid_time_fingerprint", true),
44        fingerprint_field("graph_content_fingerprint", false),
45        fingerprint_field("descriptor_fingerprint", false),
46        Field::new("source_record_uuids", uuid_list, false),
47        uuid_field("provenance_uuid", false),
48        timestamp_field("recorded_at", false),
49        Field::new("contract_version", DataType::UInt32, false),
50    ]))
51});
52
53static SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
54    fingerprint(
55        CanonicalDomain::Schema,
56        CANONICAL_CONTRACT_VERSION,
57        b"belief_projection_attachment/1|attachment_uuid:fixed[16]:required|run_uuid:fixed[16]:required|source_generation_uuid:fixed[16]:required|transaction_cutoff:timestamp_us_utc:required|valid_time:timestamp_us_utc:nullable|policy_version:u32:required|policy_bytes:binary:required|policy_fingerprint:fixed[32]:required|snapshot_fingerprint:fixed[32]:required|valid_time_fingerprint:fixed[32]:nullable|graph_content_fingerprint:fixed[32]:required|descriptor_fingerprint:fixed[32]:required|source_record_uuids:list<fixed[16]>:required|provenance_uuid:fixed[16]:required|recorded_at:timestamp_us_utc:required|contract_version:u32:required",
58    )
59    .expect("registered belief-projection attachment schema is bounded")
60});
61
62/// One immutable interpretation attachment for a completed algorithm run.
63#[derive(Clone, Debug, Eq, PartialEq)]
64pub struct BeliefProjectionAttachment {
65    /// Stable attachment identity and idempotency key.
66    pub attachment_uuid: Uuid,
67    /// Completed M20 run identity.
68    pub run_uuid: Uuid,
69    /// Source project generation pinned during resolution.
70    pub source_generation_uuid: Uuid,
71    /// Mandatory transaction-time cutoff.
72    pub transaction_cutoff_micros: i64,
73    /// Optional valid time applied after transaction reconstruction.
74    pub valid_time_micros: Option<i64>,
75    /// Explicit resolved-belief policy version.
76    pub policy_version: u32,
77    /// Canonical language-neutral policy bytes.
78    pub policy_bytes: Vec<u8>,
79    /// Canonical fingerprint of `policy_bytes`.
80    pub policy_fingerprint: [u8; 32],
81    /// Fingerprint of the transaction-time epistemic snapshot.
82    pub snapshot_fingerprint: [u8; 32],
83    /// Fingerprint of the optional valid-time interpretation.
84    pub valid_time_fingerprint: Option<[u8; 32]>,
85    /// Canonical logical graph-content fingerprint.
86    pub graph_content_fingerprint: [u8; 32],
87    /// Neutral M18 invocation-descriptor fingerprint.
88    pub descriptor_fingerprint: [u8; 32],
89    /// Sorted and deduplicated decision-relevant M21 records.
90    pub source_record_uuids: Vec<Uuid>,
91    /// Producing provenance event.
92    pub provenance_uuid: Uuid,
93    /// Mandatory attachment transaction time.
94    pub recorded_at_micros: i64,
95    /// Frozen record contract.
96    pub contract_version: u32,
97}
98
99impl BeliefProjectionAttachment {
100    /// Validate and construct an attachment, canonicalizing source UUIDs.
101    #[allow(clippy::too_many_arguments)]
102    pub fn new(
103        attachment_uuid: Uuid,
104        run_uuid: Uuid,
105        source_generation_uuid: Uuid,
106        transaction_cutoff_micros: i64,
107        valid_time_micros: Option<i64>,
108        policy_version: u32,
109        policy_bytes: Vec<u8>,
110        snapshot_fingerprint: [u8; 32],
111        valid_time_fingerprint: Option<[u8; 32]>,
112        graph_content_fingerprint: [u8; 32],
113        descriptor_fingerprint: [u8; 32],
114        mut source_record_uuids: Vec<Uuid>,
115        provenance_uuid: Uuid,
116        recorded_at_micros: i64,
117    ) -> Result<Self, KnowledgeError> {
118        source_record_uuids.sort_unstable();
119        source_record_uuids.dedup();
120        let policy_fingerprint = fingerprint(
121            CanonicalDomain::BeliefProjectionPolicy,
122            CANONICAL_CONTRACT_VERSION,
123            &policy_bytes,
124        )?;
125        let row = Self {
126            attachment_uuid,
127            run_uuid,
128            source_generation_uuid,
129            transaction_cutoff_micros,
130            valid_time_micros,
131            policy_version,
132            policy_bytes,
133            policy_fingerprint,
134            snapshot_fingerprint,
135            valid_time_fingerprint,
136            graph_content_fingerprint,
137            descriptor_fingerprint,
138            source_record_uuids,
139            provenance_uuid,
140            recorded_at_micros,
141            contract_version: BELIEF_PROJECTION_ATTACHMENT_CONTRACT_VERSION,
142        };
143        validate(&row)?;
144        Ok(row)
145    }
146}
147
148/// Validated append-only interpretation-attachment participant.
149#[derive(Clone, Debug, Default, Eq, PartialEq)]
150pub struct BeliefProjectionAttachmentLedger {
151    /// Attachments ordered by `(recorded_at, attachment_uuid)`.
152    pub attachments: Vec<BeliefProjectionAttachment>,
153}
154
155impl BeliefProjectionAttachmentLedger {
156    /// Validate, sort, and construct a complete participant.
157    pub fn new(mut attachments: Vec<BeliefProjectionAttachment>) -> Result<Self, KnowledgeError> {
158        if attachments.len() > MAX_KNOWLEDGE_ROWS {
159            return Err(KnowledgeError::Limit {
160                participant: "algorithm_interpretation_attachments",
161                observed: attachments.len(),
162                limit: MAX_KNOWLEDGE_ROWS,
163            });
164        }
165        let mut ids = HashSet::with_capacity(attachments.len());
166        for row in &attachments {
167            validate(row)?;
168            if !ids.insert(row.attachment_uuid) {
169                return Err(KnowledgeError::Duplicate("attachment_uuid"));
170            }
171        }
172        attachments.sort_by_key(|row| (row.recorded_at_micros, row.attachment_uuid));
173        Ok(Self { attachments })
174    }
175
176    /// Merge with exact replay and transaction-conflict semantics.
177    pub fn merge(&self, staged: &Self) -> Result<Self, KnowledgeError> {
178        let mut rows = self.attachments.clone();
179        let mut by_id = rows
180            .iter()
181            .cloned()
182            .map(|row| (row.attachment_uuid, row))
183            .collect::<HashMap<_, _>>();
184        for row in &staged.attachments {
185            if let Some(existing) = by_id.get(&row.attachment_uuid) {
186                if existing != row {
187                    return Err(KnowledgeError::TransactionConflict("attachment_uuid"));
188                }
189            } else {
190                rows.push(row.clone());
191                by_id.insert(row.attachment_uuid, row.clone());
192            }
193        }
194        Self::new(rows)
195    }
196
197    /// Canonical fingerprint over one exact immutable attachment.
198    pub fn attachment_fingerprint(&self, id: Uuid) -> Result<[u8; 32], KnowledgeError> {
199        let row = self
200            .attachments
201            .iter()
202            .find(|row| row.attachment_uuid == id)
203            .ok_or(KnowledgeError::Dangling("attachment_uuid"))?;
204        let mut writer = CanonicalWriter::new();
205        writer.raw(row.attachment_uuid.as_bytes())?;
206        writer.raw(row.run_uuid.as_bytes())?;
207        writer.raw(row.source_generation_uuid.as_bytes())?;
208        writer.i64(row.transaction_cutoff_micros)?;
209        optional_i64(&mut writer, row.valid_time_micros)?;
210        writer.u32(row.policy_version)?;
211        writer.binary(&row.policy_bytes)?;
212        writer.raw(&row.policy_fingerprint)?;
213        writer.raw(&row.snapshot_fingerprint)?;
214        optional_fingerprint(&mut writer, row.valid_time_fingerprint)?;
215        writer.raw(&row.graph_content_fingerprint)?;
216        writer.raw(&row.descriptor_fingerprint)?;
217        writer.u64(u64::try_from(row.source_record_uuids.len()).unwrap_or(u64::MAX))?;
218        for source in &row.source_record_uuids {
219            writer.raw(source.as_bytes())?;
220        }
221        writer.raw(row.provenance_uuid.as_bytes())?;
222        writer.i64(row.recorded_at_micros)?;
223        writer.u32(row.contract_version)?;
224        fingerprint(
225            CanonicalDomain::BeliefProjectionAttachment,
226            CANONICAL_CONTRACT_VERSION,
227            &writer.finish(),
228        )
229        .map_err(Into::into)
230    }
231
232    /// Build the authoritative Arrow batch.
233    pub fn batch(&self) -> Result<RecordBatch, KnowledgeError> {
234        let len = self.attachments.len();
235        let mut attachment = FixedSizeBinaryBuilder::with_capacity(len, 16);
236        let mut run = FixedSizeBinaryBuilder::with_capacity(len, 16);
237        let mut generation = FixedSizeBinaryBuilder::with_capacity(len, 16);
238        let mut cutoff = TimestampMicrosecondBuilder::new().with_timezone("UTC");
239        let mut valid_time = TimestampMicrosecondBuilder::new().with_timezone("UTC");
240        let mut policy_version = UInt32Builder::new();
241        let mut policy_bytes = BinaryBuilder::new();
242        let mut policy_fp = FixedSizeBinaryBuilder::with_capacity(len, 32);
243        let mut snapshot_fp = FixedSizeBinaryBuilder::with_capacity(len, 32);
244        let mut valid_fp = FixedSizeBinaryBuilder::with_capacity(len, 32);
245        let mut graph_fp = FixedSizeBinaryBuilder::with_capacity(len, 32);
246        let mut descriptor_fp = FixedSizeBinaryBuilder::with_capacity(len, 32);
247        let mut sources = ListBuilder::new(FixedSizeBinaryBuilder::new(16)).with_field(Arc::new(
248            Field::new("item", DataType::FixedSizeBinary(16), false),
249        ));
250        let mut provenance = FixedSizeBinaryBuilder::with_capacity(len, 16);
251        let mut recorded = TimestampMicrosecondBuilder::new().with_timezone("UTC");
252        let mut version = UInt32Builder::new();
253        for row in &self.attachments {
254            append(&mut attachment, row.attachment_uuid.as_bytes())?;
255            append(&mut run, row.run_uuid.as_bytes())?;
256            append(&mut generation, row.source_generation_uuid.as_bytes())?;
257            cutoff.append_value(row.transaction_cutoff_micros);
258            valid_time.append_option(row.valid_time_micros);
259            policy_version.append_value(row.policy_version);
260            policy_bytes.append_value(&row.policy_bytes);
261            append(&mut policy_fp, &row.policy_fingerprint)?;
262            append(&mut snapshot_fp, &row.snapshot_fingerprint)?;
263            append_optional(&mut valid_fp, row.valid_time_fingerprint.as_ref())?;
264            append(&mut graph_fp, &row.graph_content_fingerprint)?;
265            append(&mut descriptor_fp, &row.descriptor_fingerprint)?;
266            for source in &row.source_record_uuids {
267                append(sources.values(), source.as_bytes())?;
268            }
269            sources.append(true);
270            append(&mut provenance, row.provenance_uuid.as_bytes())?;
271            recorded.append_value(row.recorded_at_micros);
272            version.append_value(row.contract_version);
273        }
274        Ok(RecordBatch::try_new(
275            Arc::clone(&ALGORITHM_INTERPRETATION_ATTACHMENT_SCHEMA),
276            vec![
277                Arc::new(attachment.finish()),
278                Arc::new(run.finish()),
279                Arc::new(generation.finish()),
280                Arc::new(cutoff.finish()),
281                Arc::new(valid_time.finish()),
282                Arc::new(policy_version.finish()),
283                Arc::new(policy_bytes.finish()),
284                Arc::new(policy_fp.finish()),
285                Arc::new(snapshot_fp.finish()),
286                Arc::new(valid_fp.finish()),
287                Arc::new(graph_fp.finish()),
288                Arc::new(descriptor_fp.finish()),
289                Arc::new(sources.finish()),
290                Arc::new(provenance.finish()),
291                Arc::new(recorded.finish()),
292                Arc::new(version.finish()),
293            ],
294        )?)
295    }
296
297    /// Decode exact Arrow batches and re-run every invariant.
298    pub fn from_batches(batches: &[RecordBatch]) -> Result<Self, KnowledgeError> {
299        let mut rows = Vec::new();
300        for batch in batches {
301            if batch.schema().as_ref() != ALGORITHM_INTERPRETATION_ATTACHMENT_SCHEMA.as_ref() {
302                return Err(invalid(
303                    "belief_projection_attachment.schema",
304                    "schema mismatch",
305                ));
306            }
307            let attachment = fixed(batch, "attachment_uuid")?;
308            let run = fixed(batch, "run_uuid")?;
309            let generation = fixed(batch, "source_generation_uuid")?;
310            let cutoff = timestamp(batch, "transaction_cutoff")?;
311            let valid_time = timestamp(batch, "valid_time")?;
312            let policy_version = uint32(batch, "policy_version")?;
313            let policy_bytes = binary(batch, "policy_bytes")?;
314            let policy_fp = fixed(batch, "policy_fingerprint")?;
315            let snapshot_fp = fixed(batch, "snapshot_fingerprint")?;
316            let valid_fp = fixed(batch, "valid_time_fingerprint")?;
317            let graph_fp = fixed(batch, "graph_content_fingerprint")?;
318            let descriptor_fp = fixed(batch, "descriptor_fingerprint")?;
319            let sources = list(batch, "source_record_uuids")?;
320            let provenance = fixed(batch, "provenance_uuid")?;
321            let recorded = timestamp(batch, "recorded_at")?;
322            let versions = uint32(batch, "contract_version")?;
323            for row in 0..batch.num_rows() {
324                rows.push(BeliefProjectionAttachment {
325                    attachment_uuid: uuid_at(attachment, row, "attachment_uuid")?,
326                    run_uuid: uuid_at(run, row, "run_uuid")?,
327                    source_generation_uuid: uuid_at(generation, row, "source_generation_uuid")?,
328                    transaction_cutoff_micros: cutoff.value(row),
329                    valid_time_micros: optional_timestamp(valid_time, row),
330                    policy_version: policy_version.value(row),
331                    policy_bytes: policy_bytes.value(row).to_vec(),
332                    policy_fingerprint: bytes32(policy_fp, row, "policy_fingerprint")?,
333                    snapshot_fingerprint: bytes32(snapshot_fp, row, "snapshot_fingerprint")?,
334                    valid_time_fingerprint: optional_bytes32(
335                        valid_fp,
336                        row,
337                        "valid_time_fingerprint",
338                    )?,
339                    graph_content_fingerprint: bytes32(graph_fp, row, "graph_content_fingerprint")?,
340                    descriptor_fingerprint: bytes32(descriptor_fp, row, "descriptor_fingerprint")?,
341                    source_record_uuids: uuid_list_at(sources, row)?,
342                    provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
343                    recorded_at_micros: recorded.value(row),
344                    contract_version: versions.value(row),
345                });
346            }
347        }
348        Self::new(rows)
349    }
350}
351
352pub(crate) fn schema_registry_entry() -> SchemaRegistryEntry {
353    SchemaRegistryEntry {
354        capability_id: "epistemic",
355        capability_version: EPISTEMIC_CAPABILITY_VERSION,
356        record_family: "algorithm_interpretation_attachments",
357        record_version: BELIEF_PROJECTION_ATTACHMENT_CONTRACT_VERSION,
358        schema: Arc::clone(&ALGORITHM_INTERPRETATION_ATTACHMENT_SCHEMA),
359        schema_fingerprint: *SCHEMA_FINGERPRINT,
360        enum_registry_versions: &[],
361        sort_key: &["recorded_at", "attachment_uuid"],
362        diff_identity_fields: &["attachment_uuid"],
363        diff_record_uuid_field: Some("attachment_uuid"),
364        fingerprint_domain: CanonicalDomain::BeliefProjectionAttachment,
365        owner: "graphforge-knowledge",
366        implementation_issue: 2004,
367        max_rows: MAX_KNOWLEDGE_ROWS,
368    }
369}
370
371fn validate(row: &BeliefProjectionAttachment) -> Result<(), KnowledgeError> {
372    if row.contract_version != BELIEF_PROJECTION_ATTACHMENT_CONTRACT_VERSION {
373        return Err(invalid(
374            "belief_projection_attachment.contract_version",
375            "unsupported version",
376        ));
377    }
378    require_v7(row.attachment_uuid, "attachment_uuid")?;
379    require_v7(row.run_uuid, "run_uuid")?;
380    require_uuid(row.source_generation_uuid, "source_generation_uuid")?;
381    require_uuid(row.provenance_uuid, "provenance_uuid")?;
382    if row.policy_version == 0 {
383        return Err(invalid(
384            "belief_projection_attachment.policy_version",
385            "must be positive",
386        ));
387    }
388    if u64::try_from(row.policy_bytes.len()).unwrap_or(u64::MAX) > MAX_CANONICAL_BINARY_BYTES {
389        return Err(KnowledgeError::Limit {
390            participant: "policy_bytes",
391            observed: row.policy_bytes.len(),
392            limit: usize::try_from(MAX_CANONICAL_BINARY_BYTES).unwrap_or(usize::MAX),
393        });
394    }
395    let expected = fingerprint(
396        CanonicalDomain::BeliefProjectionPolicy,
397        CANONICAL_CONTRACT_VERSION,
398        &row.policy_bytes,
399    )?;
400    if row.policy_fingerprint != expected {
401        return Err(invalid(
402            "belief_projection_attachment.policy_fingerprint",
403            "does not match policy bytes",
404        ));
405    }
406    if row
407        .source_record_uuids
408        .windows(2)
409        .any(|pair| pair[0] >= pair[1])
410    {
411        return Err(invalid(
412            "belief_projection_attachment.source_record_uuids",
413            "must be sorted and deduplicated",
414        ));
415    }
416    for source in &row.source_record_uuids {
417        require_uuid(*source, "source_record_uuid")?;
418    }
419    Ok(())
420}
421
422fn optional_i64(writer: &mut CanonicalWriter, value: Option<i64>) -> Result<(), KnowledgeError> {
423    match value {
424        Some(value) => {
425            writer.u8(1)?;
426            writer.i64(value)?;
427        }
428        None => writer.u8(0)?,
429    }
430    Ok(())
431}
432fn optional_fingerprint(
433    writer: &mut CanonicalWriter,
434    value: Option<[u8; 32]>,
435) -> Result<(), KnowledgeError> {
436    match value {
437        Some(value) => {
438            writer.u8(1)?;
439            writer.raw(&value)?;
440        }
441        None => writer.u8(0)?,
442    }
443    Ok(())
444}
445fn append(builder: &mut FixedSizeBinaryBuilder, value: &[u8]) -> Result<(), KnowledgeError> {
446    builder
447        .append_value(value)
448        .map_err(|_| invalid("belief_projection_attachment", "invalid fixed-width value"))
449}
450fn append_optional(
451    builder: &mut FixedSizeBinaryBuilder,
452    value: Option<&[u8; 32]>,
453) -> Result<(), KnowledgeError> {
454    if let Some(value) = value {
455        append(builder, value)
456    } else {
457        builder.append_null();
458        Ok(())
459    }
460}
461fn uuid_field(name: &str, nullable: bool) -> Field {
462    Field::new(name, DataType::FixedSizeBinary(16), nullable)
463}
464fn fingerprint_field(name: &str, nullable: bool) -> Field {
465    Field::new(name, DataType::FixedSizeBinary(32), nullable)
466}
467fn timestamp_field(name: &str, nullable: bool) -> Field {
468    Field::new(
469        name,
470        DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
471        nullable,
472    )
473}
474fn fixed<'a>(
475    batch: &'a RecordBatch,
476    name: &'static str,
477) -> Result<&'a FixedSizeBinaryArray, KnowledgeError> {
478    batch
479        .column_by_name(name)
480        .and_then(|v| v.as_any().downcast_ref())
481        .ok_or_else(|| invalid(name, "column type mismatch"))
482}
483fn timestamp<'a>(
484    batch: &'a RecordBatch,
485    name: &'static str,
486) -> Result<&'a TimestampMicrosecondArray, KnowledgeError> {
487    batch
488        .column_by_name(name)
489        .and_then(|v| v.as_any().downcast_ref())
490        .ok_or_else(|| invalid(name, "column type mismatch"))
491}
492fn uint32<'a>(
493    batch: &'a RecordBatch,
494    name: &'static str,
495) -> Result<&'a UInt32Array, KnowledgeError> {
496    batch
497        .column_by_name(name)
498        .and_then(|v| v.as_any().downcast_ref())
499        .ok_or_else(|| invalid(name, "column type mismatch"))
500}
501fn binary<'a>(
502    batch: &'a RecordBatch,
503    name: &'static str,
504) -> Result<&'a BinaryArray, KnowledgeError> {
505    batch
506        .column_by_name(name)
507        .and_then(|v| v.as_any().downcast_ref())
508        .ok_or_else(|| invalid(name, "column type mismatch"))
509}
510fn list<'a>(batch: &'a RecordBatch, name: &'static str) -> Result<&'a ListArray, KnowledgeError> {
511    batch
512        .column_by_name(name)
513        .and_then(|v| v.as_any().downcast_ref())
514        .ok_or_else(|| invalid(name, "column type mismatch"))
515}
516fn uuid_at(
517    values: &FixedSizeBinaryArray,
518    row: usize,
519    field: &'static str,
520) -> Result<Uuid, KnowledgeError> {
521    Uuid::from_slice(values.value(row)).map_err(|_| invalid(field, "invalid UUID bytes"))
522}
523fn bytes32(
524    values: &FixedSizeBinaryArray,
525    row: usize,
526    field: &'static str,
527) -> Result<[u8; 32], KnowledgeError> {
528    values
529        .value(row)
530        .try_into()
531        .map_err(|_| invalid(field, "invalid fingerprint width"))
532}
533fn optional_bytes32(
534    values: &FixedSizeBinaryArray,
535    row: usize,
536    field: &'static str,
537) -> Result<Option<[u8; 32]>, KnowledgeError> {
538    (!values.is_null(row))
539        .then(|| bytes32(values, row, field))
540        .transpose()
541}
542fn optional_timestamp(values: &TimestampMicrosecondArray, row: usize) -> Option<i64> {
543    (!values.is_null(row)).then(|| values.value(row))
544}
545fn uuid_list_at(values: &ListArray, row: usize) -> Result<Vec<Uuid>, KnowledgeError> {
546    let value = values.value(row);
547    let fixed = value
548        .as_any()
549        .downcast_ref::<FixedSizeBinaryArray>()
550        .ok_or_else(|| invalid("source_record_uuids", "item type mismatch"))?;
551    (0..fixed.len())
552        .map(|index| uuid_at(fixed, index, "source_record_uuid"))
553        .collect()
554}
555fn require_v7(value: Uuid, field: &'static str) -> Result<(), KnowledgeError> {
556    if value.get_version() != Some(Version::SortRand) {
557        return Err(invalid(field, "must be UUIDv7"));
558    }
559    require_uuid(value, field)
560}
561fn require_uuid(value: Uuid, field: &'static str) -> Result<(), KnowledgeError> {
562    if value.is_nil() {
563        Err(invalid(field, "must not be nil"))
564    } else {
565        Ok(())
566    }
567}
568const fn invalid(field: &'static str, message: &'static str) -> KnowledgeError {
569    KnowledgeError::Invalid { field, message }
570}
571
572#[cfg(test)]
573mod tests {
574    use super::*;
575    fn uuid7(seed: u8) -> Uuid {
576        let mut bytes = [seed; 16];
577        bytes[6] = (bytes[6] & 0x0f) | 0x70;
578        bytes[8] = (bytes[8] & 0x3f) | 0x80;
579        Uuid::from_bytes(bytes)
580    }
581    fn attachment(id: u8, sources: Vec<Uuid>) -> BeliefProjectionAttachment {
582        BeliefProjectionAttachment::new(
583            uuid7(id),
584            uuid7(id + 1),
585            uuid7(id + 2),
586            10,
587            Some(11),
588            1,
589            b"policy-v1".to_vec(),
590            [3; 32],
591            Some([4; 32]),
592            [5; 32],
593            [6; 32],
594            sources,
595            uuid7(id + 3),
596            20,
597        )
598        .unwrap()
599    }
600
601    #[test]
602    fn sources_are_canonical_and_arrow_round_trip_is_exact() {
603        let row = attachment(1, vec![uuid7(9), uuid7(8), uuid7(9)]);
604        assert_eq!(row.source_record_uuids, vec![uuid7(8), uuid7(9)]);
605        let ledger = BeliefProjectionAttachmentLedger::new(vec![row]).unwrap();
606        let decoded =
607            BeliefProjectionAttachmentLedger::from_batches(&[ledger.batch().unwrap()]).unwrap();
608        assert_eq!(decoded, ledger);
609        assert_eq!(
610            decoded.attachment_fingerprint(uuid7(1)).unwrap(),
611            ledger.attachment_fingerprint(uuid7(1)).unwrap()
612        );
613    }
614
615    #[test]
616    fn replay_is_idempotent_and_conflict_has_transaction_code() {
617        let row = attachment(10, vec![]);
618        let ledger = BeliefProjectionAttachmentLedger::new(vec![row.clone()]).unwrap();
619        assert_eq!(
620            ledger
621                .merge(&BeliefProjectionAttachmentLedger::new(vec![row.clone()]).unwrap())
622                .unwrap(),
623            ledger
624        );
625        let mut different = row;
626        different.graph_content_fingerprint = [99; 32];
627        let error = ledger
628            .merge(&BeliefProjectionAttachmentLedger {
629                attachments: vec![different],
630            })
631            .unwrap_err();
632        assert_eq!(error.code(), "GF_TRANSACTION_CONFLICT");
633    }
634
635    #[test]
636    fn registry_is_epistemic_owned_and_frozen() {
637        let entry = schema_registry_entry();
638        assert_eq!(entry.capability_id, "epistemic");
639        assert_eq!(entry.record_family, "algorithm_interpretation_attachments");
640        assert_eq!(entry.sort_key, &["recorded_at", "attachment_uuid"]);
641        assert_eq!(entry.implementation_issue, 2004);
642    }
643
644    #[test]
645    fn defensive_attachment_validation_and_optional_fingerprints_are_exact() {
646        let row = attachment(30, vec![uuid7(40)]);
647        assert!(matches!(
648            BeliefProjectionAttachmentLedger::new(vec![row.clone(), row.clone()]),
649            Err(KnowledgeError::Duplicate("attachment_uuid"))
650        ));
651
652        let mut invalid = row.clone();
653        invalid.contract_version += 1;
654        assert!(
655            validate(&invalid)
656                .unwrap_err()
657                .to_string()
658                .contains("unsupported version")
659        );
660        invalid = row.clone();
661        invalid.policy_version = 0;
662        assert!(
663            validate(&invalid)
664                .unwrap_err()
665                .to_string()
666                .contains("must be positive")
667        );
668        invalid = row.clone();
669        invalid.policy_fingerprint = [0; 32];
670        assert!(
671            validate(&invalid)
672                .unwrap_err()
673                .to_string()
674                .contains("does not match policy bytes")
675        );
676        invalid = row.clone();
677        invalid.source_record_uuids = vec![uuid7(42), uuid7(41)];
678        assert!(
679            validate(&invalid)
680                .unwrap_err()
681                .to_string()
682                .contains("sorted and deduplicated")
683        );
684        invalid = row;
685        invalid.source_record_uuids = vec![Uuid::nil()];
686        assert!(
687            validate(&invalid)
688                .unwrap_err()
689                .to_string()
690                .contains("must not be nil")
691        );
692
693        let without_optional = BeliefProjectionAttachment::new(
694            uuid7(50),
695            uuid7(51),
696            uuid7(52),
697            10,
698            None,
699            1,
700            b"policy-v1".to_vec(),
701            [3; 32],
702            None,
703            [5; 32],
704            [6; 32],
705            vec![],
706            uuid7(53),
707            20,
708        )
709        .unwrap();
710        BeliefProjectionAttachmentLedger::new(vec![without_optional])
711            .unwrap()
712            .attachment_fingerprint(uuid7(50))
713            .unwrap();
714
715        let wrong = RecordBatch::new_empty(Arc::new(Schema::empty()));
716        assert!(
717            BeliefProjectionAttachmentLedger::from_batches(&[wrong])
718                .unwrap_err()
719                .to_string()
720                .contains("schema mismatch")
721        );
722    }
723}