use std::any::{Any, TypeId};
use std::cell::RefCell;
use std::collections::VecDeque;
use std::rc::Rc;
use guinea_core::guard::Verdict;
use guinea_router::router::{Router, SegmentProps};
use crate::{Iced, Nodes, UpdateCx, dispatcher};
pub(crate) type Deliver = fn(&SegmentProps<Iced>, &mut Nodes, Box<dyn Any + Send>);
pub struct Envelope(Payload);
enum Payload {
Settled,
Answer(bool),
Node {
cursor: usize,
deliver: Deliver,
message: Box<dyn Any + Send>,
},
}
impl Envelope {
pub(crate) fn settled() -> Self {
Envelope(Payload::Settled)
}
pub fn answer(allowed: bool) -> Self {
Envelope(Payload::Answer(allowed))
}
pub(crate) fn new(cursor: usize, deliver: Deliver, message: Box<dyn Any + Send>) -> Self {
Envelope(Payload::Node {
cursor,
deliver,
message,
})
}
}
pub(crate) fn deliver<Node, Message>(
props: &SegmentProps<Iced>,
nodes: &mut Nodes,
message: Box<dyn Any + Send>,
update: fn(&mut Node, Message, &mut UpdateCx<'_, Node>),
leaving: fn(&Node) -> Verdict,
) where
Node: Default + 'static,
Message: Send + 'static,
{
if (props.chain[props.cursor].type_id)() != TypeId::of::<Node>() {
tracing::debug!(
node = std::any::type_name::<Node>(),
cursor = props.cursor,
"message arrived after its node left the chain; dropped"
);
return;
}
let Ok(message) = message.downcast::<Message>() else {
tracing::warn!(
node = std::any::type_name::<Node>(),
"message of the wrong type for its node; dropped"
);
return;
};
let Some((node, verdict)) = nodes.get_mut::<Node>(props.cursor) else {
return;
};
let mut cx = UpdateCx {
props,
segment: std::marker::PhantomData,
};
update(node, *message, &mut cx);
*verdict.borrow_mut() = leaving(node);
}
thread_local! {
static PARKED: RefCell<VecDeque<Envelope>> = const { RefCell::new(VecDeque::new()) };
}
pub(crate) fn park(envelope: Envelope) {
PARKED.with(|queue| queue.borrow_mut().push_back(envelope));
}
fn take_parked() -> Vec<Envelope> {
PARKED.with(|queue| queue.borrow_mut().drain(..).collect())
}
const SETTLE_ROUNDS: usize = 16;
pub(crate) fn settle(router: &Rc<Router<Iced>>, nodes: &mut Nodes, envelope: Envelope) {
apply(router, nodes, envelope);
for round in 0.. {
let parked = take_parked();
if parked.is_empty() {
break;
}
if round == SETTLE_ROUNDS {
tracing::warn!(
dropped = parked.len(),
"observers still producing messages after {SETTLE_ROUNDS} rounds; \
something is translating its own output back into its input"
);
break;
}
for envelope in parked {
apply(router, nodes, envelope);
}
}
if let Some(chain) = router.active_chain() {
nodes.sync(chain);
}
}
fn apply(router: &Rc<Router<Iced>>, nodes: &mut Nodes, envelope: Envelope) {
match envelope.0 {
Payload::Settled => dispatcher::drain(),
Payload::Answer(allowed) => router.answer(allowed),
Payload::Node {
cursor,
deliver,
message,
} => {
let Some(props) = props_at(router, cursor) else {
return;
};
deliver(&props, nodes, message);
}
}
}
fn props_at(router: &Rc<Router<Iced>>, cursor: usize) -> Option<SegmentProps<Iced>> {
let chain = router.active_chain()?;
let scopes = router.active_scopes()?;
(cursor < chain.len() && cursor < scopes.len()).then_some(SegmentProps {
chain,
scopes,
cursor,
})
}