1use std::collections::{HashMap, HashSet};
4use std::sync::{Arc, LazyLock};
5
6use arrow::array::{
7 Array, FixedSizeBinaryArray, FixedSizeBinaryBuilder, TimestampMicrosecondArray,
8 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::{KnowledgeError, MAX_KNOWLEDGE_ROWS, SchemaRegistryEntry};
18
19pub const VALID_TIME_CAPABILITY_VERSION: u32 = 1;
21pub const ASSERTION_VALIDITY_CONTRACT_VERSION: u32 = 1;
23pub const ASSERTION_VALIDITY_POLICY_VERSION: u32 = 1;
25
26pub static ASSERTION_VALIDITY_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
28 Arc::new(Schema::new(vec![
29 uuid_field("validity_event_uuid", false),
30 uuid_field("assertion_uuid", false),
31 timestamp_field("valid_from", true),
32 timestamp_field("valid_to", true),
33 uuid_field("reasoning_uuid", true),
34 uuid_field("provenance_uuid", false),
35 timestamp_field("recorded_at", false),
36 Field::new("contract_version", DataType::UInt32, false),
37 ]))
38});
39
40static ASSERTION_VALIDITY_SCHEMA_FINGERPRINT: LazyLock<[u8; 32]> = LazyLock::new(|| {
41 fingerprint(
42 CanonicalDomain::Schema,
43 CANONICAL_CONTRACT_VERSION,
44 b"assertion_validity/1|validity_event_uuid:fixed[16]:required|assertion_uuid:fixed[16]:required|valid_from:timestamp_us_utc:nullable|valid_to:timestamp_us_utc:nullable|reasoning_uuid:fixed[16]:nullable|provenance_uuid:fixed[16]:required|recorded_at:timestamp_us_utc:required|contract_version:u32:required",
45 )
46 .expect("registered assertion-validity schema is within canonical bounds")
47});
48
49#[derive(Clone, Debug, Eq, PartialEq)]
51pub struct AssertionValidityEvent {
52 pub validity_event_uuid: Uuid,
54 pub assertion_uuid: Uuid,
56 pub valid_from_micros: Option<i64>,
58 pub valid_to_micros: Option<i64>,
60 pub reasoning_uuid: Option<Uuid>,
62 pub provenance_uuid: Uuid,
64 pub recorded_at_micros: i64,
66 pub contract_version: u32,
68}
69
70impl AssertionValidityEvent {
71 #[allow(clippy::too_many_arguments)]
73 pub fn new(
74 validity_event_uuid: Uuid,
75 assertion_uuid: Uuid,
76 valid_from_micros: Option<i64>,
77 valid_to_micros: Option<i64>,
78 reasoning_uuid: Option<Uuid>,
79 provenance_uuid: Uuid,
80 recorded_at_micros: i64,
81 ) -> Result<Self, KnowledgeError> {
82 let event = Self {
83 validity_event_uuid,
84 assertion_uuid,
85 valid_from_micros,
86 valid_to_micros,
87 reasoning_uuid,
88 provenance_uuid,
89 recorded_at_micros,
90 contract_version: ASSERTION_VALIDITY_CONTRACT_VERSION,
91 };
92 validate_event(&event)?;
93 Ok(event)
94 }
95
96 #[must_use]
100 pub fn contains(&self, valid_time_micros: i64) -> bool {
101 self.valid_from_micros
102 .is_none_or(|from| from <= valid_time_micros)
103 && self.valid_to_micros.is_none_or(|to| valid_time_micros < to)
104 }
105}
106
107#[derive(Clone, Debug, Default, Eq, PartialEq)]
109pub struct AssertionValidityLedger {
110 pub events: Vec<AssertionValidityEvent>,
112}
113
114impl AssertionValidityLedger {
115 pub fn new(mut events: Vec<AssertionValidityEvent>) -> Result<Self, KnowledgeError> {
117 if events.len() > MAX_KNOWLEDGE_ROWS {
118 return Err(KnowledgeError::Limit {
119 participant: "assertion_validity_events",
120 observed: events.len(),
121 limit: MAX_KNOWLEDGE_ROWS,
122 });
123 }
124 let mut ids = HashSet::with_capacity(events.len());
125 for event in &events {
126 validate_event(event)?;
127 if !ids.insert(event.validity_event_uuid) {
128 return Err(KnowledgeError::Duplicate("validity_event_uuid"));
129 }
130 }
131 events.sort_by_key(|row| (row.recorded_at_micros, row.validity_event_uuid));
132 Ok(Self { events })
133 }
134
135 pub fn merge(&self, staged: &Self) -> Result<Self, KnowledgeError> {
137 let mut events = self.events.clone();
138 let mut by_id = events
139 .iter()
140 .cloned()
141 .map(|row| (row.validity_event_uuid, row))
142 .collect::<HashMap<_, _>>();
143 for event in &staged.events {
144 if let Some(existing) = by_id.get(&event.validity_event_uuid) {
145 if existing != event {
146 return Err(KnowledgeError::Conflict("validity_event_uuid"));
147 }
148 } else {
149 events.push(event.clone());
150 by_id.insert(event.validity_event_uuid, event.clone());
151 }
152 }
153 Self::new(events)
154 }
155
156 #[must_use]
158 pub fn current_for_at(
159 &self,
160 assertion_uuid: Uuid,
161 transaction_cutoff_micros: i64,
162 ) -> Option<&AssertionValidityEvent> {
163 self.events
164 .iter()
165 .filter(|row| {
166 row.assertion_uuid == assertion_uuid
167 && row.recorded_at_micros <= transaction_cutoff_micros
168 })
169 .max_by_key(|row| (row.recorded_at_micros, row.validity_event_uuid))
170 }
171
172 #[must_use]
174 pub fn is_valid_at(
175 &self,
176 assertion_uuid: Uuid,
177 transaction_cutoff_micros: i64,
178 valid_time_micros: i64,
179 ) -> Option<bool> {
180 self.current_for_at(assertion_uuid, transaction_cutoff_micros)
181 .map(|event| event.contains(valid_time_micros))
182 }
183
184 pub fn event_fingerprint(&self, validity_event_uuid: Uuid) -> Result<[u8; 32], KnowledgeError> {
186 let row = self
187 .events
188 .iter()
189 .find(|row| row.validity_event_uuid == validity_event_uuid)
190 .ok_or(KnowledgeError::Dangling("validity_event_uuid"))?;
191 let mut writer = CanonicalWriter::new();
192 writer.raw(row.validity_event_uuid.as_bytes())?;
193 writer.raw(row.assertion_uuid.as_bytes())?;
194 optional_i64(&mut writer, row.valid_from_micros)?;
195 optional_i64(&mut writer, row.valid_to_micros)?;
196 optional_uuid(&mut writer, row.reasoning_uuid)?;
197 writer.raw(row.provenance_uuid.as_bytes())?;
198 writer.i64(row.recorded_at_micros)?;
199 writer.u32(row.contract_version)?;
200 fingerprint(
201 CanonicalDomain::AssertionValidity,
202 CANONICAL_CONTRACT_VERSION,
203 &writer.finish(),
204 )
205 .map_err(Into::into)
206 }
207
208 pub fn batch(&self) -> Result<RecordBatch, KnowledgeError> {
210 let mut ids = FixedSizeBinaryBuilder::with_capacity(self.events.len(), 16);
211 let mut assertions = FixedSizeBinaryBuilder::with_capacity(self.events.len(), 16);
212 let mut from = TimestampMicrosecondBuilder::new().with_timezone("UTC");
213 let mut to = TimestampMicrosecondBuilder::new().with_timezone("UTC");
214 let mut reasoning = FixedSizeBinaryBuilder::with_capacity(self.events.len(), 16);
215 let mut provenance = FixedSizeBinaryBuilder::with_capacity(self.events.len(), 16);
216 let mut times = TimestampMicrosecondBuilder::new().with_timezone("UTC");
217 let mut versions = UInt32Builder::new();
218 for row in &self.events {
219 append_uuid(&mut ids, row.validity_event_uuid, "validity_event_uuid")?;
220 append_uuid(&mut assertions, row.assertion_uuid, "assertion_uuid")?;
221 from.append_option(row.valid_from_micros);
222 to.append_option(row.valid_to_micros);
223 append_optional_uuid(&mut reasoning, row.reasoning_uuid)?;
224 append_uuid(&mut provenance, row.provenance_uuid, "provenance_uuid")?;
225 times.append_value(row.recorded_at_micros);
226 versions.append_value(row.contract_version);
227 }
228 RecordBatch::try_new(
229 Arc::clone(&ASSERTION_VALIDITY_SCHEMA),
230 vec![
231 Arc::new(ids.finish()),
232 Arc::new(assertions.finish()),
233 Arc::new(from.finish()),
234 Arc::new(to.finish()),
235 Arc::new(reasoning.finish()),
236 Arc::new(provenance.finish()),
237 Arc::new(times.finish()),
238 Arc::new(versions.finish()),
239 ],
240 )
241 .map_err(|_| invalid("assertion_validity", "Arrow batch construction failed"))
242 }
243
244 pub fn from_batches(batches: &[RecordBatch]) -> Result<Self, KnowledgeError> {
246 let mut events = Vec::new();
247 for batch in batches {
248 if batch.schema().as_ref() != ASSERTION_VALIDITY_SCHEMA.as_ref() {
249 return Err(invalid("assertion_validity.schema", "schema mismatch"));
250 }
251 let ids = fixed(batch, "validity_event_uuid")?;
252 let assertions = fixed(batch, "assertion_uuid")?;
253 let from = timestamp(batch, "valid_from")?;
254 let to = timestamp(batch, "valid_to")?;
255 let reasoning = fixed(batch, "reasoning_uuid")?;
256 let provenance = fixed(batch, "provenance_uuid")?;
257 let times = timestamp(batch, "recorded_at")?;
258 let versions = uint32(batch, "contract_version")?;
259 for row in 0..batch.num_rows() {
260 events.push(AssertionValidityEvent {
261 validity_event_uuid: uuid_at(ids, row, "validity_event_uuid")?,
262 assertion_uuid: uuid_at(assertions, row, "assertion_uuid")?,
263 valid_from_micros: optional_timestamp_at(from, row),
264 valid_to_micros: optional_timestamp_at(to, row),
265 reasoning_uuid: optional_uuid_at(reasoning, row, "reasoning_uuid")?,
266 provenance_uuid: uuid_at(provenance, row, "provenance_uuid")?,
267 recorded_at_micros: times.value(row),
268 contract_version: versions.value(row),
269 });
270 }
271 }
272 Self::new(events)
273 }
274}
275
276pub(crate) fn schema_registry_entry() -> SchemaRegistryEntry {
277 SchemaRegistryEntry {
278 capability_id: "valid_time",
279 capability_version: VALID_TIME_CAPABILITY_VERSION,
280 record_family: "assertion_validity_events",
281 record_version: ASSERTION_VALIDITY_CONTRACT_VERSION,
282 schema: Arc::clone(&ASSERTION_VALIDITY_SCHEMA),
283 schema_fingerprint: *ASSERTION_VALIDITY_SCHEMA_FINGERPRINT,
284 enum_registry_versions: &[(
285 "assertion_validity_policy",
286 ASSERTION_VALIDITY_POLICY_VERSION,
287 )],
288 sort_key: &["recorded_at", "validity_event_uuid"],
289 diff_identity_fields: &["validity_event_uuid"],
290 diff_record_uuid_field: Some("validity_event_uuid"),
291 fingerprint_domain: CanonicalDomain::AssertionValidity,
292 owner: "graphforge-knowledge",
293 implementation_issue: 781,
294 max_rows: MAX_KNOWLEDGE_ROWS,
295 }
296}
297
298fn validate_event(row: &AssertionValidityEvent) -> Result<(), KnowledgeError> {
299 if row.contract_version != ASSERTION_VALIDITY_CONTRACT_VERSION {
300 return Err(invalid(
301 "assertion_validity.contract_version",
302 "unsupported version",
303 ));
304 }
305 require_v7(row.validity_event_uuid, "validity_event_uuid")?;
306 require_v7(row.assertion_uuid, "assertion_uuid")?;
307 if let Some(reasoning_uuid) = row.reasoning_uuid {
308 require_v7(reasoning_uuid, "reasoning_uuid")?;
309 }
310 require_uuid(row.provenance_uuid, "provenance_uuid")?;
311 if matches!(
312 (row.valid_from_micros, row.valid_to_micros),
313 (Some(from), Some(to)) if from > to
314 ) {
315 return Err(invalid(
316 "assertion_validity.interval",
317 "valid_from must not exceed valid_to",
318 ));
319 }
320 Ok(())
321}
322
323fn optional_i64(writer: &mut CanonicalWriter, value: Option<i64>) -> Result<(), KnowledgeError> {
324 match value {
325 Some(value) => {
326 writer.u8(1)?;
327 writer.i64(value)?;
328 }
329 None => writer.u8(0)?,
330 }
331 Ok(())
332}
333
334fn optional_uuid(writer: &mut CanonicalWriter, value: Option<Uuid>) -> Result<(), KnowledgeError> {
335 match value {
336 Some(value) => {
337 writer.u8(1)?;
338 writer.raw(value.as_bytes())?;
339 }
340 None => writer.u8(0)?,
341 }
342 Ok(())
343}
344
345fn append_uuid(
346 builder: &mut FixedSizeBinaryBuilder,
347 value: Uuid,
348 field: &'static str,
349) -> Result<(), KnowledgeError> {
350 builder
351 .append_value(value.as_bytes())
352 .map_err(|_| invalid(field, "invalid UUID width"))
353}
354
355fn append_optional_uuid(
356 builder: &mut FixedSizeBinaryBuilder,
357 value: Option<Uuid>,
358) -> Result<(), KnowledgeError> {
359 if let Some(value) = value {
360 append_uuid(builder, value, "reasoning_uuid")?;
361 } else {
362 builder.append_null();
363 }
364 Ok(())
365}
366
367fn optional_uuid_at(
368 values: &FixedSizeBinaryArray,
369 row: usize,
370 field: &'static str,
371) -> Result<Option<Uuid>, KnowledgeError> {
372 (!values.is_null(row))
373 .then(|| uuid_at(values, row, field))
374 .transpose()
375}
376
377fn optional_timestamp_at(values: &TimestampMicrosecondArray, row: usize) -> Option<i64> {
378 (!values.is_null(row)).then(|| values.value(row))
379}
380
381fn require_v7(value: Uuid, field: &'static str) -> Result<(), KnowledgeError> {
382 if value.get_version() != Some(Version::SortRand) {
383 return Err(invalid(field, "must be UUIDv7"));
384 }
385 require_uuid(value, field)
386}
387
388fn require_uuid(value: Uuid, field: &'static str) -> Result<(), KnowledgeError> {
389 if value.is_nil() {
390 return Err(invalid(field, "must not be nil"));
391 }
392 Ok(())
393}
394
395const fn invalid(field: &'static str, message: &'static str) -> KnowledgeError {
396 KnowledgeError::Invalid { field, message }
397}
398
399fn uuid_field(name: &str, nullable: bool) -> Field {
400 Field::new(name, DataType::FixedSizeBinary(16), nullable)
401}
402
403fn timestamp_field(name: &str, nullable: bool) -> Field {
404 Field::new(
405 name,
406 DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
407 nullable,
408 )
409}
410
411fn fixed<'a>(
412 batch: &'a RecordBatch,
413 name: &'static str,
414) -> Result<&'a FixedSizeBinaryArray, KnowledgeError> {
415 batch
416 .column_by_name(name)
417 .and_then(|value| value.as_any().downcast_ref())
418 .ok_or_else(|| invalid(name, "column type mismatch"))
419}
420
421fn timestamp<'a>(
422 batch: &'a RecordBatch,
423 name: &'static str,
424) -> Result<&'a TimestampMicrosecondArray, KnowledgeError> {
425 batch
426 .column_by_name(name)
427 .and_then(|value| value.as_any().downcast_ref())
428 .ok_or_else(|| invalid(name, "column type mismatch"))
429}
430
431fn uint32<'a>(
432 batch: &'a RecordBatch,
433 name: &'static str,
434) -> Result<&'a UInt32Array, KnowledgeError> {
435 batch
436 .column_by_name(name)
437 .and_then(|value| value.as_any().downcast_ref())
438 .ok_or_else(|| invalid(name, "column type mismatch"))
439}
440
441fn uuid_at(
442 values: &FixedSizeBinaryArray,
443 row: usize,
444 field: &'static str,
445) -> Result<Uuid, KnowledgeError> {
446 Uuid::from_slice(values.value(row)).map_err(|_| invalid(field, "invalid UUID bytes"))
447}
448
449#[cfg(test)]
450mod tests {
451 use super::*;
452
453 fn uuid7(seed: u8) -> Uuid {
454 let mut bytes = [seed; 16];
455 bytes[6] = (bytes[6] & 0x0f) | 0x70;
456 bytes[8] = (bytes[8] & 0x3f) | 0x80;
457 Uuid::from_bytes(bytes)
458 }
459
460 fn event(
461 id: u8,
462 assertion: u8,
463 from: Option<i64>,
464 to: Option<i64>,
465 recorded_at: i64,
466 ) -> AssertionValidityEvent {
467 AssertionValidityEvent::new(
468 uuid7(id),
469 uuid7(assertion),
470 from,
471 to,
472 Some(uuid7(id.wrapping_add(40))),
473 uuid7(id.wrapping_add(80)),
474 recorded_at,
475 )
476 .unwrap()
477 }
478
479 #[test]
480 fn half_open_unbounded_empty_and_invalid_intervals_are_explicit() {
481 let always = event(1, 20, None, None, 1);
482 assert!(always.contains(i64::MIN));
483 assert!(always.contains(i64::MAX));
484
485 let bounded = event(2, 20, Some(10), Some(20), 2);
486 assert!(!bounded.contains(9));
487 assert!(bounded.contains(10));
488 assert!(bounded.contains(19));
489 assert!(!bounded.contains(20));
490
491 let empty = event(3, 20, Some(10), Some(10), 3);
492 assert!(!empty.contains(10));
493 assert!(
494 AssertionValidityEvent::new(
495 uuid7(4),
496 uuid7(20),
497 Some(11),
498 Some(10),
499 None,
500 uuid7(84),
501 4,
502 )
503 .is_err()
504 );
505 }
506
507 #[test]
508 fn transaction_cutoff_and_uuid_tie_breaking_preserve_prior_views() {
509 let ledger = AssertionValidityLedger::new(vec![
510 event(1, 20, Some(0), Some(10), 5),
511 event(2, 20, Some(10), None, 5),
512 event(3, 20, None, Some(0), 8),
513 ])
514 .unwrap();
515 assert_eq!(
516 ledger.current_for_at(uuid7(20), 4),
517 None,
518 "no event is visible before its transaction time"
519 );
520 assert_eq!(
521 ledger
522 .current_for_at(uuid7(20), 5)
523 .unwrap()
524 .validity_event_uuid,
525 uuid7(2)
526 );
527 assert_eq!(ledger.is_valid_at(uuid7(20), 5, 10), Some(true));
528 assert_eq!(ledger.is_valid_at(uuid7(20), 7, 10), Some(true));
529 assert_eq!(ledger.is_valid_at(uuid7(20), 8, 10), Some(false));
530 }
531
532 #[test]
533 fn round_trip_merge_replay_and_fingerprint_are_deterministic() {
534 let first = event(2, 20, Some(10), None, 2);
535 let second = event(1, 21, None, Some(10), 1);
536 let ledger = AssertionValidityLedger::new(vec![first.clone(), second]).unwrap();
537 let decoded = AssertionValidityLedger::from_batches(&[ledger.batch().unwrap()]).unwrap();
538 assert_eq!(decoded, ledger);
539 assert_eq!(decoded.events[0].validity_event_uuid, uuid7(1));
540 assert_eq!(
541 decoded.event_fingerprint(uuid7(2)).unwrap(),
542 ledger.event_fingerprint(uuid7(2)).unwrap()
543 );
544 assert_eq!(
545 ledger
546 .merge(&AssertionValidityLedger::new(vec![first.clone()]).unwrap())
547 .unwrap(),
548 ledger
549 );
550
551 let mut conflicting = first;
552 conflicting.valid_to_micros = Some(30);
553 assert!(
554 ledger
555 .merge(&AssertionValidityLedger {
556 events: vec![conflicting]
557 })
558 .is_err()
559 );
560 }
561
562 #[test]
563 fn schema_registry_freezes_capability_family_and_order() {
564 let entry = schema_registry_entry();
565 assert_eq!(entry.capability_id, "valid_time");
566 assert_eq!(entry.capability_version, 1);
567 assert_eq!(entry.record_family, "assertion_validity_events");
568 assert_eq!(entry.sort_key, &["recorded_at", "validity_event_uuid"]);
569 assert_eq!(entry.implementation_issue, 781);
570 }
571}