1use super::{
2 BoundaryPhase, BoundarySystemContract, BoundarySystemHandler, BoundaryWriteStage,
3 CORE_STATE_NAMESPACE, CanwuError, ErrorCode, IngressClass, PayloadSchema,
4 PluginActionDescriptor, PluginCommandHandler, PluginDescriptor, PluginIngressDescriptor,
5 PluginRegistrar, PluginRegistry, RandomStreamKey, RegisteredBoundarySystem, RegisteredCommand,
6 RegisteredSystem, ReservationRef, SchemaRegistry, SimulationPlugin, SimulationSystemHandler,
7 StateKey, StateVisibility, SystemCadence, SystemContract, TypeSchema, boundary_write_stage,
8 canonical_text, invalid_snapshot, invalid_snapshot_error, is_canonical_hash,
9 is_domain_record_state, validate_type_schema,
10};
11use crate::records::DomainRecordSchema;
12use canwu_event::EventAudience;
13use std::collections::{BTreeMap, BTreeSet};
14
15impl PluginRegistrar<'_> {
16 pub fn register_record_schema(
17 &mut self,
18 mut schema: DomainRecordSchema,
19 ) -> Result<(), CanwuError> {
20 schema.canonicalize();
21 schema.validate().map_err(|error| {
22 CanwuError::new(
23 ErrorCode::InvalidPluginRegistration,
24 format!("invalid domain record schema: {error}"),
25 )
26 })?;
27 let state = schema.state_key();
28 if state.namespace == CORE_STATE_NAMESPACE {
29 return Err(CanwuError::new(
30 ErrorCode::InvalidPluginRegistration,
31 "plugins cannot register domain record kinds in the core namespace",
32 ));
33 }
34 if self
35 .registry
36 .descriptors
37 .get(&self.plugin)
38 .is_some_and(|descriptor| {
39 descriptor
40 .record_schemas
41 .iter()
42 .any(|candidate| candidate.kind == schema.kind)
43 })
44 {
45 return Err(CanwuError::new(
46 ErrorCode::DuplicateDomainRecordKind,
47 format!(
48 "plugin {} registered record kind {} twice",
49 self.plugin, schema.kind
50 ),
51 ));
52 }
53 if let Some((owner, existing)) = self.registry.record_schemas.get(&schema.kind) {
54 if owner != &self.plugin {
55 return Err(CanwuError::new(
56 ErrorCode::DuplicateDomainRecordKind,
57 format!(
58 "domain record kind {} is already owned by plugin {owner}",
59 schema.kind
60 ),
61 ));
62 }
63 if existing != &schema {
64 return Err(CanwuError::new(
65 ErrorCode::PluginManifestMismatch,
66 format!(
67 "plugin {} changed the stored schema for domain record kind {}",
68 self.plugin, schema.kind
69 ),
70 ));
71 }
72 }
73 let mut candidate = self.registry.clone();
74 if candidate.immediate_write_states.contains_key(&state) {
75 return Err(CanwuError::new(
76 ErrorCode::InvalidPluginRegistration,
77 format!(
78 "domain record kind {} is already exposed as immediate component state",
79 schema.kind
80 ),
81 ));
82 }
83 register_state_owners(
84 &mut candidate.state_owners,
85 &self.plugin,
86 std::slice::from_ref(&state),
87 )?;
88 candidate
89 .record_schemas
90 .insert(schema.kind.clone(), (self.plugin.clone(), schema.clone()));
91 let descriptor = candidate
92 .descriptors
93 .entry(self.plugin.clone())
94 .or_default();
95 descriptor.name.clone_from(&self.plugin);
96 descriptor.record_schemas.push(schema);
97 descriptor
98 .record_schemas
99 .sort_by(|left, right| left.kind.cmp(&right.kind));
100 *self.registry = candidate;
101 Ok(())
102 }
103
104 pub fn register_schema(&mut self, schema: TypeSchema) -> Result<(), CanwuError> {
105 validate_type_schema(&schema)?;
106 let type_name = schema.type_name.clone();
107 let mut candidate_schema = self.schema.clone();
108 let mut candidate_registry = self.registry.clone();
109 if let Some(existing) = candidate_schema.get(&type_name) {
110 if existing != &schema {
111 return Err(CanwuError::new(
112 ErrorCode::InvalidPluginRegistration,
113 format!(
114 "schema type {type_name} is already registered with a different definition"
115 ),
116 ));
117 }
118 } else {
119 candidate_schema.register(schema);
120 }
121 let descriptor = candidate_registry
122 .descriptors
123 .entry(self.plugin.clone())
124 .or_default();
125 if descriptor.schema_types.contains(&type_name) {
126 return Err(CanwuError::new(
127 ErrorCode::InvalidPluginRegistration,
128 format!(
129 "plugin {} registered schema type {} more than once",
130 self.plugin, type_name
131 ),
132 ));
133 }
134 descriptor.name.clone_from(&self.plugin);
135 descriptor.schema_types.push(type_name);
136 descriptor.schema_types.sort();
137 *self.schema = candidate_schema;
138 *self.registry = candidate_registry;
139 Ok(())
140 }
141
142 pub fn register_event_audience(
148 &mut self,
149 event_type: impl Into<String>,
150 audience: EventAudience,
151 ) -> Result<(), CanwuError> {
152 let event_type = event_type.into();
153 validate_event_audience_name(&event_type)?;
154 validate_event_audience(&audience)?;
155 let mut candidate = self.registry.clone();
156 let descriptor = candidate
157 .descriptors
158 .entry(self.plugin.clone())
159 .or_default();
160 descriptor.name.clone_from(&self.plugin);
161 if descriptor
162 .event_audiences
163 .insert(event_type.clone(), audience)
164 .is_some()
165 {
166 return Err(CanwuError::new(
167 ErrorCode::InvalidPluginRegistration,
168 format!(
169 "plugin {} already declared event audience for {event_type}",
170 self.plugin
171 ),
172 ));
173 }
174 *self.registry = candidate;
175 Ok(())
176 }
177
178 pub fn register_system(
179 &mut self,
180 mut contract: SystemContract,
181 handler: SimulationSystemHandler,
182 ) -> Result<(), CanwuError> {
183 validate_system_contract(&self.plugin, &mut contract)?;
187 if self
188 .registry
189 .descriptors
190 .get(&self.plugin)
191 .is_some_and(|descriptor| {
192 descriptor
193 .systems
194 .iter()
195 .any(|candidate| candidate.name == contract.name)
196 || descriptor
197 .boundary_systems
198 .iter()
199 .any(|candidate| candidate.name == contract.name)
200 })
201 {
202 return Err(CanwuError::new(
203 ErrorCode::DuplicatePluginSystem,
204 format!(
205 "plugin {} already registered system {}",
206 self.plugin, contract.name
207 ),
208 ));
209 }
210 let mut candidate = self.registry.clone();
211 if contract
212 .writes
213 .iter()
214 .any(|state| is_domain_record_state(&candidate.record_schemas, state))
215 {
216 return Err(CanwuError::new(
217 ErrorCode::InvalidPluginRegistration,
218 "domain record kinds can only be mutated by phased boundary systems",
219 ));
220 }
221 register_state_owners(&mut candidate.state_owners, &self.plugin, &contract.writes)?;
222 register_immediate_write_states(
223 &mut candidate.immediate_write_states,
224 &candidate.boundary_writers,
225 &self.plugin,
226 &contract.writes,
227 )?;
228 {
229 let descriptor = candidate
230 .descriptors
231 .entry(self.plugin.clone())
232 .or_default();
233 descriptor.name.clone_from(&self.plugin);
234 descriptor.systems.push(contract.clone());
235 descriptor
236 .systems
237 .sort_by(|left, right| (left.phase, &left.name).cmp(&(right.phase, &right.name)));
238 }
239 candidate.systems.push(RegisteredSystem {
240 plugin: self.plugin.clone(),
241 contract,
242 handler,
243 });
244 candidate.systems.sort_by(|left, right| {
245 (left.contract.phase, &left.plugin, &left.contract.name).cmp(&(
246 right.contract.phase,
247 &right.plugin,
248 &right.contract.name,
249 ))
250 });
251 *self.registry = candidate;
252 Ok(())
253 }
254
255 pub fn register_boundary_system(
256 &mut self,
257 mut contract: BoundarySystemContract,
258 handler: BoundarySystemHandler,
259 ) -> Result<(), CanwuError> {
260 validate_boundary_system_contract(&mut contract)?;
261 if self
262 .registry
263 .descriptors
264 .get(&self.plugin)
265 .is_some_and(|descriptor| {
266 descriptor
267 .systems
268 .iter()
269 .any(|candidate| candidate.name == contract.name)
270 || descriptor
271 .boundary_systems
272 .iter()
273 .any(|candidate| candidate.name == contract.name)
274 })
275 {
276 return Err(CanwuError::new(
277 ErrorCode::DuplicatePluginSystem,
278 format!(
279 "plugin {} already registered system {}",
280 self.plugin, contract.name
281 ),
282 ));
283 }
284 let mut owned_state = contract.writes.clone();
285 owned_state.extend(contract.reservation_offers.iter().cloned());
286 owned_state.sort();
287 owned_state.dedup();
288 let mut candidate = self.registry.clone();
289 register_state_owners(&mut candidate.state_owners, &self.plugin, &owned_state)?;
290 register_boundary_writers(
291 &mut candidate.boundary_writers,
292 &candidate.immediate_write_states,
293 &self.plugin,
294 &contract.name,
295 contract.phase,
296 &contract.writes,
297 )?;
298 register_reservation_offerers(
299 &mut candidate.reservation_offerers,
300 &self.plugin,
301 &contract.name,
302 &contract.reservation_offers,
303 )?;
304 register_random_streams(
305 &mut candidate.random_stream_owners,
306 &self.plugin,
307 &contract.name,
308 &contract.random_streams,
309 )?;
310 {
311 let descriptor = candidate
312 .descriptors
313 .entry(self.plugin.clone())
314 .or_default();
315 descriptor.name.clone_from(&self.plugin);
316 descriptor.boundary_systems.push(contract.clone());
317 descriptor
318 .boundary_systems
319 .sort_by(|left, right| (left.phase, &left.name).cmp(&(right.phase, &right.name)));
320 }
321 candidate.boundary_systems.push(RegisteredBoundarySystem {
322 plugin: self.plugin.clone(),
323 contract,
324 handler,
325 });
326 candidate.boundary_systems.sort_by(|left, right| {
327 (left.contract.phase, &left.plugin, &left.contract.name).cmp(&(
328 right.contract.phase,
329 &right.plugin,
330 &right.contract.name,
331 ))
332 });
333 *self.registry = candidate;
334 Ok(())
335 }
336
337 pub fn register_command(
338 &mut self,
339 mut descriptor: PluginActionDescriptor,
340 handler: PluginCommandHandler,
341 ) -> Result<(), CanwuError> {
342 validate_action_descriptor(&self.plugin, &mut descriptor)?;
343 let command_key = (self.plugin.clone(), descriptor.name.clone());
344 if self.registry.commands.contains_key(&command_key) {
345 return Err(CanwuError::new(
346 ErrorCode::DuplicatePluginCommand,
347 format!(
348 "plugin {} already registered command {}",
349 self.plugin, descriptor.name
350 ),
351 ));
352 }
353 let mut candidate = self.registry.clone();
354 if descriptor
355 .writes
356 .iter()
357 .any(|state| is_domain_record_state(&candidate.record_schemas, state))
358 {
359 return Err(CanwuError::new(
360 ErrorCode::InvalidPluginRegistration,
361 "plugin commands cannot write domain record state directly",
362 ));
363 }
364 register_state_owners(
365 &mut candidate.state_owners,
366 &self.plugin,
367 &descriptor.writes,
368 )?;
369 register_immediate_write_states(
370 &mut candidate.immediate_write_states,
371 &candidate.boundary_writers,
372 &self.plugin,
373 &descriptor.writes,
374 )?;
375 {
376 let plugin_descriptor = candidate
377 .descriptors
378 .entry(self.plugin.clone())
379 .or_default();
380 plugin_descriptor.name.clone_from(&self.plugin);
381 plugin_descriptor.commands.push(descriptor.clone());
382 plugin_descriptor
383 .commands
384 .sort_by(|left, right| left.name.cmp(&right.name));
385 }
386 candidate.commands.insert(
387 command_key,
388 RegisteredCommand {
389 descriptor,
390 handler,
391 },
392 );
393 *self.registry = candidate;
394 Ok(())
395 }
396
397 pub fn register_ingress(
398 &mut self,
399 descriptor: PluginIngressDescriptor,
400 ) -> Result<(), CanwuError> {
401 validate_ingress_descriptor(&descriptor)?;
402 let key = (self.plugin.clone(), descriptor.name.clone());
403 if self
404 .registry
405 .descriptors
406 .get(&self.plugin)
407 .is_some_and(|plugin| {
408 plugin
409 .ingress
410 .iter()
411 .any(|candidate| candidate.name == descriptor.name)
412 })
413 {
414 return Err(CanwuError::new(
415 ErrorCode::DuplicatePluginIngress,
416 format!(
417 "plugin {} already registered ingress type {}",
418 self.plugin, descriptor.name
419 ),
420 ));
421 }
422 if self
423 .registry
424 .ingress
425 .get(&key)
426 .is_some_and(|existing| existing != &descriptor)
427 {
428 return Err(CanwuError::new(
429 ErrorCode::PluginManifestMismatch,
430 format!(
431 "plugin {} changed the stored ingress type {}",
432 self.plugin, descriptor.name
433 ),
434 ));
435 }
436 let mut candidate = self.registry.clone();
437 candidate.ingress.insert(key, descriptor.clone());
438 let plugin_descriptor = candidate
439 .descriptors
440 .entry(self.plugin.clone())
441 .or_default();
442 plugin_descriptor.name.clone_from(&self.plugin);
443 plugin_descriptor.ingress.push(descriptor);
444 plugin_descriptor
445 .ingress
446 .sort_by(|left, right| left.name.cmp(&right.name));
447 *self.registry = candidate;
448 Ok(())
449 }
450}
451
452impl PluginRegistry {
453 pub fn register<P: SimulationPlugin + ?Sized>(
454 &mut self,
455 plugin: &P,
456 schema: &mut SchemaRegistry,
457 ) -> Result<(), CanwuError> {
458 let raw_plugin_name = plugin.name();
459 let plugin_name = raw_plugin_name.trim();
460 if plugin_name.is_empty() || plugin_name != raw_plugin_name {
461 return Err(CanwuError::new(
462 ErrorCode::InvalidPluginRegistration,
463 "plugin name must be non-empty and have no surrounding whitespace",
464 ));
465 }
466 if self.active_plugins.contains(plugin_name) {
467 return Err(CanwuError::new(
468 ErrorCode::DuplicatePlugin,
469 format!("plugin {plugin_name} is already registered"),
470 ));
471 }
472 validate_plugin_identity(plugin_name, plugin.version(), plugin.semantic_hash())?;
473
474 let expected_descriptor = self.descriptors.get(plugin_name).cloned();
475 let mut candidate_registry = self.clone();
476 let mut candidate_schema = schema.clone();
477 candidate_registry.descriptors.insert(
478 plugin_name.to_owned(),
479 PluginDescriptor {
480 name: plugin_name.to_owned(),
481 version: plugin.version().to_owned(),
482 semantic_hash: plugin.semantic_hash().to_owned(),
483 ..PluginDescriptor::default()
484 },
485 );
486 let mut registrar = PluginRegistrar {
487 plugin: plugin_name.to_owned(),
488 registry: &mut candidate_registry,
489 schema: &mut candidate_schema,
490 };
491 plugin.register(&mut registrar)?;
492 let Some(generated_descriptor) = candidate_registry.descriptors.get(plugin_name) else {
493 return Err(CanwuError::new(
494 ErrorCode::InvalidPluginRegistration,
495 format!("plugin {plugin_name} did not produce a descriptor"),
496 ));
497 };
498 if let Some(expected) = expected_descriptor
499 && generated_descriptor != &expected
500 {
501 return Err(CanwuError::new(
502 ErrorCode::PluginManifestMismatch,
503 format!("plugin {plugin_name} registration does not match the snapshot manifest"),
504 ));
505 }
506 candidate_registry
507 .active_plugins
508 .insert(plugin_name.to_owned());
509 *self = candidate_registry;
510 *schema = candidate_schema;
511 Ok(())
512 }
513
514 pub fn descriptors(&self) -> impl Iterator<Item = &PluginDescriptor> {
515 self.descriptors.values()
516 }
517
518 pub(super) fn event_audience(&self, plugin: &str, event_type: &str) -> EventAudience {
519 self.descriptors
520 .get(plugin)
521 .and_then(|descriptor| descriptor.event_audiences.get(event_type))
522 .cloned()
523 .unwrap_or_default()
524 }
525
526 pub(super) fn from_descriptors(descriptors: Vec<PluginDescriptor>) -> Result<Self, CanwuError> {
527 let mut registry = Self {
528 descriptors: BTreeMap::new(),
529 active_plugins: BTreeSet::new(),
530 systems: Vec::new(),
531 boundary_systems: Vec::new(),
532 commands: BTreeMap::new(),
533 ingress: BTreeMap::new(),
534 state_owners: BTreeMap::new(),
535 immediate_write_states: BTreeMap::new(),
536 boundary_writers: BTreeMap::new(),
537 reservation_offerers: BTreeMap::new(),
538 random_stream_owners: BTreeMap::new(),
539 record_schemas: BTreeMap::new(),
540 };
541 let mut previous_plugin = None;
542 for mut descriptor in descriptors {
543 let plugin = descriptor.name.trim().to_owned();
544 if plugin.is_empty()
545 || descriptor.name != plugin
546 || descriptor.version.trim().is_empty()
547 || descriptor.version != descriptor.version.trim()
548 || !is_canonical_hash(&descriptor.semantic_hash)
549 || registry.descriptors.contains_key(&plugin)
550 || previous_plugin
551 .as_ref()
552 .is_some_and(|previous| previous >= &plugin)
553 {
554 return Err(CanwuError::new(
555 ErrorCode::InvalidSnapshot,
556 "snapshot contains an invalid, unversioned, or duplicate plugin descriptor",
557 ));
558 }
559 if descriptor
560 .record_schemas
561 .windows(2)
562 .any(|pair| pair[0].kind >= pair[1].kind)
563 {
564 return invalid_snapshot("plugin record schemas are not in canonical order");
565 }
566 for schema in &mut descriptor.record_schemas {
567 let original = schema.clone();
568 schema.canonicalize();
569 schema.validate().map_err(|error| {
570 invalid_snapshot_error(format!("invalid domain record schema: {error}"))
571 })?;
572 if *schema != original {
573 return invalid_snapshot(
574 "plugin record-schema declarations are not in canonical order",
575 );
576 }
577 let state = schema.state_key();
578 if state.namespace == CORE_STATE_NAMESPACE {
579 return invalid_snapshot(
580 "plugin record schemas cannot use the reserved core namespace",
581 );
582 }
583 if let Some((owner, _)) = registry.record_schemas.get(&schema.kind) {
584 return invalid_snapshot(format!(
585 "domain record kind {} is owned by both {owner} and {plugin}",
586 schema.kind
587 ));
588 }
589 register_state_owners(
590 &mut registry.state_owners,
591 &plugin,
592 std::slice::from_ref(&state),
593 )
594 .map_err(|error| {
595 invalid_snapshot_error(format!(
596 "invalid domain record state ownership descriptor: {error}"
597 ))
598 })?;
599 registry
600 .record_schemas
601 .insert(schema.kind.clone(), (plugin.clone(), schema.clone()));
602 }
603 if descriptor
604 .systems
605 .windows(2)
606 .any(|pair| (pair[0].phase, &pair[0].name) >= (pair[1].phase, &pair[1].name))
607 {
608 return invalid_snapshot("plugin systems are not in canonical order");
609 }
610 let mut system_names = BTreeSet::new();
611 for contract in &mut descriptor.systems {
612 if !system_names.insert(contract.name.clone()) {
613 return invalid_snapshot("plugin descriptor has duplicate system names");
614 }
615 let original = contract.clone();
616 validate_system_contract(&plugin, contract).map_err(|error| {
617 invalid_snapshot_error(format!("invalid plugin system descriptor: {error}"))
618 })?;
619 if *contract != original {
620 return invalid_snapshot(
621 "plugin system reads and writes are not in canonical order",
622 );
623 }
624 if contract
625 .writes
626 .iter()
627 .any(|state| is_domain_record_state(®istry.record_schemas, state))
628 {
629 return invalid_snapshot(
630 "plugin systems cannot expose domain records as immediate component state",
631 );
632 }
633 register_state_owners(&mut registry.state_owners, &plugin, &contract.writes)
634 .map_err(|error| {
635 invalid_snapshot_error(format!(
636 "invalid plugin state ownership descriptor: {error}"
637 ))
638 })?;
639 register_immediate_write_states(
640 &mut registry.immediate_write_states,
641 ®istry.boundary_writers,
642 &plugin,
643 &contract.writes,
644 )
645 .map_err(|error| {
646 invalid_snapshot_error(format!(
647 "invalid immediate state writer descriptor: {error}"
648 ))
649 })?;
650 }
651 if descriptor
652 .boundary_systems
653 .windows(2)
654 .any(|pair| (pair[0].phase, &pair[0].name) >= (pair[1].phase, &pair[1].name))
655 {
656 return invalid_snapshot("boundary systems are not in canonical order");
657 }
658 for contract in &mut descriptor.boundary_systems {
659 if !system_names.insert(contract.name.clone()) {
660 return invalid_snapshot("plugin descriptor has duplicate system names");
661 }
662 let original = contract.clone();
663 validate_boundary_system_contract(contract).map_err(|error| {
664 invalid_snapshot_error(format!("invalid boundary system descriptor: {error}"))
665 })?;
666 if *contract != original {
667 return invalid_snapshot(
668 "boundary system declarations are not in canonical order",
669 );
670 }
671 let mut owned_state = contract.writes.clone();
672 owned_state.extend(contract.reservation_offers.iter().cloned());
673 owned_state.sort();
674 owned_state.dedup();
675 register_state_owners(&mut registry.state_owners, &plugin, &owned_state).map_err(
676 |error| {
677 invalid_snapshot_error(format!(
678 "invalid boundary state ownership descriptor: {error}"
679 ))
680 },
681 )?;
682 register_boundary_writers(
683 &mut registry.boundary_writers,
684 ®istry.immediate_write_states,
685 &plugin,
686 &contract.name,
687 contract.phase,
688 &contract.writes,
689 )
690 .map_err(|error| {
691 invalid_snapshot_error(format!("invalid boundary writer descriptor: {error}"))
692 })?;
693 register_reservation_offerers(
694 &mut registry.reservation_offerers,
695 &plugin,
696 &contract.name,
697 &contract.reservation_offers,
698 )
699 .map_err(|error| {
700 invalid_snapshot_error(format!(
701 "invalid reservation offerer descriptor: {error}"
702 ))
703 })?;
704 register_random_streams(
705 &mut registry.random_stream_owners,
706 &plugin,
707 &contract.name,
708 &contract.random_streams,
709 )
710 .map_err(|error| {
711 invalid_snapshot_error(format!(
712 "invalid random stream ownership descriptor: {error}"
713 ))
714 })?;
715 }
716 if descriptor
717 .commands
718 .windows(2)
719 .any(|pair| pair[0].name >= pair[1].name)
720 {
721 return invalid_snapshot("plugin commands are not in canonical order");
722 }
723 let mut command_names = BTreeSet::new();
724 for action in &mut descriptor.commands {
725 if !command_names.insert(action.name.clone()) {
726 return invalid_snapshot("plugin descriptor has duplicate command names");
727 }
728 let original = action.clone();
729 validate_action_descriptor(&plugin, action).map_err(|error| {
730 invalid_snapshot_error(format!("invalid plugin command descriptor: {error}"))
731 })?;
732 if *action != original {
733 return invalid_snapshot(
734 "plugin command reads and writes are not in canonical order",
735 );
736 }
737 if action
738 .writes
739 .iter()
740 .any(|state| is_domain_record_state(®istry.record_schemas, state))
741 {
742 return invalid_snapshot(
743 "plugin commands cannot expose domain records as immediate component state",
744 );
745 }
746 register_state_owners(&mut registry.state_owners, &plugin, &action.writes)
747 .map_err(|error| {
748 invalid_snapshot_error(format!(
749 "invalid plugin state ownership descriptor: {error}"
750 ))
751 })?;
752 register_immediate_write_states(
753 &mut registry.immediate_write_states,
754 ®istry.boundary_writers,
755 &plugin,
756 &action.writes,
757 )
758 .map_err(|error| {
759 invalid_snapshot_error(format!(
760 "invalid immediate state writer descriptor: {error}"
761 ))
762 })?;
763 }
764 if descriptor
765 .ingress
766 .windows(2)
767 .any(|pair| pair[0].name >= pair[1].name)
768 {
769 return invalid_snapshot("plugin ingress types are not in canonical order");
770 }
771 for ingress in &descriptor.ingress {
772 validate_ingress_descriptor(ingress).map_err(|error| {
773 invalid_snapshot_error(format!("invalid plugin ingress descriptor: {error}"))
774 })?;
775 if registry
776 .ingress
777 .insert((plugin.clone(), ingress.name.clone()), ingress.clone())
778 .is_some()
779 {
780 return invalid_snapshot("plugin descriptor has duplicate ingress types");
781 }
782 }
783 for (event_type, audience) in &descriptor.event_audiences {
784 validate_event_audience_name(event_type).map_err(|error| {
785 invalid_snapshot_error(format!("invalid plugin event audience: {error}"))
786 })?;
787 validate_event_audience(audience).map_err(|error| {
788 invalid_snapshot_error(format!("invalid plugin event audience: {error}"))
789 })?;
790 }
791 let schema_types: BTreeSet<_> = descriptor.schema_types.iter().collect();
792 if schema_types.len() != descriptor.schema_types.len()
793 || descriptor
794 .schema_types
795 .windows(2)
796 .any(|pair| pair[0] >= pair[1])
797 || descriptor
798 .schema_types
799 .iter()
800 .any(|name| name.trim().is_empty() || name != name.trim())
801 {
802 return invalid_snapshot("plugin descriptor has invalid schema type names");
803 }
804 previous_plugin = Some(plugin.clone());
805 registry.descriptors.insert(plugin, descriptor);
806 }
807 Ok(registry)
808 }
809
810 pub(super) fn ensure_active(&self) -> Result<(), CanwuError> {
811 let inactive: Vec<_> = self
812 .descriptors
813 .keys()
814 .filter(|name| !self.active_plugins.contains(*name))
815 .cloned()
816 .collect();
817 if inactive.is_empty() {
818 return Ok(());
819 }
820 Err(CanwuError::new(
821 ErrorCode::PluginNotActive,
822 format!(
823 "required plugin handlers are not active: {}",
824 inactive.join(", ")
825 ),
826 ))
827 }
828}
829
830pub(super) fn validate_state_keys(keys: &mut Vec<StateKey>) -> Result<(), CanwuError> {
831 for key in keys.iter() {
832 if key.namespace.trim().is_empty()
833 || key.name.trim().is_empty()
834 || key.namespace != key.namespace.trim()
835 || key.name != key.name.trim()
836 {
837 return Err(CanwuError::new(
838 ErrorCode::InvalidPluginRegistration,
839 "state keys require non-empty canonical namespace and name values",
840 ));
841 }
842 }
843 let unique: BTreeSet<_> = keys.drain(..).collect();
844 keys.extend(unique);
845 Ok(())
846}
847
848fn validate_plugin_identity(
849 name: &str,
850 version: &str,
851 semantic_hash: &str,
852) -> Result<(), CanwuError> {
853 if name.trim().is_empty()
854 || name != name.trim()
855 || version.trim().is_empty()
856 || version != version.trim()
857 || !is_canonical_hash(semantic_hash)
858 {
859 return Err(CanwuError::new(
860 ErrorCode::InvalidPluginRegistration,
861 "plugins require canonical names, versions, and 64-character semantic hashes",
862 ));
863 }
864 Ok(())
865}
866
867fn validate_system_contract(
868 _plugin: &str,
869 contract: &mut SystemContract,
870) -> Result<(), CanwuError> {
871 if contract.name.trim().is_empty() || contract.name != contract.name.trim() {
872 return Err(CanwuError::new(
873 ErrorCode::InvalidPluginRegistration,
874 "plugin system name must be non-empty and have no surrounding whitespace",
875 ));
876 }
877 if matches!(
878 contract.phase,
879 BoundaryPhase::EventIngress
880 | BoundaryPhase::BoundarySnapshot
881 | BoundaryPhase::AtomicDomainCommit
882 | BoundaryPhase::ConditionalTransitionCommit
883 ) {
884 return Err(CanwuError::new(
885 ErrorCode::InvalidPluginRegistration,
886 format!("boundary phase {:?} is owned by the kernel", contract.phase),
887 ));
888 }
889 if contract.cadence != SystemCadence::EventDriven {
890 return Err(CanwuError::new(
891 ErrorCode::InvalidPluginRegistration,
892 format!(
893 "system {} declares {:?} cadence, but the current runtime systems are event-driven only",
894 contract.name, contract.cadence
895 ),
896 ));
897 }
898 if contract.visibility != StateVisibility::SameBoundary {
899 return Err(CanwuError::new(
900 ErrorCode::InvalidPluginRegistration,
901 format!(
902 "event-driven system {} must declare same-boundary visibility until the phased boundary runtime is active",
903 contract.name
904 ),
905 ));
906 }
907 validate_state_keys(&mut contract.reads)?;
908 validate_state_keys(&mut contract.writes)?;
909 if contract.reads.contains(&StateKey::core_ingress()) {
910 return Err(CanwuError::new(
911 ErrorCode::InvalidPluginRegistration,
912 "canonical ingress can be read only by phased boundary systems",
913 ));
914 }
915 Ok(())
916}
917
918fn validate_action_descriptor(
919 _plugin: &str,
920 descriptor: &mut PluginActionDescriptor,
921) -> Result<(), CanwuError> {
922 if descriptor.name.trim().is_empty() || descriptor.name != descriptor.name.trim() {
923 return Err(CanwuError::new(
924 ErrorCode::InvalidPluginRegistration,
925 "plugin command names must be non-empty and have no surrounding whitespace",
926 ));
927 }
928 if let PayloadSchema::Object { properties, .. } = &descriptor.payload_schema
929 && properties
930 .keys()
931 .any(|name| name.trim().is_empty() || name != name.trim())
932 {
933 return Err(CanwuError::new(
934 ErrorCode::InvalidPluginRegistration,
935 "plugin payload schema property names cannot be empty",
936 ));
937 }
938 validate_state_keys(&mut descriptor.reads)?;
939 validate_state_keys(&mut descriptor.writes)?;
940 if descriptor.reads.contains(&StateKey::core_ingress()) {
941 return Err(CanwuError::new(
942 ErrorCode::InvalidPluginRegistration,
943 "plugin commands cannot inspect the canonical ingress queue",
944 ));
945 }
946 Ok(())
947}
948
949fn validate_ingress_descriptor(descriptor: &PluginIngressDescriptor) -> Result<(), CanwuError> {
950 if descriptor.name.trim().is_empty()
951 || descriptor.name != descriptor.name.trim()
952 || descriptor.description.trim().is_empty()
953 || descriptor.description != descriptor.description.trim()
954 || descriptor.class == IngressClass::Command
955 {
956 return Err(CanwuError::new(
957 ErrorCode::InvalidPluginRegistration,
958 "plugin ingress types require canonical names/descriptions and cannot claim the core command class",
959 ));
960 }
961 if let PayloadSchema::Object { properties, .. } = &descriptor.payload_schema
962 && properties
963 .keys()
964 .any(|name| name.trim().is_empty() || name != name.trim())
965 {
966 return Err(CanwuError::new(
967 ErrorCode::InvalidPluginRegistration,
968 "plugin ingress payload property names cannot be empty",
969 ));
970 }
971 Ok(())
972}
973
974fn validate_boundary_system_contract(
975 contract: &mut BoundarySystemContract,
976) -> Result<(), CanwuError> {
977 if contract.name.trim().is_empty() || contract.name != contract.name.trim() {
978 return Err(CanwuError::new(
979 ErrorCode::InvalidPluginRegistration,
980 "boundary system name must be non-empty and canonical",
981 ));
982 }
983 validate_state_keys(&mut contract.reads)?;
984 validate_state_keys(&mut contract.writes)?;
985 validate_state_keys(&mut contract.reservation_offers)?;
986 validate_state_keys(&mut contract.reservation_requests)?;
987 validate_reservation_refs(&mut contract.reservation_reads)?;
988 validate_random_stream_keys(&mut contract.random_streams)?;
989 validate_canonical_names(&mut contract.emits, "boundary event type")?;
990
991 let may_propose_changes = matches!(
992 contract.phase,
993 BoundaryPhase::DomainDeltaProposal
994 | BoundaryPhase::HistoricalCandidateEvaluation
995 | BoundaryPhase::StrategicAggregation
996 | BoundaryPhase::PerspectiveAndReportMaterialization
997 );
998 if (!contract.writes.is_empty() || !contract.emits.is_empty()) && !may_propose_changes {
999 return Err(CanwuError::new(
1000 ErrorCode::InvalidPluginRegistration,
1001 format!(
1002 "boundary system {} declares changes in kernel-owned phase {:?}",
1003 contract.name, contract.phase
1004 ),
1005 ));
1006 }
1007 let declares_reservations =
1008 !contract.reservation_offers.is_empty() || !contract.reservation_requests.is_empty();
1009 if declares_reservations && contract.phase != BoundaryPhase::ReservationAndAllocation {
1010 return Err(CanwuError::new(
1011 ErrorCode::InvalidPluginRegistration,
1012 format!(
1013 "boundary system {} declares reservations outside reservation and allocation",
1014 contract.name
1015 ),
1016 ));
1017 }
1018 if !contract.reservation_reads.is_empty()
1019 && contract.phase <= BoundaryPhase::ReservationAndAllocation
1020 {
1021 return Err(CanwuError::new(
1022 ErrorCode::InvalidPluginRegistration,
1023 format!(
1024 "boundary system {} reads allocations before reservation commit",
1025 contract.name
1026 ),
1027 ));
1028 }
1029 Ok(())
1030}
1031
1032fn validate_reservation_refs(values: &mut Vec<ReservationRef>) -> Result<(), CanwuError> {
1033 if values.iter().any(|reservation| {
1034 reservation.plugin.trim().is_empty()
1035 || reservation.plugin != reservation.plugin.trim()
1036 || reservation.system.trim().is_empty()
1037 || reservation.system != reservation.system.trim()
1038 || reservation.request.trim().is_empty()
1039 || reservation.request != reservation.request.trim()
1040 }) {
1041 return Err(CanwuError::new(
1042 ErrorCode::InvalidPluginRegistration,
1043 "reservation read declarations must be non-empty and canonical",
1044 ));
1045 }
1046 let unique: BTreeSet<_> = values.drain(..).collect();
1047 values.extend(unique);
1048 Ok(())
1049}
1050
1051fn validate_random_stream_keys(values: &mut Vec<RandomStreamKey>) -> Result<(), CanwuError> {
1052 if values.iter().any(|stream| {
1053 stream.namespace.trim().is_empty()
1054 || stream.namespace != stream.namespace.trim()
1055 || stream.name.trim().is_empty()
1056 || stream.name != stream.name.trim()
1057 || stream.version == 0
1058 }) {
1059 return Err(CanwuError::new(
1060 ErrorCode::InvalidPluginRegistration,
1061 "random stream declarations require canonical names and a nonzero version",
1062 ));
1063 }
1064 let unique: BTreeSet<_> = values.drain(..).collect();
1065 values.extend(unique);
1066 Ok(())
1067}
1068
1069fn validate_canonical_names(values: &mut Vec<String>, label: &str) -> Result<(), CanwuError> {
1070 if values
1071 .iter()
1072 .any(|value| value.trim().is_empty() || value != value.trim())
1073 {
1074 return Err(CanwuError::new(
1075 ErrorCode::InvalidPluginRegistration,
1076 format!("{label} declarations must be non-empty and canonical"),
1077 ));
1078 }
1079 let unique: BTreeSet<_> = values.drain(..).collect();
1080 values.extend(unique);
1081 Ok(())
1082}
1083
1084fn validate_event_audience_name(event_type: &str) -> Result<(), CanwuError> {
1085 if !canonical_text(event_type) {
1086 return Err(CanwuError::new(
1087 ErrorCode::InvalidPluginRegistration,
1088 "plugin event audience names must be non-empty and canonical",
1089 ));
1090 }
1091 Ok(())
1092}
1093
1094fn validate_event_audience(audience: &EventAudience) -> Result<(), CanwuError> {
1095 match audience {
1096 EventAudience::Actor(actor) if actor.get() == 0 => {
1097 return Err(CanwuError::new(
1098 ErrorCode::InvalidPluginRegistration,
1099 "plugin event audience actors must use positive actor IDs",
1100 ));
1101 }
1102 EventAudience::Actors(actors) => {
1103 if actors.is_empty() || actors.iter().any(|actor| actor.get() == 0) {
1104 return Err(CanwuError::new(
1105 ErrorCode::InvalidPluginRegistration,
1106 "plugin event audience actor lists must contain positive actor IDs",
1107 ));
1108 }
1109 if actors.windows(2).any(|pair| pair[0] >= pair[1]) {
1110 return Err(CanwuError::new(
1111 ErrorCode::InvalidPluginRegistration,
1112 "plugin event audience actor lists must be sorted and unique",
1113 ));
1114 }
1115 }
1116 _ => {}
1117 }
1118 Ok(())
1119}
1120
1121fn register_state_owners(
1122 owners: &mut BTreeMap<StateKey, String>,
1123 plugin: &str,
1124 writes: &[StateKey],
1125) -> Result<(), CanwuError> {
1126 for key in writes {
1127 if key.namespace == CORE_STATE_NAMESPACE {
1128 return Err(CanwuError::new(
1129 ErrorCode::InvalidPluginRegistration,
1130 format!(
1131 "plugin {plugin} cannot claim reserved state {}.{}",
1132 key.namespace, key.name
1133 ),
1134 ));
1135 }
1136 if let Some(existing) = owners.get(key)
1137 && existing != plugin
1138 {
1139 return Err(CanwuError::new(
1140 ErrorCode::DuplicateStateOwner,
1141 format!(
1142 "state {}.{} is owned by both {existing} and {plugin}",
1143 key.namespace, key.name
1144 ),
1145 ));
1146 }
1147 }
1148 for key in writes {
1149 owners.insert(key.clone(), plugin.to_owned());
1150 }
1151 Ok(())
1152}
1153
1154fn register_boundary_writers(
1155 writers: &mut BTreeMap<(BoundaryWriteStage, StateKey), (String, String)>,
1156 immediate_writes: &BTreeMap<StateKey, String>,
1157 plugin: &str,
1158 system: &str,
1159 phase: BoundaryPhase,
1160 declared_states: &[StateKey],
1161) -> Result<(), CanwuError> {
1162 let Some(stage) = boundary_write_stage(phase) else {
1163 if declared_states.is_empty() {
1164 return Ok(());
1165 }
1166 return Err(CanwuError::new(
1167 ErrorCode::InvalidPluginRegistration,
1168 format!("boundary phase {phase:?} cannot own state writes"),
1169 ));
1170 };
1171 for state in declared_states {
1172 if let Some(immediate_plugin) = immediate_writes.get(state) {
1173 return Err(CanwuError::new(
1174 ErrorCode::InvalidPluginRegistration,
1175 format!(
1176 "boundary state {}.{} conflicts with immediate writes from plugin {immediate_plugin}",
1177 state.namespace, state.name
1178 ),
1179 ));
1180 }
1181 if let Some((existing_plugin, existing_system)) = writers.get(&(stage, state.clone()))
1182 && (existing_plugin != plugin || existing_system != system)
1183 {
1184 return Err(CanwuError::new(
1185 ErrorCode::DuplicateBoundaryWriter,
1186 format!(
1187 "boundary state {}.{} is written by both {existing_plugin}.{existing_system} and {plugin}.{system}",
1188 state.namespace, state.name
1189 ),
1190 ));
1191 }
1192 }
1193 for state in declared_states {
1194 writers.insert(
1195 (stage, state.clone()),
1196 (plugin.to_owned(), system.to_owned()),
1197 );
1198 }
1199 Ok(())
1200}
1201
1202fn register_immediate_write_states(
1203 immediate_writes: &mut BTreeMap<StateKey, String>,
1204 boundary_writers: &BTreeMap<(BoundaryWriteStage, StateKey), (String, String)>,
1205 plugin: &str,
1206 writes: &[StateKey],
1207) -> Result<(), CanwuError> {
1208 for state in writes {
1209 if boundary_writers
1210 .keys()
1211 .any(|(_, boundary_state)| boundary_state == state)
1212 {
1213 return Err(CanwuError::new(
1214 ErrorCode::InvalidPluginRegistration,
1215 format!(
1216 "immediate state {}.{} conflicts with a phased boundary writer",
1217 state.namespace, state.name
1218 ),
1219 ));
1220 }
1221 if immediate_writes
1222 .get(state)
1223 .is_some_and(|existing| existing != plugin)
1224 {
1225 return Err(CanwuError::new(
1226 ErrorCode::DuplicateStateOwner,
1227 format!(
1228 "immediate state {}.{} is written by multiple plugins",
1229 state.namespace, state.name
1230 ),
1231 ));
1232 }
1233 }
1234 for state in writes {
1235 immediate_writes.insert(state.clone(), plugin.to_owned());
1236 }
1237 Ok(())
1238}
1239
1240fn register_reservation_offerers(
1241 offerers: &mut BTreeMap<StateKey, (String, String)>,
1242 plugin: &str,
1243 system: &str,
1244 offered_state: &[StateKey],
1245) -> Result<(), CanwuError> {
1246 for state in offered_state {
1247 if let Some((existing_plugin, existing_system)) = offerers.get(state)
1248 && (existing_plugin != plugin || existing_system != system)
1249 {
1250 return Err(CanwuError::new(
1251 ErrorCode::DuplicateReservationOfferer,
1252 format!(
1253 "reservation state {}.{} is offered by both {existing_plugin}.{existing_system} and {plugin}.{system}",
1254 state.namespace, state.name
1255 ),
1256 ));
1257 }
1258 }
1259 for state in offered_state {
1260 offerers.insert(state.clone(), (plugin.to_owned(), system.to_owned()));
1261 }
1262 Ok(())
1263}
1264
1265fn register_random_streams(
1266 owners: &mut BTreeMap<RandomStreamKey, (String, String)>,
1267 plugin: &str,
1268 system: &str,
1269 streams: &[RandomStreamKey],
1270) -> Result<(), CanwuError> {
1271 for stream in streams {
1272 if stream.namespace != plugin || stream.namespace == CORE_STATE_NAMESPACE {
1273 return Err(CanwuError::new(
1274 ErrorCode::InvalidPluginRegistration,
1275 format!(
1276 "random stream {}.{}@{} must use its owning plugin namespace {plugin}",
1277 stream.namespace, stream.name, stream.version
1278 ),
1279 ));
1280 }
1281 if let Some((existing_plugin, existing_system)) = owners.get(stream)
1282 && (existing_plugin != plugin || existing_system != system)
1283 {
1284 return Err(CanwuError::new(
1285 ErrorCode::InvalidPluginRegistration,
1286 format!(
1287 "random stream {}.{}@{} is owned by both {existing_plugin}.{existing_system} and {plugin}.{system}",
1288 stream.namespace, stream.name, stream.version
1289 ),
1290 ));
1291 }
1292 }
1293 for stream in streams {
1294 owners.insert(stream.clone(), (plugin.to_owned(), system.to_owned()));
1295 }
1296 Ok(())
1297}