1use std::any::TypeId;
2use std::collections::HashMap;
3use std::sync::Arc;
4
5use parking_lot::{Mutex, RwLock};
6use serde::{de::DeserializeOwned, Serialize};
7
8use crate::{Context, CordisError, Fiber, FiberId, Service, Symbol};
9
10pub trait Plugin: Send + Sync + 'static {
17 type Config: Serialize + DeserializeOwned + Send + Sync + 'static;
18 type Provides: Service;
19 fn apply(
20 &self,
21 ctx: &Arc<Context>,
22 config: Self::Config,
23 ) -> Result<Arc<Self::Provides>, CordisError>;
24}
25
26type ReadinessPredicate = Arc<dyn Fn(&Arc<Context>) -> bool + Send + Sync>;
45
46#[derive(Clone)]
47pub struct ReadinessBarrier {
48 inner: ReadinessPredicate,
49 keys: Vec<TypeId>,
51}
52
53impl ReadinessBarrier {
54 pub fn new(ready: impl Fn(&Arc<Context>) -> bool + Send + Sync + 'static) -> Self {
56 Self {
57 inner: Arc::new(ready),
58 keys: Vec::new(),
59 }
60 }
61
62 pub fn watched_type_ids(&self) -> &[TypeId] {
64 &self.keys
65 }
66
67 pub fn is_ready(&self, ctx: &Arc<Context>) -> bool {
69 (self.inner)(ctx)
70 }
71
72 pub fn and(self, other: ReadinessBarrier) -> ReadinessBarrier {
76 let pair = (self.inner, other.inner);
77 ReadinessBarrier::new(move |ctx| (pair.0)(ctx) && (pair.1)(ctx))
78 .watching(self.keys.iter().chain(other.keys.iter()).copied())
79 }
80
81 pub fn watching(
86 mut self,
87 keys: impl IntoIterator<Item = TypeId>,
88 ) -> ReadinessBarrier {
89 self.keys.extend(keys);
90 self
91 }
92}
93
94pub fn with_readiness(
99 barriers: impl IntoIterator<Item = ReadinessBarrier>,
100) -> ReadinessBarrier {
101 let mut combined: Option<ReadinessBarrier> = None;
102 for barrier in barriers {
103 combined = Some(match combined {
104 None => barrier,
105 Some(acc) => acc.and(barrier),
106 });
107 }
108 combined.unwrap_or_else(|| ReadinessBarrier::new(|_ctx| true))
109}
110
111impl std::fmt::Debug for ReadinessBarrier {
112 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
113 f.debug_struct("ReadinessBarrier").finish_non_exhaustive()
114 }
115}
116
117pub struct RegistryService {
128 fibers: RwLock<HashMap<FiberId, Arc<Fiber>>>,
129 provided: RwLock<HashMap<(TypeId, Option<Symbol>), FiberId>>,
130 realms: RwLock<HashMap<FiberId, std::sync::Weak<Context>>>,
133 next_id: Mutex<FiberId>,
134}
135
136impl RegistryService {
137 pub fn new() -> Self {
138 Self {
139 fibers: RwLock::new(HashMap::new()),
140 provided: RwLock::new(HashMap::new()),
141 realms: RwLock::new(HashMap::new()),
142 next_id: Mutex::new(1),
143 }
144 }
145
146 pub fn reliance_count(&self, key: &(TypeId, Option<Symbol>)) -> usize {
157 let fibers = self.fibers.read();
158 let mut consumers = 0;
159 for (fid, fiber) in fibers.iter() {
160 if !matches!(fiber.state(), crate::FiberState::Active { .. }) {
161 continue;
162 }
163 if Some(*fid) == self.provided.read().get(key).copied() {
165 continue;
166 }
167 let Some(reg_ctx) = self
173 .realms
174 .read()
175 .get(fid)
176 .and_then(std::sync::Weak::upgrade)
177 else {
178 continue;
179 };
180 if reg_ctx.isolate_label(key.0) != key.1 {
181 continue;
182 }
183 for inject_tid in fiber.injected_type_ids() {
184 if inject_tid == key.0 {
185 consumers += 1;
186 break;
187 }
188 }
189 }
190 consumers
191 }
192
193 fn next_fiber_id(&self) -> FiberId {
194 let mut guard = self.next_id.lock();
195 let id = *guard;
196 *guard += 1;
197 id
198 }
199
200 pub fn register<P: Plugin>(
216 &self,
217 ctx: &Arc<Context>,
218 plugin: P,
219 config: P::Config,
220 ) -> Result<FiberId, CordisError> {
221 let tid = TypeId::of::<P::Provides>();
222 let isolate = ctx.isolate_label(tid);
223 let key = (tid, isolate);
224 let active_conflict = match self.provided.read().get(&key).copied() {
228 Some(existing) => self
229 .fibers
230 .read()
231 .get(&existing)
232 .map(|fiber| {
233 matches!(
234 fiber.state(),
235 crate::FiberState::Active { .. }
236 | crate::FiberState::Loading
237 | crate::FiberState::Reloading
238 )
239 })
240 .unwrap_or(false),
241 None => false,
242 };
243 if active_conflict {
244 return Err(CordisError::DuplicateProvider {
245 name: format!("{tid:?}"),
246 owner: format!("isolate realm {:?}", ctx.isolate_label(tid)),
247 });
248 }
249 self.provided.write().remove(&key);
252
253 let fid = self.next_fiber_id();
254 let fiber = Arc::new(Fiber::new());
255 fiber.set_state(crate::FiberState::Loading);
256 fiber.set_reload_context(ctx);
257 fiber.set_id(fid);
258 self.fibers.write().insert(fid, fiber.clone());
259 self.realms.write().insert(fid, Arc::downgrade(ctx));
260
261 let config_value = serde_json::to_value(&config).map_err(|error| {
262 let message = format!("cannot serialize plugin config: {error}");
263 fiber.set_state(crate::FiberState::Failed {
264 error: Some(message.clone()),
265 });
266 self.wire_failed_registration(ctx, fid, &fiber, tid);
267 CordisError::Configuration(message)
268 })?;
269 fiber.set_raw_config(config_value.clone());
272 if let Some(events) = ctx.get::<crate::EventsService>() {
277 match crate::events::blocking_intercept_config(&events, config_value.clone()) {
278 Ok(effective) => {
279 if !effective.is_null() {
280 fiber.stage_effective_config(effective);
281 }
282 }
283 Err(error) => {
284 fiber.set_state(crate::FiberState::Failed {
285 error: Some(error.to_string()),
286 });
287 self.wire_failed_registration(ctx, fid, &fiber, tid);
288 return Err(error);
289 }
290 }
291 }
292 let plugin = Arc::new(plugin);
293 let weak_fiber = Arc::downgrade(&fiber);
294 fiber.set_reload_runner(Box::new(move |ctx| {
295 let cfg_raw = weak_fiber
301 .upgrade()
302 .and_then(|owner| owner.effective_config_override())
303 .unwrap_or_else(|| config_value.clone());
304 let config =
305 serde_json::from_value::<P::Config>(cfg_raw).map_err(|error| {
306 CordisError::Configuration(format!("cannot deserialize plugin config: {error}"))
307 })?;
308 let owner = weak_fiber
309 .upgrade()
310 .ok_or_else(|| CordisError::Fiber("registration fiber was dropped".into()))?;
311 let applied = crate::hmr::catch_plugin_panic(std::panic::AssertUnwindSafe(|| {
319 ctx.with_provider_fiber(&owner, || plugin.apply(ctx, config))
320 }));
321 let applied = match applied {
322 Ok(applied) => applied,
323 Err(payload) => {
324 return Err(CordisError::Fiber(format!(
325 "plugin factory panicked: {payload}"
326 )))
327 }
328 };
329 let provides = applied?;
330 let healthy = provides.check();
331 ctx.provide_on_fiber(provides, &owner);
332 Ok(healthy)
333 }));
334
335 let healthy = match fiber.run_runner(ctx) {
336 Ok(healthy) => healthy,
337 Err(error) => {
338 fiber.set_state(crate::FiberState::Failed {
339 error: Some(error.to_string()),
340 });
341 self.wire_failed_registration(ctx, fid, &fiber, tid);
347 return Err(error);
348 }
349 };
350 let epoch = fiber.compute_epoch(ctx);
358 fiber.set_epoch(epoch.clone());
359 fiber.set_state(if healthy {
366 crate::FiberState::Active { epoch }
367 } else {
368 crate::FiberState::Failed {
369 error: Some("availability predicate rejected service".into()),
370 }
371 });
372 self.provided.write().insert(key, fid);
373
374 if let Some(reflect) = ctx.get::<crate::ReflectService>() {
375 for dependency in fiber.injected_type_ids() {
376 reflect.register_dependent(dependency, fid);
377 }
378 let _ = reflect.ensure_notifier(tid);
379 reflect.register_fiber(fid, fiber.clone(), tid);
380 }
381 Ok(fid)
382 }
383
384 pub fn register_with_readiness<P: Plugin>(
396 &self,
397 ctx: &Arc<Context>,
398 plugin: P,
399 config: P::Config,
400 ready_when: ReadinessBarrier,
401 ) -> Result<FiberId, CordisError> {
402 let fid = self.register(ctx, plugin, config)?;
403 if let Some(fiber) = self.get_fiber(fid) {
404 for tid in ready_when.watched_type_ids().iter().copied() {
410 if let Some(reflect) = ctx.get::<crate::ReflectService>() {
411 reflect.register_dependent(tid, fid);
412 let _ = reflect.ensure_notifier(tid);
413 }
414 }
415 fiber.set_readiness_gate(ready_when);
416 tokio::spawn({
421 let fiber = fiber.clone();
422 let ctx = ctx.clone();
423 async move { fiber.refresh(&ctx).await }
424 });
425 }
426 Ok(fid)
427 }
428
429 fn wire_failed_registration(
439 &self,
440 ctx: &Arc<Context>,
441 fid: FiberId,
442 fiber: &Arc<Fiber>,
443 tid: TypeId,
444 ) {
445 if let Some(reflect) = ctx.get::<crate::ReflectService>() {
446 for dependency in fiber.injected_type_ids() {
447 reflect.register_dependent(dependency, fid);
448 }
449 let _ = reflect.ensure_notifier(tid);
450 reflect.register_fiber(fid, fiber.clone(), tid);
451 reflect.notify(tid);
452 }
453 }
454
455 pub fn plugin<P: Plugin>(
459 &self,
460 ctx: &Arc<Context>,
461 plugin: P,
462 config: P::Config,
463 ) -> Result<FiberId, CordisError> {
464 self.register(ctx, plugin, config)
465 }
466
467 pub fn track_fiber(&self, fiber: Arc<Fiber>) -> FiberId {
473 let fid = self.next_fiber_id();
474 self.fibers.write().insert(fid, fiber);
475 fid
476 }
477
478 pub fn prune_disposed(&self) -> usize {
487 let disposed: Vec<FiberId> = self
488 .fibers
489 .read()
490 .iter()
491 .filter(|(_, fiber)| fiber.is_disposed())
492 .map(|(fid, _)| *fid)
493 .collect();
494 let mut removed = 0;
495 for fid in disposed {
496 if self.remove(fid).is_some() {
497 removed += 1;
498 }
499 }
500 removed
501 }
502
503 pub fn get_fiber(&self, id: FiberId) -> Option<Arc<Fiber>> {
504 self.fibers.read().get(&id).cloned()
509 }
510
511 pub fn remove(&self, id: FiberId) -> Option<Arc<Fiber>> {
512 let fiber = self.fibers.write().remove(&id)?;
513 let mut provided = self.provided.write();
514 provided.retain(|_, v| *v != id);
515 drop(provided);
516 self.realms.write().remove(&id);
517 Some(fiber)
518 }
519
520 pub fn track_fiber_in_realm(&self, fid: FiberId, ctx: &Arc<Context>) {
524 self.realms.write().insert(fid, Arc::downgrade(ctx));
525 }
526
527 pub fn provider_fibers_for(&self, ctx: &Arc<Context>, tids: &[TypeId]) -> Vec<u64> {
530 let provided = self.provided.read();
531 tids.iter()
532 .filter_map(|tid| {
533 let isolate = ctx.isolate_label(*tid);
534 provided.get(&(*tid, isolate)).copied()
535 })
536 .collect()
537 }
538
539 pub fn provided_types_of_fiber(&self, fid: FiberId) -> Vec<TypeId> {
542 self.provided
543 .read()
544 .iter()
545 .filter(|(_, owner)| **owner == fid)
546 .map(|(key, _)| key.0)
547 .collect()
548 }
549
550 pub fn tracked_ids(&self) -> Vec<FiberId> {
552 self.fibers.read().keys().copied().collect()
553 }
554
555 pub fn len(&self) -> usize {
556 self.prune_disposed();
559 self.fibers.read().len()
560 }
561
562 pub fn is_empty(&self) -> bool {
563 self.fibers.read().is_empty()
564 }
565}
566
567impl Default for RegistryService {
568 fn default() -> Self {
569 Self::new()
570 }
571}
572
573impl Service for RegistryService {}
574
575#[cfg(feature = "linkme")]
582#[linkme::distributed_slice]
583pub static REGISTRY_PLUGINS: [fn(&Arc<Context>) -> Result<FiberId, CordisError>];
584
585#[cfg(test)]
586mod tests {
587 use super::*;
588 use crate::{Context, FiberState, ReflectService, Service};
589
590 #[derive(Debug)]
591 struct FooService(pub i32);
592 impl Service for FooService {}
593
594 #[derive(Debug)]
595 struct BarService(pub i32);
596 impl Service for BarService {}
597
598 struct FooPlugin;
599 impl Plugin for FooPlugin {
600 type Config = ();
601 type Provides = FooService;
602 fn apply(
603 &self,
604 _ctx: &Arc<Context>,
605 _cfg: Self::Config,
606 ) -> Result<Arc<Self::Provides>, CordisError> {
607 Ok(Arc::new(FooService(1)))
608 }
609 }
610
611 struct FooPlugin2;
612 impl Plugin for FooPlugin2 {
613 type Config = ();
614 type Provides = FooService;
615 fn apply(
616 &self,
617 _ctx: &Arc<Context>,
618 _cfg: Self::Config,
619 ) -> Result<Arc<Self::Provides>, CordisError> {
620 Ok(Arc::new(FooService(2)))
621 }
622 }
623
624 struct BarPlugin;
625 impl Plugin for BarPlugin {
626 type Config = ();
627 type Provides = BarService;
628 fn apply(
629 &self,
630 _ctx: &Arc<Context>,
631 _cfg: Self::Config,
632 ) -> Result<Arc<Self::Provides>, CordisError> {
633 Ok(Arc::new(BarService(99)))
634 }
635 }
636
637 #[test]
638 fn duplicate_provider_rejected() {
639 let ctx = Context::new_root();
640 let registry = RegistryService::new();
641 let fid1 = registry
642 .register(&ctx, FooPlugin, ())
643 .expect("first registration should succeed");
644 assert!(registry.get_fiber(fid1).is_some());
645 let err = registry
647 .register(&ctx, FooPlugin2, ())
648 .expect_err("duplicate provider should be rejected");
649 assert!(
650 err.to_string().contains("duplicate provider for"),
651 "error should mention duplicate provider, got {err}"
652 );
653 assert!(registry.get_fiber(fid1).is_some());
655 let svc = ctx.get::<FooService>().expect("service should be present");
657 assert_eq!(svc.0, 1);
658 }
659
660 struct FailingPlugin;
661 impl Plugin for FailingPlugin {
662 type Config = ();
663 type Provides = FooService;
664 fn apply(
665 &self,
666 _ctx: &Arc<Context>,
667 _cfg: Self::Config,
668 ) -> Result<Arc<Self::Provides>, CordisError> {
669 Err(CordisError::Configuration("intentional failure".into()))
670 }
671 }
672
673 #[test]
674 fn failed_plugin_transitions_tracked_fiber_to_failed_with_error() {
675 use crate::FiberState;
676 let ctx = Context::new_root();
677 let registry = RegistryService::new();
678 let err = registry
681 .register(&ctx, FailingPlugin, ())
682 .expect_err("failing plugin should be rejected");
683 assert!(err.to_string().contains("intentional failure"));
684 let existing = registry
685 .get_fiber(1)
686 .expect("failed fiber should be tracked");
687 match existing.state() {
688 FiberState::Failed { error } => {
689 assert!(error
690 .as_deref()
691 .unwrap_or("")
692 .contains("intentional failure"));
693 }
694 other => panic!("expected Failed state, got {other:?}"),
695 }
696 assert!(ctx.get::<FooService>().is_none());
698 }
699
700 struct PanickingPlugin;
701 impl Plugin for PanickingPlugin {
702 type Config = ();
703 type Provides = FooService;
704 fn apply(
705 &self,
706 _ctx: &Arc<Context>,
707 _cfg: Self::Config,
708 ) -> Result<Arc<Self::Provides>, CordisError> {
709 panic!("factory exploded");
710 }
711 }
712
713 #[tokio::test]
718 async fn panicking_factory_registers_failed_and_notifies_dependents() {
719 let ctx = Context::new_root();
720 ctx.provide(ReflectService::new());
721 if let Some(reflect) = ctx.get::<ReflectService>() {
722 reflect.set_context(&ctx);
723 }
724 let registry = RegistryService::new();
725
726 let err = registry
727 .register(&ctx, PanickingPlugin, ())
728 .expect_err("panicking factory should be rejected");
729 match &err {
730 CordisError::Fiber(message) => {
731 assert!(
732 message.contains("plugin factory panicked"),
733 "unexpected error text: {message}"
734 );
735 assert!(message.contains("factory exploded"));
736 }
737 other => panic!("expected Fiber error, got {other:?}"),
738 }
739 let fiber = registry.get_fiber(1).expect("failed fiber is tracked");
741 match fiber.state() {
742 FiberState::Failed { error } => {
743 let error = error.as_deref().unwrap_or("");
744 assert!(error.contains("factory panicked"), "got: {error}");
745 assert!(error.contains("factory exploded"));
746 }
747 other => panic!("expected Failed state, got {other:?}"),
748 }
749 assert!(ctx.get::<FooService>().is_none());
751 let fid_retry = registry
752 .register(&ctx, FooPlugin, ())
753 .expect("fresh registration of the same key after a panic");
754 assert!(matches!(
755 registry.get_fiber(fid_retry).unwrap().state(),
756 FiberState::Active { .. }
757 ));
758
759 let dep_fid = registry
763 .register(&ctx, BarPlugin, ())
764 .expect("dependent registers without its dependency");
765 let dependent = registry.get_fiber(dep_fid).unwrap();
766 dependent.declare_inject::<FooService>();
767 dependent.refresh(&ctx).await;
768 assert!(matches!(dependent.state(), FiberState::Active { .. }));
769 }
770
771 #[test]
772 fn re_registered_good_plugin_moves_fiber_to_active() {
773 use crate::FiberState;
774 let ctx = Context::new_root();
775 let registry = RegistryService::new();
776 let _ = registry
777 .register(&ctx, FailingPlugin, ())
778 .expect_err("failing plugin should be rejected");
779 let original = registry.get_fiber(1).unwrap();
780 assert!(matches!(original.state(), FiberState::Failed { .. }));
781
782 let fid_ok = registry
784 .register(&ctx, BarPlugin, ())
785 .expect("good plugin should register");
786 let ok = registry.get_fiber(fid_ok).expect("good fiber");
787 match ok.state() {
788 FiberState::Active { .. } => {}
789 other => panic!("expected Active state, got {other:?}"),
790 }
791 }
792
793 #[test]
794 fn different_isolates_allowed() {
795 let root = Context::new_root();
796 let registry = RegistryService::new();
797 let fid_root = registry
799 .register(&root, FooPlugin, ())
800 .expect("root registration ok");
801 assert!(registry.get_fiber(fid_root).is_some());
802 let tenant_a = root.isolate::<FooService>("tenant:acme");
804 let fid_a = registry
805 .register(&tenant_a, FooPlugin2, ())
806 .expect("different isolate should be allowed");
807 assert!(registry.get_fiber(fid_a).is_some());
808 struct FooPlugin3;
810 impl Plugin for FooPlugin3 {
811 type Config = ();
812 type Provides = FooService;
813 fn apply(
814 &self,
815 _ctx: &Arc<Context>,
816 _cfg: Self::Config,
817 ) -> Result<Arc<Self::Provides>, CordisError> {
818 Ok(Arc::new(FooService(3)))
819 }
820 }
821 let err = registry
822 .register(&tenant_a, FooPlugin3, ())
823 .expect_err("duplicate in same isolate should fail");
824 assert!(err.to_string().contains("duplicate provider for"));
825 let tenant_b = root.isolate::<FooService>("tenant:other");
827 let fid_b = registry
828 .register(&tenant_b, FooPlugin3, ())
829 .expect("different isolate label should be allowed");
830 assert!(registry.get_fiber(fid_b).is_some());
831 }
832
833 fn assert_plugin_retrievable<T, P>(
837 registry: &RegistryService,
838 ctx: &Arc<Context>,
839 plugin: P,
840 expect: impl FnOnce(&T),
841 ) where
842 T: Service + std::fmt::Debug,
843 P: Plugin<Provides = T, Config = ()>,
844 {
845 let fid = registry
846 .plugin(ctx, plugin, ())
847 .expect("plugin alias should work");
848 assert!(registry.get_fiber(fid).is_some());
849 let svc = ctx
850 .get::<T>()
851 .expect("service should be retrievable via ctx.get");
852 expect(&svc);
853 }
854
855 #[test]
856 fn successful_provide_retrievable_via_ctx_get() {
857 let ctx = Context::new_root();
858 let registry = RegistryService::new();
859 assert_plugin_retrievable(®istry, &ctx, BarPlugin, |svc: &BarService| {
860 assert_eq!(svc.0, 99);
861 });
862 assert_eq!(registry.len(), 1);
864 }
865
866 #[tokio::test]
867 async fn registration_fiber_disposal_removes_only_its_service() {
868 let ctx = Context::new_root();
869 let registry = RegistryService::new();
870 let foo = registry
871 .register(&ctx, FooPlugin, ())
872 .expect("foo registration");
873 let bar = registry
874 .register(&ctx, BarPlugin, ())
875 .expect("bar registration");
876 assert!(ctx.get::<FooService>().is_some());
877 assert!(ctx.get::<BarService>().is_some());
878 let _ = registry.get_fiber(foo).unwrap().dispose().await;
879 assert!(ctx.get::<FooService>().is_none());
880 assert_eq!(ctx.get::<BarService>().unwrap().0, 99);
881 assert!(matches!(
882 registry.get_fiber(bar).unwrap().state(),
883 FiberState::Active { .. }
884 ));
885 registry.get_fiber(foo).unwrap().refresh(&ctx).await;
886 assert!(ctx.get::<FooService>().is_none());
887 }
888
889 struct Dependency;
890 impl Service for Dependency {}
891
892 struct CountingPlugin {
893 calls: std::sync::Arc<std::sync::atomic::AtomicUsize>,
894 }
895
896 impl Plugin for CountingPlugin {
897 type Config = ();
898 type Provides = FooService;
899
900 fn apply(
901 &self,
902 _ctx: &Arc<Context>,
903 _cfg: Self::Config,
904 ) -> Result<Arc<Self::Provides>, CordisError> {
905 self.calls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
906 Ok(Arc::new(FooService(7)))
907 }
908 }
909
910 #[tokio::test]
911 async fn refresh_reruns_provider_after_dependency_version_change() {
912 let ctx = Context::new_root();
913 let registry = RegistryService::new();
914 let calls = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
915 let fid = registry
916 .register(
917 &ctx,
918 CountingPlugin {
919 calls: calls.clone(),
920 },
921 (),
922 )
923 .expect("registration");
924 let fiber = registry.get_fiber(fid).unwrap();
925 fiber.declare_inject::<Dependency>();
926 ctx.provide(Dependency);
927 fiber.refresh(&ctx).await;
928 assert_eq!(calls.load(std::sync::atomic::Ordering::SeqCst), 2);
929 assert!(matches!(fiber.state(), FiberState::Active { .. }));
930 }
931
932 #[tokio::test]
933 async fn missing_dependency_deactivates_then_provide_reactivates() {
934 let ctx = Context::new_root();
935 let registry = RegistryService::new();
936 let fid = registry
937 .register(&ctx, FooPlugin, ())
938 .expect("registration");
939 let fiber = registry.get_fiber(fid).unwrap();
940 fiber.declare_inject::<Dependency>();
941 fiber.refresh(&ctx).await;
942 assert!(matches!(fiber.state(), FiberState::Inactive { .. }));
943 assert!(ctx.get::<FooService>().is_none());
944 ctx.provide(Dependency);
945 fiber.refresh(&ctx).await;
946 assert!(matches!(fiber.state(), FiberState::Active { .. }));
947 assert!(ctx.get::<FooService>().is_some());
948 }
949
950 #[test]
951 fn isolate_lookup_does_not_cross_realms() {
952 let root = Context::new_root();
953 root.provide(FooService(1));
954 let tenant = root.isolate::<FooService>("tenant:a");
955 assert!(tenant.get::<FooService>().is_none());
956 tenant.provide(FooService(2));
957 assert_eq!(tenant.get::<FooService>().unwrap().0, 2);
958 let other = root.isolate::<FooService>("tenant:b");
959 assert!(other.get::<FooService>().is_none());
960 }
961
962 struct DependentPlugin;
965 impl Plugin for DependentPlugin {
966 type Config = ();
967 type Provides = BarService;
968 fn apply(
969 &self,
970 _ctx: &Arc<Context>,
971 _cfg: Self::Config,
972 ) -> Result<Arc<Self::Provides>, CordisError> {
973 Ok(Arc::new(BarService(5)))
974 }
975 }
976
977 fn register_foo_consumer(ctx: &Arc<Context>, registry: &RegistryService) -> FiberId {
980 let fid = registry
981 .register(&ctx.clone(), DependentPlugin, ())
982 .expect("dependent registration");
983 registry
984 .get_fiber(fid)
985 .unwrap()
986 .declare_inject::<FooService>();
987 fid
988 }
989
990 #[tokio::test]
991 async fn guarded_withdrawal_blocks_removal_with_active_consumer() {
992 use crate::Context;
993 let ctx = Context::new_root();
994 let registry = RegistryService::new();
995 ctx.provide(registry);
996 let registry = ctx.get::<RegistryService>().unwrap();
997
998 let provider_fid = registry.register(&ctx, FooPlugin, ()).expect("provider");
999 let consumer_fid = register_foo_consumer(&ctx, ®istry);
1000 assert!(matches!(
1001 registry.get_fiber(consumer_fid).unwrap().state(),
1002 FiberState::Active { .. }
1003 ));
1004
1005 let key = (TypeId::of::<FooService>(), None);
1006 assert_eq!(registry.reliance_count(&key), 1);
1007
1008 let err = ctx.remove::<FooService>().expect_err("guard must block");
1010 assert!(
1011 err.to_string().contains("guarded withdrawal"),
1012 "error should mention guarded withdrawal, got {err}"
1013 );
1014 assert!(ctx.get::<FooService>().is_some(), "service stays provided");
1015 assert!(matches!(
1016 registry.get_fiber(provider_fid).unwrap().state(),
1017 FiberState::Active { .. }
1018 ));
1019 }
1020
1021 #[tokio::test]
1022 async fn guarded_withdrawal_allows_after_consumer_gone() {
1023 use crate::Context;
1024 let ctx = Context::new_root();
1025 let registry = RegistryService::new();
1026 ctx.provide(registry);
1027 let registry = ctx.get::<RegistryService>().unwrap();
1028
1029 registry.register(&ctx, FooPlugin, ()).expect("provider");
1030 let consumer_fid = register_foo_consumer(&ctx, ®istry);
1031
1032 let _ = registry.get_fiber(consumer_fid).unwrap().dispose().await;
1035 let key = (TypeId::of::<FooService>(), None);
1036 assert_eq!(registry.reliance_count(&key), 0);
1037
1038 let removed = ctx.remove::<FooService>().expect("removal now allowed");
1039 assert_eq!(removed.unwrap().0, 1);
1040 assert!(ctx.get::<FooService>().is_none());
1041 }
1042
1043 #[tokio::test]
1044 async fn internal_undo_bypasses_guard() {
1045 use crate::Context;
1046 let ctx = Context::new_root();
1047 let registry = RegistryService::new();
1048 ctx.provide(registry);
1049 let registry = ctx.get::<RegistryService>().unwrap();
1050
1051 registry.register(&ctx, FooPlugin, ()).expect("provider");
1052 register_foo_consumer(&ctx, ®istry);
1053 let key = (TypeId::of::<FooService>(), None);
1054 assert_eq!(registry.reliance_count(&key), 1);
1055
1056 let removed = ctx.remove_forced::<FooService>().expect("forced removal");
1059 assert!(removed.is_some());
1060 assert!(ctx.get::<FooService>().is_none());
1061 }
1062
1063 #[tokio::test]
1064 async fn remove_without_registry_or_consumers_still_works() {
1065 use crate::Context;
1066 let ctx = Context::new_root();
1068 ctx.provide(FooService(2));
1069 let removed = ctx.remove::<FooService>().expect("unguarded removal");
1070 assert!(removed.is_some());
1071 }
1072
1073 #[test]
1074 fn plugin_alias_behaves_like_register() {
1075 let ctx = Context::new_root();
1079 let registry = RegistryService::new();
1080 let fid = registry
1081 .plugin(&ctx, FooPlugin, ())
1082 .expect("plugin alias ok");
1083 assert!(registry.get_fiber(fid).is_some());
1084 let svc = ctx
1085 .get::<FooService>()
1086 .expect("FooService should be present");
1087 assert_eq!(svc.0, 1);
1088 assert!(ctx.get::<BarService>().is_none());
1090 assert_eq!(registry.len(), 1);
1091 }
1092
1093 #[tokio::test]
1097 async fn prune_disposed_drops_disposed_but_keeps_failed() {
1098 let ctx = Context::new_root();
1099 let registry = RegistryService::new();
1100
1101 let failed_fid = registry
1102 .register(&ctx, FailingPlugin, ())
1103 .expect_err("failing registration is rejected but tracked");
1104 let failed_fid = match failed_fid {
1105 CordisError::Configuration(_) => 1, other => panic!("unexpected error shape: {other:?}"),
1107 };
1108 let live_fid = registry
1109 .register(&ctx, FooPlugin, ())
1110 .expect("live registration");
1111 assert!(matches!(
1112 registry.get_fiber(failed_fid).unwrap().state(),
1113 FiberState::Failed { .. }
1114 ));
1115
1116 let _ = registry.get_fiber(live_fid).unwrap().dispose().await;
1119 assert!(registry.get_fiber(live_fid).unwrap().is_disposed());
1120
1121 let pruned = registry.prune_disposed();
1122 assert_eq!(pruned, 1, "exactly the disposed fiber is pruned");
1123 assert!(
1124 registry.get_fiber(live_fid).is_none(),
1125 "disposed fiber must be gone after prune"
1126 );
1127 assert!(matches!(
1129 registry.get_fiber(failed_fid).unwrap().state(),
1130 FiberState::Failed { .. }
1131 ));
1132 let bar_fid = registry
1134 .register(&ctx, BarPlugin, ())
1135 .expect("bar registration");
1136 assert!(registry.get_fiber(bar_fid).is_some());
1137
1138 let _ = registry.get_fiber(bar_fid).unwrap().dispose().await;
1141 assert_eq!(registry.len(), 1, "only the Failed fiber remains");
1142 assert!(registry.get_fiber(bar_fid).is_none());
1143 }
1144
1145 #[tokio::test]
1149 async fn pending_fiber_survives_prune_disposed() {
1150 let ctx = Context::new_root();
1151 ctx.provide(ReflectService::new());
1152 let reflect = ctx.get::<ReflectService>().unwrap();
1153 reflect.set_context(&ctx);
1154 let registry = RegistryService::new();
1155 ctx.provide(registry);
1156 let registry = ctx.get::<RegistryService>().unwrap();
1157
1158 let provider_fid = registry.register(&ctx, FooPlugin, ()).expect("provider");
1160 let fid = registry
1161 .register(&ctx, DependentPlugin, ())
1162 .expect("consumer registration");
1163 let fiber = registry.get_fiber(fid).unwrap();
1164 fiber.declare_inject::<FooService>();
1165 fiber.refresh(&ctx).await;
1166 assert!(matches!(fiber.state(), FiberState::Active { .. }));
1167
1168 let _ = registry.get_fiber(provider_fid).unwrap().dispose().await;
1172 reflect.notify_with_ctx(TypeId::of::<FooService>(), &ctx).await;
1173 assert!(
1174 matches!(fiber.state(), FiberState::Pending),
1175 "consumer must rest Pending after dep loss, got {:?}",
1176 fiber.state()
1177 );
1178
1179 let pruned = registry.prune_disposed();
1182 assert_eq!(pruned, 1, "exactly the disposed provider is pruned");
1183 assert!(
1184 matches!(fiber.state(), FiberState::Pending),
1185 "Pending fiber must survive prune"
1186 );
1187 assert!(registry.get_fiber(fid).is_some(), "still tracked");
1188
1189 registry.register(&ctx, FooPlugin, ()).expect("provider back");
1191 reflect.notify_with_ctx(TypeId::of::<FooService>(), &ctx).await;
1192 assert!(
1193 matches!(fiber.state(), FiberState::Active { .. }),
1194 "reactivation after prune-pass, got {:?}",
1195 fiber.state()
1196 );
1197 assert_eq!(
1198 registry.len(),
1199 2,
1200 "both live fibers remain tracked through the cycle"
1201 );
1202 }
1203
1204 #[tokio::test]
1207 async fn prune_clears_provided_slot_for_fresh_registration() {
1208 let ctx = Context::new_root();
1209 let registry = RegistryService::new();
1210 let fid = registry
1211 .register(&ctx, FooPlugin, ())
1212 .expect("first registration");
1213 let _ = registry.get_fiber(fid).unwrap().dispose().await;
1214
1215 let fid2 = registry
1218 .register(&ctx, FooPlugin2, ())
1219 .expect("re-registration works even pre-prune");
1220 let _ = registry.get_fiber(fid).unwrap().dispose().await;
1223 let _ = registry.get_fiber(fid2).unwrap().dispose().await;
1224 assert_eq!(registry.prune_disposed(), 2);
1225 assert!(registry.get_fiber(fid).is_none());
1226 assert!(registry.get_fiber(fid2).is_none());
1227 assert_eq!(registry.len(), 0);
1228 }
1229
1230 #[derive(Debug)]
1235 struct GatedService(bool);
1236 impl Service for GatedService {
1237 fn check(&self) -> bool {
1238 self.0
1239 }
1240 }
1241
1242 struct GatedPlugin {
1243 ready: bool,
1244 }
1245
1246 impl Plugin for GatedPlugin {
1247 type Config = ();
1248 type Provides = GatedService;
1249
1250 fn apply(
1251 &self,
1252 _ctx: &Arc<Context>,
1253 _config: Self::Config,
1254 ) -> Result<Arc<Self::Provides>, CordisError> {
1255 Ok(Arc::new(GatedService(self.ready)))
1256 }
1257 }
1258
1259 #[tokio::test]
1266 async fn availability_predicate_rejection_registers_failed() {
1267 let ctx = Context::new_root();
1268 ctx.provide(ReflectService::new());
1269 if let Some(reflect) = ctx.get::<crate::ReflectService>() {
1270 reflect.set_context(&ctx);
1271 }
1272 let registry = RegistryService::new();
1273
1274 let fid = registry
1275 .register(&ctx, GatedPlugin { ready: false }, ())
1276 .expect("predicate rejection must NOT fail registration");
1277 let fiber = registry.get_fiber(fid).expect("tracked");
1278 match fiber.state() {
1279 crate::FiberState::Failed { error } => {
1280 assert!(error
1281 .as_deref()
1282 .unwrap_or("")
1283 .contains("availability predicate rejected service"));
1284 }
1285 other => panic!("expected Failed state, got {other:?}"),
1286 }
1287 assert!(ctx.get::<GatedService>().is_none());
1289 }
1292
1293 #[tokio::test]
1297 async fn predicate_passing_reregistration_activates_dependents() {
1298 let ctx = Context::new_root();
1299 ctx.provide(ReflectService::new());
1300 if let Some(reflect) = ctx.get::<crate::ReflectService>() {
1301 reflect.set_context(&ctx);
1302 }
1303 let registry = RegistryService::new();
1304
1305 let bad_fid = registry
1308 .register(&ctx, GatedPlugin { ready: false }, ())
1309 .expect("rejection is non-throwing");
1310 assert!(matches!(
1311 registry.get_fiber(bad_fid).unwrap().state(),
1312 crate::FiberState::Failed { .. }
1313 ));
1314
1315 struct GatedConsumer;
1317 impl Plugin for GatedConsumer {
1318 type Config = ();
1319 type Provides = DerivedProbe;
1320
1321 fn apply(
1322 &self,
1323 _ctx: &Arc<Context>,
1324 _config: Self::Config,
1325 ) -> Result<Arc<Self::Provides>, CordisError> {
1326 Ok(Arc::new(DerivedProbe))
1327 }
1328 }
1329
1330 #[derive(Debug)]
1331 struct DerivedProbe;
1332 impl Service for DerivedProbe {}
1333
1334 let dep_fid = registry
1335 .register(&ctx, GatedConsumer, ())
1336 .expect("consumer registers even without its dependency");
1337 let dependent = registry.get_fiber(dep_fid).unwrap();
1338 dependent.declare_inject::<GatedService>();
1339 dependent.refresh(&ctx).await;
1340 assert!(
1341 matches!(dependent.state(), crate::FiberState::Inactive { .. }),
1342 "dependent must stay Inactive while the provider is rejected, got {:?}",
1343 dependent.state()
1344 );
1345
1346 registry
1349 .register(&ctx, GatedPlugin { ready: true }, ())
1350 .expect("passing provider must register");
1351 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1352 dependent.refresh(&ctx).await;
1353 match dependent.state() {
1354 crate::FiberState::Active { .. } => {}
1355 other => {
1356 panic!("dependent should activate after a passing re-registration, got {other:?}")
1357 }
1358 }
1359 assert!(ctx.get::<GatedService>().is_some());
1360 }
1361
1362 #[derive(Debug)]
1370 struct ReadinessProbe;
1371 impl Service for ReadinessProbe {}
1372
1373 struct ReadinessPlugin;
1374 impl Plugin for ReadinessPlugin {
1375 type Config = ();
1376 type Provides = ReadinessProbe;
1377
1378 fn apply(
1379 &self,
1380 _ctx: &Arc<Context>,
1381 _config: Self::Config,
1382 ) -> Result<Arc<Self::Provides>, CordisError> {
1383 Ok(Arc::new(ReadinessProbe))
1384 }
1385 }
1386
1387 #[tokio::test]
1392 async fn ready_when_holds_pending_until_true_then_activates() {
1393 let ctx = Context::new_root();
1394 let registry = RegistryService::new();
1395
1396 let open = Arc::new(std::sync::atomic::AtomicBool::new(false));
1397 let gate_flag = open.clone();
1398 let fid = registry
1399 .register_with_readiness(
1400 &ctx,
1401 ReadinessPlugin,
1402 (),
1403 ReadinessBarrier::new(move |_ctx| {
1404 gate_flag.load(std::sync::atomic::Ordering::Acquire)
1405 }),
1406 )
1407 .expect("gated registration is non-throwing");
1408
1409 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1412 let fiber = registry.get_fiber(fid).expect("tracked");
1413 assert!(
1414 matches!(fiber.state(), FiberState::Pending),
1415 "closed gate must rest Pending, got {:?}",
1416 fiber.state()
1417 );
1418 assert!(
1419 ctx.get::<ReadinessProbe>().is_none(),
1420 "strict get refuses values owned by non-Active fibers"
1421 );
1422
1423 open.store(true, std::sync::atomic::Ordering::Release);
1425 fiber.refresh(&ctx).await;
1426 match fiber.state() {
1427 FiberState::Active { .. } => {}
1428 other => panic!("open gate must activate, got {other:?}"),
1429 }
1430 assert!(ctx.get::<ReadinessProbe>().is_some(), "now served");
1431 }
1432
1433 #[tokio::test]
1437 async fn readiness_composes_and_semantics() {
1438 let ctx = Context::new_root();
1439 let registry = RegistryService::new();
1440
1441 let a_open = Arc::new(std::sync::atomic::AtomicBool::new(false));
1442 let b_open = Arc::new(std::sync::atomic::AtomicBool::new(false));
1443 let combined = with_readiness([
1444 {
1445 let flag = a_open.clone();
1446 ReadinessBarrier::new(move |_ctx| {
1447 flag.load(std::sync::atomic::Ordering::Acquire)
1448 })
1449 },
1450 {
1451 let flag = b_open.clone();
1452 ReadinessBarrier::new(move |_ctx| {
1453 flag.load(std::sync::atomic::Ordering::Acquire)
1454 })
1455 },
1456 ]);
1457 assert!(!combined.is_ready(&ctx), "both closed must not be ready");
1459 a_open.store(true, std::sync::atomic::Ordering::Release);
1460 assert!(
1461 !combined.is_ready(&ctx),
1462 "AND semantics: one open half is not enough"
1463 );
1464 b_open.store(true, std::sync::atomic::Ordering::Release);
1465 assert!(combined.is_ready(&ctx), "both open must be ready");
1466 assert!(with_readiness([]).is_ready(&ctx));
1468
1469 b_open.store(false, std::sync::atomic::Ordering::Release);
1473 a_open.store(true, std::sync::atomic::Ordering::Release);
1474
1475 let fid = registry
1476 .register_with_readiness(&ctx, ReadinessPlugin, (), combined)
1477 .expect("composed-gate registration");
1478 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1479 let fiber = registry.get_fiber(fid).expect("tracked");
1480 assert!(
1481 matches!(fiber.state(), FiberState::Pending),
1482 "one-open-half must still wait Pending, got {:?}",
1483 fiber.state()
1484 );
1485
1486 b_open.store(true, std::sync::atomic::Ordering::Release);
1489 fiber.refresh(&ctx).await;
1490 match fiber.state() {
1491 FiberState::Active { .. } => {}
1492 other => panic!("fully-open AND gate must activate, got {other:?}"),
1493 }
1494 }
1495
1496 #[tokio::test]
1501 async fn external_rekick_reactivates_waiting_fiber() {
1502 let ctx = Context::new_root();
1503 ctx.provide(ReflectService::new());
1504 if let Some(reflect) = ctx.get::<crate::ReflectService>() {
1505 reflect.set_context(&ctx);
1506 }
1507 let registry = RegistryService::new();
1508 ctx.provide(registry);
1509 let registry = ctx.get::<RegistryService>().unwrap();
1510
1511 let fid = registry
1515 .register_with_readiness(
1516 &ctx,
1517 ReadinessPlugin,
1518 (),
1519 ReadinessBarrier::new(|ctx: &Arc<Context>| ctx.get::<Dependency>().is_some())
1520 .watching([TypeId::of::<Dependency>()]),
1521 )
1522 .expect("fact-gated registration");
1523 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1524 let fiber = registry.get_fiber(fid).expect("tracked");
1525 assert!(matches!(fiber.state(), FiberState::Pending));
1526
1527 let dep_fid = registry
1531 .register(&ctx, BarPlugin, ())
1532 .expect("dependency provider registers");
1533 assert!(dep_fid > 0);
1534 ctx.provide(Dependency);
1535 if let Some(reflect) = ctx.get::<crate::ReflectService>() {
1536 reflect.notify_with_ctx(TypeId::of::<Dependency>(), &ctx).await;
1537 }
1538 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1539
1540 match fiber.state() {
1541 FiberState::Active { .. } => {}
1542 other => panic!(
1543 "external re-kick must reactivate the waiting fiber, got {other:?}"
1544 ),
1545 }
1546 assert!(ctx.get::<ReadinessProbe>().is_some());
1547
1548 let _removed = ctx.remove::<Dependency>();
1552 if let Some(reflect) = ctx.get::<crate::ReflectService>() {
1553 reflect.notify_with_ctx(TypeId::of::<Dependency>(), &ctx).await;
1554 }
1555 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
1556 assert!(
1557 matches!(fiber.state(), FiberState::Pending),
1558 "gate closing again must rest Pending quietly, got {:?}",
1559 fiber.state()
1560 );
1561 }
1562}
1563
1564#[cfg(test)]
1565mod config_waterfall_tests {
1566 use super::*;
1567 use crate::events::{EventsService, INTERNAL_CONFIG_EVENT};
1568 use crate::{Context, FiberState};
1569
1570 struct CapturePlugin;
1572 #[derive(Debug)]
1573 struct CapturedConfig(pub serde_json::Value);
1574 impl Service for CapturedConfig {}
1575 impl Plugin for CapturePlugin {
1576 type Config = serde_json::Value;
1577 type Provides = CapturedConfig;
1578 fn apply(
1579 &self,
1580 _ctx: &Arc<Context>,
1581 cfg: Self::Config,
1582 ) -> Result<Arc<Self::Provides>, CordisError> {
1583 Ok(Arc::new(CapturedConfig(cfg)))
1584 }
1585 }
1586
1587 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
1592 async fn config_waterfall_covers_activation_path() {
1593 let ctx = Context::new_root();
1594 let events = Arc::new(EventsService::new());
1595 ctx.provide_arc(events.clone());
1596 let _gate = events.on(INTERNAL_CONFIG_EVENT.into(), |raw| async move {
1597 let mut effective = raw;
1598 if let Some(obj) = effective.as_object_mut() {
1599 obj.insert("rewritten_at_activation".into(), serde_json::json!(true));
1600 }
1601 Ok(effective)
1602 });
1603
1604 let registry = RegistryService::new();
1605 let fid = registry
1606 .register(
1607 &ctx,
1608 CapturePlugin,
1609 serde_json::json!({ "model": "base" }),
1610 )
1611 .expect("registration with an intercept-config listener");
1612 assert!(matches!(
1613 registry.get_fiber(fid).unwrap().state(),
1614 FiberState::Active { .. }
1615 ));
1616 let captured = ctx.get::<CapturedConfig>().expect("provider active");
1617 assert_eq!(captured.0["model"], "base");
1618 assert_eq!(
1619 captured.0["rewritten_at_activation"], true,
1620 "activation pass must consume the EFFECTIVE config"
1621 );
1622 }
1623}