1use 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
15pub 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#[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 first_seen: i64,
92 last_seen: i64,
94 owner_label: Option<String>,
96}
97
98#[derive(Debug, Clone, Default)]
107pub struct RuntimeCatalogData {
108 entity_types: HashMap<String, usize>,
110 relation_types: HashMap<String, usize>,
112 properties: HashMap<(String, Option<String>), usize>,
114 entries: Vec<CatalogEntry>,
116 next_type_id: u32,
118 next_prop_id: u32,
120}
121
122impl RuntimeCatalogData {
123 #[must_use]
125 pub fn new() -> Self {
126 Self::default()
127 }
128
129 #[must_use]
131 pub fn entry_count(&self) -> usize {
132 self.entries.len()
133 }
134
135 #[must_use]
138 pub fn latest_observation_micros(&self) -> Option<i64> {
139 self.entries.iter().map(|entry| entry.last_seen).max()
140 }
141
142 #[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 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 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 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 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 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 #[must_use]
283 pub fn contains_entity_type(&self, name: &str) -> bool {
284 self.entity_types.contains_key(name)
285 }
286
287 #[must_use]
289 pub fn contains_relation_type(&self, name: &str) -> bool {
290 self.relation_types.contains_key(name)
291 }
292
293 #[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 #[must_use]
302 pub fn entity_types(&self) -> Vec<&str> {
303 self.entity_types.keys().map(String::as_str).collect()
304 }
305
306 #[must_use]
308 pub fn relation_types(&self) -> Vec<&str> {
309 self.relation_types.keys().map(String::as_str).collect()
310 }
311
312 #[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 #[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 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 #[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 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 #[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 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 #[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 pub fn from_record_batch(batch: &RecordBatch) -> Result<Self, GfError> {
449 Self::from_record_batches(std::iter::once(batch))
450 }
451
452 #[allow(clippy::too_many_lines)] 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 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 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}