use crate::effect::{AsyncDisposer, EffectHandle};
use crate::events::{Event, EventOptions, EventResult, EventValue, EventsRoot, EventsService};
use crate::fiber::{Fiber, FiberInner};
use crate::logger::{LogArg, Logger, LoggerRoot, LoggerService};
use crate::reflect::{Accessor, ReflectRoot, ReflectService};
use crate::registry::{Inject, IntoPlugin, PluginOutput, RegistryRoot, RegistryService};
use crate::{Config, Result, Value};
use std::collections::HashMap;
use std::fmt::{self, Debug, Formatter};
use std::future::Future;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, OnceLock, Weak};
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
pub struct Isolation(pub(crate) u64);
impl Isolation {
pub const fn from_raw(value: u64) -> Self {
Self(value)
}
pub const fn as_raw(self) -> u64 {
self.0
}
}
pub type ContextFilter = Arc<dyn Fn(&Context) -> bool + Send + Sync + 'static>;
#[derive(Clone, Default)]
pub struct ContextMeta {
pub(crate) isolates: Arc<HashMap<String, Isolation>>,
pub(crate) intercepts: Arc<Vec<(String, Value)>>,
pub(crate) values: Arc<HashMap<String, Value>>,
pub(crate) filter: Option<ContextFilter>,
pub(crate) base_url: Option<Arc<str>>,
}
impl Debug for ContextMeta {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
f.debug_struct("ContextMeta")
.field("isolates", &self.isolates)
.field("intercept_count", &self.intercepts.len())
.field("value_keys", &self.values.keys().collect::<Vec<_>>())
.field("has_filter", &self.filter.is_some())
.field("base_url", &self.base_url)
.finish()
}
}
pub(crate) struct RootInner {
pub(crate) reflect: ReflectRoot,
pub(crate) registry: RegistryRoot,
pub(crate) events: EventsRoot,
pub(crate) logger: LoggerRoot,
pub(crate) next_scope: AtomicU64,
pub(crate) next_fiber: AtomicU64,
pub(crate) next_effect: AtomicU64,
pub(crate) root_fiber: OnceLock<Fiber>,
}
impl RootInner {
fn new() -> Self {
Self {
reflect: ReflectRoot::new(),
registry: RegistryRoot::new(),
events: EventsRoot::new(),
logger: LoggerRoot::new(),
next_scope: AtomicU64::new(0),
next_fiber: AtomicU64::new(0),
next_effect: AtomicU64::new(0),
root_fiber: OnceLock::new(),
}
}
pub(crate) fn scope(&self) -> Isolation {
Isolation(self.next_scope.fetch_add(1, Ordering::Relaxed) + 1)
}
pub(crate) fn fiber_id(&self) -> u64 {
self.next_fiber.fetch_add(1, Ordering::Relaxed) + 1
}
pub(crate) fn effect_id(&self) -> u64 {
self.next_effect.fetch_add(1, Ordering::Relaxed) + 1
}
}
#[derive(Clone)]
pub struct Context {
pub(crate) root: Arc<RootInner>,
pub(crate) fiber: Weak<FiberInner>,
pub(crate) meta: ContextMeta,
}
impl Context {
pub fn new() -> Self {
let root = Arc::new(RootInner::new());
let fiber = Fiber::new_root(Arc::downgrade(&root), ContextMeta::default());
root.root_fiber
.set(fiber.clone())
.unwrap_or_else(|_| unreachable!("new root has no fiber"));
fiber.context()
}
pub fn same_root(&self, other: &Context) -> bool {
Arc::ptr_eq(&self.root, &other.root)
}
pub fn root(&self) -> Context {
self.root
.root_fiber
.get()
.expect("root fiber initialized")
.context()
}
pub fn fiber(&self) -> Result<Fiber> {
self.fiber
.upgrade()
.map(Fiber::from_inner)
.ok_or_else(|| crate::CordisError::new(crate::ErrorCode::InactiveEffect))
}
pub fn base_url(&self) -> Option<&str> {
self.meta.base_url.as_deref()
}
pub fn with_base_url(&self, base_url: impl Into<Arc<str>>) -> Context {
let mut child = self.clone();
child.meta.base_url = Some(base_url.into());
child
}
pub fn extend<T>(&self, name: impl Into<String>, value: T) -> Context
where
T: Send + Sync + 'static,
{
self.extend_value(name, Value::new(value))
}
pub fn extend_value(&self, name: impl Into<String>, value: Value) -> Context {
let mut values = (*self.meta.values).clone();
values.insert(name.into(), value);
let mut child = self.clone();
child.meta.values = Arc::new(values);
child
}
pub fn metadata<T>(&self, name: &str) -> Result<Option<Arc<T>>>
where
T: Send + Sync + 'static,
{
self.meta.values.get(name).map(Value::downcast).transpose()
}
pub fn new_isolation(&self) -> Isolation {
self.root.scope()
}
pub fn isolate(&self, name: impl Into<String>) -> Context {
let label = self.new_isolation();
self.isolate_with(name, label)
}
pub fn isolate_with(&self, name: impl Into<String>, label: Isolation) -> Context {
let mut isolates = (*self.meta.isolates).clone();
isolates.insert(name.into(), label);
let mut child = self.clone();
child.meta.isolates = Arc::new(isolates);
child
}
pub fn isolate_many(
&self,
names: impl IntoIterator<Item = impl Into<String>>,
label: Option<Isolation>,
) -> Context {
let label = label.unwrap_or_else(|| self.new_isolation());
let mut child = self.clone();
let mut isolates = (*child.meta.isolates).clone();
for name in names {
isolates.insert(name.into(), label);
}
child.meta.isolates = Arc::new(isolates);
child
}
pub fn intercept<T>(&self, name: impl Into<String>, config: T) -> Context
where
T: Send + Sync + 'static,
{
self.intercept_value(name, Value::new(config))
}
pub fn intercept_value(&self, name: impl Into<String>, config: Value) -> Context {
let mut intercepts = (*self.meta.intercepts).clone();
intercepts.push((name.into(), config));
let mut child = self.clone();
child.meta.intercepts = Arc::new(intercepts);
child
}
pub fn intercepts<T>(&self, name: &str) -> Result<Vec<Arc<T>>>
where
T: Send + Sync + 'static,
{
self.meta
.intercepts
.iter()
.filter(|(entry, _)| entry == name)
.map(|(_, value)| value.downcast())
.collect()
}
pub fn with_filter<F>(&self, filter: F) -> Context
where
F: Fn(&Context) -> bool + Send + Sync + 'static,
{
let mut child = self.clone();
child.meta.filter = Some(Arc::new(filter));
child
}
pub fn events(&self) -> EventsService {
EventsService::new(self.clone())
}
pub fn reflect(&self) -> ReflectService {
ReflectService::new(self.clone())
}
pub fn registry(&self) -> RegistryService {
RegistryService::new(self.clone())
}
pub fn logger(&self) -> Logger {
LoggerService::new(self.clone()).logger(None)
}
pub fn named_logger(&self, name: impl Into<String>) -> Logger {
LoggerService::new(self.clone()).logger(Some(name.into()))
}
pub fn logger_service(&self) -> LoggerService {
LoggerService::new(self.clone())
}
pub fn effect<F>(&self, label: impl Into<String>, dispose: F) -> Result<EffectHandle>
where
F: FnOnce() -> Result<()> + Send + 'static,
{
self.fiber()?
.register_effect(label, AsyncDisposer::from_sync(dispose))
}
pub fn effect_infallible<F>(&self, label: impl Into<String>, dispose: F) -> Result<EffectHandle>
where
F: FnOnce() + Send + 'static,
{
self.fiber()?
.register_effect(label, AsyncDisposer::infallible(dispose))
}
pub fn effect_async<F, Fut>(&self, label: impl Into<String>, dispose: F) -> Result<EffectHandle>
where
F: FnOnce() -> Fut + Send + 'static,
Fut: Future<Output = Result<()>> + Send + 'static,
{
self.fiber()?
.register_effect(label, AsyncDisposer::from_async(dispose))
}
pub fn provide<T>(&self, name: impl Into<String>, value: T) -> Result<EffectHandle>
where
T: Send + Sync + 'static,
{
self.reflect()
.provide_value(name.into(), Value::new(value), None)
}
pub fn provide_arc<T>(&self, name: impl Into<String>, value: Arc<T>) -> Result<EffectHandle>
where
T: Send + Sync + 'static,
{
self.reflect()
.provide_value(name.into(), Value::from_arc(value), None)
}
pub fn get<T>(&self, name: &str) -> Result<Option<Arc<T>>>
where
T: Send + Sync + 'static,
{
self.reflect().get(name, true)
}
pub fn get_unchecked<T>(&self, name: &str) -> Result<Option<Arc<T>>>
where
T: Send + Sync + 'static,
{
self.reflect().get(name, false)
}
pub fn require<T>(&self, name: &str) -> Result<Arc<T>>
where
T: Send + Sync + 'static,
{
self.reflect().require(name)
}
pub fn set<T>(&self, name: &str, value: T) -> Result<()>
where
T: Send + Sync + 'static,
{
self.reflect().set_value(name, Value::new(value))
}
pub fn notify<I, S>(&self, names: I) -> Vec<Fiber>
where
I: IntoIterator<Item = S>,
S: AsRef<str>,
{
self.reflect().notify(names)
}
pub fn accessor(&self, name: impl Into<String>, accessor: Accessor) -> Result<EffectHandle> {
self.reflect().accessor(name.into(), accessor)
}
pub fn on<F>(&self, name: impl Into<String>, listener: F) -> Result<EffectHandle>
where
F: Fn(Event) -> EventResult + Send + Sync + 'static,
{
self.events().on(name, listener, EventOptions::default())
}
pub fn on_with<F>(
&self,
name: impl Into<String>,
listener: F,
options: EventOptions,
) -> Result<EffectHandle>
where
F: Fn(Event) -> EventResult + Send + Sync + 'static,
{
self.events().on(name, listener, options)
}
pub fn on_async<F, Fut>(&self, name: impl Into<String>, listener: F) -> Result<EffectHandle>
where
F: Fn(Event) -> Fut + Send + Sync + 'static,
Fut: Future<Output = EventResult> + Send + 'static,
{
self.events()
.on_async(name, listener, EventOptions::default())
}
pub fn once<F>(&self, name: impl Into<String>, listener: F) -> Result<EffectHandle>
where
F: Fn(Event) -> EventResult + Send + Sync + 'static,
{
self.events().once(name, listener, EventOptions::default())
}
pub fn emit(
&self,
name: impl Into<String>,
args: impl IntoIterator<Item = EventValue>,
) -> Result<()> {
self.events().emit(name, args)
}
pub async fn parallel(
&self,
name: impl Into<String>,
args: impl IntoIterator<Item = EventValue>,
) -> Result<()> {
self.events().parallel(name, args).await
}
pub async fn serial(
&self,
name: impl Into<String>,
args: impl IntoIterator<Item = EventValue>,
) -> EventResult {
self.events().serial(name, args).await
}
pub fn bail(
&self,
name: impl Into<String>,
args: impl IntoIterator<Item = EventValue>,
) -> EventResult {
self.events().bail(name, args)
}
pub fn waterfall<F>(
&self,
name: impl Into<String>,
args: impl IntoIterator<Item = EventValue>,
inner: F,
) -> EventResult
where
F: Fn() -> EventResult + Send + Sync + 'static,
{
self.events().waterfall(name, args, inner)
}
pub fn plugin<P, C>(&self, plugin: P, config: C) -> Fiber
where
P: IntoPlugin,
C: Send + Sync + 'static,
{
self.registry()
.plugin_value(plugin.into_plugin(), Config::new(config))
}
pub fn plugin_default<P>(&self, plugin: P) -> Fiber
where
P: IntoPlugin,
{
self.registry()
.plugin_value(plugin.into_plugin(), Config::default())
}
pub fn plugin_object<P, C>(&self, plugin: P, config: C) -> Fiber
where
P: crate::Plugin,
C: Send + Sync + 'static,
{
self.registry()
.plugin_value(crate::PluginHandle::new(plugin), Config::new(config))
}
pub fn plugin_object_default<P>(&self, plugin: P) -> Fiber
where
P: crate::Plugin,
{
self.registry()
.plugin_value(crate::PluginHandle::new(plugin), Config::default())
}
pub fn inject<F>(&self, inject: Inject, callback: F) -> Fiber
where
F: Fn(Context) -> Result<PluginOutput> + Send + Sync + 'static,
{
let plugin = crate::plugin_sync::<(), _>("anonymous", inject, move |ctx, _| callback(ctx));
self.plugin_default(plugin)
}
pub fn log_error(&self, error: impl ToString) {
self.logger().error(error.to_string(), Vec::<LogArg>::new());
}
pub(crate) fn scope_override(&self, name: &str) -> Option<Isolation> {
self.meta.isolates.get(name).copied()
}
pub(crate) fn filter(&self) -> Option<&ContextFilter> {
self.meta.filter.as_ref()
}
pub(crate) fn root_arc(&self) -> &Arc<RootInner> {
&self.root
}
}
impl Default for Context {
fn default() -> Self {
Self::new()
}
}
impl Debug for Context {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
let name = self
.fiber()
.map(|fiber| fiber.name())
.unwrap_or_else(|_| "disposed".to_owned());
f.debug_tuple("Context").field(&name).finish()
}
}