use crate::{
disposable::Disposable,
observable::{Observable, Subscription},
observer::{Event, Flow, Observer, Termination, boxed_observer::BoxedObserver},
utils::{
mutable::{Mutable, MutableBool, MutableBoolHelper, MutableExt, MutableHelper},
on_panic::on_panic,
pending_events::PendingEvents,
types::{MaybeSend, Shared},
},
};
use educe::Educe;
pub fn unicast_subject<'or, T, E>() -> (UnicastSender<'or, T, E>, UnicastObservable<'or, T, E>) {
new_pair(PendingEvents::new())
}
pub fn unicast_subject_with_capacity<'or, T, E>(
capacity: usize,
) -> (UnicastSender<'or, T, E>, UnicastObservable<'or, T, E>) {
new_pair(PendingEvents::with_capacity(capacity))
}
fn new_pair<'or, T, E>(
pending: PendingEvents<T, E>,
) -> (UnicastSender<'or, T, E>, UnicastObservable<'or, T, E>) {
let pipe = Shared::new(Pipe {
is_disposed: MutableBool::new(false),
state: Mutable::new(State::Pending(pending)),
});
(
UnicastSender {
pipe: pipe.clone(),
observer: None,
},
UnicastObservable(Some(pipe)),
)
}
#[derive(Educe)]
#[educe(Debug)]
enum State<'or, T, E> {
Pending(PendingEvents<T, E>),
Attached(BoxedObserver<'or, T, E>),
Held,
Closed,
}
#[derive(Educe)]
#[educe(Debug)]
struct Pipe<'or, T, E> {
is_disposed: MutableBool,
state: Mutable<State<'or, T, E>>,
}
type SharedPipe<'or, T, E> = Shared<Pipe<'or, T, E>>;
#[derive(Educe)]
#[educe(Debug)]
pub struct UnicastSender<'or, T, E> {
pipe: SharedPipe<'or, T, E>,
observer: Option<BoxedObserver<'or, T, E>>,
}
impl<T, E> Drop for UnicastSender<'_, T, E> {
fn drop(&mut self) {
let observer = self.observer.take();
let previous_state = self.pipe.state.with_mut(|current| match current {
State::Pending(pending) if pending.is_terminated() => None,
state => Some(std::mem::replace(state, State::Closed)),
});
drop(observer); drop(previous_state); }
}
impl<T, E> UnicastSender<'_, T, E> {
pub fn is_disposed(&self) -> bool {
self.pipe.is_disposed.read()
}
}
impl<T, E> Observer<T, E> for UnicastSender<'_, T, E> {
fn on_next(&mut self, value: T) -> Flow {
if self.observer.is_some() {
return self.send_held(value);
}
let (delivery, rejected, discarded) = self.pipe.state.with_mut(|current| match current {
State::Pending(pending) => {
(None, pending.push(Event::Next(value)), None)
}
state @ State::Attached(_) => {
let State::Attached(observer) = std::mem::replace(state, State::Held) else {
unreachable!()
};
(Some((observer, value)), None, None)
}
State::Held => unreachable!(),
State::Closed => (None, None, Some(value)),
});
debug_assert!(rejected.is_none());
let flow = if discarded.is_some() {
Flow::Stop
} else {
Flow::Continue
};
drop((rejected, discarded)); if let Some((observer, value)) = delivery {
self.observer = Some(observer);
return self.send_held(value);
}
flow
}
fn on_termination(mut self, termination: Termination<E>) {
if let Some(observer) = self.observer.take() {
let previous_state = self.pipe.state.replace_value(State::Closed);
let is_disposed = matches!(previous_state, State::Closed);
drop(previous_state); if is_disposed {
drop(observer);
drop(termination);
} else {
observer.on_termination(termination); }
return;
}
let (delivery, rejected, discarded) = self.pipe.state.with_mut(|current| match current {
State::Pending(pending) => {
(None, pending.push(Event::Termination(termination)), None)
}
state @ State::Attached(_) => {
let State::Attached(observer) = std::mem::replace(state, State::Closed) else {
unreachable!()
};
(Some((observer, termination)), None, None)
}
State::Held => unreachable!(),
State::Closed => (None, None, Some(termination)),
});
debug_assert!(rejected.is_none());
drop((rejected, discarded)); if let Some((observer, termination)) = delivery {
observer.on_termination(termination); }
}
}
impl<T, E> UnicastSender<'_, T, E> {
fn send_held(&mut self, value: T) -> Flow {
debug_assert!(self.observer.is_some());
if self.pipe.is_disposed.read() {
let observer = self.observer.take();
drop(observer);
drop(value);
return Flow::Stop;
}
let mut flow = Flow::Continue;
if let Some(observer) = &mut self.observer {
flow = observer.on_next(value); }
if flow.is_stop() || self.pipe.is_disposed.read() {
let observer = self.observer.take();
close(&self.pipe);
drop(observer); return Flow::Stop;
}
Flow::Continue
}
}
#[derive(Educe)]
#[educe(Debug)]
pub struct UnicastObservable<'or, T, E>(Option<SharedPipe<'or, T, E>>);
impl<T, E> Drop for UnicastObservable<'_, T, E> {
fn drop(&mut self) {
if let Some(pipe) = self.0.take() {
close(&pipe);
}
}
}
impl<'or, T, E> Observable<'or, T, E> for UnicastObservable<'or, T, E> {
type D = Disposal<'or, T, E>;
fn subscribe(
mut self,
observer: impl Observer<T, E> + MaybeSend + 'or,
) -> Subscription<Self::D> {
let pipe = self
.0
.take()
.expect("the shared state is taken by either subscribing or dropping");
let is_live = deliver(&pipe, BoxedObserver::new(observer));
Subscription::new(Disposal(is_live.then_some(pipe)))
}
}
#[derive(Educe)]
#[educe(Debug)]
pub struct Disposal<'or, T, E>(Option<SharedPipe<'or, T, E>>);
impl<T, E> Disposable for Disposal<'_, T, E> {
fn dispose(self) {
if let Some(pipe) = self.0 {
close(&pipe);
}
}
}
fn close<T, E>(pipe: &SharedPipe<'_, T, E>) {
let previous_state = pipe.state.with_mut(|state| {
pipe.is_disposed.write(true);
std::mem::replace(state, State::Closed)
});
drop(previous_state); }
enum Step<'or, T, E> {
Next(BoxedObserver<'or, T, E>, T),
Terminate(BoxedObserver<'or, T, E>, Termination<E>),
Park,
Close(BoxedObserver<'or, T, E>),
}
fn deliver<'or, T, E>(
pipe: &SharedPipe<'or, T, E>,
mut observer: BoxedObserver<'or, T, E>,
) -> bool {
loop {
let step = pipe.state.with_mut(|current| {
let pending = match &mut *current {
State::Pending(pending) => pending,
State::Closed => return Step::Close(observer),
State::Attached(_) | State::Held => unreachable!(),
};
match pending.pop() {
Some(Event::Next(value)) => Step::Next(observer, value),
Some(Event::Termination(termination)) => {
*current = State::Closed;
Step::Terminate(observer, termination)
}
None => {
*current = State::Attached(observer);
Step::Park
}
}
});
match step {
Step::Next(next_observer, value) => {
observer = next_observer;
let close_on_panic = on_panic(|| close(pipe));
let flow = observer.on_next(value); drop(close_on_panic);
if flow.is_stop() {
close(pipe);
drop(observer); return false;
}
}
Step::Terminate(next_observer, termination) => {
next_observer.on_termination(termination); return false;
}
Step::Park => return true,
Step::Close(next_observer) => {
drop(next_observer); return false;
}
}
}
}