Skip to main content

graphforge_value/
catalog.rs

1//! Persisted runtime catalog schema, checked identities, and atomic batch merging.
2
3use std::collections::HashMap;
4use std::sync::{Arc, LazyLock};
5
6use arrow::array::{
7    Array, ArrayRef, RecordBatch, StringArray, StringBuilder, TimestampMicrosecondArray,
8    TimestampMicrosecondBuilder, UInt32Array, UInt32Builder, UInt64Array, UInt64Builder,
9};
10use arrow::datatypes::{DataType, Field, Schema, SchemaRef, TimeUnit};
11
12use crate::{RuntimeEntityId, RuntimePropId, RuntimeRelationId};
13use graphforge_core::GfError;
14
15// ---------------------------------------------------------------------------
16// Arrow schema
17// ---------------------------------------------------------------------------
18
19/// Arrow schema for `topology/runtime_catalog.parquet`.
20pub static RUNTIME_CATALOG_SCHEMA: LazyLock<SchemaRef> = LazyLock::new(|| {
21    Arc::new(Schema::new(vec![
22        Field::new("entry_kind", DataType::Utf8, false),
23        Field::new("name", DataType::Utf8, false),
24        Field::new("runtime_id", DataType::UInt32, false),
25        Field::new("observation_count", DataType::UInt64, false),
26        Field::new(
27            "first_seen",
28            DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
29            false,
30        ),
31        Field::new(
32            "last_seen",
33            DataType::Timestamp(TimeUnit::Microsecond, Some("UTC".into())),
34            false,
35        ),
36        Field::new("owner_label", DataType::Utf8, true),
37    ]))
38});
39
40// ---------------------------------------------------------------------------
41// Private internals
42// ---------------------------------------------------------------------------
43
44#[derive(Debug, Clone, Copy, PartialEq, Eq)]
45enum EntryKind {
46    EntityType,
47    RelationType,
48    Property,
49}
50
51#[derive(Debug, Clone, Copy, PartialEq, Eq)]
52enum CatalogIdentity {
53    Entity(RuntimeEntityId),
54    Relation(RuntimeRelationId),
55    Property(RuntimePropId),
56}
57
58impl CatalogIdentity {
59    fn checked(kind: EntryKind, raw: u32) -> Result<Self, GfError> {
60        match kind {
61            EntryKind::EntityType => RuntimeEntityId::new(raw).map(Self::Entity),
62            EntryKind::RelationType => RuntimeRelationId::new(raw).map(Self::Relation),
63            EntryKind::Property => RuntimePropId::new(raw).map(Self::Property),
64        }
65        .map_err(|error| GfError::Storage(format!("runtime_catalog invalid identity: {error}")))
66    }
67
68    fn kind(self) -> EntryKind {
69        match self {
70            Self::Entity(_) => EntryKind::EntityType,
71            Self::Relation(_) => EntryKind::RelationType,
72            Self::Property(_) => EntryKind::Property,
73        }
74    }
75
76    fn get(self) -> u32 {
77        match self {
78            Self::Entity(id) => id.get(),
79            Self::Relation(id) => id.get(),
80            Self::Property(id) => id.get(),
81        }
82    }
83}
84
85#[derive(Debug, Clone)]
86struct CatalogEntry {
87    identity: CatalogIdentity,
88    name: String,
89    observation_count: u64,
90    /// Microseconds since Unix epoch (UTC).
91    first_seen: i64,
92    /// Microseconds since Unix epoch (UTC).
93    last_seen: i64,
94    /// For `Property` entries: the label this property was observed on.
95    owner_label: Option<String>,
96}
97
98// ---------------------------------------------------------------------------
99// RuntimeCatalogData
100// ---------------------------------------------------------------------------
101
102/// Checked runtime catalog carrier and its encoding-preserving Arrow codec.
103///
104/// The caller chooses observation timestamps and whether a name may be interned.
105/// Entries preserve their stable local identity and insertion order across batches.
106#[derive(Debug, Clone, Default)]
107pub struct RuntimeCatalogData {
108    /// name → index into `entries`
109    entity_types: HashMap<String, usize>,
110    /// name → index into `entries`
111    relation_types: HashMap<String, usize>,
112    /// (name, owner_label) → index into `entries`
113    properties: HashMap<(String, Option<String>), usize>,
114    /// All entries in insertion order.
115    entries: Vec<CatalogEntry>,
116    /// Next ID to assign for entity types and relation types (shared space).
117    next_type_id: u32,
118    /// Next ID to assign for properties.
119    next_prop_id: u32,
120}
121
122impl RuntimeCatalogData {
123    /// Creates an empty catalog.
124    #[must_use]
125    pub fn new() -> Self {
126        Self::default()
127    }
128
129    /// Number of retained catalog entries.
130    #[must_use]
131    pub fn entry_count(&self) -> usize {
132        self.entries.len()
133    }
134
135    /// Latest retained observation across all entry kinds, in UTC microseconds.
136    /// Returns `None` only for an empty catalog; pre-epoch values are preserved.
137    #[must_use]
138    pub fn latest_observation_micros(&self) -> Option<i64> {
139        self.entries.iter().map(|entry| entry.last_seen).max()
140    }
141
142    /// Exact UTF-8 bytes retained by entry names and property owners.
143    #[must_use]
144    pub fn retained_identifier_bytes(&self) -> usize {
145        self.entries.iter().fold(0_usize, |bytes, entry| {
146            bytes
147                .saturating_add(entry.name.len())
148                .saturating_add(entry.owner_label.as_ref().map_or(0, String::len))
149        })
150    }
151
152    /// Interns an entity label at a caller-authoritative timestamp.
153    pub fn intern_label_at(&mut self, name: &str, now: i64) -> Result<RuntimeEntityId, GfError> {
154        self.intern_label_observed_at(name, now, 1)
155    }
156
157    /// Interns an entity label and records `observations` sightings of it at
158    /// once: the catalog `observation_count` an equal number of single
159    /// [`Self::intern_label_at`] calls would leave. `observations` is at least one.
160    pub fn intern_label_observed_at(
161        &mut self,
162        name: &str,
163        now: i64,
164        observations: u64,
165    ) -> Result<RuntimeEntityId, GfError> {
166        if let Some(&idx) = self.entity_types.get(name) {
167            let entry = &mut self.entries[idx];
168            entry.observation_count = entry
169                .observation_count
170                .checked_add(observations)
171                .ok_or_else(|| {
172                    GfError::Storage("runtime_catalog observation count overflow".to_owned())
173                })?;
174            entry.last_seen = now;
175            let CatalogIdentity::Entity(id) = entry.identity else {
176                unreachable!("catalog index matches identity kind")
177            };
178            return Ok(id);
179        }
180        let id = RuntimeEntityId::new(self.next_type_id).map_err(|error| {
181            GfError::Storage(format!("runtime_catalog exhausted ID range: {error}"))
182        })?;
183        self.next_type_id += 1;
184        let idx = self.entries.len();
185        self.entries.push(CatalogEntry {
186            identity: CatalogIdentity::Entity(id),
187            name: name.to_owned(),
188            observation_count: observations,
189            first_seen: now,
190            last_seen: now,
191            owner_label: None,
192        });
193        self.entity_types.insert(name.to_owned(), idx);
194        Ok(id)
195    }
196
197    /// Interns a relation type at a caller-authoritative timestamp.
198    pub fn intern_relation_type_at(
199        &mut self,
200        name: &str,
201        now: i64,
202    ) -> Result<RuntimeRelationId, GfError> {
203        self.intern_relation_type_observed_at(name, now, 1)
204    }
205
206    /// Interns a relation type and records `observations` sightings of it at
207    /// once, as an equal number of [`Self::intern_relation_type_at`] calls would.
208    pub fn intern_relation_type_observed_at(
209        &mut self,
210        name: &str,
211        now: i64,
212        observations: u64,
213    ) -> Result<RuntimeRelationId, GfError> {
214        if let Some(&idx) = self.relation_types.get(name) {
215            let entry = &mut self.entries[idx];
216            entry.observation_count = entry
217                .observation_count
218                .checked_add(observations)
219                .ok_or_else(|| {
220                    GfError::Storage("runtime_catalog observation count overflow".to_owned())
221                })?;
222            entry.last_seen = now;
223            let CatalogIdentity::Relation(id) = entry.identity else {
224                unreachable!("catalog index matches identity kind")
225            };
226            return Ok(id);
227        }
228        let id = RuntimeRelationId::new(self.next_type_id).map_err(|error| {
229            GfError::Storage(format!("runtime_catalog exhausted ID range: {error}"))
230        })?;
231        self.next_type_id += 1;
232        let idx = self.entries.len();
233        self.entries.push(CatalogEntry {
234            identity: CatalogIdentity::Relation(id),
235            name: name.to_owned(),
236            observation_count: observations,
237            first_seen: now,
238            last_seen: now,
239            owner_label: None,
240        });
241        self.relation_types.insert(name.to_owned(), idx);
242        Ok(id)
243    }
244
245    /// Interns a property at a caller-authoritative timestamp.
246    pub fn intern_property_at(
247        &mut self,
248        name: &str,
249        owner_label: Option<&str>,
250        now: i64,
251    ) -> Result<RuntimePropId, GfError> {
252        let key = (name.to_owned(), owner_label.map(str::to_owned));
253        if let Some(&idx) = self.properties.get(&key) {
254            let entry = &mut self.entries[idx];
255            entry.observation_count = entry.observation_count.checked_add(1).ok_or_else(|| {
256                GfError::Storage("runtime_catalog observation count overflow".to_owned())
257            })?;
258            entry.last_seen = now;
259            let CatalogIdentity::Property(id) = entry.identity else {
260                unreachable!("catalog index matches identity kind")
261            };
262            return Ok(id);
263        }
264        let id = RuntimePropId::new(self.next_prop_id).map_err(|error| {
265            GfError::Storage(format!("runtime_catalog exhausted ID range: {error}"))
266        })?;
267        self.next_prop_id += 1;
268        let idx = self.entries.len();
269        self.entries.push(CatalogEntry {
270            identity: CatalogIdentity::Property(id),
271            name: name.to_owned(),
272            observation_count: 1,
273            first_seen: now,
274            last_seen: now,
275            owner_label: owner_label.map(str::to_owned),
276        });
277        self.properties.insert(key, idx);
278        Ok(id)
279    }
280
281    /// Returns `true` if `name` has been interned as an entity type.
282    #[must_use]
283    pub fn contains_entity_type(&self, name: &str) -> bool {
284        self.entity_types.contains_key(name)
285    }
286
287    /// Returns `true` if `name` has been interned as a relation type.
288    #[must_use]
289    pub fn contains_relation_type(&self, name: &str) -> bool {
290        self.relation_types.contains_key(name)
291    }
292
293    /// Returns `true` if the owner-scoped property has been interned.
294    #[must_use]
295    pub fn contains_property(&self, name: &str, owner_label: Option<&str>) -> bool {
296        self.properties
297            .contains_key(&(name.to_owned(), owner_label.map(str::to_owned)))
298    }
299
300    /// Returns all interned entity type names (order unspecified).
301    #[must_use]
302    pub fn entity_types(&self) -> Vec<&str> {
303        self.entity_types.keys().map(String::as_str).collect()
304    }
305
306    /// Returns all interned relation type names (order unspecified).
307    #[must_use]
308    pub fn relation_types(&self) -> Vec<&str> {
309        self.relation_types.keys().map(String::as_str).collect()
310    }
311
312    /// Returns all property names observed on `label` (order unspecified).
313    #[must_use]
314    pub fn properties_for(&self, label: &str) -> Vec<&str> {
315        self.properties
316            .iter()
317            .filter(|((_, owner), _)| owner.as_deref() == Some(label))
318            .map(|((name, _), _)| name.as_str())
319            .collect()
320    }
321
322    /// Resolves a [`RuntimePropId`] back to the property name it was interned
323    /// under, or `None` if no property entry carries that ID.
324    ///
325    /// Used by the relational lowering layer to turn a numeric `PropertyAccess`
326    /// ID back into the real column name when reading exploratory property
327    /// tables.
328    #[must_use]
329    pub fn property_name(&self, id: RuntimePropId) -> Option<&str> {
330        self.entries
331            .iter()
332            .find(|e| e.identity.kind() == EntryKind::Property && e.identity.get() == id.get())
333            .map(|e| e.name.as_str())
334    }
335
336    /// Returns `(RuntimePropId, name)` for every interned property (order
337    /// unspecified). Convenient for building a `PropId → name` map in one pass.
338    pub fn property_names(&self) -> impl Iterator<Item = (RuntimePropId, &str)> + '_ {
339        self.entries.iter().filter_map(|e| match e.identity {
340            CatalogIdentity::Property(id) => Some((id, e.name.as_str())),
341            _ => None,
342        })
343    }
344
345    /// Resolves a relation-type [`RuntimeRelationId`] back to the name it was
346    /// interned under, or `None` if no relation-type entry carries that ID.
347    ///
348    /// Used by the relational lowering layer to resolve a `TypedEdgeScan`'s
349    /// relation name when reading exploratory edge tables (no ontology present).
350    #[must_use]
351    pub fn relation_type_name(&self, id: RuntimeRelationId) -> Option<&str> {
352        self.entries
353            .iter()
354            .find(|e| e.identity.kind() == EntryKind::RelationType && e.identity.get() == id.get())
355            .map(|e| e.name.as_str())
356    }
357
358    /// Returns `(RuntimeRelationId, name)` for every interned relation type (order
359    /// unspecified). Convenient for building a `TypeId → relation-name` map.
360    pub fn relation_type_names_with_ids(
361        &self,
362    ) -> impl Iterator<Item = (RuntimeRelationId, &str)> + '_ {
363        self.entries.iter().filter_map(|e| match e.identity {
364            CatalogIdentity::Relation(id) => Some((id, e.name.as_str())),
365            _ => None,
366        })
367    }
368
369    /// Resolves an entity-type (node label) [`RuntimeEntityId`] back to the name
370    /// it was interned under, or `None` if no entity-type entry carries that ID.
371    ///
372    /// Mirror of [`relation_type_name`](Self::relation_type_name) for labels —
373    /// used to render a real label name for an unlabelled `MATCH (n) RETURN n`
374    /// in exploratory mode, where the ontology map is empty (#889).
375    #[must_use]
376    pub fn entity_type_name(&self, id: RuntimeEntityId) -> Option<&str> {
377        self.entries
378            .iter()
379            .find(|e| e.identity.kind() == EntryKind::EntityType && e.identity.get() == id.get())
380            .map(|e| e.name.as_str())
381    }
382
383    /// Returns `(RuntimeEntityId, name)` for every interned entity type (node
384    /// label), order unspecified. Convenient for building a
385    /// `TypeId → label-name` map.
386    pub fn entity_type_names_with_ids(&self) -> impl Iterator<Item = (RuntimeEntityId, &str)> + '_ {
387        self.entries.iter().filter_map(|e| match e.identity {
388            CatalogIdentity::Entity(id) => Some((id, e.name.as_str())),
389            _ => None,
390        })
391    }
392
393    /// Serialises the catalog to an Arrow [`RecordBatch`] using [`RUNTIME_CATALOG_SCHEMA`].
394    ///
395    /// The resulting batch can be written to `topology/runtime_catalog.parquet`
396    /// and later restored via [`from_record_batch`](Self::from_record_batch).
397    #[must_use]
398    pub fn to_record_batch(&self) -> RecordBatch {
399        let n = self.entries.len();
400        let mut kind_b = StringBuilder::with_capacity(n, n * 12);
401        let mut name_b = StringBuilder::with_capacity(n, n * 32);
402        let mut id_b = UInt32Builder::with_capacity(n);
403        let mut count_b = UInt64Builder::with_capacity(n);
404        let mut first_b = TimestampMicrosecondBuilder::with_capacity(n);
405        let mut last_b = TimestampMicrosecondBuilder::with_capacity(n);
406        let mut owner_b = StringBuilder::with_capacity(n, n * 16);
407
408        for entry in &self.entries {
409            kind_b.append_value(match entry.identity.kind() {
410                EntryKind::EntityType => "entity_type",
411                EntryKind::RelationType => "relation_type",
412                EntryKind::Property => "property",
413            });
414            name_b.append_value(&entry.name);
415            id_b.append_value(entry.identity.get());
416            count_b.append_value(entry.observation_count);
417            first_b.append_value(entry.first_seen);
418            last_b.append_value(entry.last_seen);
419            match &entry.owner_label {
420                Some(label) => owner_b.append_value(label),
421                None => owner_b.append_null(),
422            }
423        }
424
425        let first_arr = first_b.finish().with_timezone_opt(Some(Arc::from("UTC")));
426        let last_arr = last_b.finish().with_timezone_opt(Some(Arc::from("UTC")));
427
428        let columns: Vec<ArrayRef> = vec![
429            Arc::new(kind_b.finish()),
430            Arc::new(name_b.finish()),
431            Arc::new(id_b.finish()),
432            Arc::new(count_b.finish()),
433            Arc::new(first_arr),
434            Arc::new(last_arr),
435            Arc::new(owner_b.finish()),
436        ];
437
438        RecordBatch::try_new(RUNTIME_CATALOG_SCHEMA.clone(), columns)
439            .expect("schema and array lengths must be consistent")
440    }
441
442    /// Restores a `RuntimeCatalogData` from an Arrow [`RecordBatch`] previously
443    /// produced by [`to_record_batch`](Self::to_record_batch).
444    ///
445    /// # Errors
446    /// Returns [`GfError::Storage`] if any column has the wrong type or an
447    /// unknown `entry_kind` value is encountered.
448    pub fn from_record_batch(batch: &RecordBatch) -> Result<Self, GfError> {
449        Self::from_record_batches(std::iter::once(batch))
450    }
451
452    /// Restores a catalog from a bounded stream of persisted Arrow batches
453    /// without concatenating or retaining those batches.
454    #[allow(clippy::too_many_lines)] // One persisted schema decoder; splitting obscures row authority.
455    pub fn from_record_batches<'a>(
456        batches: impl IntoIterator<Item = &'a RecordBatch>,
457    ) -> Result<Self, GfError> {
458        Self::decode_record_batches_at(batches, 0, 0)
459    }
460
461    #[allow(clippy::too_many_lines)]
462    fn decode_record_batches_at<'a>(
463        batches: impl IntoIterator<Item = &'a RecordBatch>,
464        initial_type_id: u32,
465        initial_prop_id: u32,
466    ) -> Result<Self, GfError> {
467        let storage_err = |msg: &str| GfError::Storage(msg.to_owned());
468        let mut catalog = Self::new();
469        let mut next_type_id = initial_type_id;
470        let mut next_prop_id = initial_prop_id;
471        for batch in batches {
472            if batch.schema().as_ref() != RUNTIME_CATALOG_SCHEMA.as_ref() {
473                return Err(storage_err("runtime_catalog schema is not canonical"));
474            }
475            let kinds = batch
476                .column(0)
477                .as_any()
478                .downcast_ref::<StringArray>()
479                .ok_or_else(|| storage_err("runtime_catalog col 0 (entry_kind) not Utf8"))?;
480            let names = batch
481                .column(1)
482                .as_any()
483                .downcast_ref::<StringArray>()
484                .ok_or_else(|| storage_err("runtime_catalog col 1 (name) not Utf8"))?;
485            let ids = batch
486                .column(2)
487                .as_any()
488                .downcast_ref::<UInt32Array>()
489                .ok_or_else(|| storage_err("runtime_catalog col 2 (runtime_id) not UInt32"))?;
490            let counts = batch
491                .column(3)
492                .as_any()
493                .downcast_ref::<UInt64Array>()
494                .ok_or_else(|| {
495                    storage_err("runtime_catalog col 3 (observation_count) not UInt64")
496                })?;
497            let first_seens = batch
498                .column(4)
499                .as_any()
500                .downcast_ref::<TimestampMicrosecondArray>()
501                .ok_or_else(|| {
502                    storage_err("runtime_catalog col 4 (first_seen) not TimestampMicrosecond")
503                })?;
504            let last_seens = batch
505                .column(5)
506                .as_any()
507                .downcast_ref::<TimestampMicrosecondArray>()
508                .ok_or_else(|| {
509                    storage_err("runtime_catalog col 5 (last_seen) not TimestampMicrosecond")
510                })?;
511            let owners = batch
512                .column(6)
513                .as_any()
514                .downcast_ref::<StringArray>()
515                .ok_or_else(|| storage_err("runtime_catalog col 6 (owner_label) not Utf8"))?;
516
517            if kinds.null_count() != 0
518                || names.null_count() != 0
519                || ids.null_count() != 0
520                || counts.null_count() != 0
521                || first_seens.null_count() != 0
522                || last_seens.null_count() != 0
523            {
524                return Err(storage_err("runtime_catalog required column contains null"));
525            }
526
527            for row in 0..batch.num_rows() {
528                let kind = match kinds.value(row) {
529                    "entity_type" => EntryKind::EntityType,
530                    "relation_type" => EntryKind::RelationType,
531                    "property" => EntryKind::Property,
532                    other => {
533                        return Err(GfError::Storage(format!(
534                            "runtime_catalog: unknown entry_kind '{other}'"
535                        )));
536                    }
537                };
538                let name = names.value(row).to_owned();
539                let runtime_id = ids.value(row);
540                let observation_count = counts.value(row);
541                let first_seen = first_seens.value(row);
542                let last_seen = last_seens.value(row);
543                let owner_label = if owners.is_null(row) {
544                    None
545                } else {
546                    Some(owners.value(row).to_owned())
547                };
548
549                if name.is_empty()
550                    || observation_count == 0
551                    || first_seen > last_seen
552                    || (kind != EntryKind::Property && owner_label.is_some())
553                {
554                    return Err(storage_err("runtime_catalog row is not canonical"));
555                }
556                let expected_id = match kind {
557                    EntryKind::EntityType | EntryKind::RelationType => &mut next_type_id,
558                    EntryKind::Property => &mut next_prop_id,
559                };
560                if runtime_id != *expected_id {
561                    return Err(storage_err(
562                        "runtime_catalog IDs are not unique and contiguous in insertion order",
563                    ));
564                }
565                *expected_id = expected_id.checked_add(1).ok_or_else(|| {
566                    storage_err("runtime_catalog persisted ID exceeds supported range")
567                })?;
568
569                let idx = catalog.entries.len();
570                catalog.entries.push(CatalogEntry {
571                    identity: CatalogIdentity::checked(kind, runtime_id)?,
572                    name: name.clone(),
573                    observation_count,
574                    first_seen,
575                    last_seen,
576                    owner_label: owner_label.clone(),
577                });
578
579                match kind {
580                    EntryKind::EntityType => {
581                        if catalog.entity_types.insert(name, idx).is_some() {
582                            return Err(storage_err(
583                                "runtime_catalog contains duplicate entity type",
584                            ));
585                        }
586                    }
587                    EntryKind::RelationType => {
588                        if catalog.relation_types.insert(name, idx).is_some() {
589                            return Err(storage_err(
590                                "runtime_catalog contains duplicate relation type",
591                            ));
592                        }
593                    }
594                    EntryKind::Property => {
595                        if catalog
596                            .properties
597                            .insert((name, owner_label), idx)
598                            .is_some()
599                        {
600                            return Err(storage_err("runtime_catalog contains duplicate property"));
601                        }
602                    }
603                }
604            }
605        }
606
607        catalog.next_type_id = next_type_id;
608        catalog.next_prop_id = next_prop_id;
609        Ok(catalog)
610    }
611
612    /// Appends one persisted catalog batch while preserving its stable IDs and
613    /// observations. The caller may therefore decode a Parquet catalog in
614    /// bounded windows rather than concatenating it in memory.
615    pub fn extend_from_record_batch(&mut self, batch: &RecordBatch) -> Result<(), GfError> {
616        let incoming = Self::decode_record_batches_at(
617            std::iter::once(batch),
618            self.next_type_id,
619            self.next_prop_id,
620        )?;
621        for entry in &incoming.entries {
622            let duplicate = match entry.identity.kind() {
623                EntryKind::EntityType => self.entity_types.contains_key(&entry.name),
624                EntryKind::RelationType => self.relation_types.contains_key(&entry.name),
625                EntryKind::Property => self
626                    .properties
627                    .contains_key(&(entry.name.clone(), entry.owner_label.clone())),
628            };
629            if duplicate {
630                return Err(GfError::Storage(
631                    "runtime_catalog contains a duplicate persisted entry".to_owned(),
632                ));
633            }
634        }
635        self.next_type_id = incoming.next_type_id;
636        self.next_prop_id = incoming.next_prop_id;
637        for entry in incoming.entries {
638            let idx = self.entries.len();
639            match entry.identity.kind() {
640                EntryKind::EntityType => {
641                    self.entity_types.insert(entry.name.clone(), idx);
642                }
643                EntryKind::RelationType => {
644                    self.relation_types.insert(entry.name.clone(), idx);
645                }
646                EntryKind::Property => {
647                    let key = (entry.name.clone(), entry.owner_label.clone());
648                    self.properties.insert(key, idx);
649                }
650            }
651            self.entries.push(entry);
652        }
653        Ok(())
654    }
655}
656
657#[cfg(test)]
658mod tests {
659    use super::*;
660
661    fn catalog() -> RuntimeCatalogData {
662        let mut catalog = RuntimeCatalogData::new();
663        catalog.intern_label_at("Person", 1).unwrap();
664        catalog.intern_relation_type_at("Person", 2).unwrap();
665        catalog
666            .intern_property_at("name", Some("Person"), 3)
667            .unwrap();
668        catalog
669            .intern_property_at("name", Some("Company"), 4)
670            .unwrap();
671        catalog
672    }
673
674    #[test]
675    fn bounded_batches_preserve_catalog_bytes_and_namespace_overlap() {
676        let original = catalog().to_record_batch();
677        let batches: Vec<_> = (0..original.num_rows())
678            .map(|row| original.slice(row, 1))
679            .collect();
680        let restored = RuntimeCatalogData::from_record_batches(&batches).unwrap();
681        assert_eq!(restored.to_record_batch(), original);
682        let mut appended = RuntimeCatalogData::new();
683        for batch in &batches {
684            appended.extend_from_record_batch(batch).unwrap();
685        }
686        assert_eq!(appended.to_record_batch(), original);
687    }
688
689    #[test]
690    fn failed_batch_append_preserves_existing_catalog_authority() {
691        let mut existing = catalog();
692        let before = existing.to_record_batch();
693        let mut next = existing.clone();
694        next.intern_label_at("Company", 5).unwrap();
695        let valid = next.to_record_batch().slice(before.num_rows(), 1);
696        let mut columns = valid.columns().to_vec();
697        columns[1] = Arc::new(StringArray::from(vec!["Person"]));
698        let duplicate_name = RecordBatch::try_new(RUNTIME_CATALOG_SCHEMA.clone(), columns).unwrap();
699        assert!(existing.extend_from_record_batch(&duplicate_name).is_err());
700        assert_eq!(existing.to_record_batch(), before);
701        assert!(RuntimeCatalogData::from_record_batches([&before, &duplicate_name]).is_err());
702
703        let mut columns = valid.columns().to_vec();
704        columns[2] = Arc::new(UInt32Array::from(vec![0]));
705        let duplicate_id = RecordBatch::try_new(RUNTIME_CATALOG_SCHEMA.clone(), columns).unwrap();
706        assert!(existing.extend_from_record_batch(&duplicate_id).is_err());
707        assert_eq!(existing.to_record_batch(), before);
708        assert!(RuntimeCatalogData::from_record_batches([&before, &duplicate_id]).is_err());
709
710        existing.extend_from_record_batch(&valid).unwrap();
711        assert_eq!(existing.to_record_batch(), next.to_record_batch());
712    }
713    #[test]
714    fn duplicate_ids_in_each_catalog_kind_fail_across_batch_boundaries() {
715        for kind in [
716            EntryKind::EntityType,
717            EntryKind::RelationType,
718            EntryKind::Property,
719        ] {
720            let mut existing = RuntimeCatalogData::new();
721            let intern = |catalog: &mut RuntimeCatalogData, name| match kind {
722                EntryKind::EntityType => catalog.intern_label_at(name, 1).map(|_| ()),
723                EntryKind::RelationType => catalog.intern_relation_type_at(name, 1).map(|_| ()),
724                EntryKind::Property => catalog.intern_property_at(name, None, 1).map(|_| ()),
725            };
726            intern(&mut existing, "first").unwrap();
727            let before = existing.to_record_batch();
728            let mut next = existing.clone();
729            intern(&mut next, "second").unwrap();
730            let good = next.to_record_batch().slice(1, 1);
731            let mut columns = good.columns().to_vec();
732            columns[2] = Arc::new(UInt32Array::from(vec![0]));
733            let duplicate = RecordBatch::try_new(RUNTIME_CATALOG_SCHEMA.clone(), columns).unwrap();
734            assert!(
735                existing.extend_from_record_batch(&duplicate).is_err(),
736                "{kind:?}"
737            );
738            assert_eq!(existing.to_record_batch(), before, "{kind:?}");
739            assert!(
740                RuntimeCatalogData::from_record_batches([&before, &duplicate]).is_err(),
741                "{kind:?}"
742            );
743            existing.extend_from_record_batch(&good).unwrap();
744            assert_eq!(existing.to_record_batch(), next.to_record_batch());
745        }
746    }
747
748    fn intern_kind(
749        catalog: &mut RuntimeCatalogData,
750        kind: EntryKind,
751        name: &str,
752    ) -> Result<(), GfError> {
753        match kind {
754            EntryKind::EntityType => catalog.intern_label_at(name, 9).map(|_| ()),
755            EntryKind::RelationType => catalog.intern_relation_type_at(name, 9).map(|_| ()),
756            EntryKind::Property => catalog
757                .intern_property_at(name, Some("Person"), 9)
758                .map(|_| ()),
759        }
760    }
761
762    #[test]
763    fn duplicate_names_fail_without_changing_existing_authority() {
764        for kind in [EntryKind::RelationType, EntryKind::Property] {
765            let mut existing = RuntimeCatalogData::new();
766            intern_kind(&mut existing, kind, "first").unwrap();
767            let before = existing.to_record_batch();
768            let mut next = existing.clone();
769            intern_kind(&mut next, kind, "second").unwrap();
770            let valid = next.to_record_batch().slice(1, 1);
771            let mut columns = valid.columns().to_vec();
772            columns[1] = Arc::new(StringArray::from(vec!["first"]));
773            let duplicate = RecordBatch::try_new(RUNTIME_CATALOG_SCHEMA.clone(), columns).unwrap();
774            let expected = match kind {
775                EntryKind::RelationType => "runtime_catalog contains duplicate relation type",
776                EntryKind::Property => "runtime_catalog contains duplicate property",
777                EntryKind::EntityType => unreachable!(),
778            };
779            assert!(matches!(
780                RuntimeCatalogData::from_record_batches([&before, &duplicate]),
781                Err(GfError::Storage(message)) if message == expected
782            ));
783            assert!(matches!(
784                existing.extend_from_record_batch(&duplicate),
785                Err(GfError::Storage(message))
786                    if message == "runtime_catalog contains a duplicate persisted entry"
787            ));
788            assert_eq!(existing.to_record_batch(), before);
789            existing.extend_from_record_batch(&valid).unwrap();
790            assert_eq!(existing.to_record_batch(), next.to_record_batch());
791        }
792    }
793
794    #[test]
795    fn observation_overflow_preserves_counts_timestamps_and_identity() {
796        for kind in [
797            EntryKind::EntityType,
798            EntryKind::RelationType,
799            EntryKind::Property,
800        ] {
801            let mut seed = RuntimeCatalogData::new();
802            intern_kind(&mut seed, kind, "first").unwrap();
803            let batch = seed.to_record_batch();
804            let mut columns = batch.columns().to_vec();
805            columns[3] = Arc::new(UInt64Array::from(vec![u64::MAX]));
806            let persisted = RecordBatch::try_new(RUNTIME_CATALOG_SCHEMA.clone(), columns).unwrap();
807            let mut restored = RuntimeCatalogData::from_record_batch(&persisted).unwrap();
808            let result = match kind {
809                EntryKind::EntityType => restored.intern_label_at("first", 10).map(|_| ()),
810                EntryKind::RelationType => {
811                    restored.intern_relation_type_at("first", 10).map(|_| ())
812                }
813                EntryKind::Property => restored
814                    .intern_property_at("first", Some("Person"), 10)
815                    .map(|_| ()),
816            };
817            assert!(matches!(result, Err(GfError::Storage(message))
818                if message == "runtime_catalog observation count overflow"));
819            assert_eq!(restored.to_record_batch(), persisted);
820            assert_eq!(restored.next_type_id, seed.next_type_id);
821            assert_eq!(restored.next_prop_id, seed.next_prop_id);
822            intern_kind(&mut restored, kind, "second").unwrap();
823            assert_eq!(restored.entries[1].identity.get(), 1);
824        }
825    }
826
827    #[test]
828    fn final_valid_allocation_is_followed_by_atomic_exhaustion() {
829        for kind in [
830            EntryKind::EntityType,
831            EntryKind::RelationType,
832            EntryKind::Property,
833        ] {
834            let mut catalog = RuntimeCatalogData::new();
835            // Reach the allocator boundary without allocating billions of entries.
836            let limit = match kind {
837                EntryKind::EntityType | EntryKind::RelationType => 1_u32 << 30,
838                EntryKind::Property => u32::MAX,
839            };
840            match kind {
841                EntryKind::EntityType | EntryKind::RelationType => catalog.next_type_id = limit - 1,
842                EntryKind::Property => catalog.next_prop_id = limit - 1,
843            }
844            intern_kind(&mut catalog, kind, "last").unwrap();
845            assert_eq!(catalog.entries[0].identity.get(), limit - 1);
846            let before = catalog.to_record_batch();
847            let counters = (catalog.next_type_id, catalog.next_prop_id);
848            assert!(matches!(intern_kind(&mut catalog, kind, "overflow"),
849                Err(GfError::Storage(message))
850                    if message.starts_with("runtime_catalog exhausted ID range:")));
851            assert_eq!(catalog.to_record_batch(), before);
852            assert_eq!((catalog.next_type_id, catalog.next_prop_id), counters);
853            assert!(!catalog.entity_types.contains_key("overflow"));
854            assert!(!catalog.relation_types.contains_key("overflow"));
855            assert!(
856                !catalog
857                    .properties
858                    .contains_key(&("overflow".to_owned(), Some("Person".to_owned())))
859            );
860            intern_kind(&mut catalog, kind, "last").unwrap();
861            assert_eq!(catalog.entries[0].identity.get(), limit - 1);
862            assert_eq!(catalog.entries[0].observation_count, 2);
863        }
864    }
865}