1use super::{
2 CanwuError, ErrorCode, MAX_STATE_DELTA_PAGES, PayloadSchema, StateKey, StatePageBlob,
3 StatePageProvider, StateVisibility, canonical_byte_hash, canonical_hash, state_page_id,
4};
5use canwu_core::{
6 CoreEntityKind, DomainEntityType, DomainKindClass, DomainRecordKind, DomainRecordRef,
7 DomainRecordType, DomainValueType, EntityRef, KnowledgeHolderPolicy, TypedDomainRecordRef,
8};
9use canwu_time::SimTime;
10use im::{HashMap as PersistentHashMap, HashSet as PersistentHashSet, Vector};
11use serde::de::DeserializeOwned;
12use serde::{Deserialize, Serialize};
13use serde_json::Value;
14use std::collections::{BTreeMap, BTreeSet};
15use std::fmt::Write as _;
16use std::mem::size_of;
17use std::sync::Arc;
18
19#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
25#[serde(deny_unknown_fields)]
26pub struct DomainRecordCommitmentRoots {
27 pub primary: String,
28 pub reverse_references: String,
29 pub successor_of: String,
30 pub predecessors_of: String,
31}
32
33#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
34#[serde(deny_unknown_fields)]
35pub struct DomainRecordPageRoots {
36 #[serde(default, skip_serializing_if = "Option::is_none")]
37 pub primary: Option<String>,
38 #[serde(default, skip_serializing_if = "Option::is_none")]
39 pub reverse_references: Option<String>,
40 #[serde(default, skip_serializing_if = "Option::is_none")]
41 pub successor_of: Option<String>,
42 #[serde(default, skip_serializing_if = "Option::is_none")]
43 pub predecessors_of: Option<String>,
44 pub commitment_roots: DomainRecordCommitmentRoots,
45}
46
47#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
48#[serde(deny_unknown_fields)]
49pub struct PatriciaStoreMetrics {
50 pub entries: u64,
51 pub logical_nodes: u64,
52 pub leaf_nodes: u64,
53 pub branch_nodes: u64,
54 pub structural_bytes: u64,
55 pub estimated_resident_bytes: u64,
56 pub depth_p50: u16,
57 pub depth_p95: u16,
58 pub depth_p99: u16,
59 pub max_depth: u16,
60 pub root_hash: String,
61}
62
63#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
64#[serde(deny_unknown_fields)]
65pub struct DomainStoreScaleMetrics {
66 pub records: u64,
67 pub key_pages: u64,
68 pub record_hamt_entries: u64,
69 pub reverse_lookup_keys: u64,
70 pub reverse_reference_edges: u64,
71 pub successor_lookup_entries: u64,
72 pub predecessor_lookup_keys: u64,
73 pub predecessor_edges: u64,
74 pub total_patricia_nodes: u64,
75 pub total_patricia_structural_bytes: u64,
76 pub total_patricia_estimated_resident_bytes: u64,
77 pub primary: PatriciaStoreMetrics,
78 pub reverse_references: PatriciaStoreMetrics,
79 pub successor_of: PatriciaStoreMetrics,
80 pub predecessors_of: PatriciaStoreMetrics,
81}
82
83#[derive(Clone, Debug)]
84pub struct PersistentDomainRecordStore {
85 records: PersistentHashMap<DomainRecordRef, Arc<DomainRecord>>,
86 record_key_pages: Vector<Arc<Vec<DomainRecordRef>>>,
87 reverse_lookup: PersistentHashMap<DomainRecordRef, PersistentHashSet<DomainRecordRef>>,
88 successor_lookup: PersistentHashMap<DomainRecordRef, DomainRecordRef>,
89 predecessor_lookup: PersistentHashMap<DomainRecordRef, PersistentHashSet<DomainRecordRef>>,
90 primary: PersistentPatricia,
91 reverse_references: PersistentPatricia,
92 successor_of: PersistentPatricia,
93 predecessors_of: PersistentPatricia,
94 roots: DomainRecordCommitmentRoots,
95 primary_leaf_count: usize,
96}
97
98const DOMAIN_RECORD_KEY_PAGE_CAPACITY: usize = 256;
99const MAX_AFFECTED_RECORD_VALIDATION_CLOSURE: usize = 16_384;
100
101impl PersistentDomainRecordStore {
102 pub fn from_records(
103 records: BTreeMap<DomainRecordRef, DomainRecord>,
104 ) -> Result<Self, CanwuError> {
105 let mut store = Self::empty();
106 for (reference, record) in &records {
107 store.primary.insert(reference, record)?;
108 store.add_indexes(record)?;
109 }
110 let sorted = records.into_iter().collect::<Vec<_>>();
111 for chunk in sorted.chunks(DOMAIN_RECORD_KEY_PAGE_CAPACITY) {
112 store.record_key_pages.push_back(Arc::new(
113 chunk
114 .iter()
115 .map(|(reference, _)| reference.clone())
116 .collect(),
117 ));
118 for (reference, record) in chunk {
119 store
120 .records
121 .insert(reference.clone(), Arc::new(record.clone()));
122 store.add_lookup_indexes(record);
123 }
124 }
125 store.refresh_roots();
126 Ok(store)
127 }
128
129 #[must_use]
130 pub fn get(&self, reference: &DomainRecordRef) -> Option<&DomainRecord> {
131 self.records.get(reference).map(Arc::as_ref)
132 }
133
134 #[must_use]
135 pub fn contains_key(&self, reference: &DomainRecordRef) -> bool {
136 self.records.contains_key(reference)
137 }
138
139 #[must_use]
140 pub fn len(&self) -> usize {
141 self.records.len()
142 }
143
144 #[must_use]
145 pub fn is_empty(&self) -> bool {
146 self.records.is_empty()
147 }
148
149 pub fn iter(&self) -> impl Iterator<Item = (&DomainRecordRef, &DomainRecord)> {
150 self.record_key_pages.iter().flat_map(|page| {
151 page.iter().filter_map(|reference| {
152 self.records
153 .get(reference)
154 .map(|record| (reference, record.as_ref()))
155 })
156 })
157 }
158
159 pub fn values(&self) -> impl Iterator<Item = &DomainRecord> {
160 self.iter().map(|(_, record)| record)
161 }
162
163 #[must_use]
164 pub fn roots(&self) -> &DomainRecordCommitmentRoots {
165 &self.roots
166 }
167
168 #[must_use]
171 pub const fn primary_leaf_count(&self) -> usize {
172 self.primary_leaf_count
173 }
174
175 #[must_use]
176 pub fn materialize(&self) -> BTreeMap<DomainRecordRef, DomainRecord> {
177 self.iter()
178 .map(|(reference, record)| (reference.clone(), record.clone()))
179 .collect()
180 }
181
182 #[cfg(test)]
185 #[must_use]
186 pub(crate) fn shares_root_with(&self, other: &Self) -> bool {
187 match (&self.primary.root, &other.primary.root) {
188 (None, None) => true,
189 (Some(left), Some(right)) => Arc::ptr_eq(left, right),
190 _ => false,
191 }
192 }
193
194 pub(crate) fn commitment_root(&self) -> Result<String, CanwuError> {
195 canonical_hash("canwu.format8.domain-record.roots.v1", &self.roots)
196 }
197
198 pub fn state_pages(&self) -> Result<(DomainRecordPageRoots, Vec<StatePageBlob>), CanwuError> {
199 let mut pages = BTreeMap::new();
200 let roots = DomainRecordPageRoots {
201 primary: collect_patricia_pages(&self.primary, &mut pages)?,
202 reverse_references: collect_patricia_pages(&self.reverse_references, &mut pages)?,
203 successor_of: collect_patricia_pages(&self.successor_of, &mut pages)?,
204 predecessors_of: collect_patricia_pages(&self.predecessors_of, &mut pages)?,
205 commitment_roots: self.roots.clone(),
206 };
207 if pages.len() > MAX_STATE_DELTA_PAGES {
208 return Err(invalid_patricia_page(
209 "domain record store exceeds the bounded state-page count",
210 ));
211 }
212 Ok((roots, pages.into_values().collect()))
213 }
214
215 pub fn missing_state_pages(
219 &self,
220 provider: &dyn StatePageProvider,
221 ) -> Result<(DomainRecordPageRoots, Vec<StatePageBlob>), CanwuError> {
222 let mut pages = BTreeMap::new();
223 let roots = DomainRecordPageRoots {
224 primary: collect_missing_patricia_pages(&self.primary, provider, &mut pages)?,
225 reverse_references: collect_missing_patricia_pages(
226 &self.reverse_references,
227 provider,
228 &mut pages,
229 )?,
230 successor_of: collect_missing_patricia_pages(&self.successor_of, provider, &mut pages)?,
231 predecessors_of: collect_missing_patricia_pages(
232 &self.predecessors_of,
233 provider,
234 &mut pages,
235 )?,
236 commitment_roots: self.roots.clone(),
237 };
238 if pages.len() > MAX_STATE_DELTA_PAGES {
239 return Err(invalid_patricia_page(
240 "domain record delta exceeds the bounded state-page count",
241 ));
242 }
243 Ok((roots, pages.into_values().collect()))
244 }
245
246 #[must_use]
247 pub fn primary_metrics(&self) -> PatriciaStoreMetrics {
248 self.primary.metrics()
249 }
250
251 pub fn from_state_pages(
252 roots: &DomainRecordPageRoots,
253 provider: &dyn StatePageProvider,
254 ) -> Result<Self, CanwuError> {
255 let mut cache = BTreeMap::new();
256 let primary = load_patricia_root(
257 "canwu.format8.domain-record.primary.v1",
258 roots.primary.as_deref(),
259 provider,
260 &mut cache,
261 )?;
262 let reverse_references = load_patricia_root(
263 "canwu.format8.domain-record.reverse.v1",
264 roots.reverse_references.as_deref(),
265 provider,
266 &mut cache,
267 )?;
268 let successor_of = load_patricia_root(
269 "canwu.format8.domain-record.successor.v1",
270 roots.successor_of.as_deref(),
271 provider,
272 &mut cache,
273 )?;
274 let predecessors_of = load_patricia_root(
275 "canwu.format8.domain-record.predecessors.v1",
276 roots.predecessors_of.as_deref(),
277 provider,
278 &mut cache,
279 )?;
280 let mut materialized = BTreeMap::new();
281 if let Some(root) = &primary.root {
282 collect_primary_records(root, &mut materialized)?;
283 }
284 let rebuilt = Self::from_records(materialized)?;
285 let loaded = Self {
286 records: rebuilt.records,
287 record_key_pages: rebuilt.record_key_pages,
288 reverse_lookup: rebuilt.reverse_lookup,
289 successor_lookup: rebuilt.successor_lookup,
290 predecessor_lookup: rebuilt.predecessor_lookup,
291 primary,
292 reverse_references,
293 successor_of,
294 predecessors_of,
295 roots: roots.commitment_roots.clone(),
296 primary_leaf_count: rebuilt.primary_leaf_count,
297 };
298 let actual_roots = DomainRecordCommitmentRoots {
299 primary: loaded.primary.root_hash(),
300 reverse_references: loaded.reverse_references.root_hash(),
301 successor_of: loaded.successor_of.root_hash(),
302 predecessors_of: loaded.predecessors_of.root_hash(),
303 };
304 if actual_roots != roots.commitment_roots || rebuilt.roots != roots.commitment_roots {
305 return Err(invalid_patricia_page(
306 "state pages do not reconstruct the committed domain-record indexes",
307 ));
308 }
309 loaded.validate_internal_indexes()?;
310 Ok(loaded)
311 }
312
313 pub(crate) fn validate_internal_indexes(&self) -> Result<(), CanwuError> {
317 if self.records.len() != self.primary_leaf_count {
318 return Err(invalid_patricia_page(
319 "domain record HAMT count disagrees with the primary leaf count",
320 ));
321 }
322 let mut paged_keys = BTreeSet::new();
323 let mut previous: Option<&DomainRecordRef> = None;
324 for page in &self.record_key_pages {
325 if page.is_empty()
326 || page.len() > DOMAIN_RECORD_KEY_PAGE_CAPACITY
327 || page.windows(2).any(|pair| pair[0] >= pair[1])
328 {
329 return Err(invalid_patricia_page(
330 "domain record ordered key page is malformed",
331 ));
332 }
333 for reference in page.iter() {
334 if previous.is_some_and(|prior| prior >= reference)
335 || !paged_keys.insert(reference.clone())
336 || self
337 .records
338 .get(reference)
339 .is_none_or(|record| record.reference != *reference)
340 {
341 return Err(invalid_patricia_page(
342 "domain record ordered key pages disagree with the HAMT",
343 ));
344 }
345 previous = Some(reference);
346 }
347 }
348 let hamt_keys = self.records.keys().cloned().collect::<BTreeSet<_>>();
349 if paged_keys != hamt_keys {
350 return Err(invalid_patricia_page(
351 "domain record ordered key pages omit or invent HAMT entries",
352 ));
353 }
354
355 let mut expected_reverse = BTreeMap::<DomainRecordRef, BTreeSet<DomainRecordRef>>::new();
356 let mut expected_successors = BTreeMap::<DomainRecordRef, DomainRecordRef>::new();
357 let mut expected_predecessors =
358 BTreeMap::<DomainRecordRef, BTreeSet<DomainRecordRef>>::new();
359 for record in self.records.values().map(Arc::as_ref) {
360 for target in record.references.iter().filter_map(|reference| {
361 if let DomainReferenceTarget::Domain(target) = &reference.target {
362 Some(target)
363 } else {
364 None
365 }
366 }) {
367 expected_reverse
368 .entry(target.clone())
369 .or_default()
370 .insert(record.reference.clone());
371 }
372 if let DomainRecordLifecycle::Retired {
373 successor: Some(successor),
374 ..
375 } = &record.lifecycle
376 {
377 expected_successors.insert(record.reference.clone(), successor.clone());
378 expected_predecessors
379 .entry(successor.clone())
380 .or_default()
381 .insert(record.reference.clone());
382 }
383 }
384 let actual_reverse = self
385 .reverse_lookup
386 .iter()
387 .map(|(target, sources)| {
388 (
389 target.clone(),
390 sources.iter().cloned().collect::<BTreeSet<_>>(),
391 )
392 })
393 .collect::<BTreeMap<_, _>>();
394 let actual_successors = self
395 .successor_lookup
396 .iter()
397 .map(|(source, successor)| (source.clone(), successor.clone()))
398 .collect::<BTreeMap<_, _>>();
399 let actual_predecessors = self
400 .predecessor_lookup
401 .iter()
402 .map(|(target, sources)| {
403 (
404 target.clone(),
405 sources.iter().cloned().collect::<BTreeSet<_>>(),
406 )
407 })
408 .collect::<BTreeMap<_, _>>();
409 if actual_reverse != expected_reverse
410 || actual_successors != expected_successors
411 || actual_predecessors != expected_predecessors
412 {
413 return Err(invalid_patricia_page(
414 "domain record COW lookup indexes disagree with authoritative records",
415 ));
416 }
417 let actual_roots = DomainRecordCommitmentRoots {
418 primary: self.primary.root_hash(),
419 reverse_references: self.reverse_references.root_hash(),
420 successor_of: self.successor_of.root_hash(),
421 predecessors_of: self.predecessors_of.root_hash(),
422 };
423 if actual_roots != self.roots {
424 return Err(invalid_patricia_page(
425 "domain record Patricia roots disagree with cached commitments",
426 ));
427 }
428 Ok(())
429 }
430
431 fn empty() -> Self {
432 let primary = PersistentPatricia::new("canwu.format8.domain-record.primary.v1");
433 let reverse_references = PersistentPatricia::new("canwu.format8.domain-record.reverse.v1");
434 let successor_of = PersistentPatricia::new("canwu.format8.domain-record.successor.v1");
435 let predecessors_of =
436 PersistentPatricia::new("canwu.format8.domain-record.predecessors.v1");
437 let roots = DomainRecordCommitmentRoots {
438 primary: primary.root_hash(),
439 reverse_references: reverse_references.root_hash(),
440 successor_of: successor_of.root_hash(),
441 predecessors_of: predecessors_of.root_hash(),
442 };
443 Self {
444 records: PersistentHashMap::new(),
445 record_key_pages: Vector::new(),
446 reverse_lookup: PersistentHashMap::new(),
447 successor_lookup: PersistentHashMap::new(),
448 predecessor_lookup: PersistentHashMap::new(),
449 primary,
450 reverse_references,
451 successor_of,
452 predecessors_of,
453 roots,
454 primary_leaf_count: 0,
455 }
456 }
457
458 fn insert_record(
459 &mut self,
460 reference: DomainRecordRef,
461 record: DomainRecord,
462 ) -> Result<(), CanwuError> {
463 let previous = self.records.get(&reference).cloned();
464 if let Some(previous) = previous.as_deref() {
465 self.remove_indexes(previous)?;
466 self.remove_lookup_indexes(previous);
467 }
468 self.primary.insert(&reference, &record)?;
469 self.add_indexes(&record)?;
470 self.add_lookup_indexes(&record);
471 if previous.is_none() {
472 self.insert_record_key(reference.clone());
473 }
474 self.records.insert(reference, Arc::new(record));
475 self.refresh_roots();
476 Ok(())
477 }
478
479 fn insert_record_key(&mut self, reference: DomainRecordRef) {
480 if self.record_key_pages.is_empty() {
481 self.record_key_pages.push_back(Arc::new(vec![reference]));
482 return;
483 }
484 let page_index = self.record_key_page_for(&reference);
485 let mut page = self.record_key_pages[page_index].as_ref().clone();
486 let position = page
487 .binary_search(&reference)
488 .expect_err("new domain record keys are unique");
489 page.insert(position, reference);
490 if page.len() <= DOMAIN_RECORD_KEY_PAGE_CAPACITY {
491 self.record_key_pages.set(page_index, Arc::new(page));
492 return;
493 }
494 let right = page.split_off(page.len() / 2);
495 self.record_key_pages.set(page_index, Arc::new(page));
496 self.record_key_pages
497 .insert(page_index.saturating_add(1), Arc::new(right));
498 }
499
500 fn record_key_page_for(&self, reference: &DomainRecordRef) -> usize {
501 let mut lower = 0usize;
502 let mut upper = self.record_key_pages.len();
503 while lower < upper {
504 let middle = lower + (upper - lower) / 2;
505 let last = self.record_key_pages[middle]
506 .last()
507 .expect("record key pages are never empty");
508 if last < reference {
509 lower = middle.saturating_add(1);
510 } else {
511 upper = middle;
512 }
513 }
514 lower.min(self.record_key_pages.len().saturating_sub(1))
515 }
516
517 fn record_key_start(&self, reference: &DomainRecordRef, excluded: bool) -> (usize, usize) {
518 if self.record_key_pages.is_empty() {
519 return (0, 0);
520 }
521 let page_index = self.record_key_page_for(reference);
522 let page = &self.record_key_pages[page_index];
523 let key_index = match page.binary_search(reference) {
524 Ok(index) if excluded => index.saturating_add(1),
525 Ok(index) | Err(index) => index,
526 };
527 if key_index == page.len() {
528 (page_index.saturating_add(1), 0)
529 } else {
530 (page_index, key_index)
531 }
532 }
533
534 fn add_lookup_indexes(&mut self, record: &DomainRecord) {
535 for target in record.references.iter().filter_map(|reference| {
536 if let DomainReferenceTarget::Domain(target) = &reference.target {
537 Some(target)
538 } else {
539 None
540 }
541 }) {
542 self.reverse_lookup
543 .entry(target.clone())
544 .or_default()
545 .insert(record.reference.clone());
546 }
547 if let DomainRecordLifecycle::Retired {
548 successor: Some(successor),
549 ..
550 } = &record.lifecycle
551 {
552 self.successor_lookup
553 .insert(record.reference.clone(), successor.clone());
554 self.predecessor_lookup
555 .entry(successor.clone())
556 .or_default()
557 .insert(record.reference.clone());
558 }
559 }
560
561 fn remove_lookup_indexes(&mut self, record: &DomainRecord) {
562 for target in record.references.iter().filter_map(|reference| {
563 if let DomainReferenceTarget::Domain(target) = &reference.target {
564 Some(target)
565 } else {
566 None
567 }
568 }) {
569 if let Some(mut sources) = self.reverse_lookup.get(target).cloned() {
570 sources.remove(&record.reference);
571 if sources.is_empty() {
572 self.reverse_lookup.remove(target);
573 } else {
574 self.reverse_lookup.insert(target.clone(), sources);
575 }
576 }
577 }
578 if let Some(successor) = self.successor_lookup.remove(&record.reference)
579 && let Some(mut predecessors) = self.predecessor_lookup.get(&successor).cloned()
580 {
581 predecessors.remove(&record.reference);
582 if predecessors.is_empty() {
583 self.predecessor_lookup.remove(&successor);
584 } else {
585 self.predecessor_lookup.insert(successor, predecessors);
586 }
587 }
588 }
589
590 fn add_indexes(&mut self, record: &DomainRecord) -> Result<(), CanwuError> {
591 let mut reverse = BTreeSet::new();
592 for reference in &record.references {
593 if let DomainReferenceTarget::Domain(target) = &reference.target {
594 reverse.insert(target.clone());
595 }
596 }
597 for target in reverse {
598 self.reverse_references
599 .insert(&(target, record.reference.clone()), &())?;
600 }
601 if let DomainRecordLifecycle::Retired {
602 successor: Some(target),
603 ..
604 } = &record.lifecycle
605 {
606 self.successor_of.insert(&record.reference, target)?;
607 self.predecessors_of
608 .insert(&(target.clone(), record.reference.clone()), &())?;
609 }
610 Ok(())
611 }
612
613 fn remove_indexes(&mut self, record: &DomainRecord) -> Result<(), CanwuError> {
614 let mut reverse = BTreeSet::new();
615 for reference in &record.references {
616 if let DomainReferenceTarget::Domain(target) = &reference.target {
617 reverse.insert(target.clone());
618 }
619 }
620 for target in reverse {
621 self.reverse_references
622 .remove(&(target, record.reference.clone()))?;
623 }
624 if let DomainRecordLifecycle::Retired {
625 successor: Some(target),
626 ..
627 } = &record.lifecycle
628 {
629 self.successor_of.remove(&record.reference)?;
630 self.predecessors_of
631 .remove(&(target.clone(), record.reference.clone()))?;
632 }
633 Ok(())
634 }
635
636 fn refresh_roots(&mut self) {
637 self.primary_leaf_count = self.primary.entry_count();
638 self.roots = DomainRecordCommitmentRoots {
639 primary: self.primary.root_hash(),
640 reverse_references: self.reverse_references.root_hash(),
641 successor_of: self.successor_of.root_hash(),
642 predecessors_of: self.predecessors_of.root_hash(),
643 };
644 }
645}
646
647#[derive(Clone, Debug)]
648struct PersistentPatricia {
649 domain: &'static str,
650 root: Option<Arc<PatriciaNode>>,
651}
652
653#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
654#[serde(deny_unknown_fields)]
655struct StoredPatriciaEntry {
656 key_hash: String,
657 key_bytes: Vec<u8>,
658 value_hash: String,
659 value_bytes: Vec<u8>,
660}
661
662#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
663#[serde(tag = "node", rename_all = "snake_case")]
664enum StoredPatriciaPage {
665 Leaf {
666 domain: String,
667 key_hash: String,
668 entries: Vec<StoredPatriciaEntry>,
669 entry_count: u64,
670 structural_bytes: u64,
671 },
672 Branch {
673 domain: String,
674 bit: u16,
675 left_page: String,
676 right_page: String,
677 entry_count: u64,
678 structural_bytes: u64,
679 },
680}
681
682impl PersistentPatricia {
683 const fn new(domain: &'static str) -> Self {
684 Self { domain, root: None }
685 }
686
687 fn insert<K: Serialize, V: Serialize>(&mut self, key: &K, value: &V) -> Result<(), CanwuError> {
688 let entry = PatriciaEntry::new(self.domain, key, value)?;
689 self.root = Some(insert_patricia(self.domain, self.root.as_ref(), entry));
690 Ok(())
691 }
692
693 fn remove<K: Serialize>(&mut self, key: &K) -> Result<(), CanwuError> {
694 let key_bytes = serde_json::to_vec(key).map_err(|error| {
695 CanwuError::new(
696 ErrorCode::InvalidSnapshot,
697 format!("cannot encode Patricia key: {error}"),
698 )
699 })?;
700 let key_hash = compact_canonical_byte_hash(&format!("{}.key", self.domain), &key_bytes);
701 self.root = remove_patricia(self.domain, self.root.as_ref(), key_hash, &key_bytes);
702 Ok(())
703 }
704
705 fn root_hash(&self) -> String {
706 self.root.as_ref().map_or_else(
707 || canonical_byte_hash(&format!("{}.empty", self.domain), &[]),
708 |root| root.hash().to_hex(),
709 )
710 }
711
712 fn entry_count(&self) -> usize {
713 self.root.as_ref().map_or(0, |root| root.entry_count())
714 }
715
716 fn metrics(&self) -> PatriciaStoreMetrics {
717 let Some(root) = &self.root else {
718 return PatriciaStoreMetrics {
719 entries: 0,
720 logical_nodes: 0,
721 leaf_nodes: 0,
722 branch_nodes: 0,
723 structural_bytes: 0,
724 estimated_resident_bytes: 0,
725 depth_p50: 0,
726 depth_p95: 0,
727 depth_p99: 0,
728 max_depth: 0,
729 root_hash: self.root_hash(),
730 };
731 };
732 let mut stack = vec![(Arc::clone(root), 0_u16)];
733 let mut leaf_nodes = 0_u64;
734 let mut branch_nodes = 0_u64;
735 let mut depths = Vec::with_capacity(root.entry_count());
736 let mut resident_bytes = 0_u64;
737 while let Some((node, depth)) = stack.pop() {
738 resident_bytes = resident_bytes.saturating_add(match node.as_ref() {
739 PatriciaNode::Leaf { entries, .. } => {
740 leaf_nodes += 1;
741 depths.extend(std::iter::repeat_n(depth, entries.len()));
742 size_of::<PatriciaNode>() as u64
743 + size_of::<Vec<PatriciaEntry>>() as u64
744 + entries
745 .capacity()
746 .saturating_mul(size_of::<PatriciaEntry>())
747 as u64
748 + entries
749 .iter()
750 .map(|entry| entry.key_bytes.capacity() as u64)
751 .sum::<u64>()
752 }
753 PatriciaNode::Branch { left, right, .. } => {
754 branch_nodes += 1;
755 let next_depth = depth.saturating_add(1);
756 stack.push((Arc::clone(left), next_depth));
757 stack.push((Arc::clone(right), next_depth));
758 size_of::<PatriciaNode>() as u64
759 }
760 });
761 }
762 depths.sort_unstable();
763 let percentile = |numerator: usize, denominator: usize| {
764 if depths.is_empty() {
765 return 0;
766 }
767 let index = depths
768 .len()
769 .saturating_mul(numerator)
770 .saturating_add(denominator - 1)
771 / denominator;
772 depths[index.saturating_sub(1).min(depths.len() - 1)]
773 };
774 PatriciaStoreMetrics {
775 entries: root.entry_count() as u64,
776 logical_nodes: leaf_nodes.saturating_add(branch_nodes),
777 leaf_nodes,
778 branch_nodes,
779 structural_bytes: root.structural_bytes() as u64,
780 estimated_resident_bytes: resident_bytes,
781 depth_p50: percentile(50, 100),
782 depth_p95: percentile(95, 100),
783 depth_p99: percentile(99, 100),
784 max_depth: depths.last().copied().unwrap_or(0),
785 root_hash: root.hash().to_hex(),
786 }
787 }
788}
789
790pub fn format8_patricia_scale_probe(
795 key_count: usize,
796) -> Result<DomainStoreScaleMetrics, CanwuError> {
797 let mut store = PersistentDomainRecordStore::empty();
798 let key_count_u64 = u64::try_from(key_count)
799 .map_err(|_| invalid_patricia_page("scale probe key count exceeds u64"))?;
800 for ordinal in 0..key_count {
801 let ordinal = u64::try_from(ordinal)
802 .map_err(|_| invalid_patricia_page("scale probe key exceeds u64"))?;
803 let reference =
804 DomainRecordRef::new("canwu.format8.scale", "record", format!("{ordinal:016x}"));
805 let predecessor = (ordinal > 0).then(|| {
806 DomainRecordRef::new(
807 "canwu.format8.scale",
808 "record",
809 format!("{:016x}", ordinal - 1),
810 )
811 });
812 let successor = (ordinal % 1_024 == 0 && ordinal + 1 < key_count_u64).then(|| {
813 DomainRecordRef::new(
814 "canwu.format8.scale",
815 "record",
816 format!("{:016x}", ordinal + 1),
817 )
818 });
819 let record = DomainRecord {
820 reference: reference.clone(),
821 owner: "canwu-format8-scale".to_owned(),
822 class: DomainRecordClass::Record,
823 version: 1,
824 lifecycle: successor.map_or(DomainRecordLifecycle::Active, |successor| {
825 DomainRecordLifecycle::Retired {
826 at: SimTime::EPOCH,
827 successor: Some(successor),
828 }
829 }),
830 payload: serde_json::json!({ "ordinal": ordinal }),
831 references: predecessor
832 .map(|target| {
833 vec![DomainReference {
834 role: "previous".to_owned(),
835 target: DomainReferenceTarget::Domain(target),
836 }]
837 })
838 .unwrap_or_default(),
839 };
840 store.insert_record(reference, record)?;
841 }
842 store.validate_internal_indexes()?;
843 let primary = store.primary.metrics();
844 let reverse_references = store.reverse_references.metrics();
845 let successor_of = store.successor_of.metrics();
846 let predecessors_of = store.predecessors_of.metrics();
847 let patricia = [
848 &primary,
849 &reverse_references,
850 &successor_of,
851 &predecessors_of,
852 ];
853 Ok(DomainStoreScaleMetrics {
854 records: store.records.len() as u64,
855 key_pages: store.record_key_pages.len() as u64,
856 record_hamt_entries: store.records.len() as u64,
857 reverse_lookup_keys: store.reverse_lookup.len() as u64,
858 reverse_reference_edges: store
859 .reverse_lookup
860 .iter()
861 .map(|(_, sources)| sources.len() as u64)
862 .sum(),
863 successor_lookup_entries: store.successor_lookup.len() as u64,
864 predecessor_lookup_keys: store.predecessor_lookup.len() as u64,
865 predecessor_edges: store
866 .predecessor_lookup
867 .iter()
868 .map(|(_, predecessors)| predecessors.len() as u64)
869 .sum(),
870 total_patricia_nodes: patricia.iter().map(|metrics| metrics.logical_nodes).sum(),
871 total_patricia_structural_bytes: patricia
872 .iter()
873 .map(|metrics| metrics.structural_bytes)
874 .sum(),
875 total_patricia_estimated_resident_bytes: patricia
876 .iter()
877 .map(|metrics| metrics.estimated_resident_bytes)
878 .sum(),
879 primary,
880 reverse_references,
881 successor_of,
882 predecessors_of,
883 })
884}
885
886#[derive(Clone, Debug)]
887struct PatriciaEntry {
888 key_bytes: Vec<u8>,
889 value_hash: CompactHash,
890 value_bytes: Vec<u8>,
891}
892
893impl PatriciaEntry {
894 fn new<K: Serialize, V: Serialize>(
895 domain: &str,
896 key: &K,
897 value: &V,
898 ) -> Result<Self, CanwuError> {
899 let mut key_bytes = serde_json::to_vec(key).map_err(|error| {
900 CanwuError::new(
901 ErrorCode::InvalidSnapshot,
902 format!("cannot encode Patricia key: {error}"),
903 )
904 })?;
905 key_bytes.shrink_to_fit();
906 let mut value_bytes = serde_json::to_vec(value).map_err(|error| {
907 CanwuError::new(
908 ErrorCode::InvalidSnapshot,
909 format!("cannot encode Patricia value: {error}"),
910 )
911 })?;
912 value_bytes.shrink_to_fit();
913 Ok(Self {
914 key_bytes,
915 value_hash: compact_canonical_byte_hash(&format!("{domain}.value"), &value_bytes),
916 value_bytes,
917 })
918 }
919}
920
921#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
922struct CompactHash([u8; 32]);
923
924impl CompactHash {
925 fn parse(value: &str) -> Result<Self, CanwuError> {
926 if value.len() != 64 {
927 return Err(invalid_patricia_page(
928 "Patricia hash must contain exactly 64 hexadecimal characters",
929 ));
930 }
931 let mut bytes = [0_u8; 32];
932 let (pairs, remainder) = value.as_bytes().as_chunks::<2>();
933 if !remainder.is_empty() {
934 return Err(invalid_patricia_page(
935 "Patricia hash must contain complete hexadecimal byte pairs",
936 ));
937 }
938 for (index, pair) in pairs.iter().enumerate() {
939 let high = decode_hex_nibble(pair[0])?;
940 let low = decode_hex_nibble(pair[1])?;
941 bytes[index] = (high << 4) | low;
942 }
943 Ok(Self(bytes))
944 }
945
946 fn to_hex(self) -> String {
947 let mut encoded = String::with_capacity(64);
948 for byte in self.0 {
949 write!(&mut encoded, "{byte:02x}").expect("writing to a String cannot fail");
950 }
951 encoded
952 }
953}
954
955fn decode_hex_nibble(byte: u8) -> Result<u8, CanwuError> {
956 match byte {
957 b'0'..=b'9' => Ok(byte - b'0'),
958 b'a'..=b'f' => Ok(byte - b'a' + 10),
959 _ => Err(invalid_patricia_page(
960 "Patricia hash must use lowercase hexadecimal",
961 )),
962 }
963}
964
965fn compact_canonical_byte_hash(domain: &str, bytes: &[u8]) -> CompactHash {
966 CompactHash::parse(&canonical_byte_hash(domain, bytes))
967 .expect("canonical_byte_hash always returns a lowercase 32-byte digest")
968}
969
970fn compact_state_page_id(stored: &StoredPatriciaPage) -> CompactHash {
971 let bytes =
972 serde_json::to_vec(stored).expect("canonical Patricia page values always serialize");
973 CompactHash::parse(&state_page_id(&bytes))
974 .expect("state_page_id always returns a lowercase 32-byte digest")
975}
976
977#[derive(Clone, Debug)]
978enum PatriciaNode {
979 Leaf {
980 key_hash: CompactHash,
981 entries: Arc<Vec<PatriciaEntry>>,
982 hash: CompactHash,
983 entry_count: usize,
984 structural_bytes: usize,
985 },
986 Branch {
987 bit: u16,
988 left: Arc<Self>,
989 right: Arc<Self>,
990 hash: CompactHash,
991 entry_count: usize,
992 structural_bytes: usize,
993 },
994}
995
996impl PatriciaNode {
997 const fn hash(&self) -> CompactHash {
998 match self {
999 Self::Leaf { hash, .. } | Self::Branch { hash, .. } => *hash,
1000 }
1001 }
1002
1003 fn sample_key_hash(&self) -> CompactHash {
1004 match self {
1005 Self::Leaf { key_hash, .. } => *key_hash,
1006 Self::Branch { left, .. } => left.sample_key_hash(),
1007 }
1008 }
1009
1010 const fn entry_count(&self) -> usize {
1011 match self {
1012 Self::Leaf { entry_count, .. } | Self::Branch { entry_count, .. } => *entry_count,
1013 }
1014 }
1015
1016 const fn structural_bytes(&self) -> usize {
1017 match self {
1018 Self::Leaf {
1019 structural_bytes, ..
1020 }
1021 | Self::Branch {
1022 structural_bytes, ..
1023 } => *structural_bytes,
1024 }
1025 }
1026}
1027
1028fn patricia_leaf(
1029 domain: &str,
1030 key_hash: CompactHash,
1031 mut entries: Vec<PatriciaEntry>,
1032) -> Arc<PatriciaNode> {
1033 entries.sort_by(|left, right| left.key_bytes.cmp(&right.key_bytes));
1034 let structural_bytes = 32
1035 + entries
1036 .iter()
1037 .map(|entry| 8 + entry.key_bytes.len() + 32)
1038 .sum::<usize>();
1039 let entry_count = entries.len();
1040 let stored = StoredPatriciaPage::Leaf {
1041 domain: domain.to_owned(),
1042 key_hash: key_hash.to_hex(),
1043 entries: entries
1044 .iter()
1045 .map(|entry| StoredPatriciaEntry {
1046 key_hash: key_hash.to_hex(),
1047 key_bytes: entry.key_bytes.clone(),
1048 value_hash: entry.value_hash.to_hex(),
1049 value_bytes: entry.value_bytes.clone(),
1050 })
1051 .collect(),
1052 entry_count: entry_count as u64,
1053 structural_bytes: structural_bytes as u64,
1054 };
1055 let hash = compact_state_page_id(&stored);
1056 Arc::new(PatriciaNode::Leaf {
1057 key_hash,
1058 entries: Arc::new(entries),
1059 hash,
1060 entry_count,
1061 structural_bytes,
1062 })
1063}
1064
1065fn patricia_branch(
1066 domain: &str,
1067 bit: usize,
1068 left: Arc<PatriciaNode>,
1069 right: Arc<PatriciaNode>,
1070) -> Arc<PatriciaNode> {
1071 let bit = u16::try_from(bit).expect("Patricia bit index is bounded to 256 bits");
1072 let entry_count = left.entry_count() + right.entry_count();
1073 let structural_bytes = 2 + 64 + 16 + left.structural_bytes() + right.structural_bytes();
1074 let stored = StoredPatriciaPage::Branch {
1075 domain: domain.to_owned(),
1076 bit,
1077 left_page: left.hash().to_hex(),
1078 right_page: right.hash().to_hex(),
1079 entry_count: entry_count as u64,
1080 structural_bytes: structural_bytes as u64,
1081 };
1082 let hash = compact_state_page_id(&stored);
1083 Arc::new(PatriciaNode::Branch {
1084 bit,
1085 left,
1086 right,
1087 hash,
1088 entry_count,
1089 structural_bytes,
1090 })
1091}
1092
1093fn insert_patricia(
1094 domain: &str,
1095 node: Option<&Arc<PatriciaNode>>,
1096 entry: PatriciaEntry,
1097) -> Arc<PatriciaNode> {
1098 let entry_key_hash = compact_canonical_byte_hash(&format!("{domain}.key"), &entry.key_bytes);
1099 let Some(root) = node else {
1100 return patricia_leaf(domain, entry_key_hash, vec![entry]);
1101 };
1102 let mut current = Arc::clone(root);
1103 let mut path = Vec::new();
1104 let mut entry = Some(entry);
1105 let mut replacement = loop {
1106 let difference = first_discriminating_bit(current.sample_key_hash(), entry_key_hash);
1107 if let PatriciaNode::Branch {
1108 bit, left, right, ..
1109 } = current.as_ref()
1110 {
1111 let branch_bit = usize::from(*bit);
1112 if difference.is_none_or(|difference| difference >= branch_bit) {
1113 if bit_at(entry_key_hash, branch_bit) == 0 {
1114 path.push((branch_bit, Arc::clone(right), true));
1115 current = Arc::clone(left);
1116 } else {
1117 path.push((branch_bit, Arc::clone(left), false));
1118 current = Arc::clone(right);
1119 }
1120 continue;
1121 }
1122 }
1123 let Some(difference) = difference else {
1124 let PatriciaNode::Leaf { entries, .. } = current.as_ref() else {
1125 unreachable!("a branch with an equal sample hash was descended above");
1126 };
1127 let mut next = entries.as_ref().clone();
1128 let entry = entry
1129 .take()
1130 .expect("Patricia insertion consumes its entry exactly once");
1131 match next.binary_search_by(|candidate| candidate.key_bytes.cmp(&entry.key_bytes)) {
1132 Ok(index) => next[index] = entry,
1133 Err(index) => next.insert(index, entry),
1134 }
1135 break patricia_leaf(domain, current.sample_key_hash(), next);
1136 };
1137 let leaf = patricia_leaf(
1138 domain,
1139 entry_key_hash,
1140 vec![
1141 entry
1142 .take()
1143 .expect("Patricia insertion consumes its entry exactly once"),
1144 ],
1145 );
1146 break if bit_at(entry_key_hash, difference) == 0 {
1147 patricia_branch(domain, difference, leaf, Arc::clone(¤t))
1148 } else {
1149 patricia_branch(domain, difference, Arc::clone(¤t), leaf)
1150 };
1151 };
1152 while let Some((bit, sibling, descended_left)) = path.pop() {
1153 replacement = if descended_left {
1154 patricia_branch(domain, bit, replacement, sibling)
1155 } else {
1156 patricia_branch(domain, bit, sibling, replacement)
1157 };
1158 }
1159 replacement
1160}
1161
1162fn remove_patricia(
1163 domain: &str,
1164 node: Option<&Arc<PatriciaNode>>,
1165 key_hash: CompactHash,
1166 key_bytes: &[u8],
1167) -> Option<Arc<PatriciaNode>> {
1168 let node = node?;
1169 match node.as_ref() {
1170 PatriciaNode::Leaf {
1171 key_hash: leaf_hash,
1172 entries,
1173 ..
1174 } => {
1175 if *leaf_hash != key_hash {
1176 return Some(Arc::clone(node));
1177 }
1178 let mut next = entries.as_ref().clone();
1179 let Ok(index) =
1180 next.binary_search_by(|entry| entry.key_bytes.as_slice().cmp(key_bytes))
1181 else {
1182 return Some(Arc::clone(node));
1183 };
1184 next.remove(index);
1185 (!next.is_empty()).then(|| patricia_leaf(domain, *leaf_hash, next))
1186 }
1187 PatriciaNode::Branch {
1188 bit, left, right, ..
1189 } => {
1190 let bit = usize::from(*bit);
1191 if bit_at(key_hash, bit) == 0 {
1192 let next_left = remove_patricia(domain, Some(left), key_hash, key_bytes);
1193 next_left.map_or_else(
1194 || Some(Arc::clone(right)),
1195 |next_left| Some(patricia_branch(domain, bit, next_left, Arc::clone(right))),
1196 )
1197 } else {
1198 let next_right = remove_patricia(domain, Some(right), key_hash, key_bytes);
1199 next_right.map_or_else(
1200 || Some(Arc::clone(left)),
1201 |next_right| Some(patricia_branch(domain, bit, Arc::clone(left), next_right)),
1202 )
1203 }
1204 }
1205 }
1206}
1207
1208fn first_discriminating_bit(left: CompactHash, right: CompactHash) -> Option<usize> {
1209 (0..256).find(|bit| bit_at(left, *bit) != bit_at(right, *bit))
1210}
1211
1212fn bit_at(hash: CompactHash, bit: usize) -> u8 {
1213 let byte = hash.0.get(bit / 8).copied().unwrap_or(0);
1214 (byte >> (7 - (bit % 8))) & 1
1215}
1216
1217fn collect_patricia_pages(
1218 tree: &PersistentPatricia,
1219 pages: &mut BTreeMap<String, StatePageBlob>,
1220) -> Result<Option<String>, CanwuError> {
1221 tree.root
1222 .as_ref()
1223 .map(|root| collect_patricia_node_pages(tree.domain, root, pages))
1224 .transpose()
1225}
1226
1227fn collect_missing_patricia_pages(
1228 tree: &PersistentPatricia,
1229 provider: &dyn StatePageProvider,
1230 pages: &mut BTreeMap<String, StatePageBlob>,
1231) -> Result<Option<String>, CanwuError> {
1232 tree.root
1233 .as_ref()
1234 .map(|root| collect_missing_patricia_node_pages(tree.domain, root, provider, pages))
1235 .transpose()
1236}
1237
1238fn collect_missing_patricia_node_pages(
1239 domain: &str,
1240 node: &Arc<PatriciaNode>,
1241 provider: &dyn StatePageProvider,
1242 pages: &mut BTreeMap<String, StatePageBlob>,
1243) -> Result<String, CanwuError> {
1244 let page_id = node.hash().to_hex();
1245 if let Some(existing) = provider.load_state_page(&page_id)? {
1246 existing.validate()?;
1247 if existing.page_id != page_id {
1248 return Err(invalid_patricia_page(
1249 "state-page provider returned the wrong Patricia page",
1250 ));
1251 }
1252 return Ok(page_id);
1253 }
1254 if let PatriciaNode::Branch { left, right, .. } = node.as_ref() {
1255 collect_missing_patricia_node_pages(domain, left, provider, pages)?;
1256 collect_missing_patricia_node_pages(domain, right, provider, pages)?;
1257 }
1258 let page = patricia_node_page(domain, node)?;
1259 if page.page_id != page_id {
1260 return Err(invalid_patricia_page(
1261 "Patricia node commitment disagrees with its state page",
1262 ));
1263 }
1264 pages.insert(page_id.clone(), page);
1265 if pages.len() > MAX_STATE_DELTA_PAGES {
1266 return Err(invalid_patricia_page(
1267 "Patricia delta exceeds the bounded page count",
1268 ));
1269 }
1270 Ok(page_id)
1271}
1272
1273fn collect_patricia_node_pages(
1274 domain: &str,
1275 node: &Arc<PatriciaNode>,
1276 pages: &mut BTreeMap<String, StatePageBlob>,
1277) -> Result<String, CanwuError> {
1278 if let PatriciaNode::Branch { left, right, .. } = node.as_ref() {
1279 collect_patricia_node_pages(domain, left, pages)?;
1280 collect_patricia_node_pages(domain, right, pages)?;
1281 }
1282 let page = patricia_node_page(domain, node)?;
1283 let page_id = page.page_id.clone();
1284 pages.insert(page_id.clone(), page);
1285 if pages.len() > MAX_STATE_DELTA_PAGES {
1286 return Err(invalid_patricia_page(
1287 "Patricia page graph exceeds the bounded page count",
1288 ));
1289 }
1290 Ok(page_id)
1291}
1292
1293fn patricia_node_page(domain: &str, node: &Arc<PatriciaNode>) -> Result<StatePageBlob, CanwuError> {
1294 let stored = match node.as_ref() {
1295 PatriciaNode::Leaf {
1296 key_hash,
1297 entries,
1298 entry_count,
1299 structural_bytes,
1300 ..
1301 } => StoredPatriciaPage::Leaf {
1302 domain: domain.to_owned(),
1303 key_hash: key_hash.to_hex(),
1304 entries: entries
1305 .iter()
1306 .map(|entry| StoredPatriciaEntry {
1307 key_hash: key_hash.to_hex(),
1308 key_bytes: entry.key_bytes.clone(),
1309 value_hash: entry.value_hash.to_hex(),
1310 value_bytes: entry.value_bytes.clone(),
1311 })
1312 .collect(),
1313 entry_count: u64::try_from(*entry_count)
1314 .map_err(|_| invalid_patricia_page("Patricia entry count exceeds u64"))?,
1315 structural_bytes: u64::try_from(*structural_bytes)
1316 .map_err(|_| invalid_patricia_page("Patricia byte count exceeds u64"))?,
1317 },
1318 PatriciaNode::Branch {
1319 bit,
1320 left,
1321 right,
1322 entry_count,
1323 structural_bytes,
1324 ..
1325 } => StoredPatriciaPage::Branch {
1326 domain: domain.to_owned(),
1327 bit: *bit,
1328 left_page: left.hash().to_hex(),
1329 right_page: right.hash().to_hex(),
1330 entry_count: u64::try_from(*entry_count)
1331 .map_err(|_| invalid_patricia_page("Patricia entry count exceeds u64"))?,
1332 structural_bytes: u64::try_from(*structural_bytes)
1333 .map_err(|_| invalid_patricia_page("Patricia byte count exceeds u64"))?,
1334 },
1335 };
1336 let bytes = serde_json::to_vec(&stored).map_err(|error| {
1337 invalid_patricia_page(format!("cannot encode canonical Patricia page: {error}"))
1338 })?;
1339 StatePageBlob::new(bytes)
1340}
1341
1342fn load_patricia_root(
1343 domain: &'static str,
1344 page_id: Option<&str>,
1345 provider: &dyn StatePageProvider,
1346 cache: &mut BTreeMap<String, Arc<PatriciaNode>>,
1347) -> Result<PersistentPatricia, CanwuError> {
1348 let Some(page_id) = page_id else {
1349 return Ok(PersistentPatricia::new(domain));
1350 };
1351 let mut active = BTreeSet::new();
1352 let root = load_patricia_node(domain, page_id, provider, cache, &mut active, None, 0)?;
1353 Ok(PersistentPatricia {
1354 domain,
1355 root: Some(root),
1356 })
1357}
1358
1359fn load_patricia_node(
1360 domain: &'static str,
1361 page_id: &str,
1362 provider: &dyn StatePageProvider,
1363 cache: &mut BTreeMap<String, Arc<PatriciaNode>>,
1364 active: &mut BTreeSet<String>,
1365 parent_bit: Option<u16>,
1366 depth: usize,
1367) -> Result<Arc<PatriciaNode>, CanwuError> {
1368 if depth > 256 || cache.len() >= MAX_STATE_DELTA_PAGES {
1369 return Err(invalid_patricia_page(
1370 "Patricia page graph exceeds its depth or page budget",
1371 ));
1372 }
1373 if let Some(node) = cache.get(page_id) {
1374 validate_child_bit(node, parent_bit)?;
1375 return Ok(Arc::clone(node));
1376 }
1377 if !active.insert(page_id.to_owned()) {
1378 return Err(invalid_patricia_page(
1379 "Patricia page graph contains a cycle",
1380 ));
1381 }
1382 let page = provider.load_state_page(page_id)?.ok_or_else(|| {
1383 CanwuError::new(
1384 ErrorCode::StatePageUnavailable,
1385 format!("state page {page_id} is unavailable"),
1386 )
1387 })?;
1388 page.validate()?;
1389 if page.page_id != page_id {
1390 return Err(invalid_patricia_page(
1391 "state-page provider returned the wrong content address",
1392 ));
1393 }
1394 let stored: StoredPatriciaPage = serde_json::from_slice(&page.bytes)
1395 .map_err(|error| invalid_patricia_page(format!("invalid Patricia state page: {error}")))?;
1396 let node = match stored {
1397 StoredPatriciaPage::Leaf {
1398 domain: stored_domain,
1399 key_hash,
1400 entries,
1401 entry_count,
1402 structural_bytes,
1403 } => {
1404 if stored_domain != domain || entries.is_empty() {
1405 return Err(invalid_patricia_page(
1406 "Patricia leaf domain or collision membership is invalid",
1407 ));
1408 }
1409 let mut decoded = Vec::with_capacity(entries.len());
1410 for entry in entries {
1411 if entry.key_hash != key_hash
1412 || entry.key_hash
1413 != canonical_byte_hash(&format!("{domain}.key"), &entry.key_bytes)
1414 || entry.value_hash
1415 != canonical_byte_hash(&format!("{domain}.value"), &entry.value_bytes)
1416 {
1417 return Err(invalid_patricia_page(
1418 "Patricia leaf entry hash is inconsistent",
1419 ));
1420 }
1421 decoded.push(PatriciaEntry {
1422 key_bytes: entry.key_bytes,
1423 value_hash: CompactHash::parse(&entry.value_hash)?,
1424 value_bytes: entry.value_bytes,
1425 });
1426 }
1427 if decoded
1428 .windows(2)
1429 .any(|pair| pair[0].key_bytes >= pair[1].key_bytes)
1430 {
1431 return Err(invalid_patricia_page(
1432 "Patricia collision entries are not strictly ordered",
1433 ));
1434 }
1435 let rebuilt = patricia_leaf(domain, CompactHash::parse(&key_hash)?, decoded);
1436 if rebuilt.hash().to_hex() != page_id
1437 || rebuilt.entry_count() as u64 != entry_count
1438 || rebuilt.structural_bytes() as u64 != structural_bytes
1439 {
1440 return Err(invalid_patricia_page(
1441 "Patricia leaf metadata or commitment is inconsistent",
1442 ));
1443 }
1444 rebuilt
1445 }
1446 StoredPatriciaPage::Branch {
1447 domain: stored_domain,
1448 bit,
1449 left_page,
1450 right_page,
1451 entry_count,
1452 structural_bytes,
1453 } => {
1454 if stored_domain != domain || left_page == right_page || usize::from(bit) >= 256 {
1455 return Err(invalid_patricia_page(
1456 "Patricia branch domain, bit, or child identity is invalid",
1457 ));
1458 }
1459 if parent_bit.is_some_and(|parent| bit <= parent) {
1460 return Err(invalid_patricia_page(
1461 "Patricia discriminating bits must increase down a path",
1462 ));
1463 }
1464 let left = load_patricia_node(
1465 domain,
1466 &left_page,
1467 provider,
1468 cache,
1469 active,
1470 Some(bit),
1471 depth + 1,
1472 )?;
1473 let right = load_patricia_node(
1474 domain,
1475 &right_page,
1476 provider,
1477 cache,
1478 active,
1479 Some(bit),
1480 depth + 1,
1481 )?;
1482 if bit_at(left.sample_key_hash(), usize::from(bit)) != 0
1483 || bit_at(right.sample_key_hash(), usize::from(bit)) != 1
1484 {
1485 return Err(invalid_patricia_page(
1486 "Patricia children occupy the wrong discriminating slots",
1487 ));
1488 }
1489 let rebuilt = patricia_branch(domain, usize::from(bit), left, right);
1490 if rebuilt.hash().to_hex() != page_id
1491 || rebuilt.entry_count() as u64 != entry_count
1492 || rebuilt.structural_bytes() as u64 != structural_bytes
1493 {
1494 return Err(invalid_patricia_page(
1495 "Patricia branch metadata or commitment is inconsistent",
1496 ));
1497 }
1498 rebuilt
1499 }
1500 };
1501 validate_child_bit(&node, parent_bit)?;
1502 active.remove(page_id);
1503 cache.insert(page_id.to_owned(), Arc::clone(&node));
1504 Ok(node)
1505}
1506
1507fn validate_child_bit(node: &PatriciaNode, parent_bit: Option<u16>) -> Result<(), CanwuError> {
1508 if let (Some(parent), PatriciaNode::Branch { bit: child_bit, .. }) = (parent_bit, node)
1509 && *child_bit <= parent
1510 {
1511 return Err(invalid_patricia_page(
1512 "Patricia discriminating bits must increase down a path",
1513 ));
1514 }
1515 Ok(())
1516}
1517
1518fn collect_primary_records(
1519 node: &PatriciaNode,
1520 records: &mut BTreeMap<DomainRecordRef, DomainRecord>,
1521) -> Result<(), CanwuError> {
1522 match node {
1523 PatriciaNode::Leaf { entries, .. } => {
1524 for entry in entries.iter() {
1525 let reference: DomainRecordRef =
1526 serde_json::from_slice(&entry.key_bytes).map_err(|error| {
1527 invalid_patricia_page(format!(
1528 "cannot decode domain-record key page: {error}"
1529 ))
1530 })?;
1531 let record: DomainRecord =
1532 serde_json::from_slice(&entry.value_bytes).map_err(|error| {
1533 invalid_patricia_page(format!(
1534 "cannot decode domain-record value page: {error}"
1535 ))
1536 })?;
1537 if record.reference != reference || records.insert(reference, record).is_some() {
1538 return Err(invalid_patricia_page(
1539 "domain-record primary page contains a duplicate or mismatched key",
1540 ));
1541 }
1542 }
1543 }
1544 PatriciaNode::Branch { left, right, .. } => {
1545 collect_primary_records(left, records)?;
1546 collect_primary_records(right, records)?;
1547 }
1548 }
1549 Ok(())
1550}
1551
1552fn invalid_patricia_page(message: impl Into<String>) -> CanwuError {
1553 CanwuError::new(ErrorCode::InvalidArchive, message)
1554}
1555
1556#[derive(Clone, Copy, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
1557#[serde(rename_all = "snake_case")]
1558pub enum DomainRecordClass {
1559 Entity,
1560 Record,
1561}
1562
1563#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
1564#[serde(tag = "type", content = "kind", rename_all = "snake_case")]
1565pub enum DomainReferenceTargetKind {
1566 Core(CoreEntityKind),
1567 Domain(DomainRecordKind),
1568 AnyEntity,
1569}
1570
1571impl DomainReferenceTargetKind {
1572 #[must_use]
1573 pub fn for_domain<T: DomainRecordType>() -> Self {
1574 Self::Domain(DomainRecordKind::for_type::<T>())
1575 }
1576}
1577
1578#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
1579pub struct DomainReferenceSchema {
1580 pub role: String,
1581 pub targets: Vec<DomainReferenceTargetKind>,
1582 pub required: bool,
1583 pub multiple: bool,
1584 pub allow_retired: bool,
1585}
1586
1587#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq, Serialize)]
1588#[serde(rename_all = "snake_case")]
1589pub enum DomainRecordMutationPolicy {
1590 #[default]
1591 Versioned,
1592 CreateOnly,
1593}
1594
1595#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
1596pub struct DomainRecordSchema {
1597 pub kind: DomainRecordKind,
1598 pub class: DomainRecordClass,
1599 #[serde(default)]
1600 pub holder_policy: KnowledgeHolderPolicy,
1601 #[serde(default)]
1602 pub mutation_policy: DomainRecordMutationPolicy,
1603 pub payload_schema: PayloadSchema,
1604 pub references: Vec<DomainReferenceSchema>,
1605}
1606
1607impl DomainRecordSchema {
1608 #[must_use]
1609 pub fn new(kind: DomainRecordKind, class: DomainRecordClass) -> Self {
1610 Self {
1611 kind,
1612 class,
1613 holder_policy: KnowledgeHolderPolicy::Disallowed,
1614 mutation_policy: DomainRecordMutationPolicy::Versioned,
1615 payload_schema: PayloadSchema::Any,
1616 references: Vec::new(),
1617 }
1618 }
1619
1620 #[must_use]
1621 pub fn for_type<T: DomainRecordType>() -> Self {
1622 let class = if T::Class::IS_ENTITY {
1623 DomainRecordClass::Entity
1624 } else {
1625 DomainRecordClass::Record
1626 };
1627 Self::new(DomainRecordKind::for_type::<T>(), class)
1628 }
1629
1630 #[must_use]
1631 pub fn for_entity<T: DomainEntityType>() -> Self {
1632 Self::for_type::<T>()
1633 }
1634
1635 #[must_use]
1636 pub fn for_record<T: DomainValueType>() -> Self {
1637 Self::for_type::<T>()
1638 }
1639
1640 #[must_use]
1641 pub fn state_key(&self) -> StateKey {
1642 record_state_key(&self.kind)
1643 }
1644
1645 pub(crate) fn canonicalize(&mut self) {
1646 for reference in &mut self.references {
1647 reference.targets.sort();
1648 reference.targets.dedup();
1649 }
1650 self.references
1651 .sort_by(|left, right| left.role.cmp(&right.role));
1652 }
1653
1654 pub(crate) fn validate(&self) -> Result<(), CanwuError> {
1655 validate_kind(&self.kind)?;
1656 if self.holder_policy == KnowledgeHolderPolicy::Allowed
1657 && self.class != DomainRecordClass::Entity
1658 {
1659 return invalid_record("only domain entity schemas may allow knowledge holders");
1660 }
1661 if let PayloadSchema::Object { properties, .. } = &self.payload_schema
1662 && properties.keys().any(|name| !canonical_text(name))
1663 {
1664 return invalid_record(
1665 "record payload-schema property names must be non-empty and canonical",
1666 );
1667 }
1668 if self
1669 .references
1670 .windows(2)
1671 .any(|pair| pair[0].role >= pair[1].role)
1672 {
1673 return invalid_record("record-schema reference roles must be unique and sorted");
1674 }
1675 for reference in &self.references {
1676 if !canonical_text(&reference.role)
1677 || reference.targets.is_empty()
1678 || reference.targets.windows(2).any(|pair| pair[0] >= pair[1])
1679 {
1680 return invalid_record(
1681 "record-schema references require canonical roles and unique sorted targets",
1682 );
1683 }
1684 for target in &reference.targets {
1685 if let DomainReferenceTargetKind::Domain(kind) = target {
1686 validate_kind(kind)?;
1687 }
1688 }
1689 }
1690 Ok(())
1691 }
1692}
1693
1694#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
1695#[serde(tag = "type", content = "reference", rename_all = "snake_case")]
1696pub enum DomainReferenceTarget {
1697 Core(EntityRef),
1698 Domain(DomainRecordRef),
1699}
1700
1701impl DomainReferenceTarget {
1702 #[must_use]
1703 pub fn from_typed<T: DomainRecordType>(reference: TypedDomainRecordRef<T>) -> Self {
1704 Self::Domain(reference.into_untyped())
1705 }
1706}
1707
1708#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
1709pub struct DomainReference {
1710 pub role: String,
1711 pub target: DomainReferenceTarget,
1712}
1713
1714impl DomainReference {
1715 #[must_use]
1716 pub fn from_typed<T: DomainRecordType>(
1717 role: impl Into<String>,
1718 reference: TypedDomainRecordRef<T>,
1719 ) -> Self {
1720 Self {
1721 role: role.into(),
1722 target: DomainReferenceTarget::from_typed(reference),
1723 }
1724 }
1725}
1726
1727#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
1728pub struct DomainRecordDraft {
1729 pub reference: DomainRecordRef,
1730 pub payload: Value,
1731 pub references: Vec<DomainReference>,
1732}
1733
1734impl DomainRecordDraft {
1735 #[must_use]
1736 pub fn new(reference: DomainRecordRef, payload: Value) -> Self {
1737 Self {
1738 reference,
1739 payload,
1740 references: Vec::new(),
1741 }
1742 }
1743
1744 pub fn from_typed<T: DomainRecordType>(
1745 reference: TypedDomainRecordRef<T>,
1746 payload: &T::Payload,
1747 ) -> Result<Self, CanwuError>
1748 where
1749 T::Payload: Serialize,
1750 {
1751 let payload = serde_json::to_value(payload).map_err(|error| {
1752 CanwuError::new(
1753 ErrorCode::InvalidDomainRecord,
1754 format!(
1755 "typed domain payload for {} could not be encoded: {error}",
1756 DomainRecordKind::for_type::<T>()
1757 ),
1758 )
1759 })?;
1760 Ok(Self::new(reference.into_untyped(), payload))
1761 }
1762
1763 pub(crate) fn canonicalize(&mut self) {
1764 self.references.sort();
1765 self.references.dedup();
1766 }
1767}
1768
1769#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
1770#[serde(tag = "state", rename_all = "snake_case")]
1771pub enum DomainRecordLifecycle {
1772 Active,
1773 Retired {
1774 at: SimTime,
1775 successor: Option<DomainRecordRef>,
1776 },
1777 Deleted {
1778 at: SimTime,
1779 },
1780}
1781
1782#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
1783pub struct DomainRecord {
1784 pub reference: DomainRecordRef,
1785 pub owner: String,
1786 pub class: DomainRecordClass,
1787 pub version: u64,
1788 pub lifecycle: DomainRecordLifecycle,
1789 pub payload: Value,
1790 pub references: Vec<DomainReference>,
1791}
1792
1793impl DomainRecord {
1794 #[must_use]
1795 pub const fn is_deleted(&self) -> bool {
1796 matches!(self.lifecycle, DomainRecordLifecycle::Deleted { .. })
1797 }
1798
1799 #[must_use]
1800 pub const fn is_active(&self) -> bool {
1801 matches!(self.lifecycle, DomainRecordLifecycle::Active)
1802 }
1803
1804 #[must_use]
1805 pub fn typed_reference<T: DomainRecordType>(&self) -> Option<TypedDomainRecordRef<T>> {
1806 TypedDomainRecordRef::from_untyped(self.reference.clone()).ok()
1807 }
1808
1809 pub fn decode_payload<T: DomainRecordType>(&self) -> Result<T::Payload, CanwuError>
1810 where
1811 T::Payload: DeserializeOwned,
1812 {
1813 if !self.reference.kind.matches_type::<T>() {
1814 return Err(CanwuError::new(
1815 ErrorCode::InvalidDomainRecord,
1816 format!(
1817 "domain record {} cannot be decoded as kind {}",
1818 self.reference,
1819 DomainRecordKind::for_type::<T>()
1820 ),
1821 ));
1822 }
1823 T::Payload::deserialize(&self.payload).map_err(|error| {
1824 CanwuError::new(
1825 ErrorCode::InvalidDomainRecord,
1826 format!(
1827 "domain record {} has an incompatible typed payload: {error}",
1828 self.reference
1829 ),
1830 )
1831 })
1832 }
1833}
1834
1835#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
1836#[serde(tag = "operation", rename_all = "snake_case")]
1837pub enum DomainRecordMutation {
1838 Create {
1839 record: DomainRecordDraft,
1840 },
1841 Update {
1842 record: DomainRecordDraft,
1843 expected_version: u64,
1844 },
1845 Retire {
1846 record: DomainRecordRef,
1847 expected_version: u64,
1848 successor: Option<DomainRecordRef>,
1849 },
1850 Delete {
1851 record: DomainRecordRef,
1852 expected_version: u64,
1853 },
1854}
1855
1856impl DomainRecordMutation {
1857 #[must_use]
1858 pub const fn target(&self) -> &DomainRecordRef {
1859 match self {
1860 Self::Create { record } | Self::Update { record, .. } => &record.reference,
1861 Self::Retire { record, .. } | Self::Delete { record, .. } => record,
1862 }
1863 }
1864}
1865
1866#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
1867#[serde(rename_all = "snake_case")]
1868pub enum DomainRecordOperation {
1869 Created,
1870 Updated,
1871 Retired,
1872 Deleted,
1873}
1874
1875impl DomainRecordOperation {
1876 pub(crate) const fn event_type(self) -> &'static str {
1877 match self {
1878 Self::Created => "domain_record_created",
1879 Self::Updated => "domain_record_updated",
1880 Self::Retired => "domain_record_retired",
1881 Self::Deleted => "domain_record_deleted",
1882 }
1883 }
1884}
1885
1886#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
1887pub struct DomainRecordChange {
1888 pub plugin: String,
1889 pub system: String,
1890 pub operation: DomainRecordOperation,
1891 pub previous: Option<DomainRecord>,
1892 pub current: DomainRecord,
1893 pub visibility: StateVisibility,
1894 pub summary: String,
1895}
1896
1897pub(crate) type DomainRecordSchemas = BTreeMap<DomainRecordKind, (String, DomainRecordSchema)>;
1898
1899pub(crate) trait DomainRecordRead {
1900 fn get(&self, reference: &DomainRecordRef) -> Option<&DomainRecord>;
1901 fn iter(&self) -> Box<dyn Iterator<Item = (&DomainRecordRef, &DomainRecord)> + '_>;
1902 fn range_from(
1903 &self,
1904 lower: DomainRecordRef,
1905 excluded: bool,
1906 ) -> Box<dyn Iterator<Item = (&DomainRecordRef, &DomainRecord)> + '_>;
1907
1908 fn contains_key(&self, reference: &DomainRecordRef) -> bool {
1909 self.get(reference).is_some()
1910 }
1911}
1912
1913trait DomainRecordWrite: DomainRecordRead {
1914 fn insert(&mut self, reference: DomainRecordRef, record: DomainRecord);
1915}
1916
1917impl DomainRecordRead for BTreeMap<DomainRecordRef, DomainRecord> {
1918 fn get(&self, reference: &DomainRecordRef) -> Option<&DomainRecord> {
1919 BTreeMap::get(self, reference)
1920 }
1921
1922 fn iter(&self) -> Box<dyn Iterator<Item = (&DomainRecordRef, &DomainRecord)> + '_> {
1923 Box::new(BTreeMap::iter(self))
1924 }
1925
1926 fn range_from(
1927 &self,
1928 lower: DomainRecordRef,
1929 excluded: bool,
1930 ) -> Box<dyn Iterator<Item = (&DomainRecordRef, &DomainRecord)> + '_> {
1931 use std::ops::Bound::{Excluded, Included, Unbounded};
1932 let lower = if excluded {
1933 Excluded(lower)
1934 } else {
1935 Included(lower)
1936 };
1937 Box::new(self.range((lower, Unbounded)))
1938 }
1939}
1940
1941impl DomainRecordWrite for BTreeMap<DomainRecordRef, DomainRecord> {
1942 fn insert(&mut self, reference: DomainRecordRef, record: DomainRecord) {
1943 let _ = BTreeMap::insert(self, reference, record);
1944 }
1945}
1946
1947impl DomainRecordRead for PersistentDomainRecordStore {
1948 fn get(&self, reference: &DomainRecordRef) -> Option<&DomainRecord> {
1949 self.get(reference)
1950 }
1951
1952 fn iter(&self) -> Box<dyn Iterator<Item = (&DomainRecordRef, &DomainRecord)> + '_> {
1953 Box::new(self.iter())
1954 }
1955
1956 fn range_from(
1957 &self,
1958 lower: DomainRecordRef,
1959 excluded: bool,
1960 ) -> Box<dyn Iterator<Item = (&DomainRecordRef, &DomainRecord)> + '_> {
1961 let (start_page, start_key) = self.record_key_start(&lower, excluded);
1962 Box::new(
1963 self.record_key_pages
1964 .iter()
1965 .enumerate()
1966 .skip(start_page)
1967 .flat_map(move |(page_index, page)| {
1968 page.iter()
1969 .skip(if page_index == start_page {
1970 start_key
1971 } else {
1972 0
1973 })
1974 .filter_map(|reference| {
1975 self.records
1976 .get(reference)
1977 .map(|record| (reference, record.as_ref()))
1978 })
1979 }),
1980 )
1981 }
1982}
1983
1984impl DomainRecordWrite for PersistentDomainRecordStore {
1985 fn insert(&mut self, reference: DomainRecordRef, record: DomainRecord) {
1986 self.insert_record(reference, record)
1987 .expect("validated domain records must have canonical Patricia encodings");
1988 }
1989}
1990
1991pub(crate) struct DomainMutationRequest<'a> {
1992 pub plugin: &'a str,
1993 pub system: &'a str,
1994 pub visibility: StateVisibility,
1995 pub mutation: &'a DomainRecordMutation,
1996 pub summary: &'a str,
1997}
1998
1999pub(crate) fn record_state_key(kind: &DomainRecordKind) -> StateKey {
2000 StateKey::new(&kind.namespace, &kind.name)
2001}
2002
2003pub(crate) fn validate_initial_records(
2004 records: &[DomainRecord],
2005 now: SimTime,
2006 core_exists: &dyn Fn(&EntityRef) -> bool,
2007) -> Result<(), CanwuError> {
2008 let mut store = BTreeMap::new();
2009 for record in records {
2010 validate_record_shape(record, now)?;
2011 if store
2012 .insert(record.reference.clone(), record.clone())
2013 .is_some()
2014 {
2015 return Err(CanwuError::new(
2016 ErrorCode::DuplicateDomainRecord,
2017 format!("domain record {} is duplicated", record.reference),
2018 ));
2019 }
2020 }
2021 for record in store.values().filter(|record| !record.is_deleted()) {
2022 validate_reference_targets_basic(record, &store, core_exists)?;
2023 validate_successor(record, &store)?;
2024 }
2025 validate_successor_graph(&store)?;
2026 Ok(())
2027}
2028
2029pub(crate) fn validate_record_store(
2030 records: &impl DomainRecordRead,
2031 schemas: &DomainRecordSchemas,
2032 now: SimTime,
2033 core_exists: &dyn Fn(&EntityRef) -> bool,
2034) -> Result<(), CanwuError> {
2035 for (reference, record) in records.iter() {
2036 if reference != &record.reference {
2037 return invalid_record("domain record map key disagrees with its stable reference");
2038 }
2039 validate_record_shape(record, now)?;
2040 let Some((owner, schema)) = schemas.get(&record.reference.kind) else {
2041 return Err(CanwuError::new(
2042 ErrorCode::PluginNotActive,
2043 format!(
2044 "domain record kind {} has no active schema owner",
2045 record.reference.kind
2046 ),
2047 ));
2048 };
2049 if &record.owner != owner || record.class != schema.class {
2050 return invalid_record(format!(
2051 "domain record {} disagrees with its schema owner or class",
2052 record.reference
2053 ));
2054 }
2055 schema.payload_schema.validate(&record.payload)?;
2056 validate_record_references(record, schema, records, core_exists)?;
2057 validate_successor(record, records)?;
2058 }
2059 validate_successor_graph(records)?;
2060 Ok(())
2061}
2062
2063pub(crate) fn validate_records_for_owner(
2064 records: &impl DomainRecordRead,
2065 schemas: &DomainRecordSchemas,
2066 owner: &str,
2067 now: SimTime,
2068 core_exists: &dyn Fn(&EntityRef) -> bool,
2069) -> Result<(), CanwuError> {
2070 for (_, record) in records.iter().filter(|(_, record)| record.owner == owner) {
2071 validate_record_shape(record, now)?;
2072 let Some((schema_owner, schema)) = schemas.get(&record.reference.kind) else {
2073 return invalid_record(format!(
2074 "plugin {owner} did not register schema for owned record kind {}",
2075 record.reference.kind
2076 ));
2077 };
2078 if schema_owner != owner || record.class != schema.class {
2079 return invalid_record(format!(
2080 "domain record {} disagrees with its registered owner or class",
2081 record.reference
2082 ));
2083 }
2084 schema.payload_schema.validate(&record.payload)?;
2085 validate_record_references(record, schema, records, core_exists)?;
2086 validate_successor(record, records)?;
2087 }
2088 validate_successor_graph(records)?;
2089 Ok(())
2090}
2091
2092pub(crate) fn apply_mutation_bundle(
2093 records: &BTreeMap<DomainRecordRef, DomainRecord>,
2094 schemas: &DomainRecordSchemas,
2095 now: SimTime,
2096 core_exists: &dyn Fn(&EntityRef) -> bool,
2097 requests: Vec<DomainMutationRequest<'_>>,
2098) -> Result<
2099 (
2100 BTreeMap<DomainRecordRef, DomainRecord>,
2101 Vec<DomainRecordChange>,
2102 ),
2103 CanwuError,
2104> {
2105 let mut next = records.clone();
2106 let changes =
2107 apply_mutation_bundle_in_place(records, &mut next, schemas, now, core_exists, requests)?;
2108 validate_record_store(&next, schemas, now, core_exists)?;
2109 Ok((next, changes))
2110}
2111
2112pub(crate) fn apply_mutation_bundle_cow(
2116 records: &PersistentDomainRecordStore,
2117 schemas: &DomainRecordSchemas,
2118 now: SimTime,
2119 core_exists: &dyn Fn(&EntityRef) -> bool,
2120 requests: Vec<DomainMutationRequest<'_>>,
2121) -> Result<(PersistentDomainRecordStore, Vec<DomainRecordChange>), CanwuError> {
2122 let mut next = records.clone();
2123 let changes =
2124 apply_mutation_bundle_in_place(records, &mut next, schemas, now, core_exists, requests)?;
2125 validate_affected_record_store(&next, schemas, now, core_exists, &changes)?;
2126 Ok((next, changes))
2127}
2128
2129pub(crate) fn apply_mutation_bundle_cow_with_overlay(
2137 records: &PersistentDomainRecordStore,
2138 overlay: &BTreeMap<DomainRecordRef, DomainRecord>,
2139 schemas: &DomainRecordSchemas,
2140 now: SimTime,
2141 core_exists: &dyn Fn(&EntityRef) -> bool,
2142 requests: Vec<DomainMutationRequest<'_>>,
2143) -> Result<(PersistentDomainRecordStore, Vec<DomainRecordChange>), CanwuError> {
2144 let mut base = records.clone();
2145 for (reference, record) in overlay {
2146 base.insert_record(reference.clone(), record.clone())?;
2147 }
2148 apply_mutation_bundle_cow(&base, schemas, now, core_exists, requests)
2149}
2150
2151fn apply_mutation_bundle_in_place<R, N>(
2152 records: &R,
2153 next: &mut N,
2154 schemas: &DomainRecordSchemas,
2155 now: SimTime,
2156 _core_exists: &dyn Fn(&EntityRef) -> bool,
2157 mut requests: Vec<DomainMutationRequest<'_>>,
2158) -> Result<Vec<DomainRecordChange>, CanwuError>
2159where
2160 R: DomainRecordRead,
2161 N: DomainRecordWrite,
2162{
2163 requests.sort_by(|left, right| left.mutation.target().cmp(right.mutation.target()));
2164 if requests
2165 .windows(2)
2166 .any(|pair| pair[0].mutation.target() == pair[1].mutation.target())
2167 {
2168 return invalid_record("a boundary cannot mutate the same domain record twice");
2169 }
2170 let created_records = requests
2171 .iter()
2172 .filter_map(|request| match request.mutation {
2173 DomainRecordMutation::Create { record } => Some(record.reference.clone()),
2174 DomainRecordMutation::Update { .. }
2175 | DomainRecordMutation::Retire { .. }
2176 | DomainRecordMutation::Delete { .. } => None,
2177 })
2178 .collect::<BTreeSet<_>>();
2179
2180 let mut changes = Vec::with_capacity(requests.len());
2181 for request in requests {
2182 if !canonical_text(request.summary) {
2183 return invalid_record("domain record mutations require a canonical summary");
2184 }
2185 let target = request.mutation.target();
2186 validate_reference(target)?;
2187 let Some((owner, schema)) = schemas.get(&target.kind) else {
2188 return invalid_record(format!(
2189 "domain record kind {} has no registered schema",
2190 target.kind
2191 ));
2192 };
2193 if owner != request.plugin {
2194 return Err(CanwuError::new(
2195 ErrorCode::UndeclaredStateWrite,
2196 format!(
2197 "plugin {} cannot mutate domain record kind {} owned by {owner}",
2198 request.plugin, target.kind
2199 ),
2200 ));
2201 }
2202 if schema.mutation_policy == DomainRecordMutationPolicy::CreateOnly
2203 && !matches!(request.mutation, DomainRecordMutation::Create { .. })
2204 {
2205 return invalid_record(format!("domain record kind {} is create-only", target.kind));
2206 }
2207
2208 let (operation, previous, current) = match request.mutation {
2209 DomainRecordMutation::Create { record } => {
2210 let mut record = record.clone();
2211 record.canonicalize();
2212 if next.contains_key(&record.reference) {
2213 return Err(CanwuError::new(
2214 ErrorCode::DuplicateDomainRecord,
2215 format!("domain record {} already exists", record.reference),
2216 ));
2217 }
2218 let current = DomainRecord {
2219 reference: record.reference,
2220 owner: owner.clone(),
2221 class: schema.class,
2222 version: 1,
2223 lifecycle: DomainRecordLifecycle::Active,
2224 payload: record.payload,
2225 references: record.references,
2226 };
2227 next.insert(current.reference.clone(), current.clone());
2228 (DomainRecordOperation::Created, None, current)
2229 }
2230 DomainRecordMutation::Update {
2231 record,
2232 expected_version,
2233 } => {
2234 let mut draft = record.clone();
2235 draft.canonicalize();
2236 let previous = require_mutable_record(next, target, *expected_version)?.clone();
2237 let version = next_record_version(previous.version)?;
2238 let current = DomainRecord {
2239 reference: draft.reference,
2240 owner: previous.owner.clone(),
2241 class: previous.class,
2242 version,
2243 lifecycle: DomainRecordLifecycle::Active,
2244 payload: draft.payload,
2245 references: draft.references,
2246 };
2247 next.insert(current.reference.clone(), current.clone());
2248 (DomainRecordOperation::Updated, Some(previous), current)
2249 }
2250 DomainRecordMutation::Retire {
2251 record,
2252 expected_version,
2253 successor,
2254 } => {
2255 let previous = require_mutable_record(next, record, *expected_version)?.clone();
2256 validate_new_successor(record, successor.as_ref(), records, &created_records)?;
2257 let mut current = previous.clone();
2258 current.version = next_record_version(previous.version)?;
2259 current.lifecycle = DomainRecordLifecycle::Retired {
2260 at: now,
2261 successor: successor.clone(),
2262 };
2263 next.insert(record.clone(), current.clone());
2264 (DomainRecordOperation::Retired, Some(previous), current)
2265 }
2266 DomainRecordMutation::Delete {
2267 record,
2268 expected_version,
2269 } => {
2270 let previous = next.get(record).ok_or_else(|| {
2271 CanwuError::new(
2272 ErrorCode::DomainRecordNotFound,
2273 format!("domain record {record} was not found"),
2274 )
2275 })?;
2276 if previous.version != *expected_version {
2277 return Err(version_conflict(
2278 record,
2279 *expected_version,
2280 previous.version,
2281 ));
2282 }
2283 if !matches!(previous.lifecycle, DomainRecordLifecycle::Retired { .. }) {
2284 return invalid_record("domain records must be retired before deletion");
2285 }
2286 let previous = previous.clone();
2287 let mut current = previous.clone();
2288 current.version = next_record_version(previous.version)?;
2289 current.lifecycle = DomainRecordLifecycle::Deleted { at: now };
2290 current.references.clear();
2291 next.insert(record.clone(), current.clone());
2292 (DomainRecordOperation::Deleted, Some(previous), current)
2293 }
2294 };
2295 changes.push(DomainRecordChange {
2296 plugin: request.plugin.to_owned(),
2297 system: request.system.to_owned(),
2298 operation,
2299 previous,
2300 current,
2301 visibility: request.visibility,
2302 summary: request.summary.to_owned(),
2303 });
2304 }
2305 Ok(changes)
2306}
2307
2308fn validate_affected_record_store(
2309 records: &PersistentDomainRecordStore,
2310 schemas: &DomainRecordSchemas,
2311 now: SimTime,
2312 core_exists: &dyn Fn(&EntityRef) -> bool,
2313 changes: &[DomainRecordChange],
2314) -> Result<(), CanwuError> {
2315 let mut affected = changes
2316 .iter()
2317 .map(|change| change.current.reference.clone())
2318 .collect::<BTreeSet<_>>();
2319 let mut pending = affected.iter().cloned().collect::<Vec<_>>();
2320 while let Some(reference) = pending.pop() {
2321 for related in records
2322 .reverse_lookup
2323 .get(&reference)
2324 .into_iter()
2325 .flatten()
2326 .chain(
2327 records
2328 .predecessor_lookup
2329 .get(&reference)
2330 .into_iter()
2331 .flatten(),
2332 )
2333 {
2334 if affected.insert(related.clone()) {
2335 if affected.len() > MAX_AFFECTED_RECORD_VALIDATION_CLOSURE {
2336 return invalid_record(
2337 "domain record mutation exceeded the affected-reference closure budget",
2338 );
2339 }
2340 pending.push(related.clone());
2341 }
2342 }
2343 }
2344 for reference in &affected {
2345 let record = records.get(reference).ok_or_else(|| {
2346 CanwuError::new(
2347 ErrorCode::InvalidDomainRecord,
2348 "affected domain record disappeared during validation",
2349 )
2350 })?;
2351 validate_record_shape(record, now)?;
2352 let Some((owner, schema)) = schemas.get(&record.reference.kind) else {
2353 return Err(CanwuError::new(
2354 ErrorCode::PluginNotActive,
2355 format!(
2356 "domain record kind {} has no active schema owner",
2357 record.reference.kind
2358 ),
2359 ));
2360 };
2361 if &record.owner != owner || record.class != schema.class {
2362 return invalid_record(format!(
2363 "domain record {} disagrees with its schema owner or class",
2364 record.reference
2365 ));
2366 }
2367 schema.payload_schema.validate(&record.payload)?;
2368 validate_record_references(record, schema, records, core_exists)?;
2369 validate_successor(record, records)?;
2370 }
2371 for change in changes {
2372 validate_successor_chain_from(records, &change.current.reference)?;
2373 }
2374 Ok(())
2375}
2376
2377fn validate_successor_chain_from(
2378 records: &PersistentDomainRecordStore,
2379 start: &DomainRecordRef,
2380) -> Result<(), CanwuError> {
2381 let mut visited = BTreeSet::new();
2382 let mut current = start;
2383 loop {
2384 if !visited.insert(current.clone()) {
2385 return invalid_record("domain record successor chains cannot contain cycles");
2386 }
2387 if visited.len() > MAX_AFFECTED_RECORD_VALIDATION_CLOSURE {
2388 return invalid_record(
2389 "domain record successor validation exceeded the affected-closure budget",
2390 );
2391 }
2392 let Some(successor) = records.successor_lookup.get(current) else {
2393 return Ok(());
2394 };
2395 current = successor;
2396 }
2397}
2398
2399pub(crate) fn domain_entity_exists(
2400 records: &impl DomainRecordRead,
2401 reference: &DomainRecordRef,
2402) -> bool {
2403 records
2404 .get(reference)
2405 .is_some_and(|record| record.class == DomainRecordClass::Entity && !record.is_deleted())
2406}
2407
2408pub(crate) fn mutation_from_change(change: &DomainRecordChange) -> DomainRecordMutation {
2409 match change.operation {
2410 DomainRecordOperation::Created => DomainRecordMutation::Create {
2411 record: DomainRecordDraft {
2412 reference: change.current.reference.clone(),
2413 payload: change.current.payload.clone(),
2414 references: change.current.references.clone(),
2415 },
2416 },
2417 DomainRecordOperation::Updated => DomainRecordMutation::Update {
2418 record: DomainRecordDraft {
2419 reference: change.current.reference.clone(),
2420 payload: change.current.payload.clone(),
2421 references: change.current.references.clone(),
2422 },
2423 expected_version: change.previous.as_ref().map_or(0, |record| record.version),
2424 },
2425 DomainRecordOperation::Retired => DomainRecordMutation::Retire {
2426 record: change.current.reference.clone(),
2427 expected_version: change.previous.as_ref().map_or(0, |record| record.version),
2428 successor: match &change.current.lifecycle {
2429 DomainRecordLifecycle::Retired { successor, .. } => successor.clone(),
2430 DomainRecordLifecycle::Active | DomainRecordLifecycle::Deleted { .. } => None,
2431 },
2432 },
2433 DomainRecordOperation::Deleted => DomainRecordMutation::Delete {
2434 record: change.current.reference.clone(),
2435 expected_version: change.previous.as_ref().map_or(0, |record| record.version),
2436 },
2437 }
2438}
2439
2440fn validate_record_shape(record: &DomainRecord, now: SimTime) -> Result<(), CanwuError> {
2441 validate_reference(&record.reference)?;
2442 if !canonical_text(&record.owner)
2443 || record.version == 0
2444 || record.references.windows(2).any(|pair| pair[0] >= pair[1])
2445 {
2446 return invalid_record(format!(
2447 "domain record {} has noncanonical identity, owner, version, or references",
2448 record.reference
2449 ));
2450 }
2451 for reference in &record.references {
2452 if !canonical_text(&reference.role) {
2453 return invalid_record("domain record reference roles must be canonical");
2454 }
2455 validate_target(&reference.target)?;
2456 }
2457 match &record.lifecycle {
2458 DomainRecordLifecycle::Active => {}
2459 DomainRecordLifecycle::Retired { at, successor } => {
2460 if *at > now {
2461 return invalid_record("domain record retirement cannot be future-dated");
2462 }
2463 if let Some(successor) = successor {
2464 validate_reference(successor)?;
2465 }
2466 }
2467 DomainRecordLifecycle::Deleted { at } => {
2468 if *at > now || !record.references.is_empty() {
2469 return invalid_record(
2470 "deleted domain-record tombstones cannot be future-dated or retain references",
2471 );
2472 }
2473 }
2474 }
2475 Ok(())
2476}
2477
2478fn validate_record_references(
2479 record: &DomainRecord,
2480 schema: &DomainRecordSchema,
2481 records: &impl DomainRecordRead,
2482 core_exists: &dyn Fn(&EntityRef) -> bool,
2483) -> Result<(), CanwuError> {
2484 if record.is_deleted() {
2485 return Ok(());
2486 }
2487 let mut by_role = BTreeMap::<&str, Vec<&DomainReferenceTarget>>::new();
2488 for reference in &record.references {
2489 by_role
2490 .entry(reference.role.as_str())
2491 .or_default()
2492 .push(&reference.target);
2493 }
2494 for reference_schema in &schema.references {
2495 let targets = by_role
2496 .remove(reference_schema.role.as_str())
2497 .unwrap_or_default();
2498 if (reference_schema.required && targets.is_empty())
2499 || (!reference_schema.multiple && targets.len() > 1)
2500 {
2501 return invalid_record(format!(
2502 "domain record {} violates reference cardinality for role {}",
2503 record.reference, reference_schema.role
2504 ));
2505 }
2506 for target in targets {
2507 validate_reference_target(
2508 target,
2509 reference_schema,
2510 records,
2511 core_exists,
2512 &record.reference,
2513 )?;
2514 }
2515 }
2516 if !by_role.is_empty() {
2517 return invalid_record(format!(
2518 "domain record {} contains undeclared reference roles",
2519 record.reference
2520 ));
2521 }
2522 Ok(())
2523}
2524
2525fn validate_reference_target(
2526 target: &DomainReferenceTarget,
2527 schema: &DomainReferenceSchema,
2528 records: &impl DomainRecordRead,
2529 core_exists: &dyn Fn(&EntityRef) -> bool,
2530 source: &DomainRecordRef,
2531) -> Result<(), CanwuError> {
2532 let (kind, is_entity) = match target {
2533 DomainReferenceTarget::Core(entity) => {
2534 let Some(kind) = entity.core_kind() else {
2535 return invalid_record(
2536 "domain entities use domain references rather than core-reference aliases",
2537 );
2538 };
2539 if !core_exists(entity) {
2540 return invalid_record(format!(
2541 "domain record {source} references missing core entity {entity}"
2542 ));
2543 }
2544 (DomainReferenceTargetKind::Core(kind), true)
2545 }
2546 DomainReferenceTarget::Domain(reference) => {
2547 let target_record = records.get(reference).ok_or_else(|| {
2548 CanwuError::new(
2549 ErrorCode::DomainRecordNotFound,
2550 format!("domain record {source} references missing record {reference}"),
2551 )
2552 })?;
2553 if target_record.is_deleted() || (!schema.allow_retired && !target_record.is_active()) {
2554 return Err(CanwuError::new(
2555 ErrorCode::DomainRecordReferenced,
2556 format!("domain record {source} references unavailable record {reference}"),
2557 ));
2558 }
2559 (
2560 DomainReferenceTargetKind::Domain(reference.kind.clone()),
2561 target_record.class == DomainRecordClass::Entity,
2562 )
2563 }
2564 };
2565 if !(schema.targets.contains(&kind)
2566 || is_entity
2567 && schema
2568 .targets
2569 .contains(&DomainReferenceTargetKind::AnyEntity))
2570 {
2571 return invalid_record(format!(
2572 "domain record {source} reference target does not match role {}",
2573 schema.role
2574 ));
2575 }
2576 Ok(())
2577}
2578
2579fn validate_reference_targets_basic(
2580 record: &DomainRecord,
2581 records: &BTreeMap<DomainRecordRef, DomainRecord>,
2582 core_exists: &dyn Fn(&EntityRef) -> bool,
2583) -> Result<(), CanwuError> {
2584 for reference in &record.references {
2585 match &reference.target {
2586 DomainReferenceTarget::Core(entity) => {
2587 if entity.core_kind().is_none() || !core_exists(entity) {
2588 return invalid_record(format!(
2589 "domain record {} references missing core entity {entity}",
2590 record.reference
2591 ));
2592 }
2593 }
2594 DomainReferenceTarget::Domain(target) => {
2595 if records.get(target).is_none_or(DomainRecord::is_deleted) {
2596 return invalid_record(format!(
2597 "domain record {} references unavailable record {target}",
2598 record.reference
2599 ));
2600 }
2601 }
2602 }
2603 }
2604 Ok(())
2605}
2606
2607fn validate_successor(
2608 record: &DomainRecord,
2609 records: &impl DomainRecordRead,
2610) -> Result<(), CanwuError> {
2611 let DomainRecordLifecycle::Retired {
2612 successor: Some(successor),
2613 ..
2614 } = &record.lifecycle
2615 else {
2616 return Ok(());
2617 };
2618 let Some(target) = records.get(successor) else {
2619 return invalid_record(format!(
2620 "retired domain record {} has a missing successor {successor}",
2621 record.reference
2622 ));
2623 };
2624 if successor == &record.reference
2625 || successor.kind != record.reference.kind
2626 || target.is_deleted()
2627 {
2628 return invalid_record(
2629 "domain record successors must be distinct available records of the same kind",
2630 );
2631 }
2632 Ok(())
2633}
2634
2635fn validate_new_successor(
2636 record: &DomainRecordRef,
2637 successor: Option<&DomainRecordRef>,
2638 records: &impl DomainRecordRead,
2639 created_records: &BTreeSet<DomainRecordRef>,
2640) -> Result<(), CanwuError> {
2641 let Some(successor) = successor else {
2642 return Ok(());
2643 };
2644 if successor == record || successor.kind != record.kind {
2645 return invalid_record(
2646 "domain record successors must be distinct active records of the same kind",
2647 );
2648 }
2649 if created_records.contains(successor) {
2650 return Ok(());
2651 }
2652 let Some(target) = records.get(successor) else {
2653 return invalid_record(format!(
2654 "retired domain record {record} has a missing successor {successor}",
2655 ));
2656 };
2657 if !target.is_active() {
2658 return invalid_record("new domain record successors must be active when admitted");
2659 }
2660 Ok(())
2661}
2662
2663fn validate_successor_graph(records: &impl DomainRecordRead) -> Result<(), CanwuError> {
2664 let mut complete = BTreeSet::new();
2665 for (start, _) in records.iter() {
2666 if complete.contains(start) {
2667 continue;
2668 }
2669 let mut visited = BTreeSet::new();
2670 let mut path = Vec::new();
2671 let mut current = start;
2672 loop {
2673 if complete.contains(current) {
2674 break;
2675 }
2676 if !visited.insert(current.clone()) {
2677 return invalid_record("domain record successor chains cannot contain cycles");
2678 }
2679 path.push(current.clone());
2680 let Some(DomainRecord {
2681 lifecycle:
2682 DomainRecordLifecycle::Retired {
2683 successor: Some(successor),
2684 ..
2685 },
2686 ..
2687 }) = records.get(current)
2688 else {
2689 break;
2690 };
2691 current = successor;
2692 }
2693 complete.extend(path);
2694 }
2695 Ok(())
2696}
2697
2698fn require_mutable_record<'a>(
2699 records: &'a impl DomainRecordRead,
2700 reference: &DomainRecordRef,
2701 expected_version: u64,
2702) -> Result<&'a DomainRecord, CanwuError> {
2703 let record = records.get(reference).ok_or_else(|| {
2704 CanwuError::new(
2705 ErrorCode::DomainRecordNotFound,
2706 format!("domain record {reference} was not found"),
2707 )
2708 })?;
2709 if record.version != expected_version {
2710 return Err(version_conflict(
2711 reference,
2712 expected_version,
2713 record.version,
2714 ));
2715 }
2716 if !record.is_active() {
2717 return invalid_record("only active domain records can be updated or retired");
2718 }
2719 Ok(record)
2720}
2721
2722fn next_record_version(version: u64) -> Result<u64, CanwuError> {
2723 version.checked_add(1).ok_or_else(|| {
2724 CanwuError::new(
2725 ErrorCode::IdentifierExhausted,
2726 "domain record version space is exhausted",
2727 )
2728 })
2729}
2730
2731fn version_conflict(reference: &DomainRecordRef, expected: u64, actual: u64) -> CanwuError {
2732 CanwuError::new(
2733 ErrorCode::DomainRecordVersionConflict,
2734 format!(
2735 "domain record {reference} expected version {expected}, but current version is {actual}"
2736 ),
2737 )
2738}
2739
2740fn validate_target(target: &DomainReferenceTarget) -> Result<(), CanwuError> {
2741 match target {
2742 DomainReferenceTarget::Core(entity) => {
2743 if entity.core_kind().is_none() {
2744 return invalid_record(
2745 "domain entities must use domain-record references in record fields",
2746 );
2747 }
2748 Ok(())
2749 }
2750 DomainReferenceTarget::Domain(reference) => validate_reference(reference),
2751 }
2752}
2753
2754fn validate_reference(reference: &DomainRecordRef) -> Result<(), CanwuError> {
2755 validate_kind(&reference.kind)?;
2756 if !canonical_text(&reference.id) {
2757 return invalid_record("domain record IDs must be non-empty canonical strings");
2758 }
2759 Ok(())
2760}
2761
2762fn validate_kind(kind: &DomainRecordKind) -> Result<(), CanwuError> {
2763 if !canonical_text(&kind.namespace) || !canonical_text(&kind.name) {
2764 return invalid_record("domain record kinds require canonical namespace and name values");
2765 }
2766 Ok(())
2767}
2768
2769fn canonical_text(value: &str) -> bool {
2770 !value.is_empty() && value == value.trim()
2771}
2772
2773fn invalid_record<T>(message: impl Into<String>) -> Result<T, CanwuError> {
2774 Err(CanwuError::new(ErrorCode::InvalidDomainRecord, message))
2775}
2776
2777#[cfg(test)]
2778mod tests {
2779 use super::*;
2780 use serde_json::json;
2781
2782 fn fixture_kind(name: &str) -> DomainRecordKind {
2783 DomainRecordKind::new("fixture.records", name)
2784 }
2785
2786 fn fixture_record(reference: DomainRecordRef, class: DomainRecordClass) -> DomainRecord {
2787 DomainRecord {
2788 reference,
2789 owner: "fixture".to_owned(),
2790 class,
2791 version: 1,
2792 lifecycle: DomainRecordLifecycle::Active,
2793 payload: json!(null),
2794 references: vec![],
2795 }
2796 }
2797
2798 fn mutation_request(mutation: &DomainRecordMutation) -> DomainMutationRequest<'_> {
2799 DomainMutationRequest {
2800 plugin: "fixture",
2801 system: "fixture-system",
2802 visibility: StateVisibility::NextBoundary,
2803 mutation,
2804 summary: "fixture mutation",
2805 }
2806 }
2807
2808 #[test]
2809 fn create_only_policy_rejects_raw_mutations_and_preserves_create() {
2810 let kind = fixture_kind("immutable");
2811 let existing_ref = DomainRecordRef::new(&kind.namespace, &kind.name, "existing");
2812 let existing = fixture_record(existing_ref.clone(), DomainRecordClass::Record);
2813 let records = BTreeMap::from([(existing_ref.clone(), existing.clone())]);
2814 let mut schema = DomainRecordSchema::new(kind.clone(), DomainRecordClass::Record);
2815 schema.mutation_policy = DomainRecordMutationPolicy::CreateOnly;
2816 let schemas = BTreeMap::from([(kind.clone(), ("fixture".to_owned(), schema))]);
2817 let core_exists = |_: &EntityRef| true;
2818
2819 let update = DomainRecordMutation::Update {
2820 record: DomainRecordDraft::new(existing_ref.clone(), json!("changed")),
2821 expected_version: 1,
2822 };
2823 let retire = DomainRecordMutation::Retire {
2824 record: existing_ref.clone(),
2825 expected_version: 1,
2826 successor: None,
2827 };
2828 let delete = DomainRecordMutation::Delete {
2829 record: existing_ref.clone(),
2830 expected_version: 1,
2831 };
2832 for mutation in [&update, &retire, &delete] {
2833 let error = apply_mutation_bundle(
2834 &records,
2835 &schemas,
2836 SimTime::EPOCH,
2837 &core_exists,
2838 vec![mutation_request(mutation)],
2839 )
2840 .expect_err("create-only kinds must reject non-create raw mutations");
2841 assert_eq!(error.code, ErrorCode::InvalidDomainRecord);
2842 assert_eq!(records.get(&existing_ref), Some(&existing));
2843 }
2844
2845 let new_ref = DomainRecordRef::new(&kind.namespace, &kind.name, "new");
2846 let create = DomainRecordMutation::Create {
2847 record: DomainRecordDraft::new(new_ref.clone(), json!(null)),
2848 };
2849 let (next, changes) = apply_mutation_bundle(
2850 &records,
2851 &schemas,
2852 SimTime::EPOCH,
2853 &core_exists,
2854 vec![mutation_request(&create)],
2855 )
2856 .expect("create-only kinds must still accept creates");
2857 assert!(next.contains_key(&new_ref));
2858 assert_eq!(changes[0].operation, DomainRecordOperation::Created);
2859 }
2860
2861 #[test]
2862 fn any_entity_domain_reference_is_typed_by_registered_record_class() {
2863 let entity_kind = fixture_kind("future-entity");
2864 let value_kind = fixture_kind("future-value");
2865 let entity_ref = DomainRecordRef::new(&entity_kind.namespace, &entity_kind.name, "one");
2866 let value_ref = DomainRecordRef::new(&value_kind.namespace, &value_kind.name, "one");
2867 let records = BTreeMap::from([
2868 (
2869 entity_ref.clone(),
2870 fixture_record(entity_ref.clone(), DomainRecordClass::Entity),
2871 ),
2872 (
2873 value_ref.clone(),
2874 fixture_record(value_ref.clone(), DomainRecordClass::Record),
2875 ),
2876 ]);
2877 let schema = DomainReferenceSchema {
2878 role: "target".to_owned(),
2879 targets: vec![DomainReferenceTargetKind::AnyEntity],
2880 required: true,
2881 multiple: false,
2882 allow_retired: false,
2883 };
2884 let source = DomainRecordRef::new("fixture.records", "source", "one");
2885
2886 assert!(
2887 validate_reference_target(
2888 &DomainReferenceTarget::Domain(entity_ref),
2889 &schema,
2890 &records,
2891 &|_: &EntityRef| false,
2892 &source,
2893 )
2894 .is_ok()
2895 );
2896 let error = validate_reference_target(
2897 &DomainReferenceTarget::Domain(value_ref),
2898 &schema,
2899 &records,
2900 &|_: &EntityRef| false,
2901 &source,
2902 )
2903 .expect_err("AnyEntity must reject domain value records");
2904 assert_eq!(error.code, ErrorCode::InvalidDomainRecord);
2905 }
2906
2907 #[test]
2908 fn persistent_store_updates_leave_captured_root_unchanged() {
2909 let kind = fixture_kind("persistent");
2910 let first_ref = DomainRecordRef::new(&kind.namespace, &kind.name, "first");
2911 let second_ref = DomainRecordRef::new(&kind.namespace, &kind.name, "second");
2912 let records = BTreeMap::from([
2913 (
2914 first_ref.clone(),
2915 fixture_record(first_ref.clone(), DomainRecordClass::Record),
2916 ),
2917 (
2918 second_ref.clone(),
2919 fixture_record(second_ref, DomainRecordClass::Record),
2920 ),
2921 ]);
2922 let store = PersistentDomainRecordStore::from_records(records).unwrap();
2923 let captured = store.clone();
2924 assert!(captured.shares_root_with(&store));
2925
2926 let schema = DomainRecordSchema::new(kind.clone(), DomainRecordClass::Record);
2927 let schemas = BTreeMap::from([(kind, ("fixture".to_owned(), schema))]);
2928 let update = DomainRecordMutation::Update {
2929 record: DomainRecordDraft::new(first_ref.clone(), json!("changed")),
2930 expected_version: 1,
2931 };
2932 let (updated, changes) = apply_mutation_bundle_cow(
2933 &store,
2934 &schemas,
2935 SimTime::EPOCH,
2936 &|_: &EntityRef| true,
2937 vec![mutation_request(&update)],
2938 )
2939 .unwrap();
2940
2941 assert!(!updated.shares_root_with(&captured));
2942 assert_eq!(captured.get(&first_ref).unwrap().version, 1);
2943 assert_eq!(updated.get(&first_ref).unwrap().version, 2);
2944 assert_ne!(
2945 captured.commitment_root().unwrap(),
2946 updated.commitment_root().unwrap()
2947 );
2948 assert_eq!(changes.len(), 1);
2949 }
2950
2951 #[test]
2952 fn persistent_overlay_validation_shares_unrelated_record_payloads() {
2953 let kind = fixture_kind("overlay");
2954 let first_ref = DomainRecordRef::new(&kind.namespace, &kind.name, "first");
2955 let second_ref = DomainRecordRef::new(&kind.namespace, &kind.name, "second");
2956 let cold_ref = DomainRecordRef::new(&kind.namespace, &kind.name, "cold");
2957 let records = BTreeMap::from([
2958 (
2959 first_ref.clone(),
2960 fixture_record(first_ref.clone(), DomainRecordClass::Record),
2961 ),
2962 (
2963 second_ref.clone(),
2964 fixture_record(second_ref.clone(), DomainRecordClass::Record),
2965 ),
2966 (
2967 cold_ref.clone(),
2968 DomainRecord {
2969 payload: json!({ "cold": "x".repeat(1_000_000) }),
2970 ..fixture_record(cold_ref.clone(), DomainRecordClass::Record)
2971 },
2972 ),
2973 ]);
2974 let store = PersistentDomainRecordStore::from_records(records).unwrap();
2975 let mut first_overlay = store.get(&first_ref).unwrap().clone();
2976 first_overlay.version = 2;
2977 first_overlay.payload = json!("first-overlay");
2978 let overlay = BTreeMap::from([(first_ref.clone(), first_overlay)]);
2979 let schema = DomainRecordSchema::new(kind.clone(), DomainRecordClass::Record);
2980 let schemas = BTreeMap::from([(kind, ("fixture".to_owned(), schema))]);
2981 let update = DomainRecordMutation::Update {
2982 record: DomainRecordDraft::new(second_ref.clone(), json!("second-update")),
2983 expected_version: 1,
2984 };
2985
2986 let (updated, changes) = apply_mutation_bundle_cow_with_overlay(
2987 &store,
2988 &overlay,
2989 &schemas,
2990 SimTime::EPOCH,
2991 &|_: &EntityRef| true,
2992 vec![mutation_request(&update)],
2993 )
2994 .unwrap();
2995
2996 assert_eq!(store.get(&first_ref).unwrap().version, 1);
2997 assert_eq!(updated.get(&first_ref).unwrap().version, 2);
2998 assert_eq!(updated.get(&second_ref).unwrap().version, 2);
2999 assert!(Arc::ptr_eq(
3000 store.records.get(&cold_ref).unwrap(),
3001 updated.records.get(&cold_ref).unwrap()
3002 ));
3003 assert_eq!(changes.len(), 1);
3004 assert_eq!(changes[0].current.reference, second_ref);
3005 }
3006}