Skip to main content

behavior_actors/protocol/
established.rs

1//! Exact-incarnation observation and orderly-shutdown protocols.
2
3use 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
18/// Request the exact committed result of a same-action creation at one named
19/// creator-local role.
20///
21/// The interpreter commits creation before interpreting this request. A
22/// successful report carries an [`EstablishedRecipient`]; a rejected report
23/// carries no capability. `Occurrence` keeps duplicate declarations of the
24/// same child protocol distinct without becoming protocol identity or a
25/// runtime key.
26pub 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
118/// Both capabilities established by one committed named-child creation.
119///
120/// `creation` remains relative to the creating actor and selects its private
121/// child binding together with the static occurrence. `actor` is the stronger
122/// exact installed capability used by exact delivery, observation, or
123/// [`ShutdownEstablished`].
124pub 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    /// Return the creator-local ID of this child.
140    #[must_use]
141    pub const fn creation(&self) -> CreationId {
142        self.creation
143    }
144
145    /// Clone the exact installed-actor capability.
146    #[must_use]
147    pub fn actor(&self) -> EstablishedActor<C> {
148        self.actor.clone()
149    }
150
151    /// Select this child in an existing heterogeneous shutdown target sum.
152    ///
153    /// `Parent` restores the namespace in which `Occurrence` was declared;
154    /// the compiler then selects the role's exact structural position. No
155    /// role value, address reconstruction, or runtime protocol choice is
156    /// required after creation has committed.
157    #[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    /// Consume the product into its local and exact capabilities.
168    #[must_use]
169    pub fn into_parts(self) -> (CreationId, EstablishedActor<C>) {
170        (self.creation, self.actor)
171    }
172}
173
174/// Strengthen one successful named-child creation report without losing either
175/// of its routing capabilities.
176///
177/// `Role` must be the occurrence declared by `Parent`, so the returned exact
178/// actor proves the concrete child behavior while the returned ID retains
179/// creator-local correlation. This is an Actors-level
180/// construction over existing Behavior capabilities, not another creation
181/// operation or an allocation shortcut.
182///
183/// # Errors
184///
185/// Returns the creation's typed [`behavior::CreationRejection`] and produces
186/// no actor capability when installation did not commit.
187pub 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/// Behavior-owned correlation for one observation relationship.
204///
205/// This value is local relationship evidence, not actor identity, endpoint
206/// identity, or proof that an observation was accepted.
207#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
208pub struct ObservationId(pub u64);
209
210/// Non-authorizing identity of one accepted protocol-indexed relationship.
211///
212/// The protocol brand cannot be changed even when the address is shared:
213/// ```compile_fail,E0308
214/// fn wrong_protocol<P, Q>(relationship: behavior_actors::ObservationRelationship<P>, outcome: Result<behavior_actors::Exit<Q::Addr>, behavior_actors::Crash>, at: std::time::Instant) -> behavior_actors::EstablishedObservation<Q>
215/// where P: behavior::Protocol, Q: behavior::Protocol<Addr = P::Addr>, Q::Addr: behavior::RecipientAddress {
216///     behavior_actors::EstablishedObservation::Stopped { relationship, outcome, at }
217/// }
218/// ```
219/// A read-only relationship is not a cancellation permission:
220/// ```compile_fail,E0308
221/// fn permission<P: behavior::Protocol>(relationship: behavior_actors::ObservationRelationship<P>) -> behavior_actors::ObservationAuthority<P> {
222///     relationship
223/// }
224/// ```
225///
226/// Same-protocol terminal reports and read-only identity transfer remain valid.
227/// ```no_run
228/// pub fn stopped<P, Q>(relationship: behavior_actors::ObservationRelationship<P>, outcome: Result<behavior_actors::Exit<P::Addr>, behavior_actors::Crash>, at: std::time::Instant) -> behavior_actors::EstablishedObservation<P>
229/// where P: behavior::Protocol, Q: behavior::Protocol<Addr = P::Addr>, P::Addr: behavior::RecipientAddress {
230///     behavior_actors::EstablishedObservation::Stopped { relationship, outcome, at }
231/// }
232/// ```
233/// ```no_run
234/// pub fn retained_relationship<P: behavior::Protocol>(relationship: behavior_actors::ObservationRelationship<P>) -> behavior_actors::ObservationRelationship<P> {
235///     relationship
236/// }
237/// ```
238pub 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    /// Borrow the strong accepted identity for exact runtime membership checks.
259    ///
260    /// Reading or cloning this Arc confers no cancellation permission and
261    /// cannot reconstruct or rebrand a relationship.
262    #[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
285/// Affine permission to attempt cancellation of one exact relationship.
286///
287/// A permission does not assert that its registration remains live.
288/// It transfers once and cannot be reused after constructing cancellation:
289/// ```compile_fail,E0382
290/// fn duplicate<P: behavior::Protocol>(authority: behavior_actors::ObservationAuthority<P>) {
291///     let first = behavior_actors::CancelObservation::new(authority);
292///     let second = behavior_actors::CancelObservation::new(authority);
293///     drop((first, second));
294/// }
295/// ```
296/// A different protocol cannot consume the grant:
297/// ```compile_fail,E0308
298/// fn wrong_protocol<P: behavior::Protocol, Q: behavior::Protocol<Addr = P::Addr>>(authority: behavior_actors::ObservationAuthority<P>) -> behavior_actors::CancelObservation<Q> {
299///     behavior_actors::CancelObservation::new(authority)
300/// }
301/// ```
302///
303/// A same-protocol grant transfers into exactly one cancellation request.
304/// ```no_run
305/// pub fn cancellation<P: behavior::Protocol>(authority: behavior_actors::ObservationAuthority<P>) {
306///     let first = behavior_actors::CancelObservation::new(authority);
307///     drop(first);
308/// }
309/// ```
310/// ```no_run
311/// pub fn cancellation<P: behavior::Protocol, Q: behavior::Protocol<Addr = P::Addr>>(authority: behavior_actors::ObservationAuthority<P>) -> behavior_actors::CancelObservation<P> {
312///     behavior_actors::CancelObservation::new(authority)
313/// }
314/// ```
315pub 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    /// Transfer the original private request identity into its affine authority.
330    ///
331    /// This advanced-host operation consumes the original request. The host
332    /// must commit the same identity to its exact membership owner before
333    /// publishing Started. Perform issuance outside all Behavior folds.
334    /// Never issue for a rejected request or return an admitted original.
335    #[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    /// Discharge permission and retain the exact non-authorizing relationship.
363    #[must_use]
364    pub fn into_relationship(self) -> ObservationRelationship<P> {
365        self.relationship
366    }
367}
368
369/// One affine request for exact protocol-indexed termination observation.
370///
371/// The admitted original transfers once into its authority. Nested request
372/// lanes require separate report ingress at each structural path.
373/// ```no_run
374/// pub fn accepted_request<P>(request: behavior_actors::ObserveEstablished<P>) where P: behavior::Protocol, P::Addr: behavior::RecipientAddress {
375///     let first = behavior_actors::ObservationAuthority::issued(request);
376///     drop(first);
377/// }
378/// ```
379/// ```compile_fail,E0382
380/// pub fn accepted_request<P>(request: behavior_actors::ObserveEstablished<P>) where P: behavior::Protocol, P::Addr: behavior::RecipientAddress {
381///     let first = behavior_actors::ObservationAuthority::issued(request);
382///     let second = behavior_actors::ObservationAuthority::issued(request);
383///     drop((first, second));
384/// }
385/// ```
386/// ```no_run
387/// fn accepts<Event, Sends: behavior::SendsFor<Event>>() {}
388///
389/// pub fn independent_acknowledgements<P>() where P: behavior::Protocol, P::Addr: behavior::RecipientAddress {
390///     accepts::<behavior::EventLayer<behavior_actors::EstablishedObservation<P>, behavior::EventLayer<behavior_actors::EstablishedObservation<P>, behavior::User<P::Addr, P::Msg>>>, behavior::SendLayer<behavior::InterpreterRequests<behavior_actors::CancelObservation<P>>, behavior::SendLayer<behavior::InterpreterRequests<behavior_actors::ObserveEstablished<P>>, Vec<behavior::Never>>>>();
391/// }
392/// ```
393/// ```compile_fail,E0277
394/// fn accepts<Event, Sends: behavior::SendsFor<Event>>() {}
395///
396/// pub fn independent_acknowledgements<P>() where P: behavior::Protocol, P::Addr: behavior::RecipientAddress {
397///     accepts::<behavior::EventLayer<behavior_actors::EstablishedObservation<P>, behavior::User<P::Addr, P::Msg>>, behavior::SendLayer<behavior::InterpreterRequests<behavior_actors::CancelObservation<P>>, behavior::SendLayer<behavior::InterpreterRequests<behavior_actors::ObserveEstablished<P>>, Vec<behavior::Never>>>>();
398/// }
399/// ```
400/// A request remains affine and cannot be cloned:
401/// ```no_run
402/// pub fn transfer<P>(request: behavior_actors::ObserveEstablished<P>) -> behavior_actors::ObserveEstablished<P> where P: behavior::Protocol, P::Addr: behavior::RecipientAddress { request }
403/// ```
404/// ```compile_fail,E0599
405/// pub fn transfer<P>(request: behavior_actors::ObserveEstablished<P>) -> behavior_actors::ObserveEstablished<P> where P: behavior::Protocol, P::Addr: behavior::RecipientAddress { request.clone() }
406/// ```
407/// Reading a relationship cannot reconstruct its old request identity:
408/// ```no_run
409/// pub fn new_request<P>(relationship: behavior_actors::ObservationRelationship<P>, recipient: behavior::EstablishedRecipient<P>) -> behavior_actors::ObserveEstablished<P> where P: behavior::Protocol, P::Addr: behavior::RecipientAddress { behavior_actors::ObserveEstablished::new(relationship.id(), recipient) }
410/// ```
411/// ```compile_fail,E0451
412/// pub fn old_request<P>(relationship: behavior_actors::ObservationRelationship<P>, recipient: behavior::EstablishedRecipient<P>) -> behavior_actors::ObserveEstablished<P> where P: behavior::Protocol, P::Addr: behavior::RecipientAddress { behavior_actors::ObserveEstablished { correlation: relationship.identity().clone(), recipient } }
413/// ```
414pub 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    /// Construct one affine request with a fresh private correlation.
429    ///
430    /// Issue the request outside Behavior folds, then transfer the whole
431    /// original into its owning Actions lane. A never-accepted rejected
432    /// original may be transferred again; new construction has fresh identity.
433    #[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    /// Discharge this request correlation and recover the original inputs.
451    #[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
527/// One affine attempt to cancel the exact accepted relationship.
528///
529/// Cancellation acknowledgements require their own typed report ingress.
530/// ```no_run
531/// fn accepts<Event, Sends: behavior::SendsFor<Event>>() {}
532///
533/// pub fn cancellation_acknowledgement<P>() where P: behavior::Protocol, P::Addr: behavior::RecipientAddress {
534///     accepts::<behavior::EventLayer<behavior_actors::EstablishedObservation<P>, behavior::User<P::Addr, P::Msg>>, behavior::InterpreterRequests<behavior_actors::CancelObservation<P>>>();
535/// }
536/// ```
537/// ```compile_fail,E0277
538/// fn accepts<Event, Sends: behavior::SendsFor<Event>>() {}
539///
540/// pub fn cancellation_acknowledgement<P>() where P: behavior::Protocol, P::Addr: behavior::RecipientAddress {
541///     accepts::<behavior::User<P::Addr, P::Msg>, behavior::InterpreterRequests<behavior_actors::CancelObservation<P>>>();
542/// }
543/// ```
544pub 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    /// Recover the same permission from an unconsumed cancellation request.
565    #[must_use]
566    pub fn into_authority(self) -> ObservationAuthority<P> {
567        self.authority
568    }
569
570    /// Consume permission and retain the exact non-authorizing receipt.
571    #[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/// Operation correlated by an [`ObservationId`].
643#[derive(Debug, Clone, Copy, PartialEq, Eq)]
644pub enum ObservationOperation {
645    Start,
646    Cancel,
647}
648
649/// Semantic rejection of an exact observation operation.
650#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
651pub enum ObservationRejection {
652    /// This observation ID already names a live relationship.
653    #[error("the observation ID is already bound")]
654    IdAlreadyBound,
655    /// Cancellation named no exact registered relationship.
656    #[error("the exact observation relationship is not registered")]
657    NotObserved,
658}
659
660/// Whole protocol-indexed observation and cancellation receipts.
661pub 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
740/// Advanced-host transfer of the whole original requests and exact endpoint.
741pub 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/// Behavior-owned correlation for one orderly-shutdown request.
785#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
786pub struct ShutdownId(pub u64);
787
788/// Request orderly shutdown of one exact installed concrete behavior.
789///
790/// `TargetPath` proves where [`ShutdownRequested`] enters the installed
791/// behavior's closed event algebra. The interpreter receives that typed
792/// ingress together with the exact installed actor; shutdown is therefore still an
793/// explicit event/effect transformation, not a privileged runtime side
794/// channel.
795///
796/// A concrete actor whose event algebra has no shutdown ingress cannot be
797/// strengthened into an orderly-shutdown request:
798///
799/// ```compile_fail,E0277
800/// #[derive(Clone, Copy, PartialEq, Eq)]
801/// struct RuntimeAddr(u64);
802/// impl behavior::Address for RuntimeAddr { type Nonce = u64; }
803/// struct Endpoint;
804/// impl Clone for Endpoint { fn clone(&self) -> Self { Self } }
805/// struct Installed<B: behavior::Behavior>(Endpoint, std::sync::mpsc::Sender<B::Event>);
806/// impl<B: behavior::Behavior> Clone for Installed<B> {
807///     fn clone(&self) -> Self { Self(self.0.clone(), self.1.clone()) }
808/// }
809/// impl behavior::EndpointAddress for RuntimeAddr {
810///     type Established<P> = Endpoint where P: behavior::Protocol<Addr = Self>;
811///     type Installed<B> = Installed<B>
812///         where B: behavior::Behavior<Protocol: behavior::Protocol<Addr = Self>>;
813///     fn recipient<B>(installed: &Self::Installed<B>) -> Endpoint
814///     where B: behavior::Behavior<Protocol: behavior::Protocol<Addr = Self>> {
815///         installed.0.clone()
816///     }
817/// }
818/// struct Worker;
819/// impl behavior::Protocol for Worker { type Addr = RuntimeAddr; type Msg = (); }
820/// impl behavior::Behavior for Worker {
821///     type Protocol = Self;
822///     type Event = behavior::User<RuntimeAddr, ()>;
823///     type Sends = behavior::NoSends;
824///     type Ph = behavior::Never;
825///     type Error = behavior::Never;
826///     type Birth = behavior::NoBirths;
827///     fn transition(
828///         &mut self,
829///         _: behavior::ActiveTurn,
830///         _: Self::Event,
831///     ) -> behavior::BehaviorActed<Self> { Ok(behavior::Actions::cont()) }
832/// }
833/// let (control, _inbox) = std::sync::mpsc::channel::<<Worker as behavior::Behavior>::Event>();
834/// let actor = behavior::EstablishedActor::<Worker>::issued(Installed(Endpoint, control));
835/// let _ = behavior_actors::ShutdownEstablished::<Worker, behavior::Here>::new(
836///     behavior_actors::ShutdownId(1),
837///     actor,
838///     behavior::Ingress::<behavior_actors::ShutdownRequested, behavior::Here>::new(),
839/// );
840/// ```
841pub 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    /// Clone the exact actor capability without changing request ownership.
868    #[must_use]
869    pub fn actor(&self) -> EstablishedActor<B> {
870        self.actor.clone()
871    }
872
873    /// Attempt this request and return its one statically selected settlement.
874    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/// Semantic rejection of exact orderly shutdown.
949#[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
957/// Complete immediate resolution of one exact orderly-shutdown request.
958pub 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
996/// Public power-user boundary for exact orderly-shutdown transfer.
997pub 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}