use crate::events::EventStore;
use crate::fiber::Fiber;
use crate::registry::Registry;
use crate::service::ServiceStore;
use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub(crate) struct RealmKey(u64);
impl RealmKey {
pub(crate) const DEFAULT: Self = Self(0);
}
#[derive(Debug, Default)]
pub(crate) struct SnCounters {
realm: AtomicU64,
}
impl SnCounters {
pub(crate) fn next_realm_key(&self) -> RealmKey {
let previous = self
.realm
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |key| {
key.checked_add(1)
})
.expect("Service realm identity space exhausted");
RealmKey(previous + 1)
}
}
pub(crate) struct Root {
pub(crate) events: Arc<EventStore>,
pub(crate) observations: Arc<crate::observation::ObservationHub>,
pub(crate) updates: Arc<crate::update::UpdateStore>,
pub(crate) registry: Registry,
pub(crate) deps: crate::deps::DependencyIndex,
pub(crate) services: ServiceStore,
pub(crate) logger: Arc<crate::logger::LoggerService>,
pub(crate) realm_membership: Arc<crate::service::RealmMembership>,
pub(crate) scope_membership: Arc<ScopeMembership>,
pub(crate) sn: SnCounters,
pub(crate) root_fiber: Arc<Fiber>,
}
impl Root {
pub(crate) fn dependents_for_edges(&self, edges: &[(String, RealmKey)]) -> Vec<Arc<Fiber>> {
self.dependents_for_edges_with_retention(edges).0
}
pub(crate) fn dependents_for_edges_with_retention(
&self,
edges: &[(String, RealmKey)],
) -> (Vec<Arc<Fiber>>, Vec<Arc<Fiber>>) {
if let Some(indexed) = self.deps.dependents_of(edges) {
return (indexed, Vec::new());
}
let retained_snapshot = self.registry.snapshot_fibers();
let requested = edges
.iter()
.map(|(service, realm)| (service.as_str(), *realm))
.collect::<HashSet<_>>();
let affected = retained_snapshot
.iter()
.filter(|fiber| {
fiber.is_alive()
&& fiber
.dependency_edges()
.iter()
.any(|edge| requested.contains(&(edge.service.as_str(), edge.realm)))
})
.cloned()
.collect();
(affected, retained_snapshot)
}
}
#[derive(Clone)]
pub struct Context {
pub(crate) root: Arc<Root>,
pub(crate) fiber: Arc<Fiber>,
pub(crate) isolate: Option<Arc<IsolateLayer>>,
pub(crate) scope: Option<Arc<ScopeNode>>,
pub(crate) intercept: Option<Arc<InterceptLayer>>,
}
impl Context {
pub(crate) fn with_fiber(&self, fiber: Arc<Fiber>) -> Context {
Context {
root: self.root.clone(),
fiber,
isolate: self.isolate.clone(),
scope: self.scope.clone(),
intercept: self.intercept.clone(),
}
}
#[allow(
clippy::new_without_default,
reason = "v3 requires explicit Runtime construction and removes Context::default"
)]
pub fn new() -> Self {
let sn = SnCounters::default();
let root_fiber = Fiber::root();
let root = Arc::new(Root {
events: Arc::new(EventStore::new()),
observations: Arc::new(crate::observation::ObservationHub::default()),
updates: Arc::new(crate::update::UpdateStore::new()),
registry: Registry::new(),
deps: crate::deps::DependencyIndex::new(),
services: ServiceStore::new(),
logger: Arc::new(crate::logger::LoggerService::new()),
realm_membership: Arc::new(crate::service::RealmMembership::new()),
scope_membership: Arc::new(ScopeMembership),
sn,
root_fiber,
});
Self {
fiber: root.root_fiber.clone(),
isolate: None,
scope: None,
intercept: None,
root,
}
}
#[must_use = "a derived context that is dropped has no effect"]
pub fn root(&self) -> Self {
Self {
root: self.root.clone(),
fiber: self.root.root_fiber.clone(),
isolate: None,
scope: None,
intercept: None,
}
}
pub fn observe_runtime<O>(
&self,
observer: O,
) -> std::result::Result<(), crate::effect::EffectRegistrationError>
where
O: crate::observation::RuntimeObserver,
{
self.root.observations.register(self, observer)
}
pub fn runtime_snapshot(&self) -> crate::observation::RuntimeSnapshot {
use crate::fiber::{FiberRole, FiberState};
use crate::observation::{FiberSnapshot, RuntimeSnapshot, ServiceSnapshot};
let mut fibers = Vec::new();
fibers.push(FiberSnapshot {
id: self.root.root_fiber.id().clone(),
role: FiberRole::Root,
name: self.root.root_fiber.name.clone(),
state: FiberState::Active,
missing_services: Vec::new(),
});
fibers.extend(
self.root
.registry
.snapshot_fibers()
.into_iter()
.map(|fiber| crate::observation::fiber_snapshot(&self.root, &fiber)),
);
let services = self
.root
.services
.snapshot_occurrences()
.into_iter()
.map(|record| ServiceSnapshot {
id: crate::observation::ServicePublicationId(record.id),
service: record.service,
realm: crate::ServiceRealm::new(self.root.realm_membership.clone(), record.realm),
provider: record.provider,
visible: record.visible,
})
.collect();
RuntimeSnapshot { fibers, services }
}
pub(crate) fn fiber(&self) -> &Arc<Fiber> {
&self.fiber
}
pub(crate) fn isolate_key(&self, name: &str) -> RealmKey {
self.isolate
.as_deref()
.and_then(|layer| layer.lookup(name))
.unwrap_or(RealmKey::DEFAULT)
}
pub fn new_service_realm(&self) -> crate::ServiceRealm {
crate::ServiceRealm::new(
self.root.realm_membership.clone(),
self.root.sn.next_realm_key(),
)
}
#[must_use = "a derived context that is dropped has no effect"]
pub fn with_isolated_service(&self, name: &str) -> Self {
let mut map = HashMap::new();
map.insert(name.to_owned(), self.root.sn.next_realm_key());
self.derive_isolate(map)
}
#[must_use = "a derived context that is dropped has no effect"]
pub fn with_service_realms<I, K>(
&self,
realms: I,
) -> std::result::Result<Self, crate::service::RealmMappingError>
where
I: IntoIterator<Item = (K, crate::ServiceRealm)>,
K: Into<String>,
{
let pairs: Vec<(String, crate::ServiceRealm)> = realms
.into_iter()
.map(|(name, realm)| (name.into(), realm))
.collect();
let mut names = HashSet::with_capacity(pairs.len());
for (name, _) in &pairs {
if !names.insert(name.clone()) {
return Err(crate::service::RealmMappingError::DuplicateService {
service: name.clone(),
});
}
}
for (name, realm) in &pairs {
if !realm.belongs_to(&self.root.realm_membership) {
return Err(crate::service::RealmMappingError::ForeignRealm {
service: name.clone(),
});
}
}
let map = pairs
.into_iter()
.map(|(name, realm)| (name, realm.key()))
.collect();
Ok(self.derive_isolate(map))
}
fn derive_isolate(&self, map: HashMap<String, RealmKey>) -> Self {
Self {
root: self.root.clone(),
fiber: self.fiber.clone(),
isolate: Some(Arc::new(IsolateLayer {
map,
outer: self.isolate.clone(),
})),
scope: self.scope.clone(),
intercept: self.intercept.clone(),
}
}
#[must_use]
pub fn scope(&self) -> Scope {
Scope {
membership: self.root.scope_membership.clone(),
layer: self.scope.clone(),
}
}
#[must_use = "a derived context that is dropped has no effect"]
pub fn with_child_scope(&self) -> Self {
Self {
root: self.root.clone(),
fiber: self.fiber.clone(),
isolate: self.isolate.clone(),
scope: Some(ScopeNode::derive(self.scope.as_ref())),
intercept: self.intercept.clone(),
}
}
#[must_use = "a derived Context that is dropped has no effect"]
pub fn with_intercept<S: crate::ConfigurableService>(&self, layer: S::Layer) -> Context {
self.with_configured_intercept(
S::NAME.to_owned(),
crate::plugin::ConfiguredLayer::new::<S, _>(layer),
)
}
pub(crate) fn with_configured_intercept(
&self,
name: String,
configured: crate::plugin::ConfiguredLayer,
) -> Context {
Self {
root: self.root.clone(),
fiber: self.fiber.clone(),
isolate: self.isolate.clone(),
intercept: Some(Arc::new(InterceptLayer {
name,
configured,
outer: self.intercept.clone(),
})),
scope: self.scope.clone(),
}
}
pub fn resolve_config<S: crate::ConfigurableService>(
&self,
base: Option<&S::Layer>,
head: Option<&S::Layer>,
) -> std::result::Result<S::Resolved, crate::service::ConfigResolutionError<S::ComposeError>>
{
let mut configured = Vec::new();
let mut node = self.intercept.as_deref();
while let Some(current) = node {
if current.name == S::NAME {
configured.push(current.configured.clone());
}
node = current.outer.as_deref();
}
configured.reverse();
if configured.iter().any(|layer| !layer.matches::<S>()) {
return Err(crate::service::ConfigResolutionError::ContractMismatch {
service: S::NAME,
});
}
let layers = configured.iter().map(|layer| {
layer
.downcast_ref::<S::Layer>()
.expect("matching Service contracts retain their declared Layer type")
});
S::compose_config(base, layers, head)
.map_err(crate::service::ConfigResolutionError::Compose)
}
}
#[derive(Debug)]
pub(crate) struct ScopeMembership;
#[derive(Clone)]
pub struct Scope {
pub(crate) membership: Arc<ScopeMembership>,
pub(crate) layer: Option<Arc<ScopeNode>>,
}
impl Scope {
pub(crate) fn belongs_to(&self, membership: &Arc<ScopeMembership>) -> bool {
Arc::ptr_eq(&self.membership, membership)
}
}
impl std::fmt::Debug for Scope {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("Scope(..)")
}
}
#[derive(Debug)]
pub(crate) struct ScopeNode {
parent: Option<Arc<ScopeNode>>,
}
impl ScopeNode {
pub(crate) fn reaches(self: &Arc<Self>, head: &Arc<ScopeNode>) -> bool {
let mut node: Option<&Arc<ScopeNode>> = Some(self);
while let Some(current) = node {
if Arc::ptr_eq(current, head) {
return true;
}
node = current.parent.as_ref();
}
false
}
fn derive(parent: Option<&Arc<ScopeNode>>) -> Arc<ScopeNode> {
Arc::new(ScopeNode {
parent: parent.cloned(),
})
}
}
#[derive(Debug)]
pub(crate) struct IsolateLayer {
map: HashMap<String, RealmKey>,
outer: Option<Arc<IsolateLayer>>,
}
impl IsolateLayer {
fn lookup(&self, key: &str) -> Option<RealmKey> {
let mut node = Some(self);
while let Some(current) = node {
if let Some(id) = current.map.get(key) {
return Some(*id);
}
node = current.outer.as_deref();
}
None
}
}
pub(crate) struct InterceptLayer {
name: String,
configured: crate::plugin::ConfiguredLayer,
outer: Option<Arc<InterceptLayer>>,
}
#[cfg(test)]
mod tests {
use super::SnCounters;
use std::sync::atomic::Ordering;
#[test]
fn realm_key_exhaustion_refuses_to_reuse_an_identity() {
let counters = SnCounters::default();
counters.realm.store(u64::MAX, Ordering::Relaxed);
assert!(
std::panic::catch_unwind(|| counters.next_realm_key()).is_err(),
"exhaustion must stop allocation before the default or an existing realm can repeat"
);
}
}
#[cfg(test)]
mod issue26_tests {
use super::*;
use crate::fiber::DependencyEdge;
use crate::registry::PluginKey;
use crate::{InjectSpec, Plugin, PreparedPlugin, Service};
use std::convert::Infallible;
use std::future::Future;
use std::sync::atomic::{AtomicUsize, Ordering};
struct FallbackService;
impl Service for FallbackService {
const NAME: &'static str = "issue26-fallback-service";
}
struct FallbackDependent(Arc<AtomicUsize>);
impl Plugin for FallbackDependent {
type Config = ();
type Input = ();
type PrepareError = Infallible;
type ApplyError = Infallible;
fn inject(&self) -> InjectSpec {
InjectSpec::none().require(FallbackService::NAME)
}
fn prepare(&self, (): ()) -> Result<(), Infallible> {
Ok(())
}
fn apply(
&self,
_ctx: Context,
_prepared: &(),
) -> impl Future<Output = Result<(), Infallible>> + Send {
self.0.fetch_add(1, Ordering::SeqCst);
std::future::ready(Ok(()))
}
}
#[test]
fn incomplete_dependency_projection_falls_back_to_fiber_owned_edges() {
let ctx = Context::new();
let fiber = Fiber::new_with_edges(
"dependent",
vec![DependencyEdge::new("svc".to_owned(), RealmKey::DEFAULT)],
);
ctx.root
.registry
.attach_fiber(PluginKey::Anonymous(26), fiber.clone());
ctx.root.deps.register(&fiber);
ctx.root.deps.disable_for_test();
let affected = ctx
.root
.dependents_for_edges(&[("svc".to_owned(), RealmKey::DEFAULT)]);
assert_eq!(affected.len(), 1);
assert_eq!(affected[0].id(), fiber.id());
let residents = ctx.root.registry.snapshot_fibers();
ctx.root.deps.rebuild_for_test(&residents);
let rebuilt = ctx
.root
.dependents_for_edges(&[("svc".to_owned(), RealmKey::DEFAULT)]);
assert_eq!(rebuilt.len(), 1);
assert_eq!(rebuilt[0].id(), fiber.id());
}
#[tokio::test]
async fn disabled_projection_preserves_service_settlement() {
let ctx = Context::new();
let applies = Arc::new(AtomicUsize::new(0));
let dependent = ctx
.spawn(PreparedPlugin::from_input(
FallbackDependent(applies.clone()),
(),
))
.await
.unwrap();
assert_eq!(dependent.state(), crate::FiberState::Pending);
ctx.root.deps.disable_for_test();
let _publication = ctx.provide(Arc::new(FallbackService)).unwrap();
assert_eq!(dependent.ready().await.unwrap(), crate::FiberState::Active);
assert_eq!(applies.load(Ordering::SeqCst), 1);
}
}
#[cfg(test)]
mod critical_section_tests {
use super::Context;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
struct ReentrantRealmName {
ctx: Context,
converted: Arc<AtomicUsize>,
}
impl From<ReentrantRealmName> for String {
fn from(value: ReentrantRealmName) -> Self {
let snapshot = value.ctx.runtime_snapshot();
assert!(!snapshot.fibers().is_empty());
value.converted.fetch_add(1, Ordering::SeqCst);
"issue55/conversion".to_owned()
}
}
struct ReentrantConfig;
impl crate::Service for ReentrantConfig {
const NAME: &'static str = "issue55/config-compose";
}
impl crate::ConfigurableService for ReentrantConfig {
type Config = Context;
type Layer = Context;
type Resolved = ();
type PrepareError = std::convert::Infallible;
type ComposeError = std::convert::Infallible;
fn prepare_config(config: Self::Config) -> Result<Self::Layer, Self::PrepareError> {
assert!(!config.runtime_snapshot().fibers().is_empty());
Ok(config)
}
fn compose_config<'a>(
base: Option<&'a Self::Layer>,
layers: impl IntoIterator<Item = &'a Self::Layer>,
head: Option<&'a Self::Layer>,
) -> Result<Self::Resolved, Self::ComposeError> {
for ctx in base.into_iter().chain(layers).chain(head) {
assert!(!ctx.runtime_snapshot().fibers().is_empty());
}
Ok(())
}
}
#[test]
fn configuration_prepare_and_compose_can_reenter_runtime_observation() {
use crate::ConfigurableService;
let ctx = Context::new();
let prepared = ReentrantConfig::prepare_config(ctx.clone()).unwrap();
ctx.resolve_config::<ReentrantConfig>(Some(&prepared), None)
.unwrap();
}
#[test]
fn realm_name_conversion_can_reenter_runtime_observation() {
let ctx = Context::new();
let converted = Arc::new(AtomicUsize::new(0));
let realm = ctx.new_service_realm();
let mapped = ctx
.with_service_realms([(
ReentrantRealmName {
ctx: ctx.clone(),
converted: converted.clone(),
},
realm,
)])
.unwrap();
assert_eq!(converted.load(Ordering::SeqCst), 1);
drop(mapped);
}
}