1use crate::protocol::{ObservePeer, PeerStopped};
4use crate::{Crash, Exit, TerminationMonitorError, TerminationObservation};
5use behavior::Step;
6use behavior::{
7 Actions, Address, Become, Behavior, BehaviorActed, BirthMode, EventLayer, InterpreterRequests,
8 SendEffects, SendLayer, UserEvent,
9};
10
11pub type WatchEvent<E, Report = PeerStopped<<E as UserEvent>::Addr>> = EventLayer<Report, E>;
13
14pub type LinkReaction<B> = fn(
16 &mut B,
17 behavior::BehaviorAddr<B>,
18 &Result<Exit<behavior::BehaviorAddr<B>>, Crash>,
19) -> Become;
20
21pub struct Watch<B: Behavior> {
29 inner: B,
30 peer: behavior::BehaviorAddr<B>,
31 on_stopped: LinkReaction<B>,
32}
33
34type WatchActions<B> = Actions<
35 behavior::BehaviorAddr<B>,
36 <B as Behavior>::Ph,
37 SendLayer<InterpreterRequests<ObservePeer<behavior::BehaviorAddr<B>>>, <B as Behavior>::Sends>,
38 <B as Behavior>::Birth,
39>;
40
41impl<B: Behavior> Watch<B> {
42 #[must_use]
44 pub const fn new(
45 inner: B,
46 peer: behavior::BehaviorAddr<B>,
47 on_stopped: LinkReaction<B>,
48 ) -> Self {
49 Self {
50 inner,
51 peer,
52 on_stopped,
53 }
54 }
55
56 fn wrap(
57 actions: Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>,
58 observations: InterpreterRequests<ObservePeer<behavior::BehaviorAddr<B>>>,
59 ) -> WatchActions<B> {
60 actions.map_sends(|inner| SendLayer::new(observations, inner))
61 }
62}
63
64impl<B> behavior::BehaviorBase for Watch<B>
65where
66 B: Behavior + behavior::BehaviorBase,
67{
68 type Base = B::Base;
69
70 fn base(&self) -> &Self::Base {
71 self.inner.base()
72 }
73}
74
75impl<B> crate::StashStatus for Watch<B>
76where
77 B: Behavior + crate::StashStatus,
78{
79 fn stashed_messages(&self) -> usize {
80 self.inner.stashed_messages()
81 }
82}
83
84impl<B, A, Ph, Sends, Br> Behavior for Watch<B>
85where
86 A: Address,
87 Sends: SendEffects + behavior::SendsFor<B::Event>,
88 Br: BirthMode,
89 B: Behavior<Ph = Ph, Sends = Sends, Birth = Br>,
90 B::Protocol: behavior::Protocol<Addr = A>,
91{
92 type Protocol = B::Protocol;
93 type Event = WatchEvent<B::Event>;
94 type Sends = SendLayer<InterpreterRequests<ObservePeer<A>>, Sends>;
95 type Ph = Ph;
96 type Error = TerminationMonitorError<B::Error, PeerStopped<A>>;
97 type Birth = Br;
98
99 fn init(&mut self, _: behavior::InitializationTurn) -> BehaviorActed<Self> {
100 let actions =
101 behavior::initialize(&mut self.inner).map_err(TerminationMonitorError::Inner)?;
102 Ok(Self::wrap(
103 actions,
104 InterpreterRequests::one(ObservePeer::new(self.peer)),
105 ))
106 }
107
108 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
109 match event {
110 EventLayer::Owned(report) if report.peer == self.peer => {
111 let become_ = match (self.on_stopped)(&mut self.inner, report.peer, &report.outcome)
112 {
113 Step::Continue => Step::Continue,
114 Step::Goto(never) => match never {},
115 Step::Stop(stopped) => Step::Stop(stopped),
116 };
117 Ok(Self::wrap(
118 Actions::just(become_),
119 InterpreterRequests::empty(),
120 ))
121 }
122 EventLayer::Owned(report) => Err(TerminationMonitorError::UnexpectedReport {
123 observation: TerminationObservation::Observing,
124 report,
125 }),
126 EventLayer::Inner(event) => behavior::delegate_transition(&mut self.inner, event)
127 .map(|actions| Self::wrap(actions, InterpreterRequests::empty()))
128 .map_err(TerminationMonitorError::Inner),
129 }
130 }
131}
132
133pub fn stop_on_abnormal_death<B: Behavior>(
135 _behavior: &mut B,
136 _peer: behavior::BehaviorAddr<B>,
137 outcome: &Result<Exit<behavior::BehaviorAddr<B>>, Crash>,
138) -> Become {
139 if let Ok(Exit::Normal | Exit::Collected) = outcome {
140 Step::Continue
141 } else {
142 Step::Stop(behavior::Stopped)
143 }
144}