1use super::knowledge::{KnowledgeLimitsV1, PluginKnowledgeSchema, validate_schema_set};
2use super::maintenance::{
3 OwnerAuthorizedMaintenanceRequest, OwnerAuthorizedParticipantDraft,
4 OwnerAuthorizedParticipantRole,
5};
6use super::records::DomainRecordSchema;
7use super::{
8 ArchiveReachabilityManifest, BoundaryPhase, BoundarySystemContract, BoundarySystemHandler,
9 BoundaryWriteStage, CORE_STATE_NAMESPACE, CanwuError, CommandContext, DomainRecord, EntityRef,
10 ErrorCode, IngressClass, PluginArchiveRetention, PluginIngressDescriptor, PluginIngressPermit,
11 RandomStreamKey, ReservationRef, SchemaRegistry, SimDuration, SimEvent, SimulationView,
12 StateKey, StateVisibility, SystemCadence, SystemContract, TypeSchema, boundary_write_stage,
13 canonical_hash, canonical_text, invalid_snapshot, invalid_snapshot_error, is_canonical_hash,
14 is_domain_record_state, knowledge, records, validate_type_schema,
15};
16use canwu_event::EventAudience;
17use serde::{Deserialize, Serialize};
18use serde_json::Value;
19use std::collections::{BTreeMap, BTreeSet};
20
21#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
22#[serde(rename_all = "snake_case")]
23pub enum PayloadValueType {
24 Null,
25 Boolean,
26 Integer,
27 String,
28 Object,
29 Array,
30}
31
32impl PayloadValueType {
33 fn matches(&self, value: &Value) -> bool {
34 match self {
35 Self::Null => value.is_null(),
36 Self::Boolean => value.is_boolean(),
37 Self::Integer => value.as_i64().is_some() || value.as_u64().is_some(),
38 Self::String => value.is_string(),
39 Self::Object => value.is_object(),
40 Self::Array => value.is_array(),
41 }
42 }
43}
44
45#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
46pub struct PayloadProperty {
47 pub value_type: PayloadValueType,
48 pub required: bool,
49}
50
51#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
52#[serde(tag = "type", rename_all = "snake_case")]
53pub enum PayloadSchema {
54 Any,
55 Null,
56 Boolean,
57 Integer,
58 String,
59 Object {
60 properties: BTreeMap<String, PayloadProperty>,
61 allow_additional: bool,
62 },
63}
64
65impl PayloadSchema {
66 pub(super) fn validate(&self, value: &Value) -> Result<(), CanwuError> {
67 let scalar_matches = match self {
68 Self::Any => return Ok(()),
69 Self::Null => value.is_null(),
70 Self::Boolean => value.is_boolean(),
71 Self::Integer => value.as_i64().is_some() || value.as_u64().is_some(),
72 Self::String => value.is_string(),
73 Self::Object {
74 properties,
75 allow_additional,
76 } => {
77 let Some(object) = value.as_object() else {
78 return Err(CanwuError::new(
79 ErrorCode::InvalidPayload,
80 "plugin command payload must be an object",
81 ));
82 };
83 for (name, property) in properties {
84 match object.get(name) {
85 Some(field) if !property.value_type.matches(field) => {
86 return Err(CanwuError::new(
87 ErrorCode::InvalidPayload,
88 format!("payload field {name} has the wrong type"),
89 ));
90 }
91 None if property.required => {
92 return Err(CanwuError::new(
93 ErrorCode::InvalidPayload,
94 format!("payload field {name} is required"),
95 ));
96 }
97 Some(_) | None => {}
98 }
99 }
100 if !allow_additional && object.keys().any(|name| !properties.contains_key(name)) {
101 return Err(CanwuError::new(
102 ErrorCode::InvalidPayload,
103 "plugin command payload contains an undeclared field",
104 ));
105 }
106 return Ok(());
107 }
108 };
109 if scalar_matches {
110 Ok(())
111 } else {
112 Err(CanwuError::new(
113 ErrorCode::InvalidPayload,
114 "plugin command payload does not match its declared schema",
115 ))
116 }
117 }
118}
119
120#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
121pub struct PluginActionDescriptor {
122 pub name: String,
123 pub description: String,
124 pub payload_schema: PayloadSchema,
125 pub reads: Vec<StateKey>,
126 pub writes: Vec<StateKey>,
127}
128
129pub const PLUGIN_DESCRIPTOR_FORMAT_VERSION: u32 = 1;
130
131#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
136#[serde(deny_unknown_fields)]
137pub struct MaintenanceDependencyResolverDescriptor {
138 pub target_namespace: String,
139}
140
141#[derive(Clone, Debug, Deserialize, Eq, Ord, PartialEq, PartialOrd, Serialize)]
145#[serde(deny_unknown_fields)]
146pub struct PluginArchiveRetentionBinding {
147 pub packet_type: String,
148 pub namespace: String,
149 pub object_id_json_pointers: Vec<String>,
150}
151
152#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
153pub struct PluginDescriptor {
154 #[serde(default = "missing_descriptor_format")]
157 pub descriptor_format: u32,
158 pub name: String,
159 #[serde(default)]
160 pub version: String,
161 #[serde(default)]
162 pub semantic_hash: String,
163 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
166 pub event_audiences: BTreeMap<String, EventAudience>,
167 pub systems: Vec<SystemContract>,
168 #[serde(default)]
169 pub boundary_systems: Vec<BoundarySystemContract>,
170 pub commands: Vec<PluginActionDescriptor>,
171 #[serde(default, skip_serializing_if = "Vec::is_empty")]
172 pub ingress: Vec<PluginIngressDescriptor>,
173 #[serde(default, skip_serializing_if = "Vec::is_empty")]
174 pub internal_ingress: Vec<String>,
175 #[serde(default, skip_serializing_if = "Vec::is_empty")]
176 pub archive_retention_bindings: Vec<PluginArchiveRetentionBinding>,
177 pub schema_types: Vec<String>,
178 #[serde(default, skip_serializing_if = "Vec::is_empty")]
179 pub record_schemas: Vec<DomainRecordSchema>,
180 #[serde(default, skip_serializing_if = "Vec::is_empty")]
181 pub knowledge_schemas: Vec<PluginKnowledgeSchema>,
182 #[serde(default, skip_serializing_if = "Vec::is_empty")]
183 pub maintenance_dependency_resolvers: Vec<MaintenanceDependencyResolverDescriptor>,
184 #[serde(default, skip_serializing_if = "is_false")]
185 pub owner_authorized_maintenance_participant: bool,
186 #[serde(default, skip_serializing_if = "is_false")]
187 pub archive_reachability_participant: bool,
188}
189
190#[allow(clippy::trivially_copy_pass_by_ref)]
191fn is_false(value: &bool) -> bool {
192 !*value
193}
194
195fn missing_descriptor_format() -> u32 {
196 0
197}
198
199impl Default for PluginDescriptor {
200 fn default() -> Self {
201 Self {
202 descriptor_format: PLUGIN_DESCRIPTOR_FORMAT_VERSION,
203 name: String::new(),
204 version: String::new(),
205 semantic_hash: String::new(),
206 event_audiences: BTreeMap::new(),
207 systems: Vec::new(),
208 boundary_systems: Vec::new(),
209 commands: Vec::new(),
210 ingress: Vec::new(),
211 internal_ingress: Vec::new(),
212 archive_retention_bindings: Vec::new(),
213 schema_types: Vec::new(),
214 record_schemas: Vec::new(),
215 knowledge_schemas: Vec::new(),
216 maintenance_dependency_resolvers: Vec::new(),
217 owner_authorized_maintenance_participant: false,
218 archive_reachability_participant: false,
219 }
220 }
221}
222
223#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
224pub struct PluginComponentRecord {
225 pub plugin: String,
226 pub state: StateKey,
227 pub entity: EntityRef,
228 pub component: String,
229 pub value: Value,
230}
231
232#[derive(Clone, Debug, Eq, Ord, PartialEq, PartialOrd)]
233pub(super) struct PluginComponentKey {
234 pub(super) plugin: String,
235 pub(super) state: StateKey,
236 pub(super) entity: EntityRef,
237 pub(super) component: String,
238}
239
240#[derive(Clone, Debug, Deserialize, PartialEq, Serialize)]
241#[serde(tag = "type", rename_all = "snake_case")]
242pub enum SystemDirective {
243 SetComponent {
244 state: StateKey,
245 entity: EntityRef,
246 component: String,
247 value: Value,
248 summary: String,
249 },
250 Emit {
251 event_type: String,
252 summary: String,
253 affected: Vec<EntityRef>,
254 },
255 Schedule {
256 after: SimDuration,
257 directive: Box<SystemDirective>,
258 },
259 EnqueuePluginIngress {
262 after: SimDuration,
263 packet_type: String,
264 priority: i32,
265 payload: Value,
266 affected: Vec<EntityRef>,
267 },
268}
269
270pub type SimulationSystemHandler =
277 fn(&SimulationView<'_>, &SimEvent) -> Result<Vec<SystemDirective>, CanwuError>;
278
279pub type PluginCommandHandler =
280 fn(&SimulationView<'_>, &CommandContext, &Value) -> Result<Vec<SystemDirective>, CanwuError>;
281
282pub type OwnerAuthorizedMaintenanceParticipant =
286 fn(
287 &SimulationView<'_>,
288 &OwnerAuthorizedMaintenanceRequest,
289 OwnerAuthorizedParticipantRole,
290 ) -> Result<OwnerAuthorizedParticipantDraft, CanwuError>;
291
292pub trait PluginArchiveObjectProvider {
296 fn load_plugin_archive_object(
297 &self,
298 namespace: &str,
299 object_id: &str,
300 ) -> Result<Option<Vec<u8>>, CanwuError>;
301}
302
303impl PluginArchiveObjectProvider for () {
304 fn load_plugin_archive_object(
305 &self,
306 _namespace: &str,
307 _object_id: &str,
308 ) -> Result<Option<Vec<u8>>, CanwuError> {
309 Ok(None)
310 }
311}
312
313pub type PluginArchiveReachabilityParticipant = fn(
317 &SimulationView<'_>,
318 &dyn PluginArchiveObjectProvider,
319 &mut ArchiveReachabilityManifest,
320) -> Result<(), CanwuError>;
321
322pub trait SimulationPlugin {
325 fn name(&self) -> &str;
326 fn version(&self) -> &str;
328 fn semantic_hash(&self) -> &str;
333 fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError>;
334
335 fn validate_activation(&self, _records: &[DomainRecord]) -> Result<(), CanwuError> {
341 Ok(())
342 }
343}
344
345#[derive(Clone, Default)]
346pub struct PluginRegistry {
347 pub(super) descriptors: BTreeMap<String, PluginDescriptor>,
348 pub(super) active_plugins: BTreeSet<String>,
349 pub(super) systems: Vec<RegisteredSystem>,
350 pub(super) boundary_systems: Vec<RegisteredBoundarySystem>,
351 pub(super) commands: BTreeMap<(String, String), RegisteredCommand>,
352 pub(super) ingress: BTreeMap<(String, String), PluginIngressDescriptor>,
353 pub(super) internal_ingress: BTreeSet<(String, String)>,
354 pub(super) archive_retention_bindings:
355 BTreeMap<(String, String), PluginArchiveRetentionBinding>,
356 pub(super) state_owners: BTreeMap<StateKey, String>,
357 pub(super) immediate_write_states: BTreeMap<StateKey, String>,
358 pub(super) boundary_writers: BTreeMap<(BoundaryWriteStage, StateKey), (String, String)>,
359 pub(super) reservation_offerers: BTreeMap<StateKey, (String, String)>,
360 pub(super) random_stream_owners: BTreeMap<RandomStreamKey, (String, String)>,
361 pub(super) record_schemas: records::DomainRecordSchemas,
362 pub(super) knowledge_schemas: knowledge::KnowledgeSchemas,
363 pub(super) knowledge_kind_owners: knowledge::KnowledgeKindOwners,
364 pub(super) maintenance_dependency_resolvers: BTreeMap<String, BTreeSet<String>>,
365 pub(super) maintenance_participants: BTreeMap<String, OwnerAuthorizedMaintenanceParticipant>,
366 pub(super) archive_reachability_participants:
367 BTreeMap<String, PluginArchiveReachabilityParticipant>,
368}
369
370#[derive(Clone)]
371pub(super) struct RegisteredSystem {
372 pub(super) plugin: String,
373 pub(super) contract: SystemContract,
374 pub(super) handler: SimulationSystemHandler,
375}
376
377#[derive(Clone)]
378pub(super) struct RegisteredBoundarySystem {
379 pub(super) plugin: String,
380 pub(super) contract: BoundarySystemContract,
381 pub(super) handler: BoundarySystemHandler,
382}
383
384#[derive(Clone)]
385pub(super) struct RegisteredCommand {
386 pub(super) descriptor: PluginActionDescriptor,
387 pub(super) handler: PluginCommandHandler,
388}
389
390pub struct PluginRegistrar<'a> {
391 pub(super) plugin: String,
392 pub(super) registry: &'a mut PluginRegistry,
393 pub(super) schema: &'a mut SchemaRegistry,
394}
395
396impl PluginRegistrar<'_> {
397 pub fn register_archive_reachability_participant(
398 &mut self,
399 handler: PluginArchiveReachabilityParticipant,
400 ) -> Result<(), CanwuError> {
401 let mut candidate = self.registry.clone();
402 if candidate
403 .archive_reachability_participants
404 .insert(self.plugin.clone(), handler)
405 .is_some()
406 {
407 return Err(CanwuError::new(
408 ErrorCode::InvalidPluginRegistration,
409 "archive reachability participant was registered twice",
410 ));
411 }
412 candidate
413 .descriptors
414 .entry(self.plugin.clone())
415 .or_default()
416 .archive_reachability_participant = true;
417 *self.registry = candidate;
418 Ok(())
419 }
420
421 pub fn register_owner_authorized_maintenance_participant(
422 &mut self,
423 handler: OwnerAuthorizedMaintenanceParticipant,
424 ) -> Result<(), CanwuError> {
425 let mut candidate = self.registry.clone();
426 if candidate
427 .maintenance_participants
428 .insert(self.plugin.clone(), handler)
429 .is_some()
430 {
431 return Err(CanwuError::new(
432 ErrorCode::InvalidPluginRegistration,
433 "owner-authorized maintenance participant was registered twice",
434 ));
435 }
436 candidate
437 .descriptors
438 .entry(self.plugin.clone())
439 .or_default()
440 .owner_authorized_maintenance_participant = true;
441 *self.registry = candidate;
442 Ok(())
443 }
444
445 pub fn register_maintenance_dependency_resolver(
446 &mut self,
447 target_namespace: impl Into<String>,
448 ) -> Result<(), CanwuError> {
449 let target_namespace = target_namespace.into();
450 if !canonical_text(&target_namespace) || target_namespace == CORE_STATE_NAMESPACE {
451 return Err(CanwuError::new(
452 ErrorCode::InvalidPluginRegistration,
453 "maintenance dependency target namespace is invalid",
454 ));
455 }
456 let descriptor = MaintenanceDependencyResolverDescriptor {
457 target_namespace: target_namespace.clone(),
458 };
459 let mut candidate = self.registry.clone();
460 let plugin_descriptor = candidate
461 .descriptors
462 .entry(self.plugin.clone())
463 .or_default();
464 if plugin_descriptor
465 .maintenance_dependency_resolvers
466 .contains(&descriptor)
467 {
468 return Err(CanwuError::new(
469 ErrorCode::InvalidPluginRegistration,
470 "maintenance dependency resolver was registered twice",
471 ));
472 }
473 plugin_descriptor
474 .maintenance_dependency_resolvers
475 .push(descriptor);
476 plugin_descriptor.maintenance_dependency_resolvers.sort();
477 candidate
478 .maintenance_dependency_resolvers
479 .entry(target_namespace)
480 .or_default()
481 .insert(self.plugin.clone());
482 *self.registry = candidate;
483 Ok(())
484 }
485
486 pub fn register_record_schema(
487 &mut self,
488 mut schema: DomainRecordSchema,
489 ) -> Result<(), CanwuError> {
490 schema.canonicalize();
491 schema.validate().map_err(|error| {
492 CanwuError::new(
493 ErrorCode::InvalidPluginRegistration,
494 format!("invalid domain record schema: {error}"),
495 )
496 })?;
497 let state = schema.state_key();
498 if state.namespace == CORE_STATE_NAMESPACE {
499 return Err(CanwuError::new(
500 ErrorCode::InvalidPluginRegistration,
501 "plugins cannot register domain record kinds in the core namespace",
502 ));
503 }
504 if self
505 .registry
506 .descriptors
507 .get(&self.plugin)
508 .is_some_and(|descriptor| {
509 descriptor
510 .record_schemas
511 .iter()
512 .any(|candidate| candidate.kind == schema.kind)
513 })
514 {
515 return Err(CanwuError::new(
516 ErrorCode::DuplicateDomainRecordKind,
517 format!(
518 "plugin {} registered record kind {} twice",
519 self.plugin, schema.kind
520 ),
521 ));
522 }
523 if let Some((owner, existing)) = self.registry.record_schemas.get(&schema.kind) {
524 if owner != &self.plugin {
525 return Err(CanwuError::new(
526 ErrorCode::DuplicateDomainRecordKind,
527 format!(
528 "domain record kind {} is already owned by plugin {owner}",
529 schema.kind
530 ),
531 ));
532 }
533 if existing != &schema {
534 return Err(CanwuError::new(
535 ErrorCode::PluginManifestMismatch,
536 format!(
537 "plugin {} changed the stored schema for domain record kind {}",
538 self.plugin, schema.kind
539 ),
540 ));
541 }
542 }
543 let mut candidate = self.registry.clone();
544 if candidate.immediate_write_states.contains_key(&state) {
545 return Err(CanwuError::new(
546 ErrorCode::InvalidPluginRegistration,
547 format!(
548 "domain record kind {} is already exposed as immediate component state",
549 schema.kind
550 ),
551 ));
552 }
553 register_state_owners(
554 &mut candidate.state_owners,
555 &self.plugin,
556 std::slice::from_ref(&state),
557 )?;
558 candidate
559 .record_schemas
560 .insert(schema.kind.clone(), (self.plugin.clone(), schema.clone()));
561 let descriptor = candidate
562 .descriptors
563 .entry(self.plugin.clone())
564 .or_default();
565 descriptor.name.clone_from(&self.plugin);
566 descriptor.record_schemas.push(schema);
567 descriptor
568 .record_schemas
569 .sort_by(|left, right| left.kind.cmp(&right.kind));
570 *self.registry = candidate;
571 Ok(())
572 }
573
574 pub fn register_knowledge_schema(
575 &mut self,
576 mut schema: PluginKnowledgeSchema,
577 ) -> Result<(), CanwuError> {
578 schema.canonicalize();
579 schema.validate().map_err(|error| {
580 CanwuError::new(
581 ErrorCode::InvalidPluginRegistration,
582 format!("invalid knowledge schema: {error}"),
583 )
584 })?;
585 let current_count = self
586 .registry
587 .descriptors
588 .get(&self.plugin)
589 .map_or(0, |descriptor| descriptor.knowledge_schemas.len());
590 if current_count >= KnowledgeLimitsV1::CURRENT.schemas_per_plugin {
591 return Err(CanwuError::new(
592 ErrorCode::InvalidPluginRegistration,
593 "plugin knowledge schema limit exceeded",
594 ));
595 }
596 if self
597 .registry
598 .descriptors
599 .get(&self.plugin)
600 .is_some_and(|descriptor| {
601 descriptor
602 .knowledge_schemas
603 .iter()
604 .any(|candidate| candidate.id == schema.id)
605 })
606 {
607 return Err(CanwuError::new(
608 ErrorCode::InvalidPluginRegistration,
609 format!(
610 "plugin {} registered knowledge schema {:?} twice",
611 self.plugin, schema.id
612 ),
613 ));
614 }
615 if let Some(owner) = self.registry.knowledge_kind_owners.get(&schema.id.kind)
616 && owner != &self.plugin
617 {
618 return Err(CanwuError::new(
619 ErrorCode::InvalidPluginRegistration,
620 format!(
621 "knowledge kind {:?} is already owned by plugin {owner}",
622 schema.id.kind
623 ),
624 ));
625 }
626 if let Some((owner, existing)) = self.registry.knowledge_schemas.get(&schema.id) {
627 if owner != &self.plugin {
628 return Err(CanwuError::new(
629 ErrorCode::InvalidPluginRegistration,
630 format!(
631 "knowledge schema {:?} is already owned by plugin {owner}",
632 schema.id
633 ),
634 ));
635 }
636 if existing != &schema {
637 return Err(CanwuError::new(
638 ErrorCode::PluginManifestMismatch,
639 format!(
640 "plugin {} changed the stored knowledge schema {:?}",
641 self.plugin, schema.id
642 ),
643 ));
644 }
645 }
646 if schema.writable
647 && self
648 .registry
649 .knowledge_schemas
650 .values()
651 .any(|(owner, existing)| {
652 owner == &self.plugin
653 && existing.id != schema.id
654 && existing.id.kind == schema.id.kind
655 && existing.writable
656 })
657 {
658 return Err(CanwuError::new(
659 ErrorCode::InvalidPluginRegistration,
660 format!(
661 "knowledge kind {:?} already has a writable version",
662 schema.id.kind
663 ),
664 ));
665 }
666 let mut candidate = self.registry.clone();
667 candidate
668 .knowledge_kind_owners
669 .entry(schema.id.kind.clone())
670 .or_insert_with(|| self.plugin.clone());
671 candidate
672 .knowledge_schemas
673 .insert(schema.id.clone(), (self.plugin.clone(), schema.clone()));
674 let descriptor = candidate
675 .descriptors
676 .entry(self.plugin.clone())
677 .or_default();
678 descriptor.name.clone_from(&self.plugin);
679 descriptor.knowledge_schemas.push(schema);
680 descriptor
681 .knowledge_schemas
682 .sort_by(|left, right| left.id.cmp(&right.id));
683 *self.registry = candidate;
684 Ok(())
685 }
686
687 pub fn register_schema(&mut self, schema: TypeSchema) -> Result<(), CanwuError> {
688 validate_type_schema(&schema)?;
689 let type_name = schema.type_name.clone();
690 let mut candidate_schema = self.schema.clone();
691 let mut candidate_registry = self.registry.clone();
692 if let Some(existing) = candidate_schema.get(&type_name) {
693 if existing != &schema {
694 return Err(CanwuError::new(
695 ErrorCode::InvalidPluginRegistration,
696 format!(
697 "schema type {type_name} is already registered with a different definition"
698 ),
699 ));
700 }
701 } else {
702 candidate_schema.register(schema).map_err(|error| {
703 CanwuError::new(ErrorCode::InvalidPluginRegistration, error.to_string())
704 })?;
705 }
706 let descriptor = candidate_registry
707 .descriptors
708 .entry(self.plugin.clone())
709 .or_default();
710 if descriptor.schema_types.contains(&type_name) {
711 return Err(CanwuError::new(
712 ErrorCode::InvalidPluginRegistration,
713 format!(
714 "plugin {} registered schema type {} more than once",
715 self.plugin, type_name
716 ),
717 ));
718 }
719 descriptor.name.clone_from(&self.plugin);
720 descriptor.schema_types.push(type_name);
721 descriptor.schema_types.sort();
722 *self.schema = candidate_schema;
723 *self.registry = candidate_registry;
724 Ok(())
725 }
726
727 pub fn register_event_audience(
733 &mut self,
734 event_type: impl Into<String>,
735 audience: EventAudience,
736 ) -> Result<(), CanwuError> {
737 let event_type = event_type.into();
738 validate_event_audience_name(&event_type)?;
739 validate_event_audience(&audience)?;
740 let mut candidate = self.registry.clone();
741 let descriptor = candidate
742 .descriptors
743 .entry(self.plugin.clone())
744 .or_default();
745 descriptor.name.clone_from(&self.plugin);
746 if descriptor
747 .event_audiences
748 .insert(event_type.clone(), audience)
749 .is_some()
750 {
751 return Err(CanwuError::new(
752 ErrorCode::InvalidPluginRegistration,
753 format!(
754 "plugin {} already declared event audience for {event_type}",
755 self.plugin
756 ),
757 ));
758 }
759 *self.registry = candidate;
760 Ok(())
761 }
762
763 pub fn register_system(
764 &mut self,
765 mut contract: SystemContract,
766 handler: SimulationSystemHandler,
767 ) -> Result<(), CanwuError> {
768 validate_system_contract(&self.plugin, &mut contract)?;
772 if self
773 .registry
774 .descriptors
775 .get(&self.plugin)
776 .is_some_and(|descriptor| {
777 descriptor
778 .systems
779 .iter()
780 .any(|candidate| candidate.name == contract.name)
781 || descriptor
782 .boundary_systems
783 .iter()
784 .any(|candidate| candidate.name == contract.name)
785 })
786 {
787 return Err(CanwuError::new(
788 ErrorCode::DuplicatePluginSystem,
789 format!(
790 "plugin {} already registered system {}",
791 self.plugin, contract.name
792 ),
793 ));
794 }
795 let mut candidate = self.registry.clone();
796 if contract
797 .writes
798 .iter()
799 .any(|state| is_domain_record_state(&candidate.record_schemas, state))
800 {
801 return Err(CanwuError::new(
802 ErrorCode::InvalidPluginRegistration,
803 "domain record kinds can only be mutated by phased boundary systems",
804 ));
805 }
806 register_state_owners(&mut candidate.state_owners, &self.plugin, &contract.writes)?;
807 register_immediate_write_states(
808 &mut candidate.immediate_write_states,
809 &candidate.boundary_writers,
810 &self.plugin,
811 &contract.writes,
812 )?;
813 {
814 let descriptor = candidate
815 .descriptors
816 .entry(self.plugin.clone())
817 .or_default();
818 descriptor.name.clone_from(&self.plugin);
819 descriptor.systems.push(contract.clone());
820 descriptor
821 .systems
822 .sort_by(|left, right| (left.phase, &left.name).cmp(&(right.phase, &right.name)));
823 }
824 candidate.systems.push(RegisteredSystem {
825 plugin: self.plugin.clone(),
826 contract,
827 handler,
828 });
829 candidate.systems.sort_by(|left, right| {
830 (left.contract.phase, &left.plugin, &left.contract.name).cmp(&(
831 right.contract.phase,
832 &right.plugin,
833 &right.contract.name,
834 ))
835 });
836 *self.registry = candidate;
837 Ok(())
838 }
839
840 pub fn register_boundary_system(
841 &mut self,
842 mut contract: BoundarySystemContract,
843 handler: BoundarySystemHandler,
844 ) -> Result<(), CanwuError> {
845 validate_boundary_system_contract(&mut contract)?;
846 validate_knowledge_write_grants(&self.plugin, &contract, &self.registry.knowledge_schemas)?;
847 if self
848 .registry
849 .descriptors
850 .get(&self.plugin)
851 .is_some_and(|descriptor| {
852 descriptor
853 .systems
854 .iter()
855 .any(|candidate| candidate.name == contract.name)
856 || descriptor
857 .boundary_systems
858 .iter()
859 .any(|candidate| candidate.name == contract.name)
860 })
861 {
862 return Err(CanwuError::new(
863 ErrorCode::DuplicatePluginSystem,
864 format!(
865 "plugin {} already registered system {}",
866 self.plugin, contract.name
867 ),
868 ));
869 }
870 let plugin_writes = super::persons::plugin_owned_boundary_writes(&contract)?;
871 let mut owned_state = plugin_writes.clone();
872 owned_state.extend(contract.reservation_offers.iter().cloned());
873 owned_state.sort();
874 owned_state.dedup();
875 let mut candidate = self.registry.clone();
876 register_state_owners(&mut candidate.state_owners, &self.plugin, &owned_state)?;
877 register_boundary_writers(
878 &mut candidate.boundary_writers,
879 &candidate.immediate_write_states,
880 &self.plugin,
881 &contract.name,
882 contract.phase,
883 &plugin_writes,
884 )?;
885 register_reservation_offerers(
886 &mut candidate.reservation_offerers,
887 &self.plugin,
888 &contract.name,
889 &contract.reservation_offers,
890 )?;
891 register_random_streams(
892 &mut candidate.random_stream_owners,
893 &self.plugin,
894 &contract.name,
895 &contract.random_streams,
896 )?;
897 {
898 let descriptor = candidate
899 .descriptors
900 .entry(self.plugin.clone())
901 .or_default();
902 descriptor.name.clone_from(&self.plugin);
903 descriptor.boundary_systems.push(contract.clone());
904 descriptor
905 .boundary_systems
906 .sort_by(|left, right| (left.phase, &left.name).cmp(&(right.phase, &right.name)));
907 }
908 candidate.boundary_systems.push(RegisteredBoundarySystem {
909 plugin: self.plugin.clone(),
910 contract,
911 handler,
912 });
913 candidate.boundary_systems.sort_by(|left, right| {
914 (left.contract.phase, &left.plugin, &left.contract.name).cmp(&(
915 right.contract.phase,
916 &right.plugin,
917 &right.contract.name,
918 ))
919 });
920 *self.registry = candidate;
921 Ok(())
922 }
923
924 pub fn register_command(
925 &mut self,
926 mut descriptor: PluginActionDescriptor,
927 handler: PluginCommandHandler,
928 ) -> Result<(), CanwuError> {
929 validate_action_descriptor(&self.plugin, &mut descriptor)?;
930 let command_key = (self.plugin.clone(), descriptor.name.clone());
931 if self.registry.commands.contains_key(&command_key) {
932 return Err(CanwuError::new(
933 ErrorCode::DuplicatePluginCommand,
934 format!(
935 "plugin {} already registered command {}",
936 self.plugin, descriptor.name
937 ),
938 ));
939 }
940 let mut candidate = self.registry.clone();
941 if descriptor
942 .writes
943 .iter()
944 .any(|state| is_domain_record_state(&candidate.record_schemas, state))
945 {
946 return Err(CanwuError::new(
947 ErrorCode::InvalidPluginRegistration,
948 "plugin commands cannot write domain record state directly",
949 ));
950 }
951 register_state_owners(
952 &mut candidate.state_owners,
953 &self.plugin,
954 &descriptor.writes,
955 )?;
956 register_immediate_write_states(
957 &mut candidate.immediate_write_states,
958 &candidate.boundary_writers,
959 &self.plugin,
960 &descriptor.writes,
961 )?;
962 {
963 let plugin_descriptor = candidate
964 .descriptors
965 .entry(self.plugin.clone())
966 .or_default();
967 plugin_descriptor.name.clone_from(&self.plugin);
968 plugin_descriptor.commands.push(descriptor.clone());
969 plugin_descriptor
970 .commands
971 .sort_by(|left, right| left.name.cmp(&right.name));
972 }
973 candidate.commands.insert(
974 command_key,
975 RegisteredCommand {
976 descriptor,
977 handler,
978 },
979 );
980 *self.registry = candidate;
981 Ok(())
982 }
983
984 pub fn register_ingress(
985 &mut self,
986 descriptor: PluginIngressDescriptor,
987 ) -> Result<(), CanwuError> {
988 validate_ingress_descriptor(&descriptor)?;
989 let key = (self.plugin.clone(), descriptor.name.clone());
990 if self
991 .registry
992 .descriptors
993 .get(&self.plugin)
994 .is_some_and(|plugin| {
995 plugin
996 .ingress
997 .iter()
998 .any(|candidate| candidate.name == descriptor.name)
999 })
1000 {
1001 return Err(CanwuError::new(
1002 ErrorCode::DuplicatePluginIngress,
1003 format!(
1004 "plugin {} already registered ingress type {}",
1005 self.plugin, descriptor.name
1006 ),
1007 ));
1008 }
1009 if self
1010 .registry
1011 .ingress
1012 .get(&key)
1013 .is_some_and(|existing| existing != &descriptor)
1014 {
1015 return Err(CanwuError::new(
1016 ErrorCode::PluginManifestMismatch,
1017 format!(
1018 "plugin {} changed the stored ingress type {}",
1019 self.plugin, descriptor.name
1020 ),
1021 ));
1022 }
1023 let mut candidate = self.registry.clone();
1024 candidate.ingress.insert(key, descriptor.clone());
1025 let plugin_descriptor = candidate
1026 .descriptors
1027 .entry(self.plugin.clone())
1028 .or_default();
1029 plugin_descriptor.name.clone_from(&self.plugin);
1030 plugin_descriptor.ingress.push(descriptor);
1031 plugin_descriptor
1032 .ingress
1033 .sort_by(|left, right| left.name.cmp(&right.name));
1034 *self.registry = candidate;
1035 Ok(())
1036 }
1037
1038 pub fn register_internal_ingress(
1041 &mut self,
1042 descriptor: PluginIngressDescriptor,
1043 ) -> Result<PluginIngressPermit, CanwuError> {
1044 let packet_type = descriptor.name.clone();
1045 self.register_ingress(descriptor)?;
1046 let semantic_hash = self
1047 .registry
1048 .descriptors
1049 .get(&self.plugin)
1050 .map(|descriptor| descriptor.semantic_hash.clone())
1051 .ok_or_else(|| {
1052 CanwuError::new(
1053 ErrorCode::InvalidPluginRegistration,
1054 "internal ingress owner has no plugin descriptor",
1055 )
1056 })?;
1057 let token = canonical_hash(
1058 "canwu.plugin.internal-ingress-permit.v1",
1059 &(&self.plugin, &packet_type, &semantic_hash),
1060 )?;
1061 self.registry
1062 .internal_ingress
1063 .insert((self.plugin.clone(), packet_type.clone()));
1064 let plugin_descriptor =
1065 self.registry
1066 .descriptors
1067 .get_mut(&self.plugin)
1068 .ok_or_else(|| {
1069 CanwuError::new(
1070 ErrorCode::InvalidPluginRegistration,
1071 "internal ingress owner descriptor disappeared",
1072 )
1073 })?;
1074 plugin_descriptor.internal_ingress.push(packet_type.clone());
1075 plugin_descriptor.internal_ingress.sort();
1076 plugin_descriptor.internal_ingress.dedup();
1077 Ok(PluginIngressPermit {
1078 plugin: self.plugin.clone(),
1079 packet_type,
1080 semantic_hash,
1081 token,
1082 })
1083 }
1084
1085 pub fn register_internal_ingress_with_archive_retention(
1088 &mut self,
1089 descriptor: PluginIngressDescriptor,
1090 namespace: impl Into<String>,
1091 object_id_json_pointers: Vec<String>,
1092 ) -> Result<PluginIngressPermit, CanwuError> {
1093 let binding = PluginArchiveRetentionBinding {
1094 packet_type: descriptor.name.clone(),
1095 namespace: namespace.into(),
1096 object_id_json_pointers,
1097 };
1098 validate_archive_retention_binding(&binding)?;
1099 let permit = self.register_internal_ingress(descriptor)?;
1100 let key = (self.plugin.clone(), binding.packet_type.clone());
1101 if self
1102 .registry
1103 .archive_retention_bindings
1104 .get(&key)
1105 .is_some_and(|existing| existing != &binding)
1106 {
1107 return Err(CanwuError::new(
1108 ErrorCode::PluginManifestMismatch,
1109 "plugin changed its stored archive-retention binding",
1110 ));
1111 }
1112 self.registry
1113 .archive_retention_bindings
1114 .insert(key, binding.clone());
1115 let plugin_descriptor =
1116 self.registry
1117 .descriptors
1118 .get_mut(&self.plugin)
1119 .ok_or_else(|| {
1120 CanwuError::new(
1121 ErrorCode::InvalidPluginRegistration,
1122 "archive-retention owner descriptor disappeared",
1123 )
1124 })?;
1125 plugin_descriptor.archive_retention_bindings.push(binding);
1126 plugin_descriptor
1127 .archive_retention_bindings
1128 .sort_by(|left, right| left.packet_type.cmp(&right.packet_type));
1129 Ok(permit)
1130 }
1131}
1132
1133impl PluginRegistry {
1134 pub(super) fn validate_archive_retention(
1135 &self,
1136 plugin: &str,
1137 packet_type: &str,
1138 payload: &Value,
1139 retention: &[PluginArchiveRetention],
1140 ) -> Result<(), CanwuError> {
1141 let Some(binding) = self
1142 .archive_retention_bindings
1143 .get(&(plugin.to_owned(), packet_type.to_owned()))
1144 else {
1145 return if retention.is_empty() {
1146 Ok(())
1147 } else {
1148 Err(CanwuError::new(
1149 ErrorCode::InvalidPayload,
1150 "plugin ingress carries undeclared archive retention",
1151 ))
1152 };
1153 };
1154 if retention.len() != 1 || retention[0].namespace != binding.namespace {
1155 return Err(CanwuError::new(
1156 ErrorCode::InvalidPayload,
1157 "plugin ingress archive retention does not match its declared binding",
1158 ));
1159 }
1160 let object_id = &retention[0].object_id;
1161 if !is_canonical_hash(object_id)
1162 || binding.object_id_json_pointers.iter().any(|pointer| {
1163 payload.pointer(pointer).and_then(Value::as_str) != Some(object_id.as_str())
1164 })
1165 {
1166 return Err(CanwuError::new(
1167 ErrorCode::InvalidPayload,
1168 "plugin ingress archive root is not bound to its authenticated payload roots",
1169 ));
1170 }
1171 Ok(())
1172 }
1173
1174 pub fn register<P: SimulationPlugin + ?Sized>(
1175 &mut self,
1176 plugin: &P,
1177 schema: &mut SchemaRegistry,
1178 ) -> Result<(), CanwuError> {
1179 let raw_plugin_name = plugin.name();
1180 let plugin_name = raw_plugin_name.trim();
1181 if plugin_name.is_empty() || plugin_name != raw_plugin_name {
1182 return Err(CanwuError::new(
1183 ErrorCode::InvalidPluginRegistration,
1184 "plugin name must be non-empty and have no surrounding whitespace",
1185 ));
1186 }
1187 if self.active_plugins.contains(plugin_name) {
1188 return Err(CanwuError::new(
1189 ErrorCode::DuplicatePlugin,
1190 format!("plugin {plugin_name} is already registered"),
1191 ));
1192 }
1193 validate_plugin_identity(plugin_name, plugin.version(), plugin.semantic_hash())?;
1194
1195 let expected_descriptor = self.descriptors.get(plugin_name).cloned();
1196 let mut candidate_registry = self.clone();
1197 let mut candidate_schema = schema.clone();
1198 candidate_registry.descriptors.insert(
1199 plugin_name.to_owned(),
1200 PluginDescriptor {
1201 descriptor_format: PLUGIN_DESCRIPTOR_FORMAT_VERSION,
1202 name: plugin_name.to_owned(),
1203 version: plugin.version().to_owned(),
1204 semantic_hash: plugin.semantic_hash().to_owned(),
1205 ..PluginDescriptor::default()
1206 },
1207 );
1208 let mut registrar = PluginRegistrar {
1209 plugin: plugin_name.to_owned(),
1210 registry: &mut candidate_registry,
1211 schema: &mut candidate_schema,
1212 };
1213 plugin.register(&mut registrar)?;
1214 validate_schema_set(
1215 &candidate_registry.knowledge_schemas,
1216 &candidate_registry.knowledge_kind_owners,
1217 )
1218 .map_err(|error| {
1219 CanwuError::new(
1220 ErrorCode::InvalidPluginRegistration,
1221 format!("invalid knowledge schema set: {error}"),
1222 )
1223 })?;
1224 let Some(generated_descriptor) = candidate_registry.descriptors.get(plugin_name) else {
1225 return Err(CanwuError::new(
1226 ErrorCode::InvalidPluginRegistration,
1227 format!("plugin {plugin_name} did not produce a descriptor"),
1228 ));
1229 };
1230 if let Some(expected) = expected_descriptor
1231 && generated_descriptor != &expected
1232 {
1233 return Err(CanwuError::new(
1234 ErrorCode::PluginManifestMismatch,
1235 format!("plugin {plugin_name} registration does not match the snapshot manifest"),
1236 ));
1237 }
1238 candidate_registry
1239 .active_plugins
1240 .insert(plugin_name.to_owned());
1241 *self = candidate_registry;
1242 *schema = candidate_schema;
1243 Ok(())
1244 }
1245
1246 pub fn descriptors(&self) -> impl Iterator<Item = &PluginDescriptor> {
1247 self.descriptors.values()
1248 }
1249
1250 pub(super) fn event_audience(&self, plugin: &str, event_type: &str) -> EventAudience {
1251 self.descriptors
1252 .get(plugin)
1253 .and_then(|descriptor| descriptor.event_audiences.get(event_type))
1254 .cloned()
1255 .unwrap_or_default()
1256 }
1257
1258 pub(super) fn from_descriptors(descriptors: Vec<PluginDescriptor>) -> Result<Self, CanwuError> {
1259 let mut registry = Self {
1260 descriptors: BTreeMap::new(),
1261 active_plugins: BTreeSet::new(),
1262 systems: Vec::new(),
1263 boundary_systems: Vec::new(),
1264 commands: BTreeMap::new(),
1265 ingress: BTreeMap::new(),
1266 internal_ingress: BTreeSet::new(),
1267 archive_retention_bindings: BTreeMap::new(),
1268 state_owners: BTreeMap::new(),
1269 immediate_write_states: BTreeMap::new(),
1270 boundary_writers: BTreeMap::new(),
1271 reservation_offerers: BTreeMap::new(),
1272 random_stream_owners: BTreeMap::new(),
1273 record_schemas: BTreeMap::new(),
1274 knowledge_schemas: BTreeMap::new(),
1275 knowledge_kind_owners: BTreeMap::new(),
1276 maintenance_dependency_resolvers: BTreeMap::new(),
1277 maintenance_participants: BTreeMap::new(),
1278 archive_reachability_participants: BTreeMap::new(),
1279 };
1280 let mut previous_plugin = None;
1281 for mut descriptor in descriptors {
1282 let plugin = descriptor.name.trim().to_owned();
1283 if plugin.is_empty()
1284 || descriptor.descriptor_format != PLUGIN_DESCRIPTOR_FORMAT_VERSION
1285 || descriptor.name != plugin
1286 || descriptor.version.trim().is_empty()
1287 || descriptor.version != descriptor.version.trim()
1288 || !is_canonical_hash(&descriptor.semantic_hash)
1289 || registry.descriptors.contains_key(&plugin)
1290 || previous_plugin
1291 .as_ref()
1292 .is_some_and(|previous| previous >= &plugin)
1293 {
1294 return Err(CanwuError::new(
1295 ErrorCode::InvalidSnapshot,
1296 "snapshot contains an invalid, unversioned, or duplicate plugin descriptor",
1297 ));
1298 }
1299 if descriptor
1300 .record_schemas
1301 .windows(2)
1302 .any(|pair| pair[0].kind >= pair[1].kind)
1303 {
1304 return invalid_snapshot("plugin record schemas are not in canonical order");
1305 }
1306 for schema in &mut descriptor.record_schemas {
1307 let original = schema.clone();
1308 schema.canonicalize();
1309 schema.validate().map_err(|error| {
1310 invalid_snapshot_error(format!("invalid domain record schema: {error}"))
1311 })?;
1312 if *schema != original {
1313 return invalid_snapshot(
1314 "plugin record-schema declarations are not in canonical order",
1315 );
1316 }
1317 let state = schema.state_key();
1318 if state.namespace == CORE_STATE_NAMESPACE {
1319 return invalid_snapshot(
1320 "plugin record schemas cannot use the reserved core namespace",
1321 );
1322 }
1323 if let Some((owner, _)) = registry.record_schemas.get(&schema.kind) {
1324 return invalid_snapshot(format!(
1325 "domain record kind {} is owned by both {owner} and {plugin}",
1326 schema.kind
1327 ));
1328 }
1329 register_state_owners(
1330 &mut registry.state_owners,
1331 &plugin,
1332 std::slice::from_ref(&state),
1333 )
1334 .map_err(|error| {
1335 invalid_snapshot_error(format!(
1336 "invalid domain record state ownership descriptor: {error}"
1337 ))
1338 })?;
1339 registry
1340 .record_schemas
1341 .insert(schema.kind.clone(), (plugin.clone(), schema.clone()));
1342 }
1343 if descriptor.knowledge_schemas.len() > KnowledgeLimitsV1::CURRENT.schemas_per_plugin
1344 || descriptor
1345 .knowledge_schemas
1346 .windows(2)
1347 .any(|pair| pair[0].id >= pair[1].id)
1348 {
1349 return invalid_snapshot(
1350 "plugin knowledge schemas are not in canonical order or exceed their limit",
1351 );
1352 }
1353 for schema in &mut descriptor.knowledge_schemas {
1354 let original = schema.clone();
1355 schema.canonicalize();
1356 schema.validate().map_err(|error| {
1357 invalid_snapshot_error(format!("invalid knowledge schema: {error}"))
1358 })?;
1359 if *schema != original {
1360 return invalid_snapshot(
1361 "plugin knowledge-schema declarations are not in canonical order",
1362 );
1363 }
1364 if let Some(owner) = registry.knowledge_kind_owners.get(&schema.id.kind) {
1365 if owner != &plugin {
1366 return invalid_snapshot(format!(
1367 "knowledge kind {:?} is owned by both {owner} and {plugin}",
1368 schema.id.kind
1369 ));
1370 }
1371 } else {
1372 registry
1373 .knowledge_kind_owners
1374 .insert(schema.id.kind.clone(), plugin.clone());
1375 }
1376 if registry
1377 .knowledge_schemas
1378 .insert(schema.id.clone(), (plugin.clone(), schema.clone()))
1379 .is_some()
1380 {
1381 return invalid_snapshot("knowledge schema ID is duplicated");
1382 }
1383 }
1384 if descriptor
1385 .maintenance_dependency_resolvers
1386 .windows(2)
1387 .any(|pair| pair[0] >= pair[1])
1388 {
1389 return invalid_snapshot(
1390 "plugin maintenance dependency resolvers are not in canonical order",
1391 );
1392 }
1393 for resolver in &descriptor.maintenance_dependency_resolvers {
1394 if !canonical_text(&resolver.target_namespace)
1395 || resolver.target_namespace == CORE_STATE_NAMESPACE
1396 {
1397 return invalid_snapshot(
1398 "plugin maintenance dependency resolver target is invalid",
1399 );
1400 }
1401 registry
1402 .maintenance_dependency_resolvers
1403 .entry(resolver.target_namespace.clone())
1404 .or_default()
1405 .insert(plugin.clone());
1406 }
1407 if descriptor
1408 .systems
1409 .windows(2)
1410 .any(|pair| (pair[0].phase, &pair[0].name) >= (pair[1].phase, &pair[1].name))
1411 {
1412 return invalid_snapshot("plugin systems are not in canonical order");
1413 }
1414 let mut system_names = BTreeSet::new();
1415 for contract in &mut descriptor.systems {
1416 if !system_names.insert(contract.name.clone()) {
1417 return invalid_snapshot("plugin descriptor has duplicate system names");
1418 }
1419 let original = contract.clone();
1420 validate_system_contract(&plugin, contract).map_err(|error| {
1421 invalid_snapshot_error(format!("invalid plugin system descriptor: {error}"))
1422 })?;
1423 if *contract != original {
1424 return invalid_snapshot(
1425 "plugin system reads and writes are not in canonical order",
1426 );
1427 }
1428 if contract
1429 .writes
1430 .iter()
1431 .any(|state| is_domain_record_state(®istry.record_schemas, state))
1432 {
1433 return invalid_snapshot(
1434 "plugin systems cannot expose domain records as immediate component state",
1435 );
1436 }
1437 register_state_owners(&mut registry.state_owners, &plugin, &contract.writes)
1438 .map_err(|error| {
1439 invalid_snapshot_error(format!(
1440 "invalid plugin state ownership descriptor: {error}"
1441 ))
1442 })?;
1443 register_immediate_write_states(
1444 &mut registry.immediate_write_states,
1445 ®istry.boundary_writers,
1446 &plugin,
1447 &contract.writes,
1448 )
1449 .map_err(|error| {
1450 invalid_snapshot_error(format!(
1451 "invalid immediate state writer descriptor: {error}"
1452 ))
1453 })?;
1454 }
1455 if descriptor
1456 .boundary_systems
1457 .windows(2)
1458 .any(|pair| (pair[0].phase, &pair[0].name) >= (pair[1].phase, &pair[1].name))
1459 {
1460 return invalid_snapshot("boundary systems are not in canonical order");
1461 }
1462 for contract in &mut descriptor.boundary_systems {
1463 if !system_names.insert(contract.name.clone()) {
1464 return invalid_snapshot("plugin descriptor has duplicate system names");
1465 }
1466 let original = contract.clone();
1467 validate_boundary_system_contract(contract).map_err(|error| {
1468 invalid_snapshot_error(format!("invalid boundary system descriptor: {error}"))
1469 })?;
1470 validate_knowledge_write_grants(&plugin, contract, ®istry.knowledge_schemas)
1471 .map_err(|error| {
1472 invalid_snapshot_error(format!(
1473 "invalid boundary knowledge writer descriptor: {error}"
1474 ))
1475 })?;
1476 if *contract != original {
1477 return invalid_snapshot(
1478 "boundary system declarations are not in canonical order",
1479 );
1480 }
1481 let plugin_writes = super::persons::plugin_owned_boundary_writes(contract)
1482 .map_err(|error| {
1483 invalid_snapshot_error(format!(
1484 "invalid boundary core-write descriptor: {error}"
1485 ))
1486 })?;
1487 let mut owned_state = plugin_writes.clone();
1488 owned_state.extend(contract.reservation_offers.iter().cloned());
1489 owned_state.sort();
1490 owned_state.dedup();
1491 register_state_owners(&mut registry.state_owners, &plugin, &owned_state).map_err(
1492 |error| {
1493 invalid_snapshot_error(format!(
1494 "invalid boundary state ownership descriptor: {error}"
1495 ))
1496 },
1497 )?;
1498 register_boundary_writers(
1499 &mut registry.boundary_writers,
1500 ®istry.immediate_write_states,
1501 &plugin,
1502 &contract.name,
1503 contract.phase,
1504 &plugin_writes,
1505 )
1506 .map_err(|error| {
1507 invalid_snapshot_error(format!("invalid boundary writer descriptor: {error}"))
1508 })?;
1509 register_reservation_offerers(
1510 &mut registry.reservation_offerers,
1511 &plugin,
1512 &contract.name,
1513 &contract.reservation_offers,
1514 )
1515 .map_err(|error| {
1516 invalid_snapshot_error(format!(
1517 "invalid reservation offerer descriptor: {error}"
1518 ))
1519 })?;
1520 register_random_streams(
1521 &mut registry.random_stream_owners,
1522 &plugin,
1523 &contract.name,
1524 &contract.random_streams,
1525 )
1526 .map_err(|error| {
1527 invalid_snapshot_error(format!(
1528 "invalid random stream ownership descriptor: {error}"
1529 ))
1530 })?;
1531 }
1532 if descriptor
1533 .commands
1534 .windows(2)
1535 .any(|pair| pair[0].name >= pair[1].name)
1536 {
1537 return invalid_snapshot("plugin commands are not in canonical order");
1538 }
1539 let mut command_names = BTreeSet::new();
1540 for action in &mut descriptor.commands {
1541 if !command_names.insert(action.name.clone()) {
1542 return invalid_snapshot("plugin descriptor has duplicate command names");
1543 }
1544 let original = action.clone();
1545 validate_action_descriptor(&plugin, action).map_err(|error| {
1546 invalid_snapshot_error(format!("invalid plugin command descriptor: {error}"))
1547 })?;
1548 if *action != original {
1549 return invalid_snapshot(
1550 "plugin command reads and writes are not in canonical order",
1551 );
1552 }
1553 if action
1554 .writes
1555 .iter()
1556 .any(|state| is_domain_record_state(®istry.record_schemas, state))
1557 {
1558 return invalid_snapshot(
1559 "plugin commands cannot expose domain records as immediate component state",
1560 );
1561 }
1562 register_state_owners(&mut registry.state_owners, &plugin, &action.writes)
1563 .map_err(|error| {
1564 invalid_snapshot_error(format!(
1565 "invalid plugin state ownership descriptor: {error}"
1566 ))
1567 })?;
1568 register_immediate_write_states(
1569 &mut registry.immediate_write_states,
1570 ®istry.boundary_writers,
1571 &plugin,
1572 &action.writes,
1573 )
1574 .map_err(|error| {
1575 invalid_snapshot_error(format!(
1576 "invalid immediate state writer descriptor: {error}"
1577 ))
1578 })?;
1579 }
1580 if descriptor
1581 .ingress
1582 .windows(2)
1583 .any(|pair| pair[0].name >= pair[1].name)
1584 {
1585 return invalid_snapshot("plugin ingress types are not in canonical order");
1586 }
1587 for ingress in &descriptor.ingress {
1588 validate_ingress_descriptor(ingress).map_err(|error| {
1589 invalid_snapshot_error(format!("invalid plugin ingress descriptor: {error}"))
1590 })?;
1591 if registry
1592 .ingress
1593 .insert((plugin.clone(), ingress.name.clone()), ingress.clone())
1594 .is_some()
1595 {
1596 return invalid_snapshot("plugin descriptor has duplicate ingress types");
1597 }
1598 }
1599 if descriptor
1600 .internal_ingress
1601 .windows(2)
1602 .any(|pair| pair[0] >= pair[1])
1603 || descriptor.internal_ingress.iter().any(|name| {
1604 !descriptor
1605 .ingress
1606 .iter()
1607 .any(|ingress| &ingress.name == name)
1608 })
1609 {
1610 return invalid_snapshot("plugin internal ingress declarations are not canonical");
1611 }
1612 registry.internal_ingress.extend(
1613 descriptor
1614 .internal_ingress
1615 .iter()
1616 .map(|name| (plugin.clone(), name.clone())),
1617 );
1618 if descriptor
1619 .archive_retention_bindings
1620 .windows(2)
1621 .any(|pair| pair[0].packet_type >= pair[1].packet_type)
1622 {
1623 return invalid_snapshot(
1624 "plugin archive-retention bindings are not in canonical order",
1625 );
1626 }
1627 for binding in &descriptor.archive_retention_bindings {
1628 validate_archive_retention_binding(binding).map_err(|error| {
1629 invalid_snapshot_error(format!(
1630 "invalid plugin archive-retention binding: {error}"
1631 ))
1632 })?;
1633 if !descriptor
1634 .internal_ingress
1635 .iter()
1636 .any(|name| name == &binding.packet_type)
1637 || registry
1638 .archive_retention_bindings
1639 .insert(
1640 (plugin.clone(), binding.packet_type.clone()),
1641 binding.clone(),
1642 )
1643 .is_some()
1644 {
1645 return invalid_snapshot(
1646 "plugin archive-retention binding lacks one internal ingress owner",
1647 );
1648 }
1649 }
1650 for (event_type, audience) in &descriptor.event_audiences {
1651 validate_event_audience_name(event_type).map_err(|error| {
1652 invalid_snapshot_error(format!("invalid plugin event audience: {error}"))
1653 })?;
1654 validate_event_audience(audience).map_err(|error| {
1655 invalid_snapshot_error(format!("invalid plugin event audience: {error}"))
1656 })?;
1657 }
1658 let schema_types: BTreeSet<_> = descriptor.schema_types.iter().collect();
1659 if schema_types.len() != descriptor.schema_types.len()
1660 || descriptor
1661 .schema_types
1662 .windows(2)
1663 .any(|pair| pair[0] >= pair[1])
1664 || descriptor
1665 .schema_types
1666 .iter()
1667 .any(|name| name.trim().is_empty() || name != name.trim())
1668 {
1669 return invalid_snapshot("plugin descriptor has invalid schema type names");
1670 }
1671 previous_plugin = Some(plugin.clone());
1672 registry.descriptors.insert(plugin, descriptor);
1673 }
1674 validate_schema_set(®istry.knowledge_schemas, ®istry.knowledge_kind_owners).map_err(
1675 |error| invalid_snapshot_error(format!("invalid knowledge schema set: {error}")),
1676 )?;
1677 Ok(registry)
1678 }
1679
1680 pub(super) fn ensure_active(&self) -> Result<(), CanwuError> {
1681 let inactive: Vec<_> = self
1682 .descriptors
1683 .keys()
1684 .filter(|name| !self.active_plugins.contains(*name))
1685 .cloned()
1686 .collect();
1687 if inactive.is_empty() {
1688 return Ok(());
1689 }
1690 Err(CanwuError::new(
1691 ErrorCode::PluginNotActive,
1692 format!(
1693 "required plugin handlers are not active: {}",
1694 inactive.join(", ")
1695 ),
1696 ))
1697 }
1698}
1699
1700fn validate_archive_retention_binding(
1701 binding: &PluginArchiveRetentionBinding,
1702) -> Result<(), CanwuError> {
1703 if binding.packet_type.trim().is_empty()
1704 || binding.packet_type != binding.packet_type.trim()
1705 || binding.namespace.is_empty()
1706 || binding.namespace.len() > 128
1707 || !binding.namespace.bytes().all(|byte| {
1708 byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'.' | b'-' | b'_')
1709 })
1710 || binding.object_id_json_pointers.is_empty()
1711 || binding
1712 .object_id_json_pointers
1713 .windows(2)
1714 .any(|pair| pair[0] >= pair[1])
1715 || binding
1716 .object_id_json_pointers
1717 .iter()
1718 .any(|pointer| !pointer.starts_with('/') || pointer.len() > 256 || !pointer.is_ascii())
1719 {
1720 return Err(CanwuError::new(
1721 ErrorCode::InvalidPluginRegistration,
1722 "plugin archive-retention binding is malformed",
1723 ));
1724 }
1725 Ok(())
1726}
1727
1728pub(super) fn validate_state_keys(keys: &mut Vec<StateKey>) -> Result<(), CanwuError> {
1729 for key in keys.iter() {
1730 if key.namespace.trim().is_empty()
1731 || key.name.trim().is_empty()
1732 || key.namespace != key.namespace.trim()
1733 || key.name != key.name.trim()
1734 {
1735 return Err(CanwuError::new(
1736 ErrorCode::InvalidPluginRegistration,
1737 "state keys require non-empty canonical namespace and name values",
1738 ));
1739 }
1740 }
1741 let unique: BTreeSet<_> = keys.drain(..).collect();
1742 keys.extend(unique);
1743 Ok(())
1744}
1745
1746fn validate_plugin_identity(
1747 name: &str,
1748 version: &str,
1749 semantic_hash: &str,
1750) -> Result<(), CanwuError> {
1751 if name.trim().is_empty()
1752 || name != name.trim()
1753 || version.trim().is_empty()
1754 || version != version.trim()
1755 || !is_canonical_hash(semantic_hash)
1756 {
1757 return Err(CanwuError::new(
1758 ErrorCode::InvalidPluginRegistration,
1759 "plugins require canonical names, versions, and 64-character semantic hashes",
1760 ));
1761 }
1762 Ok(())
1763}
1764
1765fn validate_system_contract(
1766 _plugin: &str,
1767 contract: &mut SystemContract,
1768) -> Result<(), CanwuError> {
1769 if contract.name.trim().is_empty() || contract.name != contract.name.trim() {
1770 return Err(CanwuError::new(
1771 ErrorCode::InvalidPluginRegistration,
1772 "plugin system name must be non-empty and have no surrounding whitespace",
1773 ));
1774 }
1775 if matches!(
1776 contract.phase,
1777 BoundaryPhase::EventIngress
1778 | BoundaryPhase::BoundarySnapshot
1779 | BoundaryPhase::AtomicDomainCommit
1780 | BoundaryPhase::ConditionalTransitionCommit
1781 ) {
1782 return Err(CanwuError::new(
1783 ErrorCode::InvalidPluginRegistration,
1784 format!("boundary phase {:?} is owned by the kernel", contract.phase),
1785 ));
1786 }
1787 if contract.cadence != SystemCadence::EventDriven {
1788 return Err(CanwuError::new(
1789 ErrorCode::InvalidPluginRegistration,
1790 format!(
1791 "system {} declares {:?} cadence, but the current runtime systems are event-driven only",
1792 contract.name, contract.cadence
1793 ),
1794 ));
1795 }
1796 if contract.visibility != StateVisibility::SameBoundary {
1797 return Err(CanwuError::new(
1798 ErrorCode::InvalidPluginRegistration,
1799 format!(
1800 "event-driven system {} must declare same-boundary visibility until the phased boundary runtime is active",
1801 contract.name
1802 ),
1803 ));
1804 }
1805 validate_state_keys(&mut contract.reads)?;
1806 validate_state_keys(&mut contract.writes)?;
1807 if contract.reads.contains(&StateKey::core_ingress()) {
1808 return Err(CanwuError::new(
1809 ErrorCode::InvalidPluginRegistration,
1810 "canonical ingress can be read only by phased boundary systems",
1811 ));
1812 }
1813 Ok(())
1814}
1815
1816fn validate_action_descriptor(
1817 _plugin: &str,
1818 descriptor: &mut PluginActionDescriptor,
1819) -> Result<(), CanwuError> {
1820 if descriptor.name.trim().is_empty() || descriptor.name != descriptor.name.trim() {
1821 return Err(CanwuError::new(
1822 ErrorCode::InvalidPluginRegistration,
1823 "plugin command names must be non-empty and have no surrounding whitespace",
1824 ));
1825 }
1826 if let PayloadSchema::Object { properties, .. } = &descriptor.payload_schema
1827 && properties
1828 .keys()
1829 .any(|name| name.trim().is_empty() || name != name.trim())
1830 {
1831 return Err(CanwuError::new(
1832 ErrorCode::InvalidPluginRegistration,
1833 "plugin payload schema property names cannot be empty",
1834 ));
1835 }
1836 validate_state_keys(&mut descriptor.reads)?;
1837 validate_state_keys(&mut descriptor.writes)?;
1838 if descriptor.reads.contains(&StateKey::core_ingress()) {
1839 return Err(CanwuError::new(
1840 ErrorCode::InvalidPluginRegistration,
1841 "plugin commands cannot inspect the canonical ingress queue",
1842 ));
1843 }
1844 Ok(())
1845}
1846
1847fn validate_ingress_descriptor(descriptor: &PluginIngressDescriptor) -> Result<(), CanwuError> {
1848 if descriptor.name.trim().is_empty()
1849 || descriptor.name != descriptor.name.trim()
1850 || descriptor.description.trim().is_empty()
1851 || descriptor.description != descriptor.description.trim()
1852 || descriptor.class == IngressClass::Command
1853 {
1854 return Err(CanwuError::new(
1855 ErrorCode::InvalidPluginRegistration,
1856 "plugin ingress types require canonical names/descriptions and cannot claim the core command class",
1857 ));
1858 }
1859 if let PayloadSchema::Object { properties, .. } = &descriptor.payload_schema
1860 && properties
1861 .keys()
1862 .any(|name| name.trim().is_empty() || name != name.trim())
1863 {
1864 return Err(CanwuError::new(
1865 ErrorCode::InvalidPluginRegistration,
1866 "plugin ingress payload property names cannot be empty",
1867 ));
1868 }
1869 Ok(())
1870}
1871
1872fn validate_boundary_system_contract(
1873 contract: &mut BoundarySystemContract,
1874) -> Result<(), CanwuError> {
1875 if contract.name.trim().is_empty() || contract.name != contract.name.trim() {
1876 return Err(CanwuError::new(
1877 ErrorCode::InvalidPluginRegistration,
1878 "boundary system name must be non-empty and canonical",
1879 ));
1880 }
1881 validate_state_keys(&mut contract.reads)?;
1882 validate_state_keys(&mut contract.writes)?;
1883 validate_state_keys(&mut contract.reservation_offers)?;
1884 validate_state_keys(&mut contract.reservation_requests)?;
1885 validate_reservation_refs(&mut contract.reservation_reads)?;
1886 validate_random_stream_keys(&mut contract.random_streams)?;
1887 validate_canonical_names(&mut contract.emits, "boundary event type")?;
1888 for grant in &mut contract.knowledge_writes {
1889 grant.visibilities.sort();
1890 grant.visibilities.dedup();
1891 if grant.visibilities.is_empty() {
1892 return Err(CanwuError::new(
1893 ErrorCode::InvalidPluginRegistration,
1894 "knowledge write grants require at least one visibility",
1895 ));
1896 }
1897 }
1898 contract
1899 .knowledge_writes
1900 .sort_by(|left, right| left.schema.cmp(&right.schema));
1901 if contract
1902 .knowledge_writes
1903 .windows(2)
1904 .any(|pair| pair[0].schema >= pair[1].schema)
1905 {
1906 return Err(CanwuError::new(
1907 ErrorCode::InvalidPluginRegistration,
1908 "knowledge write grants must name unique schemas in canonical order",
1909 ));
1910 }
1911 if !contract.knowledge_writes.is_empty()
1912 && !matches!(
1913 contract.phase,
1914 BoundaryPhase::PerceptionAndAttentionRefresh
1915 | BoundaryPhase::PerspectiveAndReportMaterialization
1916 )
1917 {
1918 return Err(CanwuError::new(
1919 ErrorCode::InvalidPluginRegistration,
1920 "knowledge publication is available only in phases 4 and 13",
1921 ));
1922 }
1923 if contract.plugin_ingress_targets.iter().any(|target| {
1924 target.target_plugin.trim().is_empty()
1925 || target.target_plugin != target.target_plugin.trim()
1926 || target.packet_type.trim().is_empty()
1927 || target.packet_type != target.packet_type.trim()
1928 }) {
1929 return Err(CanwuError::new(
1930 ErrorCode::InvalidPluginRegistration,
1931 "cross-plugin ingress targets require canonical plugin and packet names",
1932 ));
1933 }
1934 contract.plugin_ingress_targets.sort();
1935 if contract
1936 .plugin_ingress_targets
1937 .windows(2)
1938 .any(|pair| pair[0] >= pair[1])
1939 {
1940 return Err(CanwuError::new(
1941 ErrorCode::InvalidPluginRegistration,
1942 "cross-plugin ingress targets must be unique",
1943 ));
1944 }
1945
1946 let may_propose_changes = matches!(
1947 contract.phase,
1948 BoundaryPhase::DomainDeltaProposal
1949 | BoundaryPhase::HistoricalCandidateEvaluation
1950 | BoundaryPhase::StrategicAggregation
1951 | BoundaryPhase::PerspectiveAndReportMaterialization
1952 );
1953 if (!contract.writes.is_empty()
1954 || !contract.emits.is_empty()
1955 || !contract.plugin_ingress_targets.is_empty())
1956 && !may_propose_changes
1957 {
1958 return Err(CanwuError::new(
1959 ErrorCode::InvalidPluginRegistration,
1960 format!(
1961 "boundary system {} declares changes in kernel-owned phase {:?}",
1962 contract.name, contract.phase
1963 ),
1964 ));
1965 }
1966 let declares_reservations =
1967 !contract.reservation_offers.is_empty() || !contract.reservation_requests.is_empty();
1968 if declares_reservations && contract.phase != BoundaryPhase::ReservationAndAllocation {
1969 return Err(CanwuError::new(
1970 ErrorCode::InvalidPluginRegistration,
1971 format!(
1972 "boundary system {} declares reservations outside reservation and allocation",
1973 contract.name
1974 ),
1975 ));
1976 }
1977 if !contract.reservation_reads.is_empty()
1978 && contract.phase <= BoundaryPhase::ReservationAndAllocation
1979 {
1980 return Err(CanwuError::new(
1981 ErrorCode::InvalidPluginRegistration,
1982 format!(
1983 "boundary system {} reads allocations before reservation commit",
1984 contract.name
1985 ),
1986 ));
1987 }
1988 Ok(())
1989}
1990
1991fn validate_knowledge_write_grants(
1992 plugin: &str,
1993 contract: &BoundarySystemContract,
1994 schemas: &super::knowledge::KnowledgeSchemas,
1995) -> Result<(), CanwuError> {
1996 for grant in &contract.knowledge_writes {
1997 let Some((owner, schema)) = schemas.get(&grant.schema) else {
1998 return Err(CanwuError::new(
1999 ErrorCode::InvalidPluginRegistration,
2000 format!(
2001 "boundary system {plugin}.{} names an unregistered knowledge schema",
2002 contract.name
2003 ),
2004 ));
2005 };
2006 if owner != plugin || !schema.writable {
2007 return Err(CanwuError::new(
2008 ErrorCode::InvalidPluginRegistration,
2009 format!(
2010 "boundary system {plugin}.{} cannot write a foreign or read-only knowledge schema",
2011 contract.name
2012 ),
2013 ));
2014 }
2015 }
2016 Ok(())
2017}
2018
2019fn validate_reservation_refs(values: &mut Vec<ReservationRef>) -> Result<(), CanwuError> {
2020 if values.iter().any(|reservation| {
2021 reservation.plugin.trim().is_empty()
2022 || reservation.plugin != reservation.plugin.trim()
2023 || reservation.system.trim().is_empty()
2024 || reservation.system != reservation.system.trim()
2025 || reservation.request.trim().is_empty()
2026 || reservation.request != reservation.request.trim()
2027 }) {
2028 return Err(CanwuError::new(
2029 ErrorCode::InvalidPluginRegistration,
2030 "reservation read declarations must be non-empty and canonical",
2031 ));
2032 }
2033 let unique: BTreeSet<_> = values.drain(..).collect();
2034 values.extend(unique);
2035 Ok(())
2036}
2037
2038fn validate_random_stream_keys(values: &mut Vec<RandomStreamKey>) -> Result<(), CanwuError> {
2039 if values.iter().any(|stream| {
2040 stream.namespace.trim().is_empty()
2041 || stream.namespace != stream.namespace.trim()
2042 || stream.name.trim().is_empty()
2043 || stream.name != stream.name.trim()
2044 || stream.version == 0
2045 }) {
2046 return Err(CanwuError::new(
2047 ErrorCode::InvalidPluginRegistration,
2048 "random stream declarations require canonical names and a nonzero version",
2049 ));
2050 }
2051 let unique: BTreeSet<_> = values.drain(..).collect();
2052 values.extend(unique);
2053 Ok(())
2054}
2055
2056fn validate_canonical_names(values: &mut Vec<String>, label: &str) -> Result<(), CanwuError> {
2057 if values
2058 .iter()
2059 .any(|value| value.trim().is_empty() || value != value.trim())
2060 {
2061 return Err(CanwuError::new(
2062 ErrorCode::InvalidPluginRegistration,
2063 format!("{label} declarations must be non-empty and canonical"),
2064 ));
2065 }
2066 let unique: BTreeSet<_> = values.drain(..).collect();
2067 values.extend(unique);
2068 Ok(())
2069}
2070
2071fn validate_event_audience_name(event_type: &str) -> Result<(), CanwuError> {
2072 if !canonical_text(event_type) {
2073 return Err(CanwuError::new(
2074 ErrorCode::InvalidPluginRegistration,
2075 "plugin event audience names must be non-empty and canonical",
2076 ));
2077 }
2078 Ok(())
2079}
2080
2081fn validate_event_audience(audience: &EventAudience) -> Result<(), CanwuError> {
2082 match audience {
2083 EventAudience::Actor(actor) if actor.get() == 0 => {
2084 return Err(CanwuError::new(
2085 ErrorCode::InvalidPluginRegistration,
2086 "plugin event audience actors must use positive actor IDs",
2087 ));
2088 }
2089 EventAudience::Actors(actors) => {
2090 if actors.is_empty() || actors.iter().any(|actor| actor.get() == 0) {
2091 return Err(CanwuError::new(
2092 ErrorCode::InvalidPluginRegistration,
2093 "plugin event audience actor lists must contain positive actor IDs",
2094 ));
2095 }
2096 if actors.windows(2).any(|pair| pair[0] >= pair[1]) {
2097 return Err(CanwuError::new(
2098 ErrorCode::InvalidPluginRegistration,
2099 "plugin event audience actor lists must be sorted and unique",
2100 ));
2101 }
2102 }
2103 _ => {}
2104 }
2105 Ok(())
2106}
2107
2108fn register_state_owners(
2109 owners: &mut BTreeMap<StateKey, String>,
2110 plugin: &str,
2111 writes: &[StateKey],
2112) -> Result<(), CanwuError> {
2113 for key in writes {
2114 if key.namespace == CORE_STATE_NAMESPACE {
2115 return Err(CanwuError::new(
2116 ErrorCode::InvalidPluginRegistration,
2117 format!(
2118 "plugin {plugin} cannot claim reserved state {}.{}",
2119 key.namespace, key.name
2120 ),
2121 ));
2122 }
2123 if let Some(existing) = owners.get(key)
2124 && existing != plugin
2125 {
2126 return Err(CanwuError::new(
2127 ErrorCode::DuplicateStateOwner,
2128 format!(
2129 "state {}.{} is owned by both {existing} and {plugin}",
2130 key.namespace, key.name
2131 ),
2132 ));
2133 }
2134 }
2135 for key in writes {
2136 owners.insert(key.clone(), plugin.to_owned());
2137 }
2138 Ok(())
2139}
2140
2141fn register_boundary_writers(
2142 writers: &mut BTreeMap<(BoundaryWriteStage, StateKey), (String, String)>,
2143 immediate_writes: &BTreeMap<StateKey, String>,
2144 plugin: &str,
2145 system: &str,
2146 phase: BoundaryPhase,
2147 declared_states: &[StateKey],
2148) -> Result<(), CanwuError> {
2149 let Some(stage) = boundary_write_stage(phase) else {
2150 if declared_states.is_empty() {
2151 return Ok(());
2152 }
2153 return Err(CanwuError::new(
2154 ErrorCode::InvalidPluginRegistration,
2155 format!("boundary phase {phase:?} cannot own state writes"),
2156 ));
2157 };
2158 for state in declared_states {
2159 if let Some(immediate_plugin) = immediate_writes.get(state) {
2160 return Err(CanwuError::new(
2161 ErrorCode::InvalidPluginRegistration,
2162 format!(
2163 "boundary state {}.{} conflicts with immediate writes from plugin {immediate_plugin}",
2164 state.namespace, state.name
2165 ),
2166 ));
2167 }
2168 if let Some((existing_plugin, existing_system)) = writers.get(&(stage, state.clone()))
2169 && (existing_plugin != plugin || existing_system != system)
2170 {
2171 return Err(CanwuError::new(
2172 ErrorCode::DuplicateBoundaryWriter,
2173 format!(
2174 "boundary state {}.{} is written by both {existing_plugin}.{existing_system} and {plugin}.{system}",
2175 state.namespace, state.name
2176 ),
2177 ));
2178 }
2179 }
2180 for state in declared_states {
2181 writers.insert(
2182 (stage, state.clone()),
2183 (plugin.to_owned(), system.to_owned()),
2184 );
2185 }
2186 Ok(())
2187}
2188
2189fn register_immediate_write_states(
2190 immediate_writes: &mut BTreeMap<StateKey, String>,
2191 boundary_writers: &BTreeMap<(BoundaryWriteStage, StateKey), (String, String)>,
2192 plugin: &str,
2193 writes: &[StateKey],
2194) -> Result<(), CanwuError> {
2195 for state in writes {
2196 if boundary_writers
2197 .keys()
2198 .any(|(_, boundary_state)| boundary_state == state)
2199 {
2200 return Err(CanwuError::new(
2201 ErrorCode::InvalidPluginRegistration,
2202 format!(
2203 "immediate state {}.{} conflicts with a phased boundary writer",
2204 state.namespace, state.name
2205 ),
2206 ));
2207 }
2208 if immediate_writes
2209 .get(state)
2210 .is_some_and(|existing| existing != plugin)
2211 {
2212 return Err(CanwuError::new(
2213 ErrorCode::DuplicateStateOwner,
2214 format!(
2215 "immediate state {}.{} is written by multiple plugins",
2216 state.namespace, state.name
2217 ),
2218 ));
2219 }
2220 }
2221 for state in writes {
2222 immediate_writes.insert(state.clone(), plugin.to_owned());
2223 }
2224 Ok(())
2225}
2226
2227fn register_reservation_offerers(
2228 offerers: &mut BTreeMap<StateKey, (String, String)>,
2229 plugin: &str,
2230 system: &str,
2231 offered_state: &[StateKey],
2232) -> Result<(), CanwuError> {
2233 for state in offered_state {
2234 if let Some((existing_plugin, existing_system)) = offerers.get(state)
2235 && (existing_plugin != plugin || existing_system != system)
2236 {
2237 return Err(CanwuError::new(
2238 ErrorCode::DuplicateReservationOfferer,
2239 format!(
2240 "reservation state {}.{} is offered by both {existing_plugin}.{existing_system} and {plugin}.{system}",
2241 state.namespace, state.name
2242 ),
2243 ));
2244 }
2245 }
2246 for state in offered_state {
2247 offerers.insert(state.clone(), (plugin.to_owned(), system.to_owned()));
2248 }
2249 Ok(())
2250}
2251
2252fn register_random_streams(
2253 owners: &mut BTreeMap<RandomStreamKey, (String, String)>,
2254 plugin: &str,
2255 system: &str,
2256 streams: &[RandomStreamKey],
2257) -> Result<(), CanwuError> {
2258 for stream in streams {
2259 if stream.namespace != plugin || stream.namespace == CORE_STATE_NAMESPACE {
2260 return Err(CanwuError::new(
2261 ErrorCode::InvalidPluginRegistration,
2262 format!(
2263 "random stream {}.{}@{} must use its owning plugin namespace {plugin}",
2264 stream.namespace, stream.name, stream.version
2265 ),
2266 ));
2267 }
2268 if let Some((existing_plugin, existing_system)) = owners.get(stream)
2269 && (existing_plugin != plugin || existing_system != system)
2270 {
2271 return Err(CanwuError::new(
2272 ErrorCode::InvalidPluginRegistration,
2273 format!(
2274 "random stream {}.{}@{} is owned by both {existing_plugin}.{existing_system} and {plugin}.{system}",
2275 stream.namespace, stream.name, stream.version
2276 ),
2277 ));
2278 }
2279 }
2280 for stream in streams {
2281 owners.insert(stream.clone(), (plugin.to_owned(), system.to_owned()));
2282 }
2283 Ok(())
2284}
2285
2286#[cfg(test)]
2287mod tests {
2288 use super::super::{KnowledgeSubjectSchema, KnowledgeSubjectTargetKind};
2289 use super::*;
2290 use canwu_core::{CoreEntityKind, KnowledgeRecordKind, KnowledgeSchemaId};
2291
2292 struct KnowledgeSchemaPlugin {
2293 name: &'static str,
2294 schemas: Vec<PluginKnowledgeSchema>,
2295 }
2296
2297 impl SimulationPlugin for KnowledgeSchemaPlugin {
2298 fn name(&self) -> &str {
2299 self.name
2300 }
2301
2302 fn version(&self) -> &'static str {
2303 "1"
2304 }
2305
2306 fn semantic_hash(&self) -> &'static str {
2307 "0000000000000000000000000000000000000000000000000000000000000001"
2308 }
2309
2310 fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
2311 for schema in &self.schemas {
2312 registrar.register_knowledge_schema(schema.clone())?;
2313 }
2314 Ok(())
2315 }
2316 }
2317
2318 fn knowledge_kind() -> KnowledgeRecordKind {
2319 KnowledgeRecordKind::new("fixture.knowledge", "assessment")
2320 }
2321
2322 fn knowledge_schema(version: u32, writable: bool) -> PluginKnowledgeSchema {
2323 PluginKnowledgeSchema {
2324 id: KnowledgeSchemaId::new(knowledge_kind(), version),
2325 schema_hash: format!("{version:064x}"),
2326 writable,
2327 payload_schema: PayloadSchema::Any,
2328 subjects: vec![],
2329 }
2330 }
2331
2332 #[test]
2333 fn archive_retention_binding_rejects_restored_root_mismatch() {
2334 let root = "a".repeat(64);
2335 let other = "b".repeat(64);
2336 let binding = PluginArchiveRetentionBinding {
2337 packet_type: "archive_commit".to_owned(),
2338 namespace: "fixture.archive.directory".to_owned(),
2339 object_id_json_pointers: vec![
2340 "/commit/archive_head/membership_root".to_owned(),
2341 "/commit/pending_reachability/directory_root".to_owned(),
2342 ],
2343 };
2344 let mut registry = PluginRegistry::default();
2345 registry
2346 .archive_retention_bindings
2347 .insert(("fixture".to_owned(), "archive_commit".to_owned()), binding);
2348 let retention = vec![PluginArchiveRetention {
2349 namespace: "fixture.archive.directory".to_owned(),
2350 object_id: root.clone(),
2351 }];
2352 let payload = serde_json::json!({
2353 "commit": {
2354 "archive_head": { "membership_root": root },
2355 "pending_reachability": { "directory_root": other },
2356 }
2357 });
2358 assert_eq!(
2359 registry
2360 .validate_archive_retention("fixture", "archive_commit", &payload, &retention,)
2361 .unwrap_err()
2362 .code,
2363 ErrorCode::InvalidPayload
2364 );
2365 }
2366
2367 #[test]
2368 fn duplicate_schema_and_writable_conflicts_roll_back_registration() {
2369 let duplicate = KnowledgeSchemaPlugin {
2370 name: "duplicate-knowledge",
2371 schemas: vec![knowledge_schema(1, true), knowledge_schema(1, true)],
2372 };
2373 let mut registry = PluginRegistry::default();
2374 let mut types = SchemaRegistry::default();
2375 assert!(registry.register(&duplicate, &mut types).is_err());
2376 assert!(registry.descriptors.is_empty());
2377 assert!(registry.knowledge_schemas.is_empty());
2378 assert!(registry.knowledge_kind_owners.is_empty());
2379
2380 let two_writable = KnowledgeSchemaPlugin {
2381 name: "two-writable-knowledge",
2382 schemas: vec![knowledge_schema(1, true), knowledge_schema(2, true)],
2383 };
2384 assert!(registry.register(&two_writable, &mut types).is_err());
2385 assert!(registry.descriptors.is_empty());
2386 assert!(registry.knowledge_schemas.is_empty());
2387
2388 let first_owner = KnowledgeSchemaPlugin {
2389 name: "first-knowledge-owner",
2390 schemas: vec![knowledge_schema(1, true)],
2391 };
2392 registry
2393 .register(&first_owner, &mut types)
2394 .expect("the first kind owner should register");
2395 let before = registry.clone();
2396 let second_owner = KnowledgeSchemaPlugin {
2397 name: "second-knowledge-owner",
2398 schemas: vec![knowledge_schema(2, true)],
2399 };
2400 assert!(registry.register(&second_owner, &mut types).is_err());
2401 assert_eq!(registry.descriptors, before.descriptors);
2402 assert_eq!(registry.knowledge_schemas, before.knowledge_schemas);
2403 assert_eq!(registry.knowledge_kind_owners, before.knowledge_kind_owners);
2404 }
2405
2406 #[test]
2407 fn schema_hash_mismatch_blocks_exact_rehydration() {
2408 let plugin = KnowledgeSchemaPlugin {
2409 name: "rehydrated-knowledge",
2410 schemas: vec![knowledge_schema(1, true)],
2411 };
2412 let mut registry = PluginRegistry::default();
2413 let mut types = SchemaRegistry::default();
2414 registry
2415 .register(&plugin, &mut types)
2416 .expect("fixture plugin should register");
2417 let mut descriptors = registry.descriptors().cloned().collect::<Vec<_>>();
2418 descriptors[0].knowledge_schemas[0].schema_hash =
2419 "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff".to_owned();
2420 let mut rehydrated = PluginRegistry::from_descriptors(descriptors)
2421 .expect("the altered descriptor remains structurally valid");
2422 let error = rehydrated
2423 .register(&plugin, &mut SchemaRegistry::default())
2424 .expect_err("exact rehydration must compare the persisted schema hash");
2425 assert_eq!(error.code, ErrorCode::PluginManifestMismatch);
2426 }
2427
2428 #[test]
2429 fn schema_limit_accepts_boundary_and_rejects_plus_one_atomically() {
2430 let boundary = KnowledgeSchemaPlugin {
2431 name: "knowledge-limit-boundary",
2432 schemas: (1..=KnowledgeLimitsV1::CURRENT.schemas_per_plugin)
2433 .map(|version| {
2434 knowledge_schema(
2435 u32::try_from(version).expect("schema limit fits u32"),
2436 version == 1,
2437 )
2438 })
2439 .collect(),
2440 };
2441 let mut registry = PluginRegistry::default();
2442 registry
2443 .register(&boundary, &mut SchemaRegistry::default())
2444 .expect("the exact schema limit should be admitted");
2445 assert_eq!(
2446 registry.knowledge_schemas.len(),
2447 KnowledgeLimitsV1::CURRENT.schemas_per_plugin
2448 );
2449
2450 let overflow = KnowledgeSchemaPlugin {
2451 name: "knowledge-limit-overflow",
2452 schemas: (1..=KnowledgeLimitsV1::CURRENT.schemas_per_plugin + 1)
2453 .map(|version| {
2454 knowledge_schema(
2455 u32::try_from(version).expect("schema limit fits u32"),
2456 version == 1,
2457 )
2458 })
2459 .collect(),
2460 };
2461 let mut rejected = PluginRegistry::default();
2462 let error = rejected
2463 .register(&overflow, &mut SchemaRegistry::default())
2464 .expect_err("schema limit plus one must reject the whole plugin");
2465 assert_eq!(error.code, ErrorCode::InvalidPluginRegistration);
2466 assert!(rejected.descriptors.is_empty());
2467 assert!(rejected.knowledge_schemas.is_empty());
2468 }
2469
2470 #[test]
2471 fn knowledge_schema_registration_canonicalizes_roles_and_targets() {
2472 let mut schema = knowledge_schema(1, true);
2473 schema.subjects = vec![
2474 KnowledgeSubjectSchema {
2475 role: "zeta".to_owned(),
2476 targets: vec![
2477 KnowledgeSubjectTargetKind::AnyEntity,
2478 KnowledgeSubjectTargetKind::Core(CoreEntityKind::Person),
2479 KnowledgeSubjectTargetKind::AnyEntity,
2480 ],
2481 required: false,
2482 multiple: true,
2483 },
2484 KnowledgeSubjectSchema {
2485 role: "alpha".to_owned(),
2486 targets: vec![KnowledgeSubjectTargetKind::Event],
2487 required: true,
2488 multiple: false,
2489 },
2490 ];
2491 let plugin = KnowledgeSchemaPlugin {
2492 name: "canonical-knowledge",
2493 schemas: vec![schema],
2494 };
2495 let mut registry = PluginRegistry::default();
2496 registry
2497 .register(&plugin, &mut SchemaRegistry::default())
2498 .expect("registrar should canonicalize declarations transactionally");
2499 let stored = ®istry
2500 .descriptors
2501 .get(plugin.name)
2502 .expect("descriptor exists")
2503 .knowledge_schemas[0];
2504 assert_eq!(stored.subjects[0].role, "alpha");
2505 assert_eq!(stored.subjects[1].role, "zeta");
2506 assert_eq!(stored.subjects[1].targets.len(), 2);
2507 assert!(stored.validate().is_ok());
2508 }
2509
2510 #[test]
2511 fn invalid_schema_version_and_hash_roll_back_registration() {
2512 let mut version_zero = knowledge_schema(0, true);
2513 version_zero.schema_hash =
2514 "0000000000000000000000000000000000000000000000000000000000000000".to_owned();
2515 let invalid_version = KnowledgeSchemaPlugin {
2516 name: "invalid-knowledge-version",
2517 schemas: vec![version_zero],
2518 };
2519 let mut registry = PluginRegistry::default();
2520 let mut types = SchemaRegistry::default();
2521 assert!(registry.register(&invalid_version, &mut types).is_err());
2522 assert!(registry.descriptors.is_empty());
2523 assert!(registry.knowledge_schemas.is_empty());
2524
2525 let mut bad_hash = knowledge_schema(1, true);
2526 bad_hash.schema_hash = "not-a-canonical-hash".to_owned();
2527 let invalid_hash = KnowledgeSchemaPlugin {
2528 name: "invalid-knowledge-hash",
2529 schemas: vec![bad_hash],
2530 };
2531 assert!(registry.register(&invalid_hash, &mut types).is_err());
2532 assert!(registry.descriptors.is_empty());
2533 assert!(registry.knowledge_schemas.is_empty());
2534 }
2535
2536 #[test]
2537 fn knowledge_write_grants_reject_invalid_phase_and_foreign_owner() {
2538 #[allow(clippy::unnecessary_wraps)]
2539 fn no_op_boundary(
2540 _view: &crate::SimulationView<'_>,
2541 _context: &crate::BoundaryContext,
2542 ) -> Result<crate::BoundaryProposal, CanwuError> {
2543 Ok(crate::BoundaryProposal::default())
2544 }
2545
2546 let owner = KnowledgeSchemaPlugin {
2547 name: "knowledge-grant-owner",
2548 schemas: vec![knowledge_schema(1, true)],
2549 };
2550 let foreign = KnowledgeSchemaPlugin {
2551 name: "knowledge-grant-foreign",
2552 schemas: vec![PluginKnowledgeSchema {
2553 id: KnowledgeSchemaId::new(
2554 KnowledgeRecordKind::new("fixture.foreign", "assessment"),
2555 1,
2556 ),
2557 schema_hash: "f000000000000000000000000000000000000000000000000000000000000000"
2558 .to_owned(),
2559 writable: true,
2560 payload_schema: PayloadSchema::Any,
2561 subjects: Vec::new(),
2562 }],
2563 };
2564 let mut registry = PluginRegistry::default();
2565 let mut types = SchemaRegistry::default();
2566 registry
2567 .register(&owner, &mut types)
2568 .expect("knowledge owner should register");
2569 registry
2570 .register(&foreign, &mut types)
2571 .expect("foreign fixture should register");
2572 let before = registry.clone();
2573
2574 let mut phase7 = BoundarySystemContract::new(
2575 "invalid-phase7-publication",
2576 crate::BoundaryPhase::DomainDeltaProposal,
2577 SystemCadence::Daily,
2578 );
2579 phase7.knowledge_writes = vec![crate::KnowledgeWriteGrant {
2580 schema: knowledge_schema(1, true).id,
2581 visibilities: vec![StateVisibility::SameBoundary],
2582 }];
2583 let mut owner_registry = registry.clone();
2584 let mut owner_types = types.clone();
2585 let mut registrar = PluginRegistrar {
2586 plugin: owner.name.to_owned(),
2587 registry: &mut owner_registry,
2588 schema: &mut owner_types,
2589 };
2590 let error = registrar
2591 .register_boundary_system(phase7, no_op_boundary)
2592 .expect_err("phase 7 must reject knowledge publication grants");
2593 assert_eq!(error.code, ErrorCode::InvalidPluginRegistration);
2594 assert_eq!(owner_registry.descriptors, before.descriptors);
2595
2596 let mut foreign_grant = BoundarySystemContract::new(
2597 "foreign-knowledge-grant",
2598 crate::BoundaryPhase::PerspectiveAndReportMaterialization,
2599 SystemCadence::Daily,
2600 );
2601 foreign_grant.knowledge_writes = vec![crate::KnowledgeWriteGrant {
2602 schema: knowledge_schema(1, true).id,
2603 visibilities: vec![StateVisibility::SameBoundary],
2604 }];
2605 let mut foreign_registry = registry;
2606 let mut registrar = PluginRegistrar {
2607 plugin: foreign.name.to_owned(),
2608 registry: &mut foreign_registry,
2609 schema: &mut types,
2610 };
2611 let error = registrar
2612 .register_boundary_system(foreign_grant, no_op_boundary)
2613 .expect_err("a plugin cannot claim another plugin's writable schema");
2614 assert_eq!(error.code, ErrorCode::InvalidPluginRegistration);
2615 assert_eq!(foreign_registry.descriptors, before.descriptors);
2616 }
2617}