1use core::fmt;
4use core::marker::PhantomData;
5use std::sync::Arc;
6use std::time::Instant;
7
8use behavior::{
9 ActionItem, Behavior, CreationId, EndpointAddress, EstablishedActor, EstablishedRecipient,
10 Ingress, InjectEvent, InterpretEstablished, InterpretInstalledActor, InterpretationProgress,
11 InterpreterRequest, ItemSettlement, Never, Protocol, RecipientAddress, ReturnsToEmitter,
12 finish_item, prepare_item,
13};
14
15use super::ShutdownRequested;
16use crate::{Crash, Exit};
17
18pub struct ObserveEstablishedCreation<C, Occurrence>
27where
28 C: Behavior,
29 behavior::BehaviorAddr<C>: EndpointAddress,
30{
31 pub creation: CreationId,
32 occurrence: PhantomData<fn() -> (C, Occurrence)>,
33}
34
35impl<C, Occurrence> ObserveEstablishedCreation<C, Occurrence>
36where
37 C: Behavior,
38 behavior::BehaviorAddr<C>: EndpointAddress,
39{
40 #[must_use]
41 pub const fn new(creation: CreationId) -> Self {
42 Self {
43 creation,
44 occurrence: PhantomData,
45 }
46 }
47}
48
49impl<C, Occurrence> Copy for ObserveEstablishedCreation<C, Occurrence>
50where
51 C: Behavior,
52 behavior::BehaviorAddr<C>: EndpointAddress,
53{
54}
55
56impl<C, Occurrence> Clone for ObserveEstablishedCreation<C, Occurrence>
57where
58 C: Behavior,
59 behavior::BehaviorAddr<C>: EndpointAddress,
60{
61 fn clone(&self) -> Self {
62 *self
63 }
64}
65
66impl<C, Occurrence> InterpreterRequest for ObserveEstablishedCreation<C, Occurrence>
67where
68 C: Behavior,
69 behavior::BehaviorAddr<C>: EndpointAddress,
70{
71 type ReturnToEmitter =
72 ReturnsToEmitter<behavior::EstablishedCreation<C, Occurrence>, behavior::Here>;
73 type LogicalProtocols = behavior::NoBirthProtocols;
74}
75
76impl<C, Occurrence> ActionItem for ObserveEstablishedCreation<C, Occurrence>
77where
78 C: Behavior,
79 behavior::BehaviorAddr<C>: EndpointAddress,
80 <behavior::BehaviorAddr<C> as behavior::Address>::Nonce: Send,
81{
82 type Custody = (Option<Self>, Option<Self::Reply>);
83 type Input<'a>
84 = &'a mut Option<Self>
85 where
86 Self: 'a;
87 type Reply = ItemSettlement<Self, Self::Accepted, Self::Rejection, Self::Prerequisite>;
88
89 fn prepare_interpretation(
90 progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
91 ) {
92 prepare_item::<Self>(progress);
93 }
94
95 fn interpretation_input<'a>(
96 custody: &'a mut Self::Custody,
97 ) -> Option<(Self::Input<'a>, &'a mut Option<Self::Reply>)>
98 where
99 Self: 'a,
100 {
101 match custody {
102 (input @ Some(_), received @ None) => Some((input, received)),
103 _ => None,
104 }
105 }
106
107 fn finish_interpretation(
108 progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
109 ) {
110 finish_item::<Self>(progress);
111 }
112
113 type Accepted = ();
114 type Rejection = Never;
115 type Prerequisite = behavior::CreationCorrelation<C::Protocol, Occurrence>;
116}
117
118pub struct EstablishedChild<C, Occurrence>
125where
126 C: Behavior,
127 behavior::BehaviorAddr<C>: EndpointAddress,
128{
129 creation: CreationId,
130 actor: EstablishedActor<C>,
131 occurrence: PhantomData<fn() -> Occurrence>,
132}
133
134impl<C, Occurrence> EstablishedChild<C, Occurrence>
135where
136 C: Behavior,
137 behavior::BehaviorAddr<C>: EndpointAddress,
138{
139 #[must_use]
141 pub const fn creation(&self) -> CreationId {
142 self.creation
143 }
144
145 #[must_use]
147 pub fn actor(&self) -> EstablishedActor<C> {
148 self.actor.clone()
149 }
150
151 #[must_use]
158 pub fn shutdown_target<Parent, Targets>(&self) -> Targets
159 where
160 Parent: Behavior,
161 Occurrence: behavior::ChildRole<Parent, Child = C>,
162 Targets: crate::ShutdownTargetAt<C, <Occurrence as behavior::ChildRole<Parent>>::Position>,
163 {
164 Targets::shutdown_target_at(self.creation)
165 }
166
167 #[must_use]
169 pub fn into_parts(self) -> (CreationId, EstablishedActor<C>) {
170 (self.creation, self.actor)
171 }
172}
173
174pub fn established_child<Parent, Role>(
188 report: behavior::EstablishedCreation<behavior::RoleChild<Parent, Role>, Role>,
189) -> Result<EstablishedChild<behavior::RoleChild<Parent, Role>, Role>, behavior::CreationRejection>
190where
191 Parent: Behavior,
192 Role: behavior::ChildRole<Parent>,
193 behavior::BehaviorAddr<behavior::RoleChild<Parent, Role>>: EndpointAddress,
194{
195 let (creation, _, actor) = report.into_committed()?.into_parts();
196 Ok(EstablishedChild {
197 creation,
198 actor,
199 occurrence: PhantomData,
200 })
201}
202
203#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
208pub struct ObservationId(pub u64);
209
210pub struct ObservationRelationship<P: Protocol> {
239 identity: Arc<ObservationId>,
240 protocol: PhantomData<fn(P) -> P>,
241}
242
243impl<P: Protocol> fmt::Debug for ObservationRelationship<P> {
244 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
245 formatter
246 .debug_struct("ObservationRelationship")
247 .field("id", &self.id())
248 .finish_non_exhaustive()
249 }
250}
251
252impl<P: Protocol> ObservationRelationship<P> {
253 #[must_use]
254 pub fn id(&self) -> ObservationId {
255 *self.identity
256 }
257
258 #[must_use]
263 pub fn identity(&self) -> &Arc<ObservationId> {
264 &self.identity
265 }
266}
267
268impl<P: Protocol> Clone for ObservationRelationship<P> {
269 fn clone(&self) -> Self {
270 Self {
271 identity: self.identity.clone(),
272 protocol: PhantomData,
273 }
274 }
275}
276
277impl<P: Protocol> PartialEq for ObservationRelationship<P> {
278 fn eq(&self, other: &Self) -> bool {
279 Arc::ptr_eq(&self.identity, &other.identity)
280 }
281}
282
283impl<P: Protocol> Eq for ObservationRelationship<P> {}
284
285pub struct ObservationAuthority<P: Protocol> {
316 relationship: ObservationRelationship<P>,
317}
318
319impl<P: Protocol> fmt::Debug for ObservationAuthority<P> {
320 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
321 formatter
322 .debug_struct("ObservationAuthority")
323 .field("relationship", &self.relationship)
324 .finish_non_exhaustive()
325 }
326}
327
328impl<P: Protocol> ObservationAuthority<P> {
329 #[must_use]
336 pub fn issued(request: ObserveEstablished<P>) -> Self
337 where
338 P::Addr: RecipientAddress,
339 {
340 let ObserveEstablished {
341 correlation,
342 recipient,
343 } = request;
344 drop(recipient);
345 Self {
346 relationship: ObservationRelationship {
347 identity: correlation,
348 protocol: PhantomData,
349 },
350 }
351 }
352
353 #[must_use]
354 pub fn relationship(&self) -> &ObservationRelationship<P> {
355 &self.relationship
356 }
357
358 pub(crate) fn matches_request(&self, correlation: &Arc<ObservationId>) -> bool {
359 Arc::ptr_eq(self.relationship.identity(), correlation)
360 }
361
362 #[must_use]
364 pub fn into_relationship(self) -> ObservationRelationship<P> {
365 self.relationship
366 }
367}
368
369pub struct ObserveEstablished<P>
415where
416 P: Protocol,
417 P::Addr: RecipientAddress,
418{
419 correlation: Arc<ObservationId>,
420 recipient: EstablishedRecipient<P>,
421}
422
423impl<P> ObserveEstablished<P>
424where
425 P: Protocol,
426 P::Addr: RecipientAddress,
427{
428 #[must_use]
434 pub fn new(id: ObservationId, recipient: EstablishedRecipient<P>) -> Self {
435 Self {
436 correlation: Arc::new(id),
437 recipient,
438 }
439 }
440
441 #[must_use]
442 pub fn id(&self) -> ObservationId {
443 *self.correlation
444 }
445
446 pub(crate) fn correlation(&self) -> &Arc<ObservationId> {
447 &self.correlation
448 }
449
450 #[must_use]
452 pub fn into_inputs(self) -> (ObservationId, EstablishedRecipient<P>) {
453 (*self.correlation, self.recipient)
454 }
455
456 pub fn interpret<I>(self, interpreter: &mut I) -> I::Output
457 where
458 I: InterpretEstablishedObservation<P>,
459 {
460 let recipient = self.recipient.clone();
461 recipient.interpret(&mut ObservationTransfer {
462 request: Some(self),
463 interpreter,
464 })
465 }
466
467 pub fn settle<I>(self, interpreter: &mut I) -> ItemSettlement<Self, (), Never, Never>
468 where
469 I: InterpretEstablishedObservation<P, Output = ()>,
470 {
471 self.interpret(interpreter);
472 ItemSettlement::Accepted(())
473 }
474}
475
476impl<P> InterpreterRequest for ObserveEstablished<P>
477where
478 P: Protocol,
479 P::Addr: RecipientAddress,
480{
481 type ReturnToEmitter = ReturnsToEmitter<EstablishedObservation<P>, behavior::Here>;
482 type LogicalProtocols = behavior::NoBirthProtocols;
483}
484
485impl<P> ActionItem for ObserveEstablished<P>
486where
487 P: Protocol,
488 P::Addr: RecipientAddress,
489 <P::Addr as RecipientAddress>::Established<P>: Send,
490{
491 type Custody = (Option<Self>, Option<Self::Reply>);
492 type Input<'a>
493 = &'a mut Option<Self>
494 where
495 Self: 'a;
496 type Reply = ItemSettlement<Self, Self::Accepted, Self::Rejection, Self::Prerequisite>;
497
498 fn prepare_interpretation(
499 progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
500 ) {
501 prepare_item::<Self>(progress);
502 }
503
504 fn interpretation_input<'a>(
505 custody: &'a mut Self::Custody,
506 ) -> Option<(Self::Input<'a>, &'a mut Option<Self::Reply>)>
507 where
508 Self: 'a,
509 {
510 match custody {
511 (input @ Some(_), received @ None) => Some((input, received)),
512 _ => None,
513 }
514 }
515
516 fn finish_interpretation(
517 progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
518 ) {
519 finish_item::<Self>(progress);
520 }
521
522 type Accepted = ();
523 type Rejection = Never;
524 type Prerequisite = Never;
525}
526
527pub struct CancelObservation<P: Protocol> {
545 authority: ObservationAuthority<P>,
546}
547
548impl<P: Protocol> CancelObservation<P> {
549 #[must_use]
550 pub fn new(authority: ObservationAuthority<P>) -> Self {
551 Self { authority }
552 }
553
554 #[must_use]
555 pub fn id(&self) -> ObservationId {
556 self.authority.relationship().id()
557 }
558
559 #[must_use]
560 pub fn relationship(&self) -> &ObservationRelationship<P> {
561 self.authority.relationship()
562 }
563
564 #[must_use]
566 pub fn into_authority(self) -> ObservationAuthority<P> {
567 self.authority
568 }
569
570 #[must_use]
572 pub fn into_relationship(self) -> ObservationRelationship<P> {
573 self.authority.into_relationship()
574 }
575
576 pub fn interpret<I>(self, interpreter: &mut I) -> I::Output
577 where
578 P::Addr: RecipientAddress,
579 I: InterpretEstablishedObservation<P>,
580 {
581 interpreter.cancel(self)
582 }
583
584 pub fn settle<I>(self, interpreter: &mut I) -> ItemSettlement<Self, (), Never, Never>
585 where
586 P::Addr: RecipientAddress,
587 I: InterpretEstablishedObservation<P, Output = ()>,
588 {
589 self.interpret(interpreter);
590 ItemSettlement::Accepted(())
591 }
592}
593
594impl<P: Protocol> InterpreterRequest for CancelObservation<P>
595where
596 P::Addr: RecipientAddress,
597{
598 type ReturnToEmitter = ReturnsToEmitter<EstablishedObservation<P>, behavior::Here>;
599 type LogicalProtocols = behavior::NoBirthProtocols;
600}
601
602impl<P: Protocol> ActionItem for CancelObservation<P>
603where
604 P::Addr: RecipientAddress,
605{
606 type Custody = (Option<Self>, Option<Self::Reply>);
607 type Input<'a>
608 = &'a mut Option<Self>
609 where
610 Self: 'a;
611 type Reply = ItemSettlement<Self, Self::Accepted, Self::Rejection, Self::Prerequisite>;
612
613 fn prepare_interpretation(
614 progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
615 ) {
616 prepare_item::<Self>(progress);
617 }
618
619 fn interpretation_input<'a>(
620 custody: &'a mut Self::Custody,
621 ) -> Option<(Self::Input<'a>, &'a mut Option<Self::Reply>)>
622 where
623 Self: 'a,
624 {
625 match custody {
626 (input @ Some(_), received @ None) => Some((input, received)),
627 _ => None,
628 }
629 }
630
631 fn finish_interpretation(
632 progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
633 ) {
634 finish_item::<Self>(progress);
635 }
636
637 type Accepted = ();
638 type Rejection = Never;
639 type Prerequisite = Never;
640}
641
642#[derive(Debug, Clone, Copy, PartialEq, Eq)]
644pub enum ObservationOperation {
645 Start,
646 Cancel,
647}
648
649#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
651pub enum ObservationRejection {
652 #[error("the observation ID is already bound")]
654 IdAlreadyBound,
655 #[error("the exact observation relationship is not registered")]
657 NotObserved,
658}
659
660pub enum EstablishedObservation<P>
662where
663 P: Protocol,
664 P::Addr: RecipientAddress,
665{
666 Started {
667 authority: ObservationAuthority<P>,
668 },
669 Stopped {
670 relationship: ObservationRelationship<P>,
671 outcome: Result<Exit<P::Addr>, Crash>,
672 at: Instant,
673 },
674 Cancelled {
675 relationship: ObservationRelationship<P>,
676 },
677 ObserveRejected {
678 request: ObserveEstablished<P>,
679 reason: ObservationRejection,
680 },
681 CancelRejected {
682 request: CancelObservation<P>,
683 reason: ObservationRejection,
684 },
685}
686
687impl<P> EstablishedObservation<P>
688where
689 P: Protocol,
690 P::Addr: RecipientAddress,
691{
692 #[must_use]
693 pub fn started(authority: ObservationAuthority<P>) -> Self {
694 Self::Started { authority }
695 }
696
697 #[must_use]
698 pub fn stopped(
699 relationship: ObservationRelationship<P>,
700 outcome: Result<Exit<P::Addr>, Crash>,
701 at: Instant,
702 ) -> Self {
703 Self::Stopped {
704 relationship,
705 outcome,
706 at,
707 }
708 }
709
710 #[must_use]
711 pub fn cancelled(request: CancelObservation<P>) -> Self {
712 Self::Cancelled {
713 relationship: request.into_relationship(),
714 }
715 }
716
717 #[must_use]
718 pub fn observe_rejected(request: ObserveEstablished<P>, reason: ObservationRejection) -> Self {
719 Self::ObserveRejected { request, reason }
720 }
721
722 #[must_use]
723 pub fn cancel_rejected(request: CancelObservation<P>, reason: ObservationRejection) -> Self {
724 Self::CancelRejected { request, reason }
725 }
726
727 #[must_use]
728 pub fn id(&self) -> ObservationId {
729 match self {
730 Self::Started { authority } => authority.relationship().id(),
731 Self::Stopped { relationship, .. } | Self::Cancelled { relationship } => {
732 relationship.id()
733 }
734 Self::ObserveRejected { request, .. } => request.id(),
735 Self::CancelRejected { request, .. } => request.id(),
736 }
737 }
738}
739
740pub trait InterpretEstablishedObservation<P>
742where
743 P: Protocol,
744 P::Addr: RecipientAddress,
745{
746 type Output;
747
748 fn observe(
749 &mut self,
750 request: ObserveEstablished<P>,
751 endpoint: <P::Addr as RecipientAddress>::Established<P>,
752 ) -> Self::Output;
753 fn cancel(&mut self, request: CancelObservation<P>) -> Self::Output;
754}
755
756struct ObservationTransfer<'a, P, I>
757where
758 P: Protocol,
759 P::Addr: RecipientAddress,
760{
761 request: Option<ObserveEstablished<P>>,
762 interpreter: &'a mut I,
763}
764
765impl<P, I> InterpretEstablished<P> for ObservationTransfer<'_, P, I>
766where
767 P: Protocol,
768 P::Addr: RecipientAddress,
769 I: InterpretEstablishedObservation<P>,
770{
771 type Output = I::Output;
772
773 fn interpret_established(
774 &mut self,
775 endpoint: <P::Addr as RecipientAddress>::Established<P>,
776 ) -> Self::Output {
777 let Some(request) = self.request.take() else {
778 unreachable!("Core transfers one owned established endpoint once");
779 };
780 self.interpreter.observe(request, endpoint)
781 }
782}
783
784#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
786pub struct ShutdownId(pub u64);
787
788pub struct ShutdownEstablished<B, TargetPath>
842where
843 B: Behavior,
844 behavior::BehaviorAddr<B>: EndpointAddress,
845 B::Event: InjectEvent<ShutdownRequested, TargetPath>,
846{
847 pub id: ShutdownId,
848 actor: EstablishedActor<B>,
849 ingress: Ingress<ShutdownRequested, TargetPath>,
850}
851
852impl<B, TargetPath> ShutdownEstablished<B, TargetPath>
853where
854 B: Behavior,
855 behavior::BehaviorAddr<B>: EndpointAddress,
856 B::Event: InjectEvent<ShutdownRequested, TargetPath>,
857{
858 #[must_use]
859 pub const fn new(
860 id: ShutdownId,
861 actor: EstablishedActor<B>,
862 ingress: Ingress<ShutdownRequested, TargetPath>,
863 ) -> Self {
864 Self { id, actor, ingress }
865 }
866
867 #[must_use]
869 pub fn actor(&self) -> EstablishedActor<B> {
870 self.actor.clone()
871 }
872
873 pub fn settle<I>(
875 self,
876 interpreter: &mut I,
877 ) -> ItemSettlement<Self, ShutdownId, ShutdownRejection, Never>
878 where
879 I: InterpretEstablishedShutdown<B, TargetPath>,
880 {
881 let admission = self.actor.clone().interpret_actor(&mut ShutdownTransfer {
882 id: self.id,
883 ingress: self.ingress,
884 interpreter,
885 behavior: PhantomData,
886 });
887 match admission {
888 Ok(()) => ItemSettlement::Accepted(self.id),
889 Err(reason) => ItemSettlement::Rejected { item: self, reason },
890 }
891 }
892}
893
894impl<B, TargetPath> InterpreterRequest for ShutdownEstablished<B, TargetPath>
895where
896 B: Behavior,
897 behavior::BehaviorAddr<B>: EndpointAddress,
898 B::Event: InjectEvent<ShutdownRequested, TargetPath>,
899{
900 type ReturnToEmitter =
901 ReturnsToEmitter<EstablishedShutdownResolved<B::Protocol>, behavior::Here>;
902 type LogicalProtocols = behavior::NoBirthProtocols;
903}
904
905impl<B, TargetPath> ActionItem for ShutdownEstablished<B, TargetPath>
906where
907 B: Behavior,
908 behavior::BehaviorAddr<B>: EndpointAddress,
909 B::Event: InjectEvent<ShutdownRequested, TargetPath>,
910 EstablishedActor<B>: Send,
911{
912 type Custody = (Option<Self>, Option<Self::Reply>);
913 type Input<'a>
914 = &'a mut Option<Self>
915 where
916 Self: 'a;
917 type Reply = ItemSettlement<Self, Self::Accepted, Self::Rejection, Self::Prerequisite>;
918
919 fn prepare_interpretation(
920 progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
921 ) {
922 prepare_item::<Self>(progress);
923 }
924
925 fn interpretation_input<'a>(
926 custody: &'a mut Self::Custody,
927 ) -> Option<(Self::Input<'a>, &'a mut Option<Self::Reply>)>
928 where
929 Self: 'a,
930 {
931 match custody {
932 (input @ Some(_), received @ None) => Some((input, received)),
933 _ => None,
934 }
935 }
936
937 fn finish_interpretation(
938 progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
939 ) {
940 finish_item::<Self>(progress);
941 }
942
943 type Accepted = ShutdownId;
944 type Rejection = ShutdownRejection;
945 type Prerequisite = Never;
946}
947
948#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
950pub enum ShutdownRejection {
951 #[error("shutdown is already in progress for the exact incarnation")]
952 AlreadyStopping,
953 #[error("the exact incarnation is already stopped")]
954 AlreadyStopped,
955}
956
957pub enum EstablishedShutdownResolved<P: Protocol> {
959 Accepted {
960 id: ShutdownId,
961 protocol: PhantomData<fn() -> P>,
962 },
963 Rejected {
964 id: ShutdownId,
965 reason: ShutdownRejection,
966 protocol: PhantomData<fn() -> P>,
967 },
968}
969
970impl<P: Protocol> EstablishedShutdownResolved<P> {
971 #[must_use]
972 pub const fn accepted(id: ShutdownId) -> Self {
973 Self::Accepted {
974 id,
975 protocol: PhantomData,
976 }
977 }
978
979 #[must_use]
980 pub const fn rejected(id: ShutdownId, reason: ShutdownRejection) -> Self {
981 Self::Rejected {
982 id,
983 reason,
984 protocol: PhantomData,
985 }
986 }
987
988 #[must_use]
989 pub const fn id(&self) -> ShutdownId {
990 match self {
991 Self::Accepted { id, .. } | Self::Rejected { id, .. } => *id,
992 }
993 }
994}
995
996pub trait InterpretEstablishedShutdown<B, TargetPath>
998where
999 B: Behavior,
1000 behavior::BehaviorAddr<B>: EndpointAddress,
1001 B::Event: InjectEvent<ShutdownRequested, TargetPath>,
1002{
1003 fn shutdown(
1004 &mut self,
1005 id: ShutdownId,
1006 installed: <behavior::BehaviorAddr<B> as EndpointAddress>::Installed<B>,
1007 ingress: Ingress<ShutdownRequested, TargetPath>,
1008 ) -> Result<(), ShutdownRejection>;
1009}
1010
1011struct ShutdownTransfer<'a, I, B, TargetPath> {
1012 id: ShutdownId,
1013 ingress: Ingress<ShutdownRequested, TargetPath>,
1014 interpreter: &'a mut I,
1015 behavior: PhantomData<fn() -> B>,
1016}
1017
1018impl<B, TargetPath, I> InterpretInstalledActor<B> for ShutdownTransfer<'_, I, B, TargetPath>
1019where
1020 B: Behavior,
1021 behavior::BehaviorAddr<B>: EndpointAddress,
1022 B::Event: InjectEvent<ShutdownRequested, TargetPath>,
1023 I: InterpretEstablishedShutdown<B, TargetPath>,
1024{
1025 type Output = Result<(), ShutdownRejection>;
1026
1027 fn interpret_actor(
1028 &mut self,
1029 installed: <behavior::BehaviorAddr<B> as EndpointAddress>::Installed<B>,
1030 ) -> Self::Output {
1031 self.interpreter.shutdown(self.id, installed, self.ingress)
1032 }
1033}
1034
1035#[cfg(test)]
1036mod observation_request_ownership {
1037 use std::sync::Arc;
1038
1039 use behavior::{
1040 Address, EstablishedRecipient, InterpretEstablished, Protocol, RecipientAddress,
1041 };
1042
1043 use super::{ObservationAuthority, ObservationId, ObserveEstablished};
1044
1045 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
1046 struct ObservationAddr(u64);
1047
1048 impl Address for ObservationAddr {
1049 type Nonce = u64;
1050 }
1051 impl RecipientAddress for ObservationAddr {
1052 type Established<P>
1053 = Arc<Vec<u64>>
1054 where
1055 P: Protocol<Addr = Self>;
1056 }
1057
1058 struct ObservedProtocol;
1059 impl Protocol for ObservedProtocol {
1060 type Addr = ObservationAddr;
1061 type Msg = ();
1062 }
1063
1064 struct ObservationEndpoint;
1065 impl InterpretEstablished<ObservedProtocol> for ObservationEndpoint {
1066 type Output = Arc<Vec<u64>>;
1067 fn interpret_established(&mut self, endpoint: Arc<Vec<u64>>) -> Self::Output {
1068 endpoint
1069 }
1070 }
1071
1072 #[test]
1073 fn original_request_inputs_return_without_reconstruction() {
1074 let values = Arc::new(vec![17, 43, 31]);
1075 let original = values.as_ptr();
1076 let id = ObservationId(91);
1077 let recipient = EstablishedRecipient::<ObservedProtocol>::issued(values);
1078 let request = ObserveEstablished::new(id, recipient);
1079 let (returned_id, recipient) = request.into_inputs();
1080 let returned = recipient.interpret(&mut ObservationEndpoint);
1081 assert_eq!(returned_id, id);
1082 assert_eq!(returned.as_ptr(), original);
1083 assert_eq!(returned.as_slice(), [17, 43, 31]);
1084 }
1085
1086 #[test]
1087 fn distinct_request_construction_cannot_alias_preacceptance() {
1088 let id = ObservationId(92);
1089 let recipient = EstablishedRecipient::<ObservedProtocol>::issued(Arc::new(vec![17]));
1090 let first = ObserveEstablished::new(id, recipient.clone());
1091 let second = ObserveEstablished::new(id, recipient.clone());
1092 let original_correlation = first.correlation().clone();
1093 let foreign_correlation = second.correlation().clone();
1094 let original_authority = ObservationAuthority::issued(first);
1095 let foreign_authority = ObservationAuthority::issued(second);
1096 assert!(original_authority.matches_request(&original_correlation));
1097 assert!(!original_authority.matches_request(&foreign_correlation));
1098 assert_ne!(
1099 original_authority.relationship(),
1100 foreign_authority.relationship()
1101 );
1102 assert!(Arc::ptr_eq(
1103 original_authority.relationship().identity(),
1104 &original_correlation
1105 ));
1106 assert!(Arc::ptr_eq(
1107 foreign_authority.relationship().identity(),
1108 &foreign_correlation
1109 ));
1110 }
1111}