use std::any::{Any, TypeId};
use std::cell::{Ref, RefCell};
use std::collections::{HashMap, HashSet};
use std::rc::Rc;
use std::sync::atomic::{AtomicU64, Ordering};
use tokio::task::JoinHandle;
use crate::actor::Addr;
use crate::actor::event_bus::subscribe::BusSubscription;
use crate::actor::shape::Declared;
static NEXT_SUBSCRIBER_ID: AtomicU64 = AtomicU64::new(0);
pub struct Subscription {
unsubscribe: Option<Box<dyn FnOnce()>>,
}
impl Subscription {
pub(crate) fn inert() -> Self {
Self { unsubscribe: None }
}
}
impl Drop for Subscription {
fn drop(&mut self) {
if let Some(unsubscribe) = self.unsubscribe.take() {
unsubscribe();
}
}
}
pub struct StateHandle<T>(Rc<RefCell<T>>);
impl<T> StateHandle<T> {
pub fn borrow(&self) -> Ref<'_, T> {
self.0.borrow()
}
}
impl<T> Clone for StateHandle<T> {
fn clone(&self) -> Self {
Self(self.0.clone())
}
}
impl<T> From<Rc<RefCell<T>>> for StateHandle<T> {
fn from(inner: Rc<RefCell<T>>) -> Self {
Self(inner)
}
}
pub trait Reducer: Clone + Default + std::fmt::Debug + 'static {
type Update: Clone + 'static;
fn reduce(&mut self, update: Self::Update);
}
pub type Slot<R> = RefCell<Rc<R>>;
struct Cell {
state: Rc<dyn Any>,
listeners: RefCell<Vec<(u64, Rc<dyn Fn()>)>>,
observers: RefCell<Vec<(u64, Rc<dyn Fn(&dyn Any)>)>>,
}
#[derive(Clone, Copy)]
struct Kind {
name: &'static str,
describe: fn(&dyn Any) -> String,
}
impl Kind {
fn of<R: Reducer>() -> Self {
Self {
name: std::any::type_name::<R>(),
describe: |state| match state.downcast_ref::<Slot<R>>() {
Some(cell) => match cell.try_borrow() {
Ok(state) => format!("{:#?}", **state),
Err(_) => "<being changed>".to_string(),
},
None => "<unknown>".to_string(),
},
}
}
}
#[derive(Clone, Debug, PartialEq)]
pub struct Listener {
pub event: &'static str,
pub actor: Option<&'static str>,
pub bus: crate::trace::Bus,
pub feature: Option<&'static str>,
}
#[derive(Clone, Debug, PartialEq)]
pub struct DescribedState {
pub type_name: &'static str,
pub state: String,
pub feature: Option<&'static str>,
pub declared: Option<Declared>,
}
#[derive(Clone, Copy, Debug, PartialEq)]
pub struct Installed {
pub name: &'static str,
pub declared: Option<Declared>,
}
impl Cell {
fn empty<R: Reducer>() -> Self {
Cell {
state: Rc::new(Slot::new(Rc::new(R::default()))),
listeners: RefCell::new(Vec::new()),
observers: RefCell::new(Vec::new()),
}
}
}
#[derive(Default)]
pub struct Scope {
cells: RefCell<HashMap<TypeId, Cell>>,
kinds: RefCell<HashMap<TypeId, Kind>>,
teardowns: RefCell<Vec<Box<dyn FnOnce()>>>,
installed_features: RefCell<HashSet<TypeId>>,
exports: RefCell<HashSet<TypeId>>,
sections: RefCell<Vec<HashMap<TypeId, Rc<dyn Any>>>>,
section_names: RefCell<Vec<Option<&'static str>>>,
section_declarations: RefCell<Vec<Option<Declared>>>,
declarations: RefCell<HashMap<TypeId, Declared>>,
listeners: RefCell<Vec<Listener>>,
owners: RefCell<HashMap<TypeId, usize>>,
installing: RefCell<Vec<usize>>,
leave_guards: RefCell<Vec<Rc<dyn Fn() -> crate::guard::Verdict>>>,
}
impl Drop for Scope {
fn drop(&mut self) {
for teardown in self.teardowns.get_mut().drain(..) {
teardown();
}
}
}
impl Scope {
pub fn new() -> Self {
Self::default()
}
pub fn mark_feature_installed<F: 'static>(&self) {
let newly_inserted = self.installed_features.borrow_mut().insert(TypeId::of::<F>());
assert!(
newly_inserted,
"feature already installed in this scope - install() called twice for the same feature"
);
}
pub fn claims<R: 'static>(&self) -> bool {
self.owners.borrow().contains_key(&TypeId::of::<R>())
}
pub fn note_reducer_owner<R: 'static>(&self) {
self.installed_features.borrow_mut().insert(TypeId::of::<R>());
self.owners
.borrow_mut()
.entry(TypeId::of::<R>())
.or_insert_with(|| self.current_section());
}
pub fn note_reducer_declared<R: 'static>(&self, declared: Declared) {
self.declarations.borrow_mut().entry(TypeId::of::<R>()).or_insert(declared);
}
pub fn current_crate_dir(&self) -> Option<&'static str> {
let section = self.current_section();
let declarations = self.section_declarations.borrow();
declarations.get(section).copied().flatten().map(|declared| declared.crate_dir)
}
pub fn open_section(&self, name: &'static str, declared: Option<Declared>) -> usize {
let mut sections = self.sections.borrow_mut();
if sections.is_empty() {
sections.push(HashMap::new());
}
sections.push(HashMap::new());
let index = sections.len() - 1;
let mut names = self.section_names.borrow_mut();
names.resize(index + 1, None);
names[index] = Some(name);
let mut declarations = self.section_declarations.borrow_mut();
declarations.resize(index + 1, None);
declarations[index] = declared;
self.installing.borrow_mut().push(index);
index
}
pub fn section_name(&self, section: usize) -> Option<&'static str> {
self.section_names.borrow().get(section).copied().flatten()
}
pub fn features(&self) -> Vec<Installed> {
let declarations = self.section_declarations.borrow();
self.section_names
.borrow()
.iter()
.enumerate()
.filter_map(|(section, name)| {
Some(Installed {
name: (*name)?,
declared: declarations.get(section).copied().flatten(),
})
})
.collect()
}
pub fn current_feature(&self) -> Option<&'static str> {
self.section_name(self.current_section())
}
pub fn note_listener(
&self,
event: &'static str,
actor: Option<&'static str>,
bus: crate::trace::Bus,
) {
self.listeners.borrow_mut().push(Listener {
event,
actor,
bus,
feature: self.current_feature(),
});
}
pub fn listeners(&self) -> Vec<Listener> {
self.listeners.borrow().clone()
}
pub fn key(self: &Rc<Self>) -> usize {
Rc::as_ptr(self) as usize
}
pub fn close_section(&self) {
self.installing.borrow_mut().pop();
}
pub fn current_section(&self) -> usize {
self.installing.borrow().last().copied().unwrap_or(0)
}
pub fn section_of<R: 'static>(&self) -> usize {
self.owners
.borrow()
.get(&TypeId::of::<R>())
.copied()
.unwrap_or(0)
}
pub fn has_feature<F: 'static>(&self) -> bool {
self.installed_features.borrow().contains(&TypeId::of::<F>())
}
pub fn note_export<R: 'static>(&self) {
self.exports.borrow_mut().insert(TypeId::of::<R>());
}
pub fn exports<R: 'static>(&self) -> bool {
self.exports.borrow().contains(&TypeId::of::<R>())
}
fn note_kind<R: Reducer>(&self) {
self.kinds
.borrow_mut()
.entry(TypeId::of::<R>())
.or_insert_with(Kind::of::<R>);
}
pub fn describe_states(&self) -> Vec<DescribedState> {
let cells = self.cells.borrow();
let kinds = self.kinds.borrow();
let owners = self.owners.borrow();
let mut described: Vec<DescribedState> = cells
.iter()
.filter_map(|(type_id, cell)| {
let kind = kinds.get(type_id)?;
Some(DescribedState {
type_name: kind.name,
state: (kind.describe)(&*cell.state),
feature: owners
.get(type_id)
.and_then(|section| self.section_name(*section)),
declared: self.declarations.borrow().get(type_id).copied(),
})
})
.collect();
described.sort_by_key(|state| state.type_name);
described
}
pub fn state<R: Reducer>(&self) -> Rc<Slot<R>> {
self.note_kind::<R>();
let mut cells = self.cells.borrow_mut();
let cell = cells.entry(TypeId::of::<R>()).or_insert_with(Cell::empty::<R>);
cell.state
.clone()
.downcast::<Slot<R>>()
.expect("Scope cell type mismatch for this TypeId - unreachable, keyed by R")
}
pub fn peek<R: Reducer>(&self) -> Option<Rc<Slot<R>>> {
self.note_kind::<R>();
let cells = self.cells.borrow();
let cell = cells.get(&TypeId::of::<R>())?;
Some(
cell.state
.clone()
.downcast::<Slot<R>>()
.expect("Scope cell type mismatch for this TypeId - unreachable, keyed by R"),
)
}
pub fn seed<R: Reducer>(&self, state: R) {
self.note_kind::<R>();
self.cells.borrow_mut().insert(
TypeId::of::<R>(),
Cell {
state: Rc::new(Slot::new(Rc::new(state))),
listeners: RefCell::new(Vec::new()),
observers: RefCell::new(Vec::new()),
},
);
}
pub fn answers<M: 'static>(&self, answer: impl Fn(M) + 'static) {
let answer: Rc<dyn Fn(M)> = Rc::new(answer);
let section = self.current_section();
let mut sections = self.sections.borrow_mut();
while sections.len() <= section {
sections.push(HashMap::new());
}
sections[section].insert(TypeId::of::<M>(), Rc::new(answer) as Rc<dyn Any>);
}
pub fn answerer<M: 'static>(
&self,
section: usize,
) -> Option<Rc<dyn Fn(M)>> {
let sections = self.sections.borrow();
let answer = sections.get(section)?.get(&TypeId::of::<M>())?.clone();
answer.downcast::<Rc<dyn Fn(M)>>().ok().map(|a| (*a).clone())
}
pub fn first_answerer<M: 'static>(&self) -> Option<Rc<dyn Fn(M)>> {
let sections = self.sections.borrow().len();
(0..sections).find_map(|section| self.answerer::<M>(section))
}
pub fn push<R: Reducer>(self: &Rc<Self>, update: R::Update) {
let carried: Option<Box<dyn Any>> = self
.is_observed(TypeId::of::<R>())
.then(|| Box::new(update.clone()) as Box<dyn Any>);
let state = self.state::<R>();
{
let mut state = state.borrow_mut();
Rc::make_mut(&mut state).reduce(update);
}
crate::notify::mark(self, TypeId::of::<R>(), carried);
}
pub fn observe<R: Reducer>(
self: &Rc<Self>,
callback: impl Fn(&R::Update) + 'static,
) -> Subscription {
let id = NEXT_SUBSCRIBER_ID.fetch_add(1, Ordering::Relaxed);
{
let mut cells = self.cells.borrow_mut();
let cell = cells.entry(TypeId::of::<R>()).or_insert_with(Cell::empty::<R>);
cell.observers.borrow_mut().push((
id,
Rc::new(move |update: &dyn Any| {
if let Some(update) = update.downcast_ref::<R::Update>() {
callback(update);
}
}),
));
}
let scope = Rc::downgrade(self);
Subscription {
unsubscribe: Some(Box::new(move || {
let Some(scope) = scope.upgrade() else { return };
if let Some(cell) = scope.cells.borrow_mut().get_mut(&TypeId::of::<R>()) {
cell.observers.borrow_mut().retain(|(oid, _)| *oid != id);
}
})),
}
}
fn is_observed(&self, cell: TypeId) -> bool {
self.cells
.borrow()
.get(&cell)
.is_some_and(|cell| !cell.observers.borrow().is_empty())
}
pub(crate) fn listeners_of(&self, cell: TypeId) -> Vec<Rc<dyn Fn()>> {
let cells = self.cells.borrow();
cells
.get(&cell)
.map(|cell| cell.listeners.borrow().iter().map(|(_, f)| f.clone()).collect())
.unwrap_or_default()
}
pub(crate) fn observers_of(&self, cell: TypeId) -> Vec<Rc<dyn Fn(&dyn Any)>> {
let cells = self.cells.borrow();
cells
.get(&cell)
.map(|cell| cell.observers.borrow().iter().map(|(_, f)| f.clone()).collect())
.unwrap_or_default()
}
pub fn subscribe<R: Reducer>(
self: &Rc<Self>,
callback: impl Fn() + 'static,
) -> Subscription {
let id = NEXT_SUBSCRIBER_ID.fetch_add(1, Ordering::Relaxed);
{
let mut cells = self.cells.borrow_mut();
let cell = cells.entry(TypeId::of::<R>()).or_insert_with(Cell::empty::<R>);
cell.listeners.borrow_mut().push((id, Rc::new(callback)));
}
let scope = Rc::downgrade(self);
Subscription {
unsubscribe: Some(Box::new(move || {
let Some(scope) = scope.upgrade() else { return };
if let Some(cell) = scope.cells.borrow_mut().get_mut(&TypeId::of::<R>()) {
cell.listeners.borrow_mut().retain(|(lid, _)| *lid != id);
}
})),
}
}
pub fn binding<R: Reducer>(self: &Rc<Self>) -> crate::binding::ReducerBinding<R> {
crate::binding::ReducerBinding::new(self)
}
pub fn own_subscription(&self, subscription: BusSubscription) {
self.teardowns
.borrow_mut()
.push(Box::new(move || drop(subscription)));
}
pub fn snapshot_states(&self) -> HashMap<TypeId, Rc<dyn Any>> {
self.cells
.borrow()
.iter()
.map(|(type_id, cell)| (*type_id, cell.state.clone()))
.collect()
}
pub fn restore_states(&self, states: HashMap<TypeId, Rc<dyn Any>>) {
let mut cells = self.cells.borrow_mut();
for (type_id, state) in states {
cells.insert(
type_id,
Cell {
state,
listeners: RefCell::new(Vec::new()),
observers: RefCell::new(Vec::new()),
},
);
}
}
pub fn own<R: Teardown>(&self, resource: R) {
self.teardowns.borrow_mut().push(Box::new(move || resource.teardown()));
}
pub fn on_leave(&self, guard: impl Fn() -> crate::guard::Verdict + 'static) {
self.leave_guards.borrow_mut().push(Rc::new(guard));
}
pub fn leave_guards(&self) -> Vec<Rc<dyn Fn() -> crate::guard::Verdict>> {
self.leave_guards.borrow().clone()
}
}
pub trait Teardown: 'static {
fn teardown(self);
}
impl Teardown for JoinHandle<()> {
fn teardown(self) {
self.abort();
}
}
impl<A: 'static> Teardown for Addr<A> {
fn teardown(self) {
self.dispose();
drop(self);
}
}
pub struct DropGuard<T: 'static>(pub T);
impl<T: 'static> Teardown for DropGuard<T> {
fn teardown(self) {
drop(self.0);
}
}
pub struct GlobalScope;
impl GlobalScope {
pub fn instance() -> Rc<Scope> {
thread_local! {
static SCOPE: Rc<Scope> = Rc::new(Scope::new());
}
SCOPE.with(|scope| scope.clone())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
#[derive(Default, Clone, PartialEq, Debug)]
struct Counter {
value: i32,
}
#[derive(Clone)]
enum CounterMsg {
Set(i32),
}
impl Reducer for Counter {
type Update = CounterMsg;
fn reduce(&mut self, update: CounterMsg) {
match update {
CounterMsg::Set(v) => self.value = v,
}
}
}
#[test]
fn a_described_state_names_the_feature_that_claimed_it() {
#[derive(Clone, Default, Debug)]
struct Loose;
impl Reducer for Loose {
type Update = ();
fn reduce(&mut self, _: ()) {}
}
let scope = Scope::new();
scope.open_section("app::CounterFeature", None);
scope.note_reducer_owner::<Counter>();
scope.state::<Counter>();
scope.close_section();
scope.state::<Loose>();
let described = scope.describe_states();
let feature = |name: &str| {
described
.iter()
.find(|state| state.type_name.ends_with(name))
.map(|state| state.feature)
};
assert_eq!(feature("Counter"), Some(Some("app::CounterFeature")));
assert_eq!(feature("Loose"), Some(None));
}
#[test]
fn state_survives_unmount_remount_within_a_live_store() {
let store = Rc::new(Scope::new());
let first_read = store.state::<Counter>();
assert_eq!(first_read.borrow().value, 0);
store.push::<Counter>(CounterMsg::Set(42));
drop(first_read);
let second_read = store.state::<Counter>();
assert_eq!(second_read.borrow().value, 42);
}
#[test]
fn push_notifies_subscribers_and_unsubscribe_stops_it() {
let store = Rc::new(Scope::new());
let seen = Rc::new(RefCell::new(Vec::new()));
let seen_for_sub = seen.clone();
let store_for_sub = store.clone();
let sub = store.subscribe::<Counter>(move || {
seen_for_sub.borrow_mut().push(store_for_sub.state::<Counter>().borrow().value);
});
store.push::<Counter>(CounterMsg::Set(1));
store.push::<Counter>(CounterMsg::Set(2));
assert_eq!(*seen.borrow(), vec![1, 2]);
drop(sub);
store.push::<Counter>(CounterMsg::Set(3));
assert_eq!(
*seen.borrow(),
vec![1, 2],
"no further notifications after the Subscription is dropped"
);
}
#[tokio::test]
async fn dropping_the_store_aborts_owned_tasks() {
let ran_to_completion = Arc::new(AtomicBool::new(false));
let flag = ran_to_completion.clone();
let store = Scope::new();
let handle = tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
flag.store(true, Ordering::SeqCst);
});
store.own(handle);
drop(store);
tokio::time::sleep(std::time::Duration::from_millis(100)).await;
assert!(
!ran_to_completion.load(Ordering::SeqCst),
"task should have been aborted when its owning cell was dropped"
);
}
#[test]
fn own_actor_disposes_the_registry_entry_on_cell_drop() {
let token = crate::actor::UiThreadToken::dangerously_create_token_unchecked();
let addr = Addr::new_scoped((), token);
let counter = addr.strong_count_ptr();
let store = Scope::new();
store.own(addr.clone());
drop(addr);
assert!(
Rc::strong_count(&counter) > 1,
"REGISTRY should still hold the actor alive while its Scope is alive"
);
drop(store);
assert_eq!(
Rc::strong_count(&counter),
1,
"dropping the Scope's cell should dispose the REGISTRY entry, \
leaving only this test's own counter handle"
);
}
}