Skip to main content

graphforge_knowledge/
valid_time.rs

1//! Optional append-only M21 assertion valid-time events.
2
3use 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
19/// Optional valid-time capability contract.
20pub const VALID_TIME_CAPABILITY_VERSION: u32 = 1;
21/// Assertion-validity record contract.
22pub const ASSERTION_VALIDITY_CONTRACT_VERSION: u32 = 1;
23/// Half-open interval interpretation policy.
24pub const ASSERTION_VALIDITY_POLICY_VERSION: u32 = 1;
25
26/// Authoritative assertion-validity event schema.
27pub 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/// One immutable correction to an assertion's valid-time interpretation.
50#[derive(Clone, Debug, Eq, PartialEq)]
51pub struct AssertionValidityEvent {
52    /// Caller-supplied UUIDv7 identity/idempotency key.
53    pub validity_event_uuid: Uuid,
54    /// Existing immutable assertion.
55    pub assertion_uuid: Uuid,
56    /// Inclusive lower bound, or unbounded when absent.
57    pub valid_from_micros: Option<i64>,
58    /// Exclusive upper bound, or unbounded when absent.
59    pub valid_to_micros: Option<i64>,
60    /// Optional existing immutable reasoning record.
61    pub reasoning_uuid: Option<Uuid>,
62    /// Existing producing provenance event.
63    pub provenance_uuid: Uuid,
64    /// Mandatory transaction time.
65    pub recorded_at_micros: i64,
66    /// Frozen record contract.
67    pub contract_version: u32,
68}
69
70impl AssertionValidityEvent {
71    /// Validate and construct one immutable half-open interval event.
72    #[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    /// Evaluate the half-open interval `[valid_from, valid_to)`.
97    ///
98    /// Equal bounds form a valid empty interval. Missing bounds are unbounded.
99    #[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/// Validated append-only assertion-validity participant.
108#[derive(Clone, Debug, Default, Eq, PartialEq)]
109pub struct AssertionValidityLedger {
110    /// Events ordered by `(recorded_at, validity_event_uuid)`.
111    pub events: Vec<AssertionValidityEvent>,
112}
113
114impl AssertionValidityLedger {
115    /// Validate, sort, and construct one complete participant.
116    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    /// Merge staged append-only events with exact replay semantics.
136    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    /// Select the validity interpretation visible at transaction cutoff.
157    #[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    /// Test validity using the interpretation visible at transaction cutoff.
173    #[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    /// Canonical fingerprint over one exact immutable event.
185    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    /// Build the authoritative Arrow batch.
209    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    /// Decode and validate exact Arrow batches.
245    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}