1use std::collections::{HashMap, HashSet};
4use std::sync::{Arc, LazyLock};
5
6use arrow::array::{
7 Array, FixedSizeBinaryArray, FixedSizeBinaryBuilder, StringArray, StringBuilder,
8 TimestampMicrosecondArray, TimestampMicrosecondBuilder, UInt32Array, UInt32Builder,
9};
10use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
11use arrow::record_batch::RecordBatch;
12use graphforge_core::canonical::{
13 CANONICAL_CONTRACT_VERSION, CanonicalDomain, CanonicalWriter, fingerprint,
14};
15use uuid::{Uuid, Version};
16
17use crate::{
18 EPISTEMIC_CAPABILITY_VERSION, KnowledgeError, MAX_KNOWLEDGE_ROWS, SchemaRegistryEntry,
19};
20
21pub const ASSERTION_STATUS_CONTRACT_VERSION: u32 = 1;
23pub const ASSERTION_STATUS_REGISTRY_VERSION: u32 = 1;
25
26pub static ASSERTION_STATUS_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
28 Arc::new(Schema::new(vec![
29 uuid_field("status_event_uuid", false),
30 uuid_field("assertion_uuid", false),
31 Field::new("status", DataType::Utf8, false),
32 uuid_field("confidence_uuid", true),
33 uuid_field("reasoning_uuid", true),
34 uuid_field("provenance_uuid", false),
35 Field::new(
36 "recorded_at",
37 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
38 false,
39 ),
40 Field::new("contract_version", DataType::UInt32, false),
41 ]))
42});
43
44static ASSERTION_STATUS_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
45 fingerprint(
46 CanonicalDomain::AssertionStatus,
47 CANONICAL_CONTRACT_VERSION,
48 b"assertion_status/1|status_event_uuid:fixed[16]:required|assertion_uuid:fixed[16]:required|status:utf8:required|confidence_uuid:fixed[16]:nullable|reasoning_uuid:fixed[16]:nullable|provenance_uuid:fixed[16]:required|recorded_at:timestamp_us_utc:required|contract_version:u32:required",
49 )
50 .expect("registered assertion-status schema is within canonical bounds")
51});
52
53#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
55pub enum AssertionStatus {
56 Hypothesis,
58 Supported,
60 Refuted,
62 Disputed,
64 Retracted,
66 Superseded,
68}
69
70impl AssertionStatus {
71 #[must_use]
73 pub const fn as_str(self) -> &'static str {
74 match self {
75 Self::Hypothesis => "hypothesis",
76 Self::Supported => "supported",
77 Self::Refuted => "refuted",
78 Self::Disputed => "disputed",
79 Self::Retracted => "retracted",
80 Self::Superseded => "superseded",
81 }
82 }
83
84 fn parse(value: &str) -> Result<Self, KnowledgeError> {
85 match value {
86 "hypothesis" => Ok(Self::Hypothesis),
87 "supported" => Ok(Self::Supported),
88 "refuted" => Ok(Self::Refuted),
89 "disputed" => Ok(Self::Disputed),
90 "retracted" => Ok(Self::Retracted),
91 "superseded" => Ok(Self::Superseded),
92 _ => Err(invalid("assertion_status.status", "unknown registry value")),
93 }
94 }
95
96 #[must_use]
98 pub const fn is_terminal(self) -> bool {
99 matches!(self, Self::Superseded)
100 }
101}
102
103#[derive(Clone, Debug, Eq, PartialEq)]
105pub struct AssertionStatusEvent {
106 pub status_event_uuid: Uuid,
108 pub assertion_uuid: Uuid,
110 pub status: AssertionStatus,
112 pub confidence_uuid: Option<Uuid>,
114 pub reasoning_uuid: Option<Uuid>,
116 pub provenance_uuid: Uuid,
118 pub recorded_at_micros: i64,
120 pub contract_version: u32,
122}
123
124impl AssertionStatusEvent {
125 #[allow(clippy::too_many_arguments)]
127 pub fn new(
128 status_event_uuid: Uuid,
129 assertion_uuid: Uuid,
130 status: AssertionStatus,
131 confidence_uuid: Option<Uuid>,
132 reasoning_uuid: Option<Uuid>,
133 provenance_uuid: Uuid,
134 recorded_at_micros: i64,
135 ) -> Result<Self, KnowledgeError> {
136 require_v7(status_event_uuid, "status_event_uuid")?;
137 require_v7(assertion_uuid, "assertion_uuid")?;
138 if let Some(value) = confidence_uuid {
139 require_v7(value, "confidence_uuid")?;
140 }
141 if let Some(value) = reasoning_uuid {
142 require_v7(value, "reasoning_uuid")?;
143 }
144 require_uuid(provenance_uuid, "provenance_uuid")?;
145 Ok(Self {
146 status_event_uuid,
147 assertion_uuid,
148 status,
149 confidence_uuid,
150 reasoning_uuid,
151 provenance_uuid,
152 recorded_at_micros,
153 contract_version: ASSERTION_STATUS_CONTRACT_VERSION,
154 })
155 }
156}
157
158#[derive(Clone, Debug, Default, Eq, PartialEq)]
160pub struct AssertionStatusLedger {
161 pub events: Vec<AssertionStatusEvent>,
163}
164
165impl AssertionStatusLedger {
166 pub fn new(mut events: Vec<AssertionStatusEvent>) -> Result<Self, KnowledgeError> {
168 if events.len() > MAX_KNOWLEDGE_ROWS {
169 return Err(KnowledgeError::Limit {
170 participant: "assertion_status_events",
171 observed: events.len(),
172 limit: MAX_KNOWLEDGE_ROWS,
173 });
174 }
175 let mut ids = HashSet::with_capacity(events.len());
176 for event in &events {
177 validate_event(event)?;
178 if !ids.insert(event.status_event_uuid) {
179 return Err(KnowledgeError::Duplicate("status_event_uuid"));
180 }
181 }
182 events.sort_by_key(|row| (row.recorded_at_micros, row.status_event_uuid));
183 let mut terminal_assertions = HashSet::new();
184 for event in &events {
185 if terminal_assertions.contains(&event.assertion_uuid) && !event.status.is_terminal() {
186 return Err(invalid("assertion_status.status", "superseded is terminal"));
187 }
188 if event.status.is_terminal() {
189 terminal_assertions.insert(event.assertion_uuid);
190 }
191 }
192 Ok(Self { events })
193 }
194
195 pub fn merge(&self, staged: &Self) -> Result<Self, KnowledgeError> {
197 let mut events = self.events.clone();
198 let terminal_assertions = self
199 .events
200 .iter()
201 .filter(|row| row.status.is_terminal())
202 .map(|row| row.assertion_uuid)
203 .collect::<HashSet<_>>();
204 let mut by_id = events
205 .iter()
206 .cloned()
207 .map(|row| (row.status_event_uuid, row))
208 .collect::<HashMap<_, _>>();
209 for event in &staged.events {
210 if let Some(existing) = by_id.get(&event.status_event_uuid) {
211 if existing != event {
212 return Err(KnowledgeError::Conflict("status_event_uuid"));
213 }
214 } else {
215 if terminal_assertions.contains(&event.assertion_uuid)
216 && !event.status.is_terminal()
217 {
218 return Err(invalid("assertion_status.status", "superseded is terminal"));
219 }
220 events.push(event.clone());
221 by_id.insert(event.status_event_uuid, event.clone());
222 }
223 }
224 Self::new(events)
225 }
226
227 pub fn event_fingerprint(&self, status_event_uuid: Uuid) -> Result<[u8; 32], KnowledgeError> {
229 let row = self
230 .events
231 .iter()
232 .find(|row| row.status_event_uuid == status_event_uuid)
233 .ok_or(KnowledgeError::Dangling("status_event_uuid"))?;
234 let mut writer = CanonicalWriter::new();
235 writer.raw(row.status_event_uuid.as_bytes())?;
236 writer.raw(row.assertion_uuid.as_bytes())?;
237 writer.text(row.status.as_str())?;
238 optional_uuid(&mut writer, row.confidence_uuid)?;
239 optional_uuid(&mut writer, row.reasoning_uuid)?;
240 writer.raw(row.provenance_uuid.as_bytes())?;
241 writer.i64(row.recorded_at_micros)?;
242 writer.u32(row.contract_version)?;
243 fingerprint(
244 CanonicalDomain::AssertionStatus,
245 CANONICAL_CONTRACT_VERSION,
246 &writer.finish(),
247 )
248 .map_err(Into::into)
249 }
250
251 #[must_use]
253 pub fn current_for(&self, assertion_uuid: Uuid) -> Option<&AssertionStatusEvent> {
254 self.events
255 .iter()
256 .filter(|row| row.assertion_uuid == assertion_uuid)
257 .max_by_key(|row| (row.recorded_at_micros, row.status_event_uuid))
258 }
259
260 pub fn batch(&self) -> Result<RecordBatch, KnowledgeError> {
262 let mut ids = FixedSizeBinaryBuilder::with_capacity(self.events.len(), 16);
263 let mut assertions = FixedSizeBinaryBuilder::with_capacity(self.events.len(), 16);
264 let mut statuses = StringBuilder::new();
265 let mut confidence = FixedSizeBinaryBuilder::with_capacity(self.events.len(), 16);
266 let mut reasoning = FixedSizeBinaryBuilder::with_capacity(self.events.len(), 16);
267 let mut provenance = FixedSizeBinaryBuilder::with_capacity(self.events.len(), 16);
268 let mut times = TimestampMicrosecondBuilder::new().with_timezone("UTC");
269 let mut versions = UInt32Builder::new();
270 for row in &self.events {
271 ids.append_value(row.status_event_uuid.as_bytes())
272 .map_err(|_| invalid("status_event_uuid", "invalid UUID width"))?;
273 assertions
274 .append_value(row.assertion_uuid.as_bytes())
275 .map_err(|_| invalid("assertion_uuid", "invalid UUID width"))?;
276 statuses.append_value(row.status.as_str());
277 append_optional_uuid(&mut confidence, row.confidence_uuid)?;
278 append_optional_uuid(&mut reasoning, row.reasoning_uuid)?;
279 provenance
280 .append_value(row.provenance_uuid.as_bytes())
281 .map_err(|_| invalid("provenance_uuid", "invalid UUID width"))?;
282 times.append_value(row.recorded_at_micros);
283 versions.append_value(row.contract_version);
284 }
285 RecordBatch::try_new(
286 Arc::clone(&ASSERTION_STATUS_SCHEMA),
287 vec![
288 Arc::new(ids.finish()),
289 Arc::new(assertions.finish()),
290 Arc::new(statuses.finish()),
291 Arc::new(confidence.finish()),
292 Arc::new(reasoning.finish()),
293 Arc::new(provenance.finish()),
294 Arc::new(times.finish()),
295 Arc::new(versions.finish()),
296 ],
297 )
298 .map_err(|_| invalid("assertion_status", "Arrow batch construction failed"))
299 }
300
301 pub fn from_batches(batches: &[RecordBatch]) -> Result<Self, KnowledgeError> {
303 let mut events = Vec::new();
304 for batch in batches {
305 if batch.schema().as_ref() != ASSERTION_STATUS_SCHEMA.as_ref() {
306 return Err(invalid("assertion_status.schema", "schema mismatch"));
307 }
308 let ids = fixed(batch, "status_event_uuid")?;
309 let assertions = fixed(batch, "assertion_uuid")?;
310 let statuses = string(batch, "status")?;
311 let confidence = fixed(batch, "confidence_uuid")?;
312 let reasoning = fixed(batch, "reasoning_uuid")?;
313 let provenance = fixed(batch, "provenance_uuid")?;
314 let times = timestamp(batch, "recorded_at")?;
315 let versions = uint32(batch, "contract_version")?;
316 for row in 0..batch.num_rows() {
317 events.push(AssertionStatusEvent {
318 status_event_uuid: uuid_at(ids, row, "status_event_uuid")?,
319 assertion_uuid: uuid_at(assertions, row, "assertion_uuid")?,
320 status: AssertionStatus::parse(statuses.value(row))?,
321 confidence_uuid: optional_uuid_at(confidence, row, "confidence_uuid")?,
322 reasoning_uuid: optional_uuid_at(reasoning, row, "reasoning_uuid")?,
323 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
324 recorded_at_micros: times.value(row),
325 contract_version: versions.value(row),
326 });
327 }
328 }
329 Self::new(events)
330 }
331}
332
333pub(crate) fn schema_registry_entry() -> SchemaRegistryEntry {
334 SchemaRegistryEntry {
335 capability_id: "epistemic",
336 capability_version: EPISTEMIC_CAPABILITY_VERSION,
337 record_family: "assertion_status_events",
338 record_version: ASSERTION_STATUS_CONTRACT_VERSION,
339 schema: Arc::clone(&ASSERTION_STATUS_SCHEMA),
340 schema_fingerprint: *ASSERTION_STATUS_SCHEMA_FINGERPRINT,
341 enum_registry_versions: &[("assertion_status", ASSERTION_STATUS_REGISTRY_VERSION)],
342 sort_key: &["recorded_at", "status_event_uuid"],
343 diff_identity_fields: &["status_event_uuid"],
344 diff_record_uuid_field: Some("status_event_uuid"),
345 fingerprint_domain: CanonicalDomain::AssertionStatus,
346 owner: "graphforge-knowledge",
347 implementation_issue: 777,
348 max_rows: MAX_KNOWLEDGE_ROWS,
349 }
350}
351
352fn validate_event(row: &AssertionStatusEvent) -> Result<(), KnowledgeError> {
353 if row.contract_version != ASSERTION_STATUS_CONTRACT_VERSION {
354 return Err(invalid(
355 "assertion_status.contract_version",
356 "unsupported version",
357 ));
358 }
359 require_v7(row.status_event_uuid, "status_event_uuid")?;
360 require_v7(row.assertion_uuid, "assertion_uuid")?;
361 if let Some(value) = row.confidence_uuid {
362 require_v7(value, "confidence_uuid")?;
363 }
364 if let Some(value) = row.reasoning_uuid {
365 require_v7(value, "reasoning_uuid")?;
366 }
367 require_uuid(row.provenance_uuid, "provenance_uuid")
368}
369
370fn optional_uuid(writer: &mut CanonicalWriter, value: Option<Uuid>) -> Result<(), KnowledgeError> {
371 match value {
372 Some(value) => {
373 writer.u8(1)?;
374 writer.raw(value.as_bytes())?;
375 }
376 None => writer.u8(0)?,
377 }
378 Ok(())
379}
380
381fn append_optional_uuid(
382 builder: &mut FixedSizeBinaryBuilder,
383 value: Option<Uuid>,
384) -> Result<(), KnowledgeError> {
385 if let Some(value) = value {
386 builder
387 .append_value(value.as_bytes())
388 .map_err(|_| invalid("assertion_status", "invalid UUID width"))?;
389 } else {
390 builder.append_null();
391 }
392 Ok(())
393}
394
395fn optional_uuid_at(
396 values: &FixedSizeBinaryArray,
397 row: usize,
398 field: &'static str,
399) -> Result<Option<Uuid>, KnowledgeError> {
400 (!values.is_null(row))
401 .then(|| uuid_at(values, row, field))
402 .transpose()
403}
404
405fn require_v7(value: Uuid, field: &'static str) -> Result<(), KnowledgeError> {
406 if value.get_version() != Some(Version::SortRand) {
407 return Err(invalid(field, "must be UUIDv7"));
408 }
409 require_uuid(value, field)
410}
411
412fn require_uuid(value: Uuid, field: &'static str) -> Result<(), KnowledgeError> {
413 if value.is_nil() {
414 return Err(invalid(field, "must not be nil"));
415 }
416 Ok(())
417}
418
419const fn invalid(field: &'static str, message: &'static str) -> KnowledgeError {
420 KnowledgeError::Invalid { field, message }
421}
422
423fn uuid_field(name: &str, nullable: bool) -> Field {
424 Field::new(name, DataType::FixedSizeBinary(16), nullable)
425}
426
427fn fixed<'a>(
428 batch: &'a RecordBatch,
429 name: &'static str,
430) -> Result<&'a FixedSizeBinaryArray, KnowledgeError> {
431 batch
432 .column_by_name(name)
433 .and_then(|value| value.as_any().downcast_ref())
434 .ok_or_else(|| invalid(name, "column type mismatch"))
435}
436
437fn string<'a>(
438 batch: &'a RecordBatch,
439 name: &'static str,
440) -> Result<&'a StringArray, KnowledgeError> {
441 batch
442 .column_by_name(name)
443 .and_then(|value| value.as_any().downcast_ref())
444 .ok_or_else(|| invalid(name, "column type mismatch"))
445}
446
447fn timestamp<'a>(
448 batch: &'a RecordBatch,
449 name: &'static str,
450) -> Result<&'a TimestampMicrosecondArray, KnowledgeError> {
451 batch
452 .column_by_name(name)
453 .and_then(|value| value.as_any().downcast_ref())
454 .ok_or_else(|| invalid(name, "column type mismatch"))
455}
456
457fn uint32<'a>(
458 batch: &'a RecordBatch,
459 name: &'static str,
460) -> Result<&'a UInt32Array, KnowledgeError> {
461 batch
462 .column_by_name(name)
463 .and_then(|value| value.as_any().downcast_ref())
464 .ok_or_else(|| invalid(name, "column type mismatch"))
465}
466
467fn uuid_at(
468 values: &FixedSizeBinaryArray,
469 row: usize,
470 field: &'static str,
471) -> Result<Uuid, KnowledgeError> {
472 Uuid::from_slice(values.value(row)).map_err(|_| invalid(field, "invalid UUID bytes"))
473}
474
475#[cfg(test)]
476mod tests {
477 use super::*;
478
479 fn uuid7(seed: u8) -> Uuid {
480 let mut bytes = [seed; 16];
481 bytes[6] = (bytes[6] & 0x0f) | 0x70;
482 bytes[8] = (bytes[8] & 0x3f) | 0x80;
483 Uuid::from_bytes(bytes)
484 }
485
486 fn event(id: u8, assertion: u8, status: AssertionStatus, time: i64) -> AssertionStatusEvent {
487 AssertionStatusEvent::new(
488 uuid7(id),
489 uuid7(assertion),
490 status,
491 None,
492 None,
493 uuid7(id.wrapping_add(100)),
494 time,
495 )
496 .unwrap()
497 }
498
499 #[test]
500 fn round_trip_fingerprint_statusless_and_identical_time_order_are_stable() {
501 assert_eq!(
502 ASSERTION_STATUS_SCHEMA_FINGERPRINT
503 .iter()
504 .map(|byte| format!("{byte:02x}"))
505 .collect::<String>(),
506 "83e696b90c9151fefe6b92c18b752c0d717fa4fd967997f60e995c79d59f54cd"
507 );
508 let ledger = AssertionStatusLedger::new(vec![
509 event(2, 20, AssertionStatus::Supported, 10),
510 event(1, 20, AssertionStatus::Hypothesis, 10),
511 ])
512 .unwrap();
513 let decoded = AssertionStatusLedger::from_batches(&[ledger.batch().unwrap()]).unwrap();
514 assert_eq!(decoded, ledger);
515 assert_eq!(
516 decoded.event_fingerprint(uuid7(2)).unwrap(),
517 ledger.event_fingerprint(uuid7(2)).unwrap()
518 );
519 assert_eq!(
520 ledger.current_for(uuid7(20)).unwrap().status_event_uuid,
521 uuid7(2)
522 );
523 assert!(ledger.current_for(uuid7(21)).is_none());
524 }
525
526 #[test]
527 fn every_nonterminal_transition_is_allowed_and_superseded_is_terminal() {
528 let nonterminal = [
529 AssertionStatus::Hypothesis,
530 AssertionStatus::Supported,
531 AssertionStatus::Refuted,
532 AssertionStatus::Disputed,
533 AssertionStatus::Retracted,
534 ];
535 let mut id = 1_u8;
536 for from in nonterminal {
537 for to in nonterminal {
538 assert!(
539 AssertionStatusLedger::new(vec![
540 event(id, id, from, 1),
541 event(id.wrapping_add(1), id, to, 2),
542 ])
543 .is_ok()
544 );
545 id = id.wrapping_add(2);
546 }
547 assert!(
548 AssertionStatusLedger::new(vec![
549 event(id, id, from, 1),
550 event(id.wrapping_add(1), id, AssertionStatus::Superseded, 2),
551 ])
552 .is_ok()
553 );
554 id = id.wrapping_add(2);
555 }
556 for to in [
557 AssertionStatus::Hypothesis,
558 AssertionStatus::Supported,
559 AssertionStatus::Refuted,
560 AssertionStatus::Disputed,
561 AssertionStatus::Retracted,
562 ] {
563 assert!(matches!(
564 AssertionStatusLedger::new(vec![
565 event(id, id, AssertionStatus::Superseded, 1),
566 event(id.wrapping_add(1), id, to, 2),
567 ]),
568 Err(KnowledgeError::Invalid {
569 field: "assertion_status.status",
570 message: "superseded is terminal",
571 })
572 ));
573 id = id.wrapping_add(2);
574 }
575 assert!(
576 AssertionStatusLedger::new(vec![
577 event(id, id, AssertionStatus::Superseded, 1),
578 event(id.wrapping_add(1), id, AssertionStatus::Superseded, 2),
579 ])
580 .is_ok(),
581 "multiple terminal events preserve explicit supersession branches"
582 );
583 let terminal =
584 AssertionStatusLedger::new(vec![event(id, id, AssertionStatus::Superseded, 10)])
585 .unwrap();
586 let second_terminal = AssertionStatusLedger::new(vec![event(
587 id.wrapping_add(2),
588 id,
589 AssertionStatus::Superseded,
590 11,
591 )])
592 .unwrap();
593 assert!(terminal.merge(&second_terminal).is_ok());
594 let backdated = AssertionStatusLedger::new(vec![event(
595 id.wrapping_add(1),
596 id,
597 AssertionStatus::Supported,
598 1,
599 )])
600 .unwrap();
601 assert!(matches!(
602 terminal.merge(&backdated),
603 Err(KnowledgeError::Invalid {
604 field: "assertion_status.status",
605 message: "superseded is terminal",
606 })
607 ));
608 }
609
610 #[test]
611 fn replay_is_idempotent_and_conflicting_identity_is_rejected() {
612 let base = AssertionStatusLedger::new(vec![event(1, 20, AssertionStatus::Hypothesis, 10)])
613 .unwrap();
614 assert_eq!(base.merge(&base).unwrap(), base);
615 let conflict =
616 AssertionStatusLedger::new(vec![event(1, 20, AssertionStatus::Refuted, 10)]).unwrap();
617 assert!(matches!(
618 base.merge(&conflict),
619 Err(KnowledgeError::Conflict("status_event_uuid"))
620 ));
621 }
622
623 #[test]
624 fn assertion_status_registry_is_closed_and_terminal_only_for_superseded() {
625 let values = [
626 AssertionStatus::Hypothesis,
627 AssertionStatus::Supported,
628 AssertionStatus::Refuted,
629 AssertionStatus::Disputed,
630 AssertionStatus::Retracted,
631 AssertionStatus::Superseded,
632 ];
633 for value in values {
634 assert_eq!(AssertionStatus::parse(value.as_str()).unwrap(), value);
635 assert_eq!(value.is_terminal(), value == AssertionStatus::Superseded);
636 }
637 assert!(AssertionStatus::parse("unknown").is_err());
638 }
639}