1use core::mem;
4use std::sync::Arc;
5
6use crate::{
7 CancelObservation, EstablishedObservation, ObservationAuthority, ObservationId,
8 ObservationOperation, ObservationRejection, ObservationRelationship, ObserveEstablished,
9 ObservePeer, PeerStopped, WatchEvent,
10};
11use behavior::{
12 Actions, Address, Behavior, BehaviorActed, BirthMode, EndpointAddress, EventLayer,
13 InterpreterRequest, InterpreterRequests, ReturnsToEmitter, SendEffects, SendLayer,
14};
15
16pub type TerminationReaction<B> = fn(
18 &mut B,
19 PeerStopped<behavior::BehaviorAddr<B>>,
20) -> Actions<
21 behavior::BehaviorAddr<B>,
22 <B as Behavior>::Ph,
23 <B as Behavior>::Sends,
24 <B as Behavior>::Birth,
25>;
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub enum TerminationObservation {
30 Requested,
32 Observing,
34 Observed,
36 Cancelled,
38 Rejected {
40 operation: ObservationOperation,
41 reason: ObservationRejection,
42 },
43}
44
45#[derive(thiserror::Error)]
47pub enum TerminationMonitorError<E, Report> {
48 #[error("wrapped behavior rejected its event")]
50 Inner(#[source] E),
51 #[error("observation report does not match the active relationship phase")]
54 UnexpectedReport {
55 observation: TerminationObservation,
56 report: Report,
57 },
58}
59
60impl<E: core::fmt::Debug, Report> core::fmt::Debug for TerminationMonitorError<E, Report> {
61 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
62 match self {
63 Self::Inner(error) => formatter.debug_tuple("Inner").field(error).finish(),
64 Self::UnexpectedReport { observation, .. } => formatter
65 .debug_struct("UnexpectedReport")
66 .field("observation", observation)
67 .field("report", &"<retained>")
68 .finish(),
69 }
70 }
71}
72
73pub(crate) mod sealed {
74 pub trait TerminationObservationTarget<B: behavior::Behavior> {}
75}
76
77pub trait TerminationObservationTarget<B: Behavior>:
79 sealed::TerminationObservationTarget<B>
80{
81 type Report;
82 type Request: InterpreterRequest<
83 ReturnToEmitter = ReturnsToEmitter<Self::Report, behavior::Here>,
84 >;
85
86 fn request(&mut self) -> Option<Self::Request>;
87 fn observation(&self) -> TerminationObservation;
88 fn react(
89 &mut self,
90 inner: &mut B,
91 report: Self::Report,
92 ) -> Result<Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>, Self::Report>;
93}
94
95enum LogicalTermination<A: Address> {
96 Observing(A),
97 Observed(A),
98}
99
100pub struct LogicalTerminationTarget<B: Behavior> {
102 observation: LogicalTermination<behavior::BehaviorAddr<B>>,
103 react: TerminationReaction<B>,
104}
105
106impl<B: Behavior> sealed::TerminationObservationTarget<B> for LogicalTerminationTarget<B> {}
107
108impl<B: Behavior> TerminationObservationTarget<B> for LogicalTerminationTarget<B> {
109 type Report = PeerStopped<behavior::BehaviorAddr<B>>;
110 type Request = ObservePeer<behavior::BehaviorAddr<B>>;
111
112 fn request(&mut self) -> Option<Self::Request> {
113 let peer = match &self.observation {
114 LogicalTermination::Observing(peer) | LogicalTermination::Observed(peer) => *peer,
115 };
116 Some(ObservePeer::new(peer))
117 }
118
119 fn observation(&self) -> TerminationObservation {
120 match self.observation {
121 LogicalTermination::Observing(_) => TerminationObservation::Observing,
122 LogicalTermination::Observed(_) => TerminationObservation::Observed,
123 }
124 }
125
126 fn react(
127 &mut self,
128 inner: &mut B,
129 report: Self::Report,
130 ) -> Result<Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>, Self::Report> {
131 let LogicalTermination::Observing(peer) = &self.observation else {
132 return Err(report);
133 };
134 if report.peer != *peer {
135 return Err(report);
136 }
137 self.observation = LogicalTermination::Observed(*peer);
138 Ok((self.react)(inner, report))
139 }
140}
141
142pub type EstablishedTerminationReaction<B, P> = fn(
148 &mut B,
149 EstablishedObservation<P>,
150) -> Actions<
151 behavior::BehaviorAddr<B>,
152 <B as Behavior>::Ph,
153 <B as Behavior>::Sends,
154 <B as Behavior>::Birth,
155>;
156
157enum EstablishedTermination<P>
158where
159 P: behavior::Protocol,
160 P::Addr: behavior::RecipientAddress,
161{
162 Unissued(ObserveEstablished<P>),
163 Requested { correlation: Arc<ObservationId> },
164 Observing(ObservationAuthority<P>),
165 CancelPending(ObservationRelationship<P>),
166 CancelRejectedWaiting(CancelObservation<P>),
167 StoppedAwaitingCancel(ObservationRelationship<P>),
168 ObserveRejected(ObserveEstablished<P>),
169 Stopped,
170 StoppedWithCancelRejected(CancelObservation<P>),
171 Cancelled,
172}
173
174pub struct EstablishedTerminationTarget<B: Behavior, P>
176where
177 P: behavior::Protocol<Addr = behavior::BehaviorAddr<B>>,
178 behavior::BehaviorAddr<B>: EndpointAddress,
179{
180 observation: EstablishedTermination<P>,
181 react: EstablishedTerminationReaction<B, P>,
182}
183
184impl<B, P> sealed::TerminationObservationTarget<B> for EstablishedTerminationTarget<B, P>
185where
186 B: Behavior,
187 P: behavior::Protocol<Addr = behavior::BehaviorAddr<B>>,
188 behavior::BehaviorAddr<B>: EndpointAddress,
189{
190}
191
192impl<B, P> TerminationObservationTarget<B> for EstablishedTerminationTarget<B, P>
193where
194 B: Behavior,
195 P: behavior::Protocol<Addr = behavior::BehaviorAddr<B>>,
196 behavior::BehaviorAddr<B>: EndpointAddress,
197{
198 type Report = EstablishedObservation<P>;
199 type Request = ObserveEstablished<P>;
200
201 fn request(&mut self) -> Option<Self::Request> {
202 self.request_once()
203 }
204 fn observation(&self) -> TerminationObservation {
205 self.current_observation()
206 }
207 fn react(
208 &mut self,
209 inner: &mut B,
210 report: Self::Report,
211 ) -> Result<Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>, Self::Report> {
212 self.accept_report(inner, report)
213 }
214}
215
216impl<B, P> EstablishedTerminationTarget<B, P>
217where
218 B: Behavior,
219 P: behavior::Protocol<Addr = behavior::BehaviorAddr<B>>,
220 behavior::BehaviorAddr<B>: EndpointAddress,
221{
222 #[must_use]
225 pub fn take_cancellation(&mut self) -> Option<CancelObservation<P>> {
226 let EstablishedTermination::Observing(authority) = &self.observation else {
227 return None;
228 };
229 let pending = EstablishedTermination::CancelPending(authority.relationship().clone());
230 match mem::replace(&mut self.observation, pending) {
231 EstablishedTermination::Observing(authority) => Some(CancelObservation::new(authority)),
232 _ => unreachable!("the exclusive target held its one observed grant"),
233 }
234 }
235
236 pub fn into_rejected_observe(
238 self,
239 ) -> Result<(ObserveEstablished<P>, ObservationRejection), Self> {
240 match self {
241 Self {
242 observation: EstablishedTermination::ObserveRejected(request),
243 ..
244 } => Ok((request, ObservationRejection::IdAlreadyBound)),
245 other => Err(other),
246 }
247 }
248
249 pub fn into_rejected_cancel(
251 self,
252 ) -> Result<(CancelObservation<P>, ObservationRejection), Self> {
253 match self {
254 Self {
255 observation:
256 EstablishedTermination::CancelRejectedWaiting(request)
257 | EstablishedTermination::StoppedWithCancelRejected(request),
258 ..
259 } => Ok((request, ObservationRejection::NotObserved)),
260 other => Err(other),
261 }
262 }
263
264 fn request_once(&mut self) -> Option<ObserveEstablished<P>> {
265 let EstablishedTermination::Unissued(request) = &self.observation else {
266 return None;
267 };
268 let requested = EstablishedTermination::Requested {
269 correlation: request.correlation().clone(),
270 };
271 match mem::replace(&mut self.observation, requested) {
272 EstablishedTermination::Unissued(request) => Some(request),
273 _ => unreachable!("the exclusive target held its unissued request"),
274 }
275 }
276
277 fn current_observation(&self) -> TerminationObservation {
278 match &self.observation {
279 EstablishedTermination::Unissued(_) | EstablishedTermination::Requested { .. } => {
280 TerminationObservation::Requested
281 }
282 EstablishedTermination::Observing(_)
283 | EstablishedTermination::CancelPending(_)
284 | EstablishedTermination::CancelRejectedWaiting(..) => {
285 TerminationObservation::Observing
286 }
287 EstablishedTermination::Stopped
288 | EstablishedTermination::StoppedAwaitingCancel(_)
289 | EstablishedTermination::StoppedWithCancelRejected(..) => {
290 TerminationObservation::Observed
291 }
292 EstablishedTermination::Cancelled => TerminationObservation::Cancelled,
293 EstablishedTermination::ObserveRejected(_) => TerminationObservation::Rejected {
294 operation: ObservationOperation::Start,
295 reason: ObservationRejection::IdAlreadyBound,
296 },
297 }
298 }
299
300 fn accept_report(
301 &mut self,
302 inner: &mut B,
303 report: EstablishedObservation<P>,
304 ) -> Result<
305 Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>,
306 EstablishedObservation<P>,
307 > {
308 match report {
309 EstablishedObservation::Started { authority } => match &self.observation {
310 EstablishedTermination::Requested { correlation }
311 if authority.matches_request(correlation) =>
312 {
313 self.observation = EstablishedTermination::Observing(authority);
314 Ok(Actions::cont())
315 }
316 _ => Err(EstablishedObservation::Started { authority }),
317 },
318 EstablishedObservation::ObserveRejected { request, reason } => {
319 match &self.observation {
320 EstablishedTermination::Requested { correlation }
321 if reason == ObservationRejection::IdAlreadyBound
322 && Arc::ptr_eq(request.correlation(), correlation) =>
323 {
324 self.observation = EstablishedTermination::ObserveRejected(request);
325 Ok(Actions::cont())
326 }
327 _ => Err(EstablishedObservation::ObserveRejected { request, reason }),
328 }
329 }
330 EstablishedObservation::CancelRejected { request, reason } => match &self.observation {
331 EstablishedTermination::CancelPending(expected)
332 if reason == ObservationRejection::NotObserved
333 && expected == request.relationship() =>
334 {
335 self.observation = EstablishedTermination::CancelRejectedWaiting(request);
336 Ok(Actions::cont())
337 }
338 EstablishedTermination::StoppedAwaitingCancel(expected)
339 if reason == ObservationRejection::NotObserved
340 && expected == request.relationship() =>
341 {
342 self.observation = EstablishedTermination::StoppedWithCancelRejected(request);
343 Ok(Actions::cont())
344 }
345 _ => Err(EstablishedObservation::CancelRejected { request, reason }),
346 },
347 EstablishedObservation::Cancelled { relationship } => match &self.observation {
348 EstablishedTermination::CancelPending(expected) if expected == &relationship => {
349 self.observation = EstablishedTermination::Cancelled;
350 drop(relationship);
351 Ok(Actions::cont())
352 }
353 _ => Err(EstablishedObservation::Cancelled { relationship }),
354 },
355 stopped @ EstablishedObservation::Stopped { .. } => {
356 let EstablishedObservation::Stopped { relationship, .. } = &stopped else {
357 unreachable!("this branch owns a complete Stopped report");
358 };
359 match &self.observation {
360 EstablishedTermination::Observing(authority)
361 if authority.relationship() == relationship =>
362 {
363 self.observation = EstablishedTermination::Stopped;
364 }
365 EstablishedTermination::CancelPending(expected) if expected == relationship => {
366 self.observation =
367 EstablishedTermination::StoppedAwaitingCancel(expected.clone());
368 }
369 EstablishedTermination::CancelRejectedWaiting(request)
370 if request.relationship() == relationship =>
371 {
372 let terminal =
373 mem::replace(&mut self.observation, EstablishedTermination::Stopped);
374 let EstablishedTermination::CancelRejectedWaiting(request) = terminal
375 else {
376 unreachable!("the exclusive target retained its rejected cancellation");
377 };
378 self.observation =
379 EstablishedTermination::StoppedWithCancelRejected(request);
380 }
381 _ => return Err(stopped),
382 }
383 Ok((self.react)(inner, stopped))
386 }
387 }
388 }
389}
390
391pub struct TerminationMonitorWith<B: Behavior, Target: TerminationObservationTarget<B>> {
417 inner: B,
418 target: Target,
419}
420
421pub type TerminationMonitor<B> = TerminationMonitorWith<B, LogicalTerminationTarget<B>>;
423
424pub type EstablishedTerminationMonitor<B, P> =
426 TerminationMonitorWith<B, EstablishedTerminationTarget<B, P>>;
427
428type TerminationMonitorActions<B, Target> = Actions<
429 behavior::BehaviorAddr<B>,
430 <B as Behavior>::Ph,
431 SendLayer<
432 InterpreterRequests<<Target as TerminationObservationTarget<B>>::Request>,
433 <B as Behavior>::Sends,
434 >,
435 <B as Behavior>::Birth,
436>;
437
438impl<B: Behavior> TerminationMonitorWith<B, LogicalTerminationTarget<B>> {
439 #[must_use]
441 pub const fn new(
442 inner: B,
443 peer: behavior::BehaviorAddr<B>,
444 on_stopped: TerminationReaction<B>,
445 ) -> Self {
446 Self {
447 inner,
448 target: LogicalTerminationTarget {
449 observation: LogicalTermination::Observing(peer),
450 react: on_stopped,
451 },
452 }
453 }
454}
455
456impl<B, P> TerminationMonitorWith<B, EstablishedTerminationTarget<B, P>>
457where
458 B: Behavior,
459 P: behavior::Protocol<Addr = behavior::BehaviorAddr<B>>,
460 behavior::BehaviorAddr<B>: EndpointAddress,
461{
462 #[must_use]
463 pub fn established(
464 inner: B,
465 request: ObserveEstablished<P>,
466 react: EstablishedTerminationReaction<B, P>,
467 ) -> Self {
468 Self {
469 inner,
470 target: EstablishedTerminationTarget {
471 observation: EstablishedTermination::Unissued(request),
472 react,
473 },
474 }
475 }
476
477 #[must_use]
479 pub fn take_cancellation(&mut self) -> Option<CancelObservation<P>> {
480 self.target.take_cancellation()
481 }
482}
483
484impl<B: Behavior, Target: TerminationObservationTarget<B>> TerminationMonitorWith<B, Target> {
485 #[must_use]
487 pub fn observation(&self) -> TerminationObservation {
488 self.target.observation()
489 }
490
491 #[must_use]
493 pub fn into_parts(self) -> (B, Target) {
494 (self.inner, self.target)
495 }
496
497 fn wrap(
498 actions: Actions<behavior::BehaviorAddr<B>, B::Ph, B::Sends, B::Birth>,
499 observations: InterpreterRequests<Target::Request>,
500 ) -> TerminationMonitorActions<B, Target> {
501 actions.map_sends(|inner| SendLayer::new(observations, inner))
502 }
503}
504
505impl<B, Target> behavior::BehaviorBase for TerminationMonitorWith<B, Target>
506where
507 B: Behavior + behavior::BehaviorBase,
508 Target: TerminationObservationTarget<B>,
509{
510 type Base = B::Base;
511
512 fn base(&self) -> &Self::Base {
513 self.inner.base()
514 }
515}
516
517impl<B, Target> crate::StashStatus for TerminationMonitorWith<B, Target>
518where
519 B: Behavior + crate::StashStatus,
520 Target: TerminationObservationTarget<B>,
521{
522 fn stashed_messages(&self) -> usize {
523 self.inner.stashed_messages()
524 }
525}
526
527impl<B, Target, A, Ph, Sends, Br> Behavior for TerminationMonitorWith<B, Target>
528where
529 A: Address,
530 Sends: SendEffects + behavior::SendsFor<B::Event>,
531 Br: BirthMode,
532 B: Behavior<Ph = Ph, Sends = Sends, Birth = Br>,
533 B::Protocol: behavior::Protocol<Addr = A>,
534 Target: TerminationObservationTarget<B>,
535{
536 type Protocol = B::Protocol;
537 type Event = WatchEvent<B::Event, Target::Report>;
538 type Sends = SendLayer<InterpreterRequests<Target::Request>, Sends>;
539 type Ph = Ph;
540 type Error = TerminationMonitorError<B::Error, Target::Report>;
541 type Birth = Br;
542
543 fn init(&mut self, _: behavior::InitializationTurn) -> BehaviorActed<Self> {
544 let actions =
545 behavior::initialize(&mut self.inner).map_err(TerminationMonitorError::Inner)?;
546 Ok(Self::wrap(
547 actions,
548 match self.target.request() {
549 Some(request) => InterpreterRequests::one(request),
550 None => InterpreterRequests::empty(),
551 },
552 ))
553 }
554
555 fn transition(&mut self, _: behavior::ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
556 match event {
557 EventLayer::Owned(report) => {
558 let actions = self
559 .target
560 .react(&mut self.inner, report)
561 .map_err(|report| TerminationMonitorError::UnexpectedReport {
562 observation: self.target.observation(),
563 report,
564 })?;
565 Ok(Self::wrap(actions, InterpreterRequests::empty()))
566 }
567 EventLayer::Inner(event) => behavior::delegate_transition(&mut self.inner, event)
568 .map(|actions| Self::wrap(actions, InterpreterRequests::empty()))
569 .map_err(TerminationMonitorError::Inner),
570 }
571 }
572}
573
574#[cfg(test)]
575mod tests {
576 use super::*;
577 use crate::Activate as _;
578 use crate::{Crash, Exit};
579 use behavior::{
580 Births, CreateChild, CreationId, CreationSequence, Creations, MailAddr, Never, Step, User,
581 };
582
583 fn worker_creation() -> CreationId {
584 CreationSequence::new()
585 .issue()
586 .expect("the first worker creation ID exists")
587 }
588
589 struct Probe;
590
591 impl behavior::BehaviorBase for Probe {
592 type Base = Self;
593
594 fn base(&self) -> &Self {
595 self
596 }
597 }
598
599 impl behavior::Protocol for Probe {
600 type Addr = MailAddr;
601 type Msg = u8;
602 }
603
604 impl Behavior for Probe {
605 type Protocol = Self;
606 type Event = User<MailAddr, u8>;
607 type Sends = Vec<u8>;
608 type Ph = Never;
609 type Error = Never;
610 type Birth = Births<()>;
611
612 fn init(&mut self, _: behavior::InitializationTurn) -> BehaviorActed<Self> {
613 Ok(Actions::send(vec![1]))
614 }
615
616 fn transition(
617 &mut self,
618 _: behavior::ActiveTurn,
619 event: Self::Event,
620 ) -> BehaviorActed<Self> {
621 Ok(Actions::send(vec![event.message]))
622 }
623 }
624
625 #[allow(
626 clippy::needless_pass_by_value,
627 reason = "the reaction contract transfers ownership of the complete terminal report"
628 )]
629 fn reap(
630 _: &mut Probe,
631 stopped: PeerStopped<MailAddr>,
632 ) -> Actions<MailAddr, Never, Vec<u8>, Births<()>> {
633 assert_eq!(stopped.peer, MailAddr(4));
634 assert_eq!(stopped.outcome, Err(Crash::Panicked));
635 Actions::new(
636 vec![9],
637 Creations::one(CreateChild::birth(worker_creation(), ())),
638 Step::Continue,
639 )
640 }
641
642 #[test]
643 fn matching_terminal_report_preserves_complete_reaction_actions() {
644 let initialized = crate::TerminationMonitor::new(Probe, MailAddr(4), reap)
645 .initialize()
646 .unwrap();
647 assert_eq!(initialized.actions.sends.inner, [1]);
648 assert_eq!(
649 initialized.actions.sends.owned,
650 InterpreterRequests::one(crate::ObservePeer::new(MailAddr(4)))
651 );
652
653 let mut active = initialized.behavior;
654 let actions = active
655 .on_path(PeerStopped::new(MailAddr(4), Err(Crash::Panicked)))
656 .unwrap();
657 assert_eq!(actions.sends.inner, [9]);
658 assert!(actions.sends.owned.is_empty());
659 assert_eq!(
660 actions.creates,
661 Creations::one(CreateChild::birth(worker_creation(), ()))
662 );
663 assert!(matches!(actions.become_, Step::Continue));
664 assert_eq!(active.observation(), TerminationObservation::Observed);
665
666 let duplicate = PeerStopped::new(MailAddr(4), Err(Crash::Panicked));
667 assert!(matches!(
668 active.on_path(duplicate.clone()),
669 Err(TerminationMonitorError::UnexpectedReport {
670 observation: TerminationObservation::Observed,
671 report,
672 }) if report == duplicate
673 ));
674 }
675
676 #[allow(
677 clippy::needless_pass_by_value,
678 reason = "the reaction explicitly consumes the complete terminal report"
679 )]
680 fn acknowledge_capability_failure(
681 _: &mut Probe,
682 stopped: PeerStopped<MailAddr>,
683 ) -> Actions<MailAddr, Never, Vec<u8>, Births<()>> {
684 assert_eq!(stopped.peer, MailAddr(4));
685 assert_eq!(stopped.outcome, Err(Crash::CapabilityFailed));
686 Actions::send(vec![9])
687 }
688
689 #[test]
690 fn capability_failure_reaches_the_reaction_once_without_reclassification() {
691 let initialized =
692 crate::TerminationMonitor::new(Probe, MailAddr(4), acknowledge_capability_failure)
693 .initialize()
694 .unwrap();
695 assert_eq!(initialized.actions.sends.inner, [1]);
696 assert_eq!(
697 initialized.actions.sends.owned,
698 InterpreterRequests::one(crate::ObservePeer::new(MailAddr(4)))
699 );
700 assert!(initialized.actions.creates.is_empty());
701 assert!(matches!(initialized.actions.become_, Step::Continue));
702 let mut active = initialized.behavior;
703 let actions = active
704 .on_path(PeerStopped::new(MailAddr(4), Err(Crash::CapabilityFailed)))
705 .unwrap();
706 assert_eq!(actions.sends.inner, [9]);
707 assert!(actions.sends.owned.is_empty());
708 assert!(actions.creates.is_empty());
709 assert!(matches!(actions.become_, Step::Continue));
710 assert_eq!(active.observation(), TerminationObservation::Observed);
711
712 let duplicate = PeerStopped::new(MailAddr(4), Err(Crash::CapabilityFailed));
713 let rejected = active.on_path(duplicate.clone());
714 assert!(matches!(
715 rejected,
716 Err(TerminationMonitorError::UnexpectedReport {
717 observation: TerminationObservation::Observed,
718 report,
719 }) if report == duplicate
720 ));
721 }
722
723 #[test]
724 fn unmatched_terminal_report_is_returned_complete_and_user_actions_still_delegate() {
725 let mut active = crate::TerminationMonitor::new(Probe, MailAddr(4), reap)
726 .initialize()
727 .unwrap()
728 .behavior;
729 let unmatched = PeerStopped::new(MailAddr(5), Ok(Exit::Normal));
730 assert!(matches!(
731 active.on_path(unmatched.clone()),
732 Err(TerminationMonitorError::UnexpectedReport {
733 observation: TerminationObservation::Observing,
734 report,
735 }) if report == unmatched
736 ));
737
738 let delegated = active.receive(MailAddr(0), 7).unwrap();
739 assert_eq!(delegated.sends.inner, [7]);
740 assert!(delegated.sends.owned.is_empty());
741 }
742}