1use 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
23pub const BELIEF_PROJECTION_ATTACHMENT_CONTRACT_VERSION: u32 = 1;
25
26pub 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#[derive(Clone, Debug, Eq, PartialEq)]
64pub struct BeliefProjectionAttachment {
65 pub attachment_uuid: Uuid,
67 pub run_uuid: Uuid,
69 pub source_generation_uuid: Uuid,
71 pub transaction_cutoff_micros: i64,
73 pub valid_time_micros: Option<i64>,
75 pub policy_version: u32,
77 pub policy_bytes: Vec<u8>,
79 pub policy_fingerprint: [u8; 32],
81 pub snapshot_fingerprint: [u8; 32],
83 pub valid_time_fingerprint: Option<[u8; 32]>,
85 pub graph_content_fingerprint: [u8; 32],
87 pub descriptor_fingerprint: [u8; 32],
89 pub source_record_uuids: Vec<Uuid>,
91 pub provenance_uuid: Uuid,
93 pub recorded_at_micros: i64,
95 pub contract_version: u32,
97}
98
99impl BeliefProjectionAttachment {
100 #[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#[derive(Clone, Debug, Default, Eq, PartialEq)]
150pub struct BeliefProjectionAttachmentLedger {
151 pub attachments: Vec<BeliefProjectionAttachment>,
153}
154
155impl BeliefProjectionAttachmentLedger {
156 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 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 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 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 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}