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 mut owned_state = contract.writes.clone();
871 owned_state.extend(contract.reservation_offers.iter().cloned());
872 owned_state.sort();
873 owned_state.dedup();
874 let mut candidate = self.registry.clone();
875 register_state_owners(&mut candidate.state_owners, &self.plugin, &owned_state)?;
876 register_boundary_writers(
877 &mut candidate.boundary_writers,
878 &candidate.immediate_write_states,
879 &self.plugin,
880 &contract.name,
881 contract.phase,
882 &contract.writes,
883 )?;
884 register_reservation_offerers(
885 &mut candidate.reservation_offerers,
886 &self.plugin,
887 &contract.name,
888 &contract.reservation_offers,
889 )?;
890 register_random_streams(
891 &mut candidate.random_stream_owners,
892 &self.plugin,
893 &contract.name,
894 &contract.random_streams,
895 )?;
896 {
897 let descriptor = candidate
898 .descriptors
899 .entry(self.plugin.clone())
900 .or_default();
901 descriptor.name.clone_from(&self.plugin);
902 descriptor.boundary_systems.push(contract.clone());
903 descriptor
904 .boundary_systems
905 .sort_by(|left, right| (left.phase, &left.name).cmp(&(right.phase, &right.name)));
906 }
907 candidate.boundary_systems.push(RegisteredBoundarySystem {
908 plugin: self.plugin.clone(),
909 contract,
910 handler,
911 });
912 candidate.boundary_systems.sort_by(|left, right| {
913 (left.contract.phase, &left.plugin, &left.contract.name).cmp(&(
914 right.contract.phase,
915 &right.plugin,
916 &right.contract.name,
917 ))
918 });
919 *self.registry = candidate;
920 Ok(())
921 }
922
923 pub fn register_command(
924 &mut self,
925 mut descriptor: PluginActionDescriptor,
926 handler: PluginCommandHandler,
927 ) -> Result<(), CanwuError> {
928 validate_action_descriptor(&self.plugin, &mut descriptor)?;
929 let command_key = (self.plugin.clone(), descriptor.name.clone());
930 if self.registry.commands.contains_key(&command_key) {
931 return Err(CanwuError::new(
932 ErrorCode::DuplicatePluginCommand,
933 format!(
934 "plugin {} already registered command {}",
935 self.plugin, descriptor.name
936 ),
937 ));
938 }
939 let mut candidate = self.registry.clone();
940 if descriptor
941 .writes
942 .iter()
943 .any(|state| is_domain_record_state(&candidate.record_schemas, state))
944 {
945 return Err(CanwuError::new(
946 ErrorCode::InvalidPluginRegistration,
947 "plugin commands cannot write domain record state directly",
948 ));
949 }
950 register_state_owners(
951 &mut candidate.state_owners,
952 &self.plugin,
953 &descriptor.writes,
954 )?;
955 register_immediate_write_states(
956 &mut candidate.immediate_write_states,
957 &candidate.boundary_writers,
958 &self.plugin,
959 &descriptor.writes,
960 )?;
961 {
962 let plugin_descriptor = candidate
963 .descriptors
964 .entry(self.plugin.clone())
965 .or_default();
966 plugin_descriptor.name.clone_from(&self.plugin);
967 plugin_descriptor.commands.push(descriptor.clone());
968 plugin_descriptor
969 .commands
970 .sort_by(|left, right| left.name.cmp(&right.name));
971 }
972 candidate.commands.insert(
973 command_key,
974 RegisteredCommand {
975 descriptor,
976 handler,
977 },
978 );
979 *self.registry = candidate;
980 Ok(())
981 }
982
983 pub fn register_ingress(
984 &mut self,
985 descriptor: PluginIngressDescriptor,
986 ) -> Result<(), CanwuError> {
987 validate_ingress_descriptor(&descriptor)?;
988 let key = (self.plugin.clone(), descriptor.name.clone());
989 if self
990 .registry
991 .descriptors
992 .get(&self.plugin)
993 .is_some_and(|plugin| {
994 plugin
995 .ingress
996 .iter()
997 .any(|candidate| candidate.name == descriptor.name)
998 })
999 {
1000 return Err(CanwuError::new(
1001 ErrorCode::DuplicatePluginIngress,
1002 format!(
1003 "plugin {} already registered ingress type {}",
1004 self.plugin, descriptor.name
1005 ),
1006 ));
1007 }
1008 if self
1009 .registry
1010 .ingress
1011 .get(&key)
1012 .is_some_and(|existing| existing != &descriptor)
1013 {
1014 return Err(CanwuError::new(
1015 ErrorCode::PluginManifestMismatch,
1016 format!(
1017 "plugin {} changed the stored ingress type {}",
1018 self.plugin, descriptor.name
1019 ),
1020 ));
1021 }
1022 let mut candidate = self.registry.clone();
1023 candidate.ingress.insert(key, descriptor.clone());
1024 let plugin_descriptor = candidate
1025 .descriptors
1026 .entry(self.plugin.clone())
1027 .or_default();
1028 plugin_descriptor.name.clone_from(&self.plugin);
1029 plugin_descriptor.ingress.push(descriptor);
1030 plugin_descriptor
1031 .ingress
1032 .sort_by(|left, right| left.name.cmp(&right.name));
1033 *self.registry = candidate;
1034 Ok(())
1035 }
1036
1037 pub fn register_internal_ingress(
1040 &mut self,
1041 descriptor: PluginIngressDescriptor,
1042 ) -> Result<PluginIngressPermit, CanwuError> {
1043 let packet_type = descriptor.name.clone();
1044 self.register_ingress(descriptor)?;
1045 let semantic_hash = self
1046 .registry
1047 .descriptors
1048 .get(&self.plugin)
1049 .map(|descriptor| descriptor.semantic_hash.clone())
1050 .ok_or_else(|| {
1051 CanwuError::new(
1052 ErrorCode::InvalidPluginRegistration,
1053 "internal ingress owner has no plugin descriptor",
1054 )
1055 })?;
1056 let token = canonical_hash(
1057 "canwu.plugin.internal-ingress-permit.v1",
1058 &(&self.plugin, &packet_type, &semantic_hash),
1059 )?;
1060 self.registry
1061 .internal_ingress
1062 .insert((self.plugin.clone(), packet_type.clone()));
1063 let plugin_descriptor =
1064 self.registry
1065 .descriptors
1066 .get_mut(&self.plugin)
1067 .ok_or_else(|| {
1068 CanwuError::new(
1069 ErrorCode::InvalidPluginRegistration,
1070 "internal ingress owner descriptor disappeared",
1071 )
1072 })?;
1073 plugin_descriptor.internal_ingress.push(packet_type.clone());
1074 plugin_descriptor.internal_ingress.sort();
1075 plugin_descriptor.internal_ingress.dedup();
1076 Ok(PluginIngressPermit {
1077 plugin: self.plugin.clone(),
1078 packet_type,
1079 semantic_hash,
1080 token,
1081 })
1082 }
1083
1084 pub fn register_internal_ingress_with_archive_retention(
1087 &mut self,
1088 descriptor: PluginIngressDescriptor,
1089 namespace: impl Into<String>,
1090 object_id_json_pointers: Vec<String>,
1091 ) -> Result<PluginIngressPermit, CanwuError> {
1092 let binding = PluginArchiveRetentionBinding {
1093 packet_type: descriptor.name.clone(),
1094 namespace: namespace.into(),
1095 object_id_json_pointers,
1096 };
1097 validate_archive_retention_binding(&binding)?;
1098 let permit = self.register_internal_ingress(descriptor)?;
1099 let key = (self.plugin.clone(), binding.packet_type.clone());
1100 if self
1101 .registry
1102 .archive_retention_bindings
1103 .get(&key)
1104 .is_some_and(|existing| existing != &binding)
1105 {
1106 return Err(CanwuError::new(
1107 ErrorCode::PluginManifestMismatch,
1108 "plugin changed its stored archive-retention binding",
1109 ));
1110 }
1111 self.registry
1112 .archive_retention_bindings
1113 .insert(key, binding.clone());
1114 let plugin_descriptor =
1115 self.registry
1116 .descriptors
1117 .get_mut(&self.plugin)
1118 .ok_or_else(|| {
1119 CanwuError::new(
1120 ErrorCode::InvalidPluginRegistration,
1121 "archive-retention owner descriptor disappeared",
1122 )
1123 })?;
1124 plugin_descriptor.archive_retention_bindings.push(binding);
1125 plugin_descriptor
1126 .archive_retention_bindings
1127 .sort_by(|left, right| left.packet_type.cmp(&right.packet_type));
1128 Ok(permit)
1129 }
1130}
1131
1132impl PluginRegistry {
1133 pub(super) fn validate_archive_retention(
1134 &self,
1135 plugin: &str,
1136 packet_type: &str,
1137 payload: &Value,
1138 retention: &[PluginArchiveRetention],
1139 ) -> Result<(), CanwuError> {
1140 let Some(binding) = self
1141 .archive_retention_bindings
1142 .get(&(plugin.to_owned(), packet_type.to_owned()))
1143 else {
1144 return if retention.is_empty() {
1145 Ok(())
1146 } else {
1147 Err(CanwuError::new(
1148 ErrorCode::InvalidPayload,
1149 "plugin ingress carries undeclared archive retention",
1150 ))
1151 };
1152 };
1153 if retention.len() != 1 || retention[0].namespace != binding.namespace {
1154 return Err(CanwuError::new(
1155 ErrorCode::InvalidPayload,
1156 "plugin ingress archive retention does not match its declared binding",
1157 ));
1158 }
1159 let object_id = &retention[0].object_id;
1160 if !is_canonical_hash(object_id)
1161 || binding.object_id_json_pointers.iter().any(|pointer| {
1162 payload.pointer(pointer).and_then(Value::as_str) != Some(object_id.as_str())
1163 })
1164 {
1165 return Err(CanwuError::new(
1166 ErrorCode::InvalidPayload,
1167 "plugin ingress archive root is not bound to its authenticated payload roots",
1168 ));
1169 }
1170 Ok(())
1171 }
1172
1173 pub fn register<P: SimulationPlugin + ?Sized>(
1174 &mut self,
1175 plugin: &P,
1176 schema: &mut SchemaRegistry,
1177 ) -> Result<(), CanwuError> {
1178 let raw_plugin_name = plugin.name();
1179 let plugin_name = raw_plugin_name.trim();
1180 if plugin_name.is_empty() || plugin_name != raw_plugin_name {
1181 return Err(CanwuError::new(
1182 ErrorCode::InvalidPluginRegistration,
1183 "plugin name must be non-empty and have no surrounding whitespace",
1184 ));
1185 }
1186 if self.active_plugins.contains(plugin_name) {
1187 return Err(CanwuError::new(
1188 ErrorCode::DuplicatePlugin,
1189 format!("plugin {plugin_name} is already registered"),
1190 ));
1191 }
1192 validate_plugin_identity(plugin_name, plugin.version(), plugin.semantic_hash())?;
1193
1194 let expected_descriptor = self.descriptors.get(plugin_name).cloned();
1195 let mut candidate_registry = self.clone();
1196 let mut candidate_schema = schema.clone();
1197 candidate_registry.descriptors.insert(
1198 plugin_name.to_owned(),
1199 PluginDescriptor {
1200 descriptor_format: PLUGIN_DESCRIPTOR_FORMAT_VERSION,
1201 name: plugin_name.to_owned(),
1202 version: plugin.version().to_owned(),
1203 semantic_hash: plugin.semantic_hash().to_owned(),
1204 ..PluginDescriptor::default()
1205 },
1206 );
1207 let mut registrar = PluginRegistrar {
1208 plugin: plugin_name.to_owned(),
1209 registry: &mut candidate_registry,
1210 schema: &mut candidate_schema,
1211 };
1212 plugin.register(&mut registrar)?;
1213 validate_schema_set(
1214 &candidate_registry.knowledge_schemas,
1215 &candidate_registry.knowledge_kind_owners,
1216 )
1217 .map_err(|error| {
1218 CanwuError::new(
1219 ErrorCode::InvalidPluginRegistration,
1220 format!("invalid knowledge schema set: {error}"),
1221 )
1222 })?;
1223 let Some(generated_descriptor) = candidate_registry.descriptors.get(plugin_name) else {
1224 return Err(CanwuError::new(
1225 ErrorCode::InvalidPluginRegistration,
1226 format!("plugin {plugin_name} did not produce a descriptor"),
1227 ));
1228 };
1229 if let Some(expected) = expected_descriptor
1230 && generated_descriptor != &expected
1231 {
1232 return Err(CanwuError::new(
1233 ErrorCode::PluginManifestMismatch,
1234 format!("plugin {plugin_name} registration does not match the snapshot manifest"),
1235 ));
1236 }
1237 candidate_registry
1238 .active_plugins
1239 .insert(plugin_name.to_owned());
1240 *self = candidate_registry;
1241 *schema = candidate_schema;
1242 Ok(())
1243 }
1244
1245 pub fn descriptors(&self) -> impl Iterator<Item = &PluginDescriptor> {
1246 self.descriptors.values()
1247 }
1248
1249 pub(super) fn event_audience(&self, plugin: &str, event_type: &str) -> EventAudience {
1250 self.descriptors
1251 .get(plugin)
1252 .and_then(|descriptor| descriptor.event_audiences.get(event_type))
1253 .cloned()
1254 .unwrap_or_default()
1255 }
1256
1257 pub(super) fn from_descriptors(descriptors: Vec<PluginDescriptor>) -> Result<Self, CanwuError> {
1258 let mut registry = Self {
1259 descriptors: BTreeMap::new(),
1260 active_plugins: BTreeSet::new(),
1261 systems: Vec::new(),
1262 boundary_systems: Vec::new(),
1263 commands: BTreeMap::new(),
1264 ingress: BTreeMap::new(),
1265 internal_ingress: BTreeSet::new(),
1266 archive_retention_bindings: BTreeMap::new(),
1267 state_owners: BTreeMap::new(),
1268 immediate_write_states: BTreeMap::new(),
1269 boundary_writers: BTreeMap::new(),
1270 reservation_offerers: BTreeMap::new(),
1271 random_stream_owners: BTreeMap::new(),
1272 record_schemas: BTreeMap::new(),
1273 knowledge_schemas: BTreeMap::new(),
1274 knowledge_kind_owners: BTreeMap::new(),
1275 maintenance_dependency_resolvers: BTreeMap::new(),
1276 maintenance_participants: BTreeMap::new(),
1277 archive_reachability_participants: BTreeMap::new(),
1278 };
1279 let mut previous_plugin = None;
1280 for mut descriptor in descriptors {
1281 let plugin = descriptor.name.trim().to_owned();
1282 if plugin.is_empty()
1283 || descriptor.descriptor_format != PLUGIN_DESCRIPTOR_FORMAT_VERSION
1284 || descriptor.name != plugin
1285 || descriptor.version.trim().is_empty()
1286 || descriptor.version != descriptor.version.trim()
1287 || !is_canonical_hash(&descriptor.semantic_hash)
1288 || registry.descriptors.contains_key(&plugin)
1289 || previous_plugin
1290 .as_ref()
1291 .is_some_and(|previous| previous >= &plugin)
1292 {
1293 return Err(CanwuError::new(
1294 ErrorCode::InvalidSnapshot,
1295 "snapshot contains an invalid, unversioned, or duplicate plugin descriptor",
1296 ));
1297 }
1298 if descriptor
1299 .record_schemas
1300 .windows(2)
1301 .any(|pair| pair[0].kind >= pair[1].kind)
1302 {
1303 return invalid_snapshot("plugin record schemas are not in canonical order");
1304 }
1305 for schema in &mut descriptor.record_schemas {
1306 let original = schema.clone();
1307 schema.canonicalize();
1308 schema.validate().map_err(|error| {
1309 invalid_snapshot_error(format!("invalid domain record schema: {error}"))
1310 })?;
1311 if *schema != original {
1312 return invalid_snapshot(
1313 "plugin record-schema declarations are not in canonical order",
1314 );
1315 }
1316 let state = schema.state_key();
1317 if state.namespace == CORE_STATE_NAMESPACE {
1318 return invalid_snapshot(
1319 "plugin record schemas cannot use the reserved core namespace",
1320 );
1321 }
1322 if let Some((owner, _)) = registry.record_schemas.get(&schema.kind) {
1323 return invalid_snapshot(format!(
1324 "domain record kind {} is owned by both {owner} and {plugin}",
1325 schema.kind
1326 ));
1327 }
1328 register_state_owners(
1329 &mut registry.state_owners,
1330 &plugin,
1331 std::slice::from_ref(&state),
1332 )
1333 .map_err(|error| {
1334 invalid_snapshot_error(format!(
1335 "invalid domain record state ownership descriptor: {error}"
1336 ))
1337 })?;
1338 registry
1339 .record_schemas
1340 .insert(schema.kind.clone(), (plugin.clone(), schema.clone()));
1341 }
1342 if descriptor.knowledge_schemas.len() > KnowledgeLimitsV1::CURRENT.schemas_per_plugin
1343 || descriptor
1344 .knowledge_schemas
1345 .windows(2)
1346 .any(|pair| pair[0].id >= pair[1].id)
1347 {
1348 return invalid_snapshot(
1349 "plugin knowledge schemas are not in canonical order or exceed their limit",
1350 );
1351 }
1352 for schema in &mut descriptor.knowledge_schemas {
1353 let original = schema.clone();
1354 schema.canonicalize();
1355 schema.validate().map_err(|error| {
1356 invalid_snapshot_error(format!("invalid knowledge schema: {error}"))
1357 })?;
1358 if *schema != original {
1359 return invalid_snapshot(
1360 "plugin knowledge-schema declarations are not in canonical order",
1361 );
1362 }
1363 if let Some(owner) = registry.knowledge_kind_owners.get(&schema.id.kind) {
1364 if owner != &plugin {
1365 return invalid_snapshot(format!(
1366 "knowledge kind {:?} is owned by both {owner} and {plugin}",
1367 schema.id.kind
1368 ));
1369 }
1370 } else {
1371 registry
1372 .knowledge_kind_owners
1373 .insert(schema.id.kind.clone(), plugin.clone());
1374 }
1375 if registry
1376 .knowledge_schemas
1377 .insert(schema.id.clone(), (plugin.clone(), schema.clone()))
1378 .is_some()
1379 {
1380 return invalid_snapshot("knowledge schema ID is duplicated");
1381 }
1382 }
1383 if descriptor
1384 .maintenance_dependency_resolvers
1385 .windows(2)
1386 .any(|pair| pair[0] >= pair[1])
1387 {
1388 return invalid_snapshot(
1389 "plugin maintenance dependency resolvers are not in canonical order",
1390 );
1391 }
1392 for resolver in &descriptor.maintenance_dependency_resolvers {
1393 if !canonical_text(&resolver.target_namespace)
1394 || resolver.target_namespace == CORE_STATE_NAMESPACE
1395 {
1396 return invalid_snapshot(
1397 "plugin maintenance dependency resolver target is invalid",
1398 );
1399 }
1400 registry
1401 .maintenance_dependency_resolvers
1402 .entry(resolver.target_namespace.clone())
1403 .or_default()
1404 .insert(plugin.clone());
1405 }
1406 if descriptor
1407 .systems
1408 .windows(2)
1409 .any(|pair| (pair[0].phase, &pair[0].name) >= (pair[1].phase, &pair[1].name))
1410 {
1411 return invalid_snapshot("plugin systems are not in canonical order");
1412 }
1413 let mut system_names = BTreeSet::new();
1414 for contract in &mut descriptor.systems {
1415 if !system_names.insert(contract.name.clone()) {
1416 return invalid_snapshot("plugin descriptor has duplicate system names");
1417 }
1418 let original = contract.clone();
1419 validate_system_contract(&plugin, contract).map_err(|error| {
1420 invalid_snapshot_error(format!("invalid plugin system descriptor: {error}"))
1421 })?;
1422 if *contract != original {
1423 return invalid_snapshot(
1424 "plugin system reads and writes are not in canonical order",
1425 );
1426 }
1427 if contract
1428 .writes
1429 .iter()
1430 .any(|state| is_domain_record_state(®istry.record_schemas, state))
1431 {
1432 return invalid_snapshot(
1433 "plugin systems cannot expose domain records as immediate component state",
1434 );
1435 }
1436 register_state_owners(&mut registry.state_owners, &plugin, &contract.writes)
1437 .map_err(|error| {
1438 invalid_snapshot_error(format!(
1439 "invalid plugin state ownership descriptor: {error}"
1440 ))
1441 })?;
1442 register_immediate_write_states(
1443 &mut registry.immediate_write_states,
1444 ®istry.boundary_writers,
1445 &plugin,
1446 &contract.writes,
1447 )
1448 .map_err(|error| {
1449 invalid_snapshot_error(format!(
1450 "invalid immediate state writer descriptor: {error}"
1451 ))
1452 })?;
1453 }
1454 if descriptor
1455 .boundary_systems
1456 .windows(2)
1457 .any(|pair| (pair[0].phase, &pair[0].name) >= (pair[1].phase, &pair[1].name))
1458 {
1459 return invalid_snapshot("boundary systems are not in canonical order");
1460 }
1461 for contract in &mut descriptor.boundary_systems {
1462 if !system_names.insert(contract.name.clone()) {
1463 return invalid_snapshot("plugin descriptor has duplicate system names");
1464 }
1465 let original = contract.clone();
1466 validate_boundary_system_contract(contract).map_err(|error| {
1467 invalid_snapshot_error(format!("invalid boundary system descriptor: {error}"))
1468 })?;
1469 validate_knowledge_write_grants(&plugin, contract, ®istry.knowledge_schemas)
1470 .map_err(|error| {
1471 invalid_snapshot_error(format!(
1472 "invalid boundary knowledge writer descriptor: {error}"
1473 ))
1474 })?;
1475 if *contract != original {
1476 return invalid_snapshot(
1477 "boundary system declarations are not in canonical order",
1478 );
1479 }
1480 let mut owned_state = contract.writes.clone();
1481 owned_state.extend(contract.reservation_offers.iter().cloned());
1482 owned_state.sort();
1483 owned_state.dedup();
1484 register_state_owners(&mut registry.state_owners, &plugin, &owned_state).map_err(
1485 |error| {
1486 invalid_snapshot_error(format!(
1487 "invalid boundary state ownership descriptor: {error}"
1488 ))
1489 },
1490 )?;
1491 register_boundary_writers(
1492 &mut registry.boundary_writers,
1493 ®istry.immediate_write_states,
1494 &plugin,
1495 &contract.name,
1496 contract.phase,
1497 &contract.writes,
1498 )
1499 .map_err(|error| {
1500 invalid_snapshot_error(format!("invalid boundary writer descriptor: {error}"))
1501 })?;
1502 register_reservation_offerers(
1503 &mut registry.reservation_offerers,
1504 &plugin,
1505 &contract.name,
1506 &contract.reservation_offers,
1507 )
1508 .map_err(|error| {
1509 invalid_snapshot_error(format!(
1510 "invalid reservation offerer descriptor: {error}"
1511 ))
1512 })?;
1513 register_random_streams(
1514 &mut registry.random_stream_owners,
1515 &plugin,
1516 &contract.name,
1517 &contract.random_streams,
1518 )
1519 .map_err(|error| {
1520 invalid_snapshot_error(format!(
1521 "invalid random stream ownership descriptor: {error}"
1522 ))
1523 })?;
1524 }
1525 if descriptor
1526 .commands
1527 .windows(2)
1528 .any(|pair| pair[0].name >= pair[1].name)
1529 {
1530 return invalid_snapshot("plugin commands are not in canonical order");
1531 }
1532 let mut command_names = BTreeSet::new();
1533 for action in &mut descriptor.commands {
1534 if !command_names.insert(action.name.clone()) {
1535 return invalid_snapshot("plugin descriptor has duplicate command names");
1536 }
1537 let original = action.clone();
1538 validate_action_descriptor(&plugin, action).map_err(|error| {
1539 invalid_snapshot_error(format!("invalid plugin command descriptor: {error}"))
1540 })?;
1541 if *action != original {
1542 return invalid_snapshot(
1543 "plugin command reads and writes are not in canonical order",
1544 );
1545 }
1546 if action
1547 .writes
1548 .iter()
1549 .any(|state| is_domain_record_state(®istry.record_schemas, state))
1550 {
1551 return invalid_snapshot(
1552 "plugin commands cannot expose domain records as immediate component state",
1553 );
1554 }
1555 register_state_owners(&mut registry.state_owners, &plugin, &action.writes)
1556 .map_err(|error| {
1557 invalid_snapshot_error(format!(
1558 "invalid plugin state ownership descriptor: {error}"
1559 ))
1560 })?;
1561 register_immediate_write_states(
1562 &mut registry.immediate_write_states,
1563 ®istry.boundary_writers,
1564 &plugin,
1565 &action.writes,
1566 )
1567 .map_err(|error| {
1568 invalid_snapshot_error(format!(
1569 "invalid immediate state writer descriptor: {error}"
1570 ))
1571 })?;
1572 }
1573 if descriptor
1574 .ingress
1575 .windows(2)
1576 .any(|pair| pair[0].name >= pair[1].name)
1577 {
1578 return invalid_snapshot("plugin ingress types are not in canonical order");
1579 }
1580 for ingress in &descriptor.ingress {
1581 validate_ingress_descriptor(ingress).map_err(|error| {
1582 invalid_snapshot_error(format!("invalid plugin ingress descriptor: {error}"))
1583 })?;
1584 if registry
1585 .ingress
1586 .insert((plugin.clone(), ingress.name.clone()), ingress.clone())
1587 .is_some()
1588 {
1589 return invalid_snapshot("plugin descriptor has duplicate ingress types");
1590 }
1591 }
1592 if descriptor
1593 .internal_ingress
1594 .windows(2)
1595 .any(|pair| pair[0] >= pair[1])
1596 || descriptor.internal_ingress.iter().any(|name| {
1597 !descriptor
1598 .ingress
1599 .iter()
1600 .any(|ingress| &ingress.name == name)
1601 })
1602 {
1603 return invalid_snapshot("plugin internal ingress declarations are not canonical");
1604 }
1605 registry.internal_ingress.extend(
1606 descriptor
1607 .internal_ingress
1608 .iter()
1609 .map(|name| (plugin.clone(), name.clone())),
1610 );
1611 if descriptor
1612 .archive_retention_bindings
1613 .windows(2)
1614 .any(|pair| pair[0].packet_type >= pair[1].packet_type)
1615 {
1616 return invalid_snapshot(
1617 "plugin archive-retention bindings are not in canonical order",
1618 );
1619 }
1620 for binding in &descriptor.archive_retention_bindings {
1621 validate_archive_retention_binding(binding).map_err(|error| {
1622 invalid_snapshot_error(format!(
1623 "invalid plugin archive-retention binding: {error}"
1624 ))
1625 })?;
1626 if !descriptor
1627 .internal_ingress
1628 .iter()
1629 .any(|name| name == &binding.packet_type)
1630 || registry
1631 .archive_retention_bindings
1632 .insert(
1633 (plugin.clone(), binding.packet_type.clone()),
1634 binding.clone(),
1635 )
1636 .is_some()
1637 {
1638 return invalid_snapshot(
1639 "plugin archive-retention binding lacks one internal ingress owner",
1640 );
1641 }
1642 }
1643 for (event_type, audience) in &descriptor.event_audiences {
1644 validate_event_audience_name(event_type).map_err(|error| {
1645 invalid_snapshot_error(format!("invalid plugin event audience: {error}"))
1646 })?;
1647 validate_event_audience(audience).map_err(|error| {
1648 invalid_snapshot_error(format!("invalid plugin event audience: {error}"))
1649 })?;
1650 }
1651 let schema_types: BTreeSet<_> = descriptor.schema_types.iter().collect();
1652 if schema_types.len() != descriptor.schema_types.len()
1653 || descriptor
1654 .schema_types
1655 .windows(2)
1656 .any(|pair| pair[0] >= pair[1])
1657 || descriptor
1658 .schema_types
1659 .iter()
1660 .any(|name| name.trim().is_empty() || name != name.trim())
1661 {
1662 return invalid_snapshot("plugin descriptor has invalid schema type names");
1663 }
1664 previous_plugin = Some(plugin.clone());
1665 registry.descriptors.insert(plugin, descriptor);
1666 }
1667 validate_schema_set(®istry.knowledge_schemas, ®istry.knowledge_kind_owners).map_err(
1668 |error| invalid_snapshot_error(format!("invalid knowledge schema set: {error}")),
1669 )?;
1670 Ok(registry)
1671 }
1672
1673 pub(super) fn ensure_active(&self) -> Result<(), CanwuError> {
1674 let inactive: Vec<_> = self
1675 .descriptors
1676 .keys()
1677 .filter(|name| !self.active_plugins.contains(*name))
1678 .cloned()
1679 .collect();
1680 if inactive.is_empty() {
1681 return Ok(());
1682 }
1683 Err(CanwuError::new(
1684 ErrorCode::PluginNotActive,
1685 format!(
1686 "required plugin handlers are not active: {}",
1687 inactive.join(", ")
1688 ),
1689 ))
1690 }
1691}
1692
1693fn validate_archive_retention_binding(
1694 binding: &PluginArchiveRetentionBinding,
1695) -> Result<(), CanwuError> {
1696 if binding.packet_type.trim().is_empty()
1697 || binding.packet_type != binding.packet_type.trim()
1698 || binding.namespace.is_empty()
1699 || binding.namespace.len() > 128
1700 || !binding.namespace.bytes().all(|byte| {
1701 byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'.' | b'-' | b'_')
1702 })
1703 || binding.object_id_json_pointers.is_empty()
1704 || binding
1705 .object_id_json_pointers
1706 .windows(2)
1707 .any(|pair| pair[0] >= pair[1])
1708 || binding
1709 .object_id_json_pointers
1710 .iter()
1711 .any(|pointer| !pointer.starts_with('/') || pointer.len() > 256 || !pointer.is_ascii())
1712 {
1713 return Err(CanwuError::new(
1714 ErrorCode::InvalidPluginRegistration,
1715 "plugin archive-retention binding is malformed",
1716 ));
1717 }
1718 Ok(())
1719}
1720
1721pub(super) fn validate_state_keys(keys: &mut Vec<StateKey>) -> Result<(), CanwuError> {
1722 for key in keys.iter() {
1723 if key.namespace.trim().is_empty()
1724 || key.name.trim().is_empty()
1725 || key.namespace != key.namespace.trim()
1726 || key.name != key.name.trim()
1727 {
1728 return Err(CanwuError::new(
1729 ErrorCode::InvalidPluginRegistration,
1730 "state keys require non-empty canonical namespace and name values",
1731 ));
1732 }
1733 }
1734 let unique: BTreeSet<_> = keys.drain(..).collect();
1735 keys.extend(unique);
1736 Ok(())
1737}
1738
1739fn validate_plugin_identity(
1740 name: &str,
1741 version: &str,
1742 semantic_hash: &str,
1743) -> Result<(), CanwuError> {
1744 if name.trim().is_empty()
1745 || name != name.trim()
1746 || version.trim().is_empty()
1747 || version != version.trim()
1748 || !is_canonical_hash(semantic_hash)
1749 {
1750 return Err(CanwuError::new(
1751 ErrorCode::InvalidPluginRegistration,
1752 "plugins require canonical names, versions, and 64-character semantic hashes",
1753 ));
1754 }
1755 Ok(())
1756}
1757
1758fn validate_system_contract(
1759 _plugin: &str,
1760 contract: &mut SystemContract,
1761) -> Result<(), CanwuError> {
1762 if contract.name.trim().is_empty() || contract.name != contract.name.trim() {
1763 return Err(CanwuError::new(
1764 ErrorCode::InvalidPluginRegistration,
1765 "plugin system name must be non-empty and have no surrounding whitespace",
1766 ));
1767 }
1768 if matches!(
1769 contract.phase,
1770 BoundaryPhase::EventIngress
1771 | BoundaryPhase::BoundarySnapshot
1772 | BoundaryPhase::AtomicDomainCommit
1773 | BoundaryPhase::ConditionalTransitionCommit
1774 ) {
1775 return Err(CanwuError::new(
1776 ErrorCode::InvalidPluginRegistration,
1777 format!("boundary phase {:?} is owned by the kernel", contract.phase),
1778 ));
1779 }
1780 if contract.cadence != SystemCadence::EventDriven {
1781 return Err(CanwuError::new(
1782 ErrorCode::InvalidPluginRegistration,
1783 format!(
1784 "system {} declares {:?} cadence, but the current runtime systems are event-driven only",
1785 contract.name, contract.cadence
1786 ),
1787 ));
1788 }
1789 if contract.visibility != StateVisibility::SameBoundary {
1790 return Err(CanwuError::new(
1791 ErrorCode::InvalidPluginRegistration,
1792 format!(
1793 "event-driven system {} must declare same-boundary visibility until the phased boundary runtime is active",
1794 contract.name
1795 ),
1796 ));
1797 }
1798 validate_state_keys(&mut contract.reads)?;
1799 validate_state_keys(&mut contract.writes)?;
1800 if contract.reads.contains(&StateKey::core_ingress()) {
1801 return Err(CanwuError::new(
1802 ErrorCode::InvalidPluginRegistration,
1803 "canonical ingress can be read only by phased boundary systems",
1804 ));
1805 }
1806 Ok(())
1807}
1808
1809fn validate_action_descriptor(
1810 _plugin: &str,
1811 descriptor: &mut PluginActionDescriptor,
1812) -> Result<(), CanwuError> {
1813 if descriptor.name.trim().is_empty() || descriptor.name != descriptor.name.trim() {
1814 return Err(CanwuError::new(
1815 ErrorCode::InvalidPluginRegistration,
1816 "plugin command names must be non-empty and have no surrounding whitespace",
1817 ));
1818 }
1819 if let PayloadSchema::Object { properties, .. } = &descriptor.payload_schema
1820 && properties
1821 .keys()
1822 .any(|name| name.trim().is_empty() || name != name.trim())
1823 {
1824 return Err(CanwuError::new(
1825 ErrorCode::InvalidPluginRegistration,
1826 "plugin payload schema property names cannot be empty",
1827 ));
1828 }
1829 validate_state_keys(&mut descriptor.reads)?;
1830 validate_state_keys(&mut descriptor.writes)?;
1831 if descriptor.reads.contains(&StateKey::core_ingress()) {
1832 return Err(CanwuError::new(
1833 ErrorCode::InvalidPluginRegistration,
1834 "plugin commands cannot inspect the canonical ingress queue",
1835 ));
1836 }
1837 Ok(())
1838}
1839
1840fn validate_ingress_descriptor(descriptor: &PluginIngressDescriptor) -> Result<(), CanwuError> {
1841 if descriptor.name.trim().is_empty()
1842 || descriptor.name != descriptor.name.trim()
1843 || descriptor.description.trim().is_empty()
1844 || descriptor.description != descriptor.description.trim()
1845 || descriptor.class == IngressClass::Command
1846 {
1847 return Err(CanwuError::new(
1848 ErrorCode::InvalidPluginRegistration,
1849 "plugin ingress types require canonical names/descriptions and cannot claim the core command class",
1850 ));
1851 }
1852 if let PayloadSchema::Object { properties, .. } = &descriptor.payload_schema
1853 && properties
1854 .keys()
1855 .any(|name| name.trim().is_empty() || name != name.trim())
1856 {
1857 return Err(CanwuError::new(
1858 ErrorCode::InvalidPluginRegistration,
1859 "plugin ingress payload property names cannot be empty",
1860 ));
1861 }
1862 Ok(())
1863}
1864
1865fn validate_boundary_system_contract(
1866 contract: &mut BoundarySystemContract,
1867) -> Result<(), CanwuError> {
1868 if contract.name.trim().is_empty() || contract.name != contract.name.trim() {
1869 return Err(CanwuError::new(
1870 ErrorCode::InvalidPluginRegistration,
1871 "boundary system name must be non-empty and canonical",
1872 ));
1873 }
1874 validate_state_keys(&mut contract.reads)?;
1875 validate_state_keys(&mut contract.writes)?;
1876 validate_state_keys(&mut contract.reservation_offers)?;
1877 validate_state_keys(&mut contract.reservation_requests)?;
1878 validate_reservation_refs(&mut contract.reservation_reads)?;
1879 validate_random_stream_keys(&mut contract.random_streams)?;
1880 validate_canonical_names(&mut contract.emits, "boundary event type")?;
1881 for grant in &mut contract.knowledge_writes {
1882 grant.visibilities.sort();
1883 grant.visibilities.dedup();
1884 if grant.visibilities.is_empty() {
1885 return Err(CanwuError::new(
1886 ErrorCode::InvalidPluginRegistration,
1887 "knowledge write grants require at least one visibility",
1888 ));
1889 }
1890 }
1891 contract
1892 .knowledge_writes
1893 .sort_by(|left, right| left.schema.cmp(&right.schema));
1894 if contract
1895 .knowledge_writes
1896 .windows(2)
1897 .any(|pair| pair[0].schema >= pair[1].schema)
1898 {
1899 return Err(CanwuError::new(
1900 ErrorCode::InvalidPluginRegistration,
1901 "knowledge write grants must name unique schemas in canonical order",
1902 ));
1903 }
1904 if !contract.knowledge_writes.is_empty()
1905 && !matches!(
1906 contract.phase,
1907 BoundaryPhase::PerceptionAndAttentionRefresh
1908 | BoundaryPhase::PerspectiveAndReportMaterialization
1909 )
1910 {
1911 return Err(CanwuError::new(
1912 ErrorCode::InvalidPluginRegistration,
1913 "knowledge publication is available only in phases 4 and 13",
1914 ));
1915 }
1916 if contract.plugin_ingress_targets.iter().any(|target| {
1917 target.target_plugin.trim().is_empty()
1918 || target.target_plugin != target.target_plugin.trim()
1919 || target.packet_type.trim().is_empty()
1920 || target.packet_type != target.packet_type.trim()
1921 }) {
1922 return Err(CanwuError::new(
1923 ErrorCode::InvalidPluginRegistration,
1924 "cross-plugin ingress targets require canonical plugin and packet names",
1925 ));
1926 }
1927 contract.plugin_ingress_targets.sort();
1928 if contract
1929 .plugin_ingress_targets
1930 .windows(2)
1931 .any(|pair| pair[0] >= pair[1])
1932 {
1933 return Err(CanwuError::new(
1934 ErrorCode::InvalidPluginRegistration,
1935 "cross-plugin ingress targets must be unique",
1936 ));
1937 }
1938
1939 let may_propose_changes = matches!(
1940 contract.phase,
1941 BoundaryPhase::DomainDeltaProposal
1942 | BoundaryPhase::HistoricalCandidateEvaluation
1943 | BoundaryPhase::StrategicAggregation
1944 | BoundaryPhase::PerspectiveAndReportMaterialization
1945 );
1946 if (!contract.writes.is_empty()
1947 || !contract.emits.is_empty()
1948 || !contract.plugin_ingress_targets.is_empty())
1949 && !may_propose_changes
1950 {
1951 return Err(CanwuError::new(
1952 ErrorCode::InvalidPluginRegistration,
1953 format!(
1954 "boundary system {} declares changes in kernel-owned phase {:?}",
1955 contract.name, contract.phase
1956 ),
1957 ));
1958 }
1959 let declares_reservations =
1960 !contract.reservation_offers.is_empty() || !contract.reservation_requests.is_empty();
1961 if declares_reservations && contract.phase != BoundaryPhase::ReservationAndAllocation {
1962 return Err(CanwuError::new(
1963 ErrorCode::InvalidPluginRegistration,
1964 format!(
1965 "boundary system {} declares reservations outside reservation and allocation",
1966 contract.name
1967 ),
1968 ));
1969 }
1970 if !contract.reservation_reads.is_empty()
1971 && contract.phase <= BoundaryPhase::ReservationAndAllocation
1972 {
1973 return Err(CanwuError::new(
1974 ErrorCode::InvalidPluginRegistration,
1975 format!(
1976 "boundary system {} reads allocations before reservation commit",
1977 contract.name
1978 ),
1979 ));
1980 }
1981 Ok(())
1982}
1983
1984fn validate_knowledge_write_grants(
1985 plugin: &str,
1986 contract: &BoundarySystemContract,
1987 schemas: &super::knowledge::KnowledgeSchemas,
1988) -> Result<(), CanwuError> {
1989 for grant in &contract.knowledge_writes {
1990 let Some((owner, schema)) = schemas.get(&grant.schema) else {
1991 return Err(CanwuError::new(
1992 ErrorCode::InvalidPluginRegistration,
1993 format!(
1994 "boundary system {plugin}.{} names an unregistered knowledge schema",
1995 contract.name
1996 ),
1997 ));
1998 };
1999 if owner != plugin || !schema.writable {
2000 return Err(CanwuError::new(
2001 ErrorCode::InvalidPluginRegistration,
2002 format!(
2003 "boundary system {plugin}.{} cannot write a foreign or read-only knowledge schema",
2004 contract.name
2005 ),
2006 ));
2007 }
2008 }
2009 Ok(())
2010}
2011
2012fn validate_reservation_refs(values: &mut Vec<ReservationRef>) -> Result<(), CanwuError> {
2013 if values.iter().any(|reservation| {
2014 reservation.plugin.trim().is_empty()
2015 || reservation.plugin != reservation.plugin.trim()
2016 || reservation.system.trim().is_empty()
2017 || reservation.system != reservation.system.trim()
2018 || reservation.request.trim().is_empty()
2019 || reservation.request != reservation.request.trim()
2020 }) {
2021 return Err(CanwuError::new(
2022 ErrorCode::InvalidPluginRegistration,
2023 "reservation read declarations must be non-empty and canonical",
2024 ));
2025 }
2026 let unique: BTreeSet<_> = values.drain(..).collect();
2027 values.extend(unique);
2028 Ok(())
2029}
2030
2031fn validate_random_stream_keys(values: &mut Vec<RandomStreamKey>) -> Result<(), CanwuError> {
2032 if values.iter().any(|stream| {
2033 stream.namespace.trim().is_empty()
2034 || stream.namespace != stream.namespace.trim()
2035 || stream.name.trim().is_empty()
2036 || stream.name != stream.name.trim()
2037 || stream.version == 0
2038 }) {
2039 return Err(CanwuError::new(
2040 ErrorCode::InvalidPluginRegistration,
2041 "random stream declarations require canonical names and a nonzero version",
2042 ));
2043 }
2044 let unique: BTreeSet<_> = values.drain(..).collect();
2045 values.extend(unique);
2046 Ok(())
2047}
2048
2049fn validate_canonical_names(values: &mut Vec<String>, label: &str) -> Result<(), CanwuError> {
2050 if values
2051 .iter()
2052 .any(|value| value.trim().is_empty() || value != value.trim())
2053 {
2054 return Err(CanwuError::new(
2055 ErrorCode::InvalidPluginRegistration,
2056 format!("{label} declarations must be non-empty and canonical"),
2057 ));
2058 }
2059 let unique: BTreeSet<_> = values.drain(..).collect();
2060 values.extend(unique);
2061 Ok(())
2062}
2063
2064fn validate_event_audience_name(event_type: &str) -> Result<(), CanwuError> {
2065 if !canonical_text(event_type) {
2066 return Err(CanwuError::new(
2067 ErrorCode::InvalidPluginRegistration,
2068 "plugin event audience names must be non-empty and canonical",
2069 ));
2070 }
2071 Ok(())
2072}
2073
2074fn validate_event_audience(audience: &EventAudience) -> Result<(), CanwuError> {
2075 match audience {
2076 EventAudience::Actor(actor) if actor.get() == 0 => {
2077 return Err(CanwuError::new(
2078 ErrorCode::InvalidPluginRegistration,
2079 "plugin event audience actors must use positive actor IDs",
2080 ));
2081 }
2082 EventAudience::Actors(actors) => {
2083 if actors.is_empty() || actors.iter().any(|actor| actor.get() == 0) {
2084 return Err(CanwuError::new(
2085 ErrorCode::InvalidPluginRegistration,
2086 "plugin event audience actor lists must contain positive actor IDs",
2087 ));
2088 }
2089 if actors.windows(2).any(|pair| pair[0] >= pair[1]) {
2090 return Err(CanwuError::new(
2091 ErrorCode::InvalidPluginRegistration,
2092 "plugin event audience actor lists must be sorted and unique",
2093 ));
2094 }
2095 }
2096 _ => {}
2097 }
2098 Ok(())
2099}
2100
2101fn register_state_owners(
2102 owners: &mut BTreeMap<StateKey, String>,
2103 plugin: &str,
2104 writes: &[StateKey],
2105) -> Result<(), CanwuError> {
2106 for key in writes {
2107 if key.namespace == CORE_STATE_NAMESPACE {
2108 return Err(CanwuError::new(
2109 ErrorCode::InvalidPluginRegistration,
2110 format!(
2111 "plugin {plugin} cannot claim reserved state {}.{}",
2112 key.namespace, key.name
2113 ),
2114 ));
2115 }
2116 if let Some(existing) = owners.get(key)
2117 && existing != plugin
2118 {
2119 return Err(CanwuError::new(
2120 ErrorCode::DuplicateStateOwner,
2121 format!(
2122 "state {}.{} is owned by both {existing} and {plugin}",
2123 key.namespace, key.name
2124 ),
2125 ));
2126 }
2127 }
2128 for key in writes {
2129 owners.insert(key.clone(), plugin.to_owned());
2130 }
2131 Ok(())
2132}
2133
2134fn register_boundary_writers(
2135 writers: &mut BTreeMap<(BoundaryWriteStage, StateKey), (String, String)>,
2136 immediate_writes: &BTreeMap<StateKey, String>,
2137 plugin: &str,
2138 system: &str,
2139 phase: BoundaryPhase,
2140 declared_states: &[StateKey],
2141) -> Result<(), CanwuError> {
2142 let Some(stage) = boundary_write_stage(phase) else {
2143 if declared_states.is_empty() {
2144 return Ok(());
2145 }
2146 return Err(CanwuError::new(
2147 ErrorCode::InvalidPluginRegistration,
2148 format!("boundary phase {phase:?} cannot own state writes"),
2149 ));
2150 };
2151 for state in declared_states {
2152 if let Some(immediate_plugin) = immediate_writes.get(state) {
2153 return Err(CanwuError::new(
2154 ErrorCode::InvalidPluginRegistration,
2155 format!(
2156 "boundary state {}.{} conflicts with immediate writes from plugin {immediate_plugin}",
2157 state.namespace, state.name
2158 ),
2159 ));
2160 }
2161 if let Some((existing_plugin, existing_system)) = writers.get(&(stage, state.clone()))
2162 && (existing_plugin != plugin || existing_system != system)
2163 {
2164 return Err(CanwuError::new(
2165 ErrorCode::DuplicateBoundaryWriter,
2166 format!(
2167 "boundary state {}.{} is written by both {existing_plugin}.{existing_system} and {plugin}.{system}",
2168 state.namespace, state.name
2169 ),
2170 ));
2171 }
2172 }
2173 for state in declared_states {
2174 writers.insert(
2175 (stage, state.clone()),
2176 (plugin.to_owned(), system.to_owned()),
2177 );
2178 }
2179 Ok(())
2180}
2181
2182fn register_immediate_write_states(
2183 immediate_writes: &mut BTreeMap<StateKey, String>,
2184 boundary_writers: &BTreeMap<(BoundaryWriteStage, StateKey), (String, String)>,
2185 plugin: &str,
2186 writes: &[StateKey],
2187) -> Result<(), CanwuError> {
2188 for state in writes {
2189 if boundary_writers
2190 .keys()
2191 .any(|(_, boundary_state)| boundary_state == state)
2192 {
2193 return Err(CanwuError::new(
2194 ErrorCode::InvalidPluginRegistration,
2195 format!(
2196 "immediate state {}.{} conflicts with a phased boundary writer",
2197 state.namespace, state.name
2198 ),
2199 ));
2200 }
2201 if immediate_writes
2202 .get(state)
2203 .is_some_and(|existing| existing != plugin)
2204 {
2205 return Err(CanwuError::new(
2206 ErrorCode::DuplicateStateOwner,
2207 format!(
2208 "immediate state {}.{} is written by multiple plugins",
2209 state.namespace, state.name
2210 ),
2211 ));
2212 }
2213 }
2214 for state in writes {
2215 immediate_writes.insert(state.clone(), plugin.to_owned());
2216 }
2217 Ok(())
2218}
2219
2220fn register_reservation_offerers(
2221 offerers: &mut BTreeMap<StateKey, (String, String)>,
2222 plugin: &str,
2223 system: &str,
2224 offered_state: &[StateKey],
2225) -> Result<(), CanwuError> {
2226 for state in offered_state {
2227 if let Some((existing_plugin, existing_system)) = offerers.get(state)
2228 && (existing_plugin != plugin || existing_system != system)
2229 {
2230 return Err(CanwuError::new(
2231 ErrorCode::DuplicateReservationOfferer,
2232 format!(
2233 "reservation state {}.{} is offered by both {existing_plugin}.{existing_system} and {plugin}.{system}",
2234 state.namespace, state.name
2235 ),
2236 ));
2237 }
2238 }
2239 for state in offered_state {
2240 offerers.insert(state.clone(), (plugin.to_owned(), system.to_owned()));
2241 }
2242 Ok(())
2243}
2244
2245fn register_random_streams(
2246 owners: &mut BTreeMap<RandomStreamKey, (String, String)>,
2247 plugin: &str,
2248 system: &str,
2249 streams: &[RandomStreamKey],
2250) -> Result<(), CanwuError> {
2251 for stream in streams {
2252 if stream.namespace != plugin || stream.namespace == CORE_STATE_NAMESPACE {
2253 return Err(CanwuError::new(
2254 ErrorCode::InvalidPluginRegistration,
2255 format!(
2256 "random stream {}.{}@{} must use its owning plugin namespace {plugin}",
2257 stream.namespace, stream.name, stream.version
2258 ),
2259 ));
2260 }
2261 if let Some((existing_plugin, existing_system)) = owners.get(stream)
2262 && (existing_plugin != plugin || existing_system != system)
2263 {
2264 return Err(CanwuError::new(
2265 ErrorCode::InvalidPluginRegistration,
2266 format!(
2267 "random stream {}.{}@{} is owned by both {existing_plugin}.{existing_system} and {plugin}.{system}",
2268 stream.namespace, stream.name, stream.version
2269 ),
2270 ));
2271 }
2272 }
2273 for stream in streams {
2274 owners.insert(stream.clone(), (plugin.to_owned(), system.to_owned()));
2275 }
2276 Ok(())
2277}
2278
2279#[cfg(test)]
2280mod tests {
2281 use super::super::{KnowledgeSubjectSchema, KnowledgeSubjectTargetKind};
2282 use super::*;
2283 use canwu_core::{CoreEntityKind, KnowledgeRecordKind, KnowledgeSchemaId};
2284
2285 struct KnowledgeSchemaPlugin {
2286 name: &'static str,
2287 schemas: Vec<PluginKnowledgeSchema>,
2288 }
2289
2290 impl SimulationPlugin for KnowledgeSchemaPlugin {
2291 fn name(&self) -> &str {
2292 self.name
2293 }
2294
2295 fn version(&self) -> &'static str {
2296 "1"
2297 }
2298
2299 fn semantic_hash(&self) -> &'static str {
2300 "0000000000000000000000000000000000000000000000000000000000000001"
2301 }
2302
2303 fn register(&self, registrar: &mut PluginRegistrar<'_>) -> Result<(), CanwuError> {
2304 for schema in &self.schemas {
2305 registrar.register_knowledge_schema(schema.clone())?;
2306 }
2307 Ok(())
2308 }
2309 }
2310
2311 fn knowledge_kind() -> KnowledgeRecordKind {
2312 KnowledgeRecordKind::new("fixture.knowledge", "assessment")
2313 }
2314
2315 fn knowledge_schema(version: u32, writable: bool) -> PluginKnowledgeSchema {
2316 PluginKnowledgeSchema {
2317 id: KnowledgeSchemaId::new(knowledge_kind(), version),
2318 schema_hash: format!("{version:064x}"),
2319 writable,
2320 payload_schema: PayloadSchema::Any,
2321 subjects: vec![],
2322 }
2323 }
2324
2325 #[test]
2326 fn archive_retention_binding_rejects_restored_root_mismatch() {
2327 let root = "a".repeat(64);
2328 let other = "b".repeat(64);
2329 let binding = PluginArchiveRetentionBinding {
2330 packet_type: "archive_commit".to_owned(),
2331 namespace: "fixture.archive.directory".to_owned(),
2332 object_id_json_pointers: vec![
2333 "/commit/archive_head/membership_root".to_owned(),
2334 "/commit/pending_reachability/directory_root".to_owned(),
2335 ],
2336 };
2337 let mut registry = PluginRegistry::default();
2338 registry
2339 .archive_retention_bindings
2340 .insert(("fixture".to_owned(), "archive_commit".to_owned()), binding);
2341 let retention = vec![PluginArchiveRetention {
2342 namespace: "fixture.archive.directory".to_owned(),
2343 object_id: root.clone(),
2344 }];
2345 let payload = serde_json::json!({
2346 "commit": {
2347 "archive_head": { "membership_root": root },
2348 "pending_reachability": { "directory_root": other },
2349 }
2350 });
2351 assert_eq!(
2352 registry
2353 .validate_archive_retention("fixture", "archive_commit", &payload, &retention,)
2354 .unwrap_err()
2355 .code,
2356 ErrorCode::InvalidPayload
2357 );
2358 }
2359
2360 #[test]
2361 fn duplicate_schema_and_writable_conflicts_roll_back_registration() {
2362 let duplicate = KnowledgeSchemaPlugin {
2363 name: "duplicate-knowledge",
2364 schemas: vec![knowledge_schema(1, true), knowledge_schema(1, true)],
2365 };
2366 let mut registry = PluginRegistry::default();
2367 let mut types = SchemaRegistry::default();
2368 assert!(registry.register(&duplicate, &mut types).is_err());
2369 assert!(registry.descriptors.is_empty());
2370 assert!(registry.knowledge_schemas.is_empty());
2371 assert!(registry.knowledge_kind_owners.is_empty());
2372
2373 let two_writable = KnowledgeSchemaPlugin {
2374 name: "two-writable-knowledge",
2375 schemas: vec![knowledge_schema(1, true), knowledge_schema(2, true)],
2376 };
2377 assert!(registry.register(&two_writable, &mut types).is_err());
2378 assert!(registry.descriptors.is_empty());
2379 assert!(registry.knowledge_schemas.is_empty());
2380
2381 let first_owner = KnowledgeSchemaPlugin {
2382 name: "first-knowledge-owner",
2383 schemas: vec![knowledge_schema(1, true)],
2384 };
2385 registry
2386 .register(&first_owner, &mut types)
2387 .expect("the first kind owner should register");
2388 let before = registry.clone();
2389 let second_owner = KnowledgeSchemaPlugin {
2390 name: "second-knowledge-owner",
2391 schemas: vec![knowledge_schema(2, true)],
2392 };
2393 assert!(registry.register(&second_owner, &mut types).is_err());
2394 assert_eq!(registry.descriptors, before.descriptors);
2395 assert_eq!(registry.knowledge_schemas, before.knowledge_schemas);
2396 assert_eq!(registry.knowledge_kind_owners, before.knowledge_kind_owners);
2397 }
2398
2399 #[test]
2400 fn schema_hash_mismatch_blocks_exact_rehydration() {
2401 let plugin = KnowledgeSchemaPlugin {
2402 name: "rehydrated-knowledge",
2403 schemas: vec![knowledge_schema(1, true)],
2404 };
2405 let mut registry = PluginRegistry::default();
2406 let mut types = SchemaRegistry::default();
2407 registry
2408 .register(&plugin, &mut types)
2409 .expect("fixture plugin should register");
2410 let mut descriptors = registry.descriptors().cloned().collect::<Vec<_>>();
2411 descriptors[0].knowledge_schemas[0].schema_hash =
2412 "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff".to_owned();
2413 let mut rehydrated = PluginRegistry::from_descriptors(descriptors)
2414 .expect("the altered descriptor remains structurally valid");
2415 let error = rehydrated
2416 .register(&plugin, &mut SchemaRegistry::default())
2417 .expect_err("exact rehydration must compare the persisted schema hash");
2418 assert_eq!(error.code, ErrorCode::PluginManifestMismatch);
2419 }
2420
2421 #[test]
2422 fn schema_limit_accepts_boundary_and_rejects_plus_one_atomically() {
2423 let boundary = KnowledgeSchemaPlugin {
2424 name: "knowledge-limit-boundary",
2425 schemas: (1..=KnowledgeLimitsV1::CURRENT.schemas_per_plugin)
2426 .map(|version| {
2427 knowledge_schema(
2428 u32::try_from(version).expect("schema limit fits u32"),
2429 version == 1,
2430 )
2431 })
2432 .collect(),
2433 };
2434 let mut registry = PluginRegistry::default();
2435 registry
2436 .register(&boundary, &mut SchemaRegistry::default())
2437 .expect("the exact schema limit should be admitted");
2438 assert_eq!(
2439 registry.knowledge_schemas.len(),
2440 KnowledgeLimitsV1::CURRENT.schemas_per_plugin
2441 );
2442
2443 let overflow = KnowledgeSchemaPlugin {
2444 name: "knowledge-limit-overflow",
2445 schemas: (1..=KnowledgeLimitsV1::CURRENT.schemas_per_plugin + 1)
2446 .map(|version| {
2447 knowledge_schema(
2448 u32::try_from(version).expect("schema limit fits u32"),
2449 version == 1,
2450 )
2451 })
2452 .collect(),
2453 };
2454 let mut rejected = PluginRegistry::default();
2455 let error = rejected
2456 .register(&overflow, &mut SchemaRegistry::default())
2457 .expect_err("schema limit plus one must reject the whole plugin");
2458 assert_eq!(error.code, ErrorCode::InvalidPluginRegistration);
2459 assert!(rejected.descriptors.is_empty());
2460 assert!(rejected.knowledge_schemas.is_empty());
2461 }
2462
2463 #[test]
2464 fn knowledge_schema_registration_canonicalizes_roles_and_targets() {
2465 let mut schema = knowledge_schema(1, true);
2466 schema.subjects = vec![
2467 KnowledgeSubjectSchema {
2468 role: "zeta".to_owned(),
2469 targets: vec![
2470 KnowledgeSubjectTargetKind::AnyEntity,
2471 KnowledgeSubjectTargetKind::Core(CoreEntityKind::Person),
2472 KnowledgeSubjectTargetKind::AnyEntity,
2473 ],
2474 required: false,
2475 multiple: true,
2476 },
2477 KnowledgeSubjectSchema {
2478 role: "alpha".to_owned(),
2479 targets: vec![KnowledgeSubjectTargetKind::Event],
2480 required: true,
2481 multiple: false,
2482 },
2483 ];
2484 let plugin = KnowledgeSchemaPlugin {
2485 name: "canonical-knowledge",
2486 schemas: vec![schema],
2487 };
2488 let mut registry = PluginRegistry::default();
2489 registry
2490 .register(&plugin, &mut SchemaRegistry::default())
2491 .expect("registrar should canonicalize declarations transactionally");
2492 let stored = ®istry
2493 .descriptors
2494 .get(plugin.name)
2495 .expect("descriptor exists")
2496 .knowledge_schemas[0];
2497 assert_eq!(stored.subjects[0].role, "alpha");
2498 assert_eq!(stored.subjects[1].role, "zeta");
2499 assert_eq!(stored.subjects[1].targets.len(), 2);
2500 assert!(stored.validate().is_ok());
2501 }
2502
2503 #[test]
2504 fn invalid_schema_version_and_hash_roll_back_registration() {
2505 let mut version_zero = knowledge_schema(0, true);
2506 version_zero.schema_hash =
2507 "0000000000000000000000000000000000000000000000000000000000000000".to_owned();
2508 let invalid_version = KnowledgeSchemaPlugin {
2509 name: "invalid-knowledge-version",
2510 schemas: vec![version_zero],
2511 };
2512 let mut registry = PluginRegistry::default();
2513 let mut types = SchemaRegistry::default();
2514 assert!(registry.register(&invalid_version, &mut types).is_err());
2515 assert!(registry.descriptors.is_empty());
2516 assert!(registry.knowledge_schemas.is_empty());
2517
2518 let mut bad_hash = knowledge_schema(1, true);
2519 bad_hash.schema_hash = "not-a-canonical-hash".to_owned();
2520 let invalid_hash = KnowledgeSchemaPlugin {
2521 name: "invalid-knowledge-hash",
2522 schemas: vec![bad_hash],
2523 };
2524 assert!(registry.register(&invalid_hash, &mut types).is_err());
2525 assert!(registry.descriptors.is_empty());
2526 assert!(registry.knowledge_schemas.is_empty());
2527 }
2528
2529 #[test]
2530 fn knowledge_write_grants_reject_invalid_phase_and_foreign_owner() {
2531 #[allow(clippy::unnecessary_wraps)]
2532 fn no_op_boundary(
2533 _view: &crate::SimulationView<'_>,
2534 _context: &crate::BoundaryContext,
2535 ) -> Result<crate::BoundaryProposal, CanwuError> {
2536 Ok(crate::BoundaryProposal::default())
2537 }
2538
2539 let owner = KnowledgeSchemaPlugin {
2540 name: "knowledge-grant-owner",
2541 schemas: vec![knowledge_schema(1, true)],
2542 };
2543 let foreign = KnowledgeSchemaPlugin {
2544 name: "knowledge-grant-foreign",
2545 schemas: vec![PluginKnowledgeSchema {
2546 id: KnowledgeSchemaId::new(
2547 KnowledgeRecordKind::new("fixture.foreign", "assessment"),
2548 1,
2549 ),
2550 schema_hash: "f000000000000000000000000000000000000000000000000000000000000000"
2551 .to_owned(),
2552 writable: true,
2553 payload_schema: PayloadSchema::Any,
2554 subjects: Vec::new(),
2555 }],
2556 };
2557 let mut registry = PluginRegistry::default();
2558 let mut types = SchemaRegistry::default();
2559 registry
2560 .register(&owner, &mut types)
2561 .expect("knowledge owner should register");
2562 registry
2563 .register(&foreign, &mut types)
2564 .expect("foreign fixture should register");
2565 let before = registry.clone();
2566
2567 let mut phase7 = BoundarySystemContract::new(
2568 "invalid-phase7-publication",
2569 crate::BoundaryPhase::DomainDeltaProposal,
2570 SystemCadence::Daily,
2571 );
2572 phase7.knowledge_writes = vec![crate::KnowledgeWriteGrant {
2573 schema: knowledge_schema(1, true).id,
2574 visibilities: vec![StateVisibility::SameBoundary],
2575 }];
2576 let mut owner_registry = registry.clone();
2577 let mut owner_types = types.clone();
2578 let mut registrar = PluginRegistrar {
2579 plugin: owner.name.to_owned(),
2580 registry: &mut owner_registry,
2581 schema: &mut owner_types,
2582 };
2583 let error = registrar
2584 .register_boundary_system(phase7, no_op_boundary)
2585 .expect_err("phase 7 must reject knowledge publication grants");
2586 assert_eq!(error.code, ErrorCode::InvalidPluginRegistration);
2587 assert_eq!(owner_registry.descriptors, before.descriptors);
2588
2589 let mut foreign_grant = BoundarySystemContract::new(
2590 "foreign-knowledge-grant",
2591 crate::BoundaryPhase::PerspectiveAndReportMaterialization,
2592 SystemCadence::Daily,
2593 );
2594 foreign_grant.knowledge_writes = vec![crate::KnowledgeWriteGrant {
2595 schema: knowledge_schema(1, true).id,
2596 visibilities: vec![StateVisibility::SameBoundary],
2597 }];
2598 let mut foreign_registry = registry;
2599 let mut registrar = PluginRegistrar {
2600 plugin: foreign.name.to_owned(),
2601 registry: &mut foreign_registry,
2602 schema: &mut types,
2603 };
2604 let error = registrar
2605 .register_boundary_system(foreign_grant, no_op_boundary)
2606 .expect_err("a plugin cannot claim another plugin's writable schema");
2607 assert_eq!(error.code, ErrorCode::InvalidPluginRegistration);
2608 assert_eq!(foreign_registry.descriptors, before.descriptors);
2609 }
2610}