Skip to main content

behavior_actors/atomic/fifo_pool/
protocol.rs

1//! FIFO submission, assignment, completion, and customer outcomes.
2
3use std::sync::Arc;
4
5use behavior::{Address, BehaviorAddr, EndpointAddress, MessageProtocol};
6
7use crate::atomic::RoleName;
8use crate::{ChildStopped, ReplyRoute};
9
10use super::super::pool::worker::{CurrentWorker, WorkerPreparationError, WorkerReplacementError};
11use super::super::pool::{Assignment, JobId, SubmissionId};
12use super::super::worker::{WorkerActivationGrant, WorkerActivationOutcome};
13use super::super::{
14    ActivationPermit, ActivationPlan, WorkerAttempt, WorkerCreationRejection, WorkerPreparation,
15    WorkerSource, WorkerSubmission,
16};
17use super::{FifoEvent, PrepareWorkers};
18
19/// Exact reason a submitted job never entered pool ownership.
20#[derive(Clone, Copy, Debug, Eq, PartialEq)]
21pub enum AdmissionRejection {
22    /// No direct worker can serve now or recover later.
23    NoRecoverableWorkers,
24    /// Waiting capacity is full and no worker is ready for immediate assignment.
25    BacklogFull,
26    /// The accepted-job sequence cannot issue another non-reused correlation.
27    JobCorrelationUnavailable,
28    /// Immediate assignment cannot issue another non-reused correlation.
29    AssignmentCorrelationUnavailable,
30    /// The pool has closed admission.
31    ShuttingDown,
32}
33
34/// Exact reason a queued job returned without assignment completion.
35#[derive(Clone, Copy, Debug, Eq, PartialEq)]
36pub enum QueuedReturnReason {
37    /// Every direct worker retired without a recoverable successor.
38    NoRecoverableWorkers,
39    /// Assignment correlation could not be issued while a worker was ready.
40    AssignmentCorrelationUnavailable,
41    /// Pool shutdown extracted the waiting job.
42    PoolShutdown,
43}
44
45/// Exact reason an assigned job returned without a completed result.
46#[derive(Clone, Copy, Debug, Eq, PartialEq)]
47pub enum AssignedReturnReason {
48    /// The exact assigned worker stopped and interruption policy selected failure.
49    WorkerStopped,
50    /// A retried job could not reserve fresh assignment authority.
51    RetryPreparationRejected,
52    /// Pool shutdown extracted the assigned customer obligation.
53    PoolShutdown,
54    /// Delivery returned authority that contradicted an already received completion.
55    ContradictoryAssignmentSettlement,
56}
57
58enum CustomerOutcome<Role, Job, WorkerResult> {
59    Accepted {
60        submission: SubmissionId,
61        job: JobId,
62    },
63    Rejected {
64        submission: SubmissionId,
65        payload: Job,
66        reason: AdmissionRejection,
67    },
68    Completed {
69        job: JobId,
70        role: RoleName<Role>,
71        worker_result: WorkerResult,
72    },
73    ReturnedQueued {
74        job: JobId,
75        payload: Job,
76        reason: QueuedReturnReason,
77    },
78    ReturnedAssigned {
79        job: JobId,
80        role: RoleName<Role>,
81        payload: Job,
82        reason: AssignedReturnReason,
83    },
84}
85
86/// Customer-visible alternative selected by one FIFO outcome.
87#[derive(Clone, Copy, Debug, Eq, PartialEq)]
88pub enum FifoOutcomeKind {
89    Accepted,
90    Rejected,
91    Completed,
92    ReturnedQueued,
93    ReturnedAssigned,
94}
95
96/// Complete customer-visible progress or terminal outcome for FIFO work.
97///
98/// The value is opaque so a pool can retain a non-`Clone` application role
99/// while an outcome borrows the same immutable role name. Exact consuming
100/// projections return owned payloads and results without exposing shared
101/// storage.
102pub struct FifoOutcome<Role, Job, WorkerResult> {
103    outcome: CustomerOutcome<Role, Job, WorkerResult>,
104}
105
106impl<Role, Job, WorkerResult> FifoOutcome<Role, Job, WorkerResult> {
107    pub(super) const fn accepted(submission: SubmissionId, job: JobId) -> Self {
108        Self {
109            outcome: CustomerOutcome::Accepted { submission, job },
110        }
111    }
112
113    pub(super) const fn rejected(
114        submission: SubmissionId,
115        payload: Job,
116        reason: AdmissionRejection,
117    ) -> Self {
118        Self {
119            outcome: CustomerOutcome::Rejected {
120                submission,
121                payload,
122                reason,
123            },
124        }
125    }
126
127    pub(super) const fn completed(
128        job: JobId,
129        role: RoleName<Role>,
130        worker_result: WorkerResult,
131    ) -> Self {
132        Self {
133            outcome: CustomerOutcome::Completed {
134                job,
135                role,
136                worker_result,
137            },
138        }
139    }
140
141    pub(super) const fn returned_queued(
142        job: JobId,
143        payload: Job,
144        reason: QueuedReturnReason,
145    ) -> Self {
146        Self {
147            outcome: CustomerOutcome::ReturnedQueued {
148                job,
149                payload,
150                reason,
151            },
152        }
153    }
154
155    pub(super) const fn returned_assigned(
156        job: JobId,
157        role: RoleName<Role>,
158        payload: Job,
159        reason: AssignedReturnReason,
160    ) -> Self {
161        Self {
162            outcome: CustomerOutcome::ReturnedAssigned {
163                job,
164                role,
165                payload,
166                reason,
167            },
168        }
169    }
170
171    /// Identify the customer-visible alternative without consuming it.
172    #[must_use]
173    pub const fn kind(&self) -> FifoOutcomeKind {
174        match &self.outcome {
175            CustomerOutcome::Accepted { .. } => FifoOutcomeKind::Accepted,
176            CustomerOutcome::Rejected { .. } => FifoOutcomeKind::Rejected,
177            CustomerOutcome::Completed { .. } => FifoOutcomeKind::Completed,
178            CustomerOutcome::ReturnedQueued { .. } => FifoOutcomeKind::ReturnedQueued,
179            CustomerOutcome::ReturnedAssigned { .. } => FifoOutcomeKind::ReturnedAssigned,
180        }
181    }
182
183    /// Borrow the semantic worker role carried by a terminal assigned outcome.
184    #[must_use]
185    pub fn role(&self) -> Option<&Role> {
186        match &self.outcome {
187            CustomerOutcome::Completed { role, .. }
188            | CustomerOutcome::ReturnedAssigned { role, .. } => Some(role.role()),
189            CustomerOutcome::Accepted { .. }
190            | CustomerOutcome::Rejected { .. }
191            | CustomerOutcome::ReturnedQueued { .. } => None,
192        }
193    }
194
195    /// Borrow a returned customer payload, when present.
196    #[must_use]
197    pub const fn payload(&self) -> Option<&Job> {
198        match &self.outcome {
199            CustomerOutcome::Rejected { payload, .. }
200            | CustomerOutcome::ReturnedQueued { payload, .. }
201            | CustomerOutcome::ReturnedAssigned { payload, .. } => Some(payload),
202            CustomerOutcome::Accepted { .. } | CustomerOutcome::Completed { .. } => None,
203        }
204    }
205
206    /// Borrow a completed worker result, when present.
207    #[must_use]
208    pub const fn worker_result(&self) -> Option<&WorkerResult> {
209        match &self.outcome {
210            CustomerOutcome::Completed { worker_result, .. } => Some(worker_result),
211            CustomerOutcome::Accepted { .. }
212            | CustomerOutcome::Rejected { .. }
213            | CustomerOutcome::ReturnedQueued { .. }
214            | CustomerOutcome::ReturnedAssigned { .. } => None,
215        }
216    }
217
218    /// Consume an accepted-admission outcome.
219    pub fn into_accepted(self) -> core::result::Result<(SubmissionId, JobId), Self> {
220        match self.outcome {
221            CustomerOutcome::Accepted { submission, job } => Ok((submission, job)),
222            outcome => Err(Self { outcome }),
223        }
224    }
225
226    /// Consume a rejected-admission outcome and recover its complete payload.
227    pub fn into_rejected(
228        self,
229    ) -> core::result::Result<(SubmissionId, Job, AdmissionRejection), Self> {
230        match self.outcome {
231            CustomerOutcome::Rejected {
232                submission,
233                payload,
234                reason,
235            } => Ok((submission, payload, reason)),
236            outcome => Err(Self { outcome }),
237        }
238    }
239
240    /// Consume a completed outcome after borrowing its role if needed.
241    pub fn into_completed(self) -> core::result::Result<(JobId, WorkerResult), Self> {
242        match self.outcome {
243            CustomerOutcome::Completed {
244                job, worker_result, ..
245            } => Ok((job, worker_result)),
246            outcome => Err(Self { outcome }),
247        }
248    }
249
250    /// Consume a returned queued outcome and recover its complete payload.
251    pub fn into_returned_queued(
252        self,
253    ) -> core::result::Result<(JobId, Job, QueuedReturnReason), Self> {
254        match self.outcome {
255            CustomerOutcome::ReturnedQueued {
256                job,
257                payload,
258                reason,
259            } => Ok((job, payload, reason)),
260            outcome => Err(Self { outcome }),
261        }
262    }
263
264    /// Consume a returned assigned outcome after borrowing its role if needed.
265    pub fn into_returned_assigned(
266        self,
267    ) -> core::result::Result<(JobId, Job, AssignedReturnReason), Self> {
268        match self.outcome {
269            CustomerOutcome::ReturnedAssigned {
270                job,
271                payload,
272                reason,
273                ..
274            } => Ok((job, payload, reason)),
275            outcome => Err(Self { outcome }),
276        }
277    }
278}
279
280/// Application commands accepted by one FIFO pool.
281pub enum FifoCommand<A, Role, Job, WorkerResult>
282where
283    A: Address + EndpointAddress,
284{
285    /// Submit one payload and its customer correlation and reply capability.
286    Submit {
287        /// Customer-authored admission correlation.
288        submission: SubmissionId,
289        /// Payload retained by the pool while accepted.
290        payload: Job,
291        /// Temporary typed customer capability.
292        customer: ReplyRoute<MessageProtocol<A, FifoOutcome<Role, Job, WorkerResult>>>,
293    },
294    /// Close admission and retire the complete owned actor graph.
295    Shutdown,
296}
297
298impl<A, Role, Job, WorkerResult> FifoCommand<A, Role, Job, WorkerResult>
299where
300    A: Address + EndpointAddress,
301{
302    /// Construct one complete submission without exposing pool internals.
303    #[must_use]
304    pub fn submit<Route>(submission: SubmissionId, payload: Job, customer: Route) -> Self
305    where
306        Route: Into<ReplyRoute<MessageProtocol<A, FifoOutcome<Role, Job, WorkerResult>>>>,
307    {
308        Self::Submit {
309            submission,
310            payload,
311            customer: customer.into(),
312        }
313    }
314
315    /// Construct shutdown without a reply placeholder.
316    #[must_use]
317    pub const fn shutdown() -> Self {
318        Self::Shutdown
319    }
320}
321
322#[expect(
323    dead_code,
324    reason = "diagnostic disposition transfers complete affine values to their next owner"
325)]
326pub(super) enum FifoDiagnosticCause<Role, W, P, Source, Job, WorkerResult>
327where
328    Role: Send + Sync,
329    W: behavior::Behavior + Send,
330    W::Protocol: behavior::Protocol<Msg = Assignment<Job>>,
331    P: ActivationPlan,
332    Source: WorkerSource<Role, W, P>,
333    BehaviorAddr<W>: EndpointAddress,
334    <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
335    Job: Send,
336{
337    Unexpected(
338        FifoEvent<
339            Role,
340            W,
341            P,
342            Job,
343            WorkerResult,
344            behavior::ActionItemResult<PrepareWorkers<Source, Role, W, P>>,
345            WorkerPreparation<Source, Role, W, P>,
346        >,
347    ),
348    WorkerReturned {
349        role: RoleName<Role>,
350        rejection: WorkerCreationRejection<W>,
351        activation: P,
352        stopped: Option<crate::ChildStopped<BehaviorAddr<W>>>,
353    },
354    WorkersReturned(behavior::CreationsSettled<BehaviorAddr<W>, crate::StopOnShutdown<W>>),
355    UnusedActivation {
356        role: RoleName<Role>,
357        permit: Option<ActivationPermit<W>>,
358        activation: P,
359    },
360    StoppedWorkerEstablished {
361        role: RoleName<Role>,
362        worker: CurrentWorker<W>,
363        activation: P,
364        stopped: crate::ChildStopped<BehaviorAddr<W>>,
365    },
366    UnmatchedWorkerReturn {
367        id: behavior::CreationId,
368        kind: behavior::CreationKind,
369        rejection: WorkerCreationRejection<W>,
370    },
371    WorkerActivationReturned {
372        role: RoleName<Role>,
373        worker: WorkerAttempt,
374        activation: WorkerActivationGrant<W, P>,
375        outcome: WorkerActivationOutcome<W, P>,
376    },
377    WorkerShutdownRejected {
378        role: RoleName<Role>,
379        shutdown: crate::EstablishedShutdownResolved<W::Protocol>,
380    },
381    WorkerPreparationFailed {
382        role: RoleName<Role>,
383        previous: super::super::WorkerAttempt,
384        stopped: crate::ChildStopped<BehaviorAddr<W>>,
385        returned_source: Option<Source>,
386        error: WorkerPreparationError<Source::WorkerRejection, Source::SourceRejection>,
387    },
388    WorkerReplacementFailed {
389        role: RoleName<Role>,
390        previous: super::super::WorkerAttempt,
391        stopped: crate::ChildStopped<BehaviorAddr<W>>,
392        submission: WorkerSubmission<W, P>,
393        returned_source: Option<Source>,
394        error: WorkerReplacementError,
395    },
396}
397
398/// Complete FIFO input or precise aggregate failure selected for diagnostics.
399///
400/// The value is opaque so runtime correlation and private role-sharing evidence
401/// cannot become application construction syntax. The selected diagnostic route
402/// or Bombay's terminal custodian receives the complete owned value. Applications
403/// can borrow the associated semantic role when one exists, or consume an original
404/// source rejection with [`Self::into_source_rejection`]. That projection retains
405/// shared role ownership and returns every other complete diagnostic unchanged.
406pub struct FifoDiagnostic<Role, W, P, Source, Job, WorkerResult>
407where
408    Role: Send + Sync,
409    W: behavior::Behavior + Send,
410    W::Protocol: behavior::Protocol<Msg = Assignment<Job>>,
411    P: ActivationPlan,
412    Source: WorkerSource<Role, W, P>,
413    BehaviorAddr<W>: EndpointAddress,
414    <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
415    Job: Send,
416{
417    cause: FifoDiagnosticCause<Role, W, P, Source, Job, WorkerResult>,
418}
419
420impl<Role, W, P, Source, Job, WorkerResult> FifoDiagnostic<Role, W, P, Source, Job, WorkerResult>
421where
422    Role: Send + Sync,
423    W: behavior::Behavior + Send,
424    W::Protocol: behavior::Protocol<Msg = Assignment<Job>>,
425    P: ActivationPlan,
426    Source: WorkerSource<Role, W, P>,
427    BehaviorAddr<W>: EndpointAddress,
428    <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
429    Job: Send,
430{
431    pub(super) const fn new(
432        cause: FifoDiagnosticCause<Role, W, P, Source, Job, WorkerResult>,
433    ) -> Self {
434        Self { cause }
435    }
436
437    /// Borrow the semantic role associated with a role-specific failure.
438    #[must_use]
439    pub fn role(&self) -> Option<&Role> {
440        match &self.cause {
441            FifoDiagnosticCause::WorkerReturned { role, .. }
442            | FifoDiagnosticCause::UnusedActivation { role, .. }
443            | FifoDiagnosticCause::StoppedWorkerEstablished { role, .. }
444            | FifoDiagnosticCause::WorkerActivationReturned { role, .. }
445            | FifoDiagnosticCause::WorkerShutdownRejected { role, .. }
446            | FifoDiagnosticCause::WorkerPreparationFailed { role, .. }
447            | FifoDiagnosticCause::WorkerReplacementFailed { role, .. } => Some(role.role()),
448            FifoDiagnosticCause::Unexpected(_)
449            | FifoDiagnosticCause::WorkersReturned(_)
450            | FifoDiagnosticCause::UnmatchedWorkerReturn { .. } => None,
451        }
452    }
453
454    /// Consume the complete source-rejection diagnostic, returning every other
455    /// diagnostic unchanged. The role retains its original shared allocation;
456    /// this operation does not clone a role or alter the pool's recovery policy.
457    /// The returned source is present only when the existing diagnostic owned it.
458    pub fn into_source_rejection(
459        self,
460    ) -> Result<
461        (
462            Arc<Role>,
463            WorkerAttempt,
464            ChildStopped<BehaviorAddr<W>>,
465            Option<Source>,
466            Source::SourceRejection,
467        ),
468        Self,
469    > {
470        match self.cause {
471            FifoDiagnosticCause::WorkerPreparationFailed {
472                role,
473                previous,
474                stopped,
475                returned_source,
476                error: WorkerPreparationError::SourceRejected(reason),
477            } => Ok((role.into_role(), previous, stopped, returned_source, reason)),
478            cause => Err(Self { cause }),
479        }
480    }
481}
482
483impl<Role, W, P, Source, Job, WorkerResult> core::fmt::Debug
484    for FifoDiagnostic<Role, W, P, Source, Job, WorkerResult>
485where
486    Role: Send + Sync,
487    W: behavior::Behavior + Send,
488    W::Protocol: behavior::Protocol<Msg = Assignment<Job>>,
489    P: ActivationPlan,
490    Source: WorkerSource<Role, W, P>,
491    BehaviorAddr<W>: EndpointAddress,
492    <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
493    Job: Send,
494{
495    fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
496        formatter
497            .debug_struct("FifoDiagnostic")
498            .finish_non_exhaustive()
499    }
500}
501
502#[cfg(test)]
503mod live_activation_correlation {
504    use super::super::requests::FifoRequests;
505    use super::super::{FifoEvent, FifoOperating, FifoPool, PoolState};
506    use super::{FifoDiagnostic, FifoDiagnosticCause};
507    use crate::atomic::pool::ShutdownSequence;
508    use crate::atomic::pool::assignment::{AcceptedJobSequence, AssignmentSequence};
509    use crate::atomic::pool::worker::activation_pool_correlation::{
510        ActivationEndpoint, ActivationWorker, assert_current_worker, initialized_activation,
511    };
512    use crate::atomic::pool::worker::{CurrentWorker, Member, MemberState, Worker, WorkerPhase};
513    use crate::atomic::restart::RecoveryCount;
514    use crate::atomic::restart::RestartBudget;
515    use crate::atomic::worker::WorkerActivationOutcome;
516    use crate::atomic::{
517        ActivationPolicy, ActorDrainPolicy, BacklogCapacity, ImmediateActivation, Interruption,
518        PoolFailureReaction, PoolRecovery, RoleName, WorkerAttempt,
519    };
520    use crate::{
521        Active, ChildStopped, DiagnosticAction, DiagnosticDisposition, Exit, StopOnShutdown,
522    };
523    use behavior::{Actions, CreationSequence, EstablishedActor, Never, Step};
524    use core::convert::Infallible;
525    use std::collections::BTreeMap;
526    use std::sync::Arc;
527    use std::time::Instant;
528
529    #[tokio::test]
530    async fn fifo_dispatched_live_returns_the_original_foreign_activation_twice() {
531        let mut creations = CreationSequence::new();
532        let creation = creations.issue().expect("one original worker creation");
533        let worker = WorkerAttempt::issued(creation);
534        let actor =
535            EstablishedActor::<StopOnShutdown<ActivationWorker>>::issued(ActivationEndpoint(19));
536        let expected_actor = actor.clone();
537        let original = initialized_activation(worker.clone(), actor.recipient());
538        let foreign = initialized_activation(worker.clone(), actor.recipient());
539        let original_attempt = original.attempt();
540        let foreign_attempt = foreign.attempt();
541        assert!(original_attempt != foreign_attempt);
542        let inputs = [
543            (
544                foreign.started(),
545                WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Started,
546            ),
547            (
548                foreign.activate().await,
549                WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Ready(()),
550            ),
551        ];
552        for (input, report) in inputs {
553            let role = RoleName::new(23_u64);
554            let expected_role = role.clone().into_role();
555            let member = Member {
556                role,
557                recoveries: RecoveryCount::default(),
558                state: MemberState::Worker(Worker {
559                    current: CurrentWorker::new(worker.clone(), expected_actor.clone()),
560                    phase: WorkerPhase::ActivationDispatched {
561                        attempt: original_attempt.clone(),
562                        stopped: None,
563                    },
564                }),
565            };
566            let pool: FifoPool<
567                u64,
568                ActivationWorker,
569                ImmediateActivation,
570                Never,
571                Infallible,
572                u8,
573                u16,
574            > = FifoPool {
575                state: PoolState::Operating(FifoOperating {
576                    members: vec![member],
577                    backlog: BTreeMap::new(),
578                    cursor: 0,
579                }),
580                activation: ActivationPolicy::new(1).expect("one actual activation slot"),
581                recovery: PoolRecovery::<Never>::temporary(PoolFailureReaction::RetireRole).into(),
582                restarts: RestartBudget::empty(),
583                backlog: BacklogCapacity::new(2),
584                interruption: Interruption::Fail,
585                actor_drain: ActorDrainPolicy::WaitForActorGraph,
586                diagnostics: DiagnosticDisposition::terminate(),
587                creations: CreationSequence::new(),
588                jobs: AcceptedJobSequence::new(),
589                assignments: AssignmentSequence::new(),
590                shutdowns: ShutdownSequence::new(),
591                next_restart_timer: 1,
592            };
593            let mut pool = Active { behavior: pool };
594            let mut returned = input;
595            for _ in 0..2 {
596                let actions = pool
597                    .transition(FifoEvent::WorkerActivationReported(returned))
598                    .unwrap_or_else(|error| {
599                        panic!("the activation fold returns its real Actions: {error}")
600                    });
601                let Actions {
602                    sends,
603                    creates,
604                    become_,
605                } = actions;
606                let FifoRequests {
607                    worker_observations,
608                    worker_initializations,
609                    worker_activations,
610                    customer_outcomes,
611                    worker_assignments,
612                    worker_preparations,
613                    restart_schedules,
614                    worker_shutdowns,
615                    diagnostics,
616                } = sends;
617                assert!(worker_observations.is_empty());
618                assert!(worker_initializations.is_empty());
619                assert!(worker_activations.is_empty());
620                assert!(customer_outcomes.as_slice().is_empty());
621                assert!(worker_assignments.is_empty());
622                assert!(worker_preparations.is_empty());
623                assert!(restart_schedules.is_empty());
624                assert!(worker_shutdowns.is_empty());
625                assert!(creates.is_empty());
626                assert!(matches!(become_, Step::Continue));
627                let mut diagnostics = diagnostics.into_requests();
628                assert_eq!(diagnostics.len(), 1);
629                let DiagnosticAction::Terminal {
630                    diagnostic:
631                        FifoDiagnostic {
632                            cause:
633                                FifoDiagnosticCause::Unexpected(FifoEvent::WorkerActivationReported(
634                                    input,
635                                )),
636                        },
637                } = diagnostics.remove(0)
638                else {
639                    panic!(
640                        "a foreign attempt returns the complete original activation through the existing Unexpected policy"
641                    )
642                };
643                assert_eq!(input.worker(), worker);
644                assert!(input.attempt() == &foreign_attempt);
645                returned = input;
646                let PoolState::Operating(operating) = &pool.state else {
647                    panic!("foreign activation must not retire the original pool")
648                };
649                assert_eq!(operating.members.len(), 1);
650                assert!(operating.backlog.is_empty());
651                assert_eq!(operating.cursor, 0);
652                let member = &operating.members[0];
653                let role = member.role.clone().into_role();
654                assert!(Arc::ptr_eq(&role, &expected_role));
655                assert_eq!(*role, 23);
656                assert_eq!(member.recoveries, RecoveryCount::default());
657                let MemberState::Worker(Worker {
658                    current,
659                    phase: WorkerPhase::ActivationDispatched { attempt, stopped },
660                }) = &member.state
661                else {
662                    panic!("foreign report preserves the exact original worker phase")
663                };
664                assert_current_worker(current, &worker, &expected_actor);
665                assert!(attempt == &original_attempt);
666                assert!(stopped.is_none());
667            }
668            let (returned_worker, returned_attempt, outcome) = returned.into_parts();
669            assert_eq!(returned_worker, worker);
670            assert!(returned_attempt == foreign_attempt);
671            match (report, outcome) {
672                (WorkerActivationOutcome::Started, WorkerActivationOutcome::Started)
673                | (WorkerActivationOutcome::Ready(()), WorkerActivationOutcome::Ready(())) => {}
674                _ => panic!("the original complete activation outcome is preserved"),
675            }
676        }
677        drop(original);
678    }
679
680    #[tokio::test]
681    async fn fifo_dispatched_stopped_returns_the_original_foreign_activation_twice() {
682        let mut creations = CreationSequence::new();
683        let creation = creations.issue().expect("one original worker creation");
684        let worker = WorkerAttempt::issued(creation);
685        let actor =
686            EstablishedActor::<StopOnShutdown<ActivationWorker>>::issued(ActivationEndpoint(19));
687        let expected_actor = actor.clone();
688        let original = initialized_activation(worker.clone(), actor.recipient());
689        let foreign = initialized_activation(worker.clone(), actor.recipient());
690        let original_attempt = original.attempt();
691        let foreign_attempt = foreign.attempt();
692        assert!(original_attempt != foreign_attempt);
693        let stopped_at = Instant::now();
694        let inputs = [
695            (
696                foreign.started(),
697                WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Started,
698            ),
699            (
700                foreign.activate().await,
701                WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Ready(()),
702            ),
703        ];
704        for (input, report) in inputs {
705            let role = RoleName::new(23_u64);
706            let expected_role = role.clone().into_role();
707            let member = Member {
708                role,
709                recoveries: RecoveryCount::default(),
710                state: MemberState::Worker(Worker {
711                    current: CurrentWorker::new(worker.clone(), expected_actor.clone()),
712                    phase: WorkerPhase::ActivationDispatched {
713                        attempt: original_attempt.clone(),
714                        stopped: Some(ChildStopped::new(creation, Ok(Exit::Normal), stopped_at)),
715                    },
716                }),
717            };
718            let pool: FifoPool<
719                u64,
720                ActivationWorker,
721                ImmediateActivation,
722                Never,
723                Infallible,
724                u8,
725                u16,
726            > = FifoPool {
727                state: PoolState::Operating(FifoOperating {
728                    members: vec![member],
729                    backlog: BTreeMap::new(),
730                    cursor: 0,
731                }),
732                activation: ActivationPolicy::new(1).expect("one actual activation slot"),
733                recovery: PoolRecovery::<Never>::temporary(PoolFailureReaction::RetireRole).into(),
734                restarts: RestartBudget::empty(),
735                backlog: BacklogCapacity::new(2),
736                interruption: Interruption::Fail,
737                actor_drain: ActorDrainPolicy::WaitForActorGraph,
738                diagnostics: DiagnosticDisposition::terminate(),
739                creations: CreationSequence::new(),
740                jobs: AcceptedJobSequence::new(),
741                assignments: AssignmentSequence::new(),
742                shutdowns: ShutdownSequence::new(),
743                next_restart_timer: 1,
744            };
745            let mut pool = Active { behavior: pool };
746            let mut returned = input;
747            for _ in 0..2 {
748                let actions = pool
749                    .transition(FifoEvent::WorkerActivationReported(returned))
750                    .unwrap_or_else(|error| {
751                        panic!("the activation fold returns its real Actions: {error}")
752                    });
753                let Actions {
754                    sends,
755                    creates,
756                    become_,
757                } = actions;
758                let FifoRequests {
759                    worker_observations,
760                    worker_initializations,
761                    worker_activations,
762                    customer_outcomes,
763                    worker_assignments,
764                    worker_preparations,
765                    restart_schedules,
766                    worker_shutdowns,
767                    diagnostics,
768                } = sends;
769                assert!(worker_observations.is_empty());
770                assert!(worker_initializations.is_empty());
771                assert!(worker_activations.is_empty());
772                assert!(customer_outcomes.as_slice().is_empty());
773                assert!(worker_assignments.is_empty());
774                assert!(worker_preparations.is_empty());
775                assert!(restart_schedules.is_empty());
776                assert!(worker_shutdowns.is_empty());
777                assert!(creates.is_empty());
778                assert!(matches!(become_, Step::Continue));
779                let mut diagnostics = diagnostics.into_requests();
780                assert_eq!(diagnostics.len(), 1);
781                let DiagnosticAction::Terminal {
782                    diagnostic:
783                        FifoDiagnostic {
784                            cause:
785                                FifoDiagnosticCause::Unexpected(FifoEvent::WorkerActivationReported(
786                                    input,
787                                )),
788                        },
789                } = diagnostics.remove(0)
790                else {
791                    panic!(
792                        "a foreign attempt returns the complete original activation through the existing Unexpected policy"
793                    )
794                };
795                assert_eq!(input.worker(), worker);
796                assert!(input.attempt() == &foreign_attempt);
797                returned = input;
798                let PoolState::Operating(operating) = &pool.state else {
799                    panic!("foreign activation must not retire the original pool")
800                };
801                assert_eq!(operating.members.len(), 1);
802                assert!(operating.backlog.is_empty());
803                assert_eq!(operating.cursor, 0);
804                let member = &operating.members[0];
805                let role = member.role.clone().into_role();
806                assert!(Arc::ptr_eq(&role, &expected_role));
807                assert_eq!(*role, 23);
808                assert_eq!(member.recoveries, RecoveryCount::default());
809                let MemberState::Worker(Worker {
810                    current,
811                    phase: WorkerPhase::ActivationDispatched { attempt, stopped },
812                }) = &member.state
813                else {
814                    panic!("foreign report preserves the exact original worker phase")
815                };
816                assert_current_worker(current, &worker, &expected_actor);
817                assert!(attempt == &original_attempt);
818                let Some(stopped) = stopped else {
819                    panic!("the original stopped child remains owned")
820                };
821                assert_eq!(stopped.child, creation);
822                assert!(matches!(&stopped.outcome, Ok(Exit::Normal)));
823                assert_eq!(stopped.at, stopped_at);
824            }
825            let (returned_worker, returned_attempt, outcome) = returned.into_parts();
826            assert_eq!(returned_worker, worker);
827            assert!(returned_attempt == foreign_attempt);
828            match (report, outcome) {
829                (WorkerActivationOutcome::Started, WorkerActivationOutcome::Started)
830                | (WorkerActivationOutcome::Ready(()), WorkerActivationOutcome::Ready(())) => {}
831                _ => panic!("the original complete activation outcome is preserved"),
832            }
833        }
834        drop(original);
835    }
836
837    #[tokio::test]
838    async fn fifo_activating_live_returns_the_original_foreign_activation_twice() {
839        let mut creations = CreationSequence::new();
840        let creation = creations.issue().expect("one original worker creation");
841        let worker = WorkerAttempt::issued(creation);
842        let actor =
843            EstablishedActor::<StopOnShutdown<ActivationWorker>>::issued(ActivationEndpoint(19));
844        let expected_actor = actor.clone();
845        let original = initialized_activation(worker.clone(), actor.recipient());
846        let foreign = initialized_activation(worker.clone(), actor.recipient());
847        let original_attempt = original.attempt();
848        let foreign_attempt = foreign.attempt();
849        assert!(original_attempt != foreign_attempt);
850        let inputs = [
851            (
852                foreign.started(),
853                WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Started,
854            ),
855            (
856                foreign.activate().await,
857                WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Ready(()),
858            ),
859        ];
860        for (input, report) in inputs {
861            let role = RoleName::new(23_u64);
862            let expected_role = role.clone().into_role();
863            let member = Member {
864                role,
865                recoveries: RecoveryCount::default(),
866                state: MemberState::Worker(Worker {
867                    current: CurrentWorker::new(worker.clone(), expected_actor.clone()),
868                    phase: WorkerPhase::Activating {
869                        attempt: original_attempt.clone(),
870                        stopped: None,
871                    },
872                }),
873            };
874            let pool: FifoPool<
875                u64,
876                ActivationWorker,
877                ImmediateActivation,
878                Never,
879                Infallible,
880                u8,
881                u16,
882            > = FifoPool {
883                state: PoolState::Operating(FifoOperating {
884                    members: vec![member],
885                    backlog: BTreeMap::new(),
886                    cursor: 0,
887                }),
888                activation: ActivationPolicy::new(1).expect("one actual activation slot"),
889                recovery: PoolRecovery::<Never>::temporary(PoolFailureReaction::RetireRole).into(),
890                restarts: RestartBudget::empty(),
891                backlog: BacklogCapacity::new(2),
892                interruption: Interruption::Fail,
893                actor_drain: ActorDrainPolicy::WaitForActorGraph,
894                diagnostics: DiagnosticDisposition::terminate(),
895                creations: CreationSequence::new(),
896                jobs: AcceptedJobSequence::new(),
897                assignments: AssignmentSequence::new(),
898                shutdowns: ShutdownSequence::new(),
899                next_restart_timer: 1,
900            };
901            let mut pool = Active { behavior: pool };
902            let mut returned = input;
903            for _ in 0..2 {
904                let actions = pool
905                    .transition(FifoEvent::WorkerActivationReported(returned))
906                    .unwrap_or_else(|error| {
907                        panic!("the activation fold returns its real Actions: {error}")
908                    });
909                let Actions {
910                    sends,
911                    creates,
912                    become_,
913                } = actions;
914                let FifoRequests {
915                    worker_observations,
916                    worker_initializations,
917                    worker_activations,
918                    customer_outcomes,
919                    worker_assignments,
920                    worker_preparations,
921                    restart_schedules,
922                    worker_shutdowns,
923                    diagnostics,
924                } = sends;
925                assert!(worker_observations.is_empty());
926                assert!(worker_initializations.is_empty());
927                assert!(worker_activations.is_empty());
928                assert!(customer_outcomes.as_slice().is_empty());
929                assert!(worker_assignments.is_empty());
930                assert!(worker_preparations.is_empty());
931                assert!(restart_schedules.is_empty());
932                assert!(worker_shutdowns.is_empty());
933                assert!(creates.is_empty());
934                assert!(matches!(become_, Step::Continue));
935                let mut diagnostics = diagnostics.into_requests();
936                assert_eq!(diagnostics.len(), 1);
937                let DiagnosticAction::Terminal {
938                    diagnostic:
939                        FifoDiagnostic {
940                            cause:
941                                FifoDiagnosticCause::Unexpected(FifoEvent::WorkerActivationReported(
942                                    input,
943                                )),
944                        },
945                } = diagnostics.remove(0)
946                else {
947                    panic!(
948                        "a foreign attempt returns the complete original activation through the existing Unexpected policy"
949                    )
950                };
951                assert_eq!(input.worker(), worker);
952                assert!(input.attempt() == &foreign_attempt);
953                returned = input;
954                let PoolState::Operating(operating) = &pool.state else {
955                    panic!("foreign activation must not retire the original pool")
956                };
957                assert_eq!(operating.members.len(), 1);
958                assert!(operating.backlog.is_empty());
959                assert_eq!(operating.cursor, 0);
960                let member = &operating.members[0];
961                let role = member.role.clone().into_role();
962                assert!(Arc::ptr_eq(&role, &expected_role));
963                assert_eq!(*role, 23);
964                assert_eq!(member.recoveries, RecoveryCount::default());
965                let MemberState::Worker(Worker {
966                    current,
967                    phase: WorkerPhase::Activating { attempt, stopped },
968                }) = &member.state
969                else {
970                    panic!("foreign report preserves the exact original worker phase")
971                };
972                assert_current_worker(current, &worker, &expected_actor);
973                assert!(attempt == &original_attempt);
974                assert!(stopped.is_none());
975            }
976            let (returned_worker, returned_attempt, outcome) = returned.into_parts();
977            assert_eq!(returned_worker, worker);
978            assert!(returned_attempt == foreign_attempt);
979            match (report, outcome) {
980                (WorkerActivationOutcome::Started, WorkerActivationOutcome::Started)
981                | (WorkerActivationOutcome::Ready(()), WorkerActivationOutcome::Ready(())) => {}
982                _ => panic!("the original complete activation outcome is preserved"),
983            }
984        }
985        drop(original);
986    }
987
988    #[tokio::test]
989    async fn fifo_activating_stopped_returns_the_original_foreign_activation_twice() {
990        let mut creations = CreationSequence::new();
991        let creation = creations.issue().expect("one original worker creation");
992        let worker = WorkerAttempt::issued(creation);
993        let actor =
994            EstablishedActor::<StopOnShutdown<ActivationWorker>>::issued(ActivationEndpoint(19));
995        let expected_actor = actor.clone();
996        let original = initialized_activation(worker.clone(), actor.recipient());
997        let foreign = initialized_activation(worker.clone(), actor.recipient());
998        let original_attempt = original.attempt();
999        let foreign_attempt = foreign.attempt();
1000        assert!(original_attempt != foreign_attempt);
1001        let stopped_at = Instant::now();
1002        let inputs = [
1003            (
1004                foreign.started(),
1005                WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Started,
1006            ),
1007            (
1008                foreign.activate().await,
1009                WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Ready(()),
1010            ),
1011        ];
1012        for (input, report) in inputs {
1013            let role = RoleName::new(23_u64);
1014            let expected_role = role.clone().into_role();
1015            let member = Member {
1016                role,
1017                recoveries: RecoveryCount::default(),
1018                state: MemberState::Worker(Worker {
1019                    current: CurrentWorker::new(worker.clone(), expected_actor.clone()),
1020                    phase: WorkerPhase::Activating {
1021                        attempt: original_attempt.clone(),
1022                        stopped: Some(ChildStopped::new(creation, Ok(Exit::Normal), stopped_at)),
1023                    },
1024                }),
1025            };
1026            let pool: FifoPool<
1027                u64,
1028                ActivationWorker,
1029                ImmediateActivation,
1030                Never,
1031                Infallible,
1032                u8,
1033                u16,
1034            > = FifoPool {
1035                state: PoolState::Operating(FifoOperating {
1036                    members: vec![member],
1037                    backlog: BTreeMap::new(),
1038                    cursor: 0,
1039                }),
1040                activation: ActivationPolicy::new(1).expect("one actual activation slot"),
1041                recovery: PoolRecovery::<Never>::temporary(PoolFailureReaction::RetireRole).into(),
1042                restarts: RestartBudget::empty(),
1043                backlog: BacklogCapacity::new(2),
1044                interruption: Interruption::Fail,
1045                actor_drain: ActorDrainPolicy::WaitForActorGraph,
1046                diagnostics: DiagnosticDisposition::terminate(),
1047                creations: CreationSequence::new(),
1048                jobs: AcceptedJobSequence::new(),
1049                assignments: AssignmentSequence::new(),
1050                shutdowns: ShutdownSequence::new(),
1051                next_restart_timer: 1,
1052            };
1053            let mut pool = Active { behavior: pool };
1054            let mut returned = input;
1055            for _ in 0..2 {
1056                let actions = pool
1057                    .transition(FifoEvent::WorkerActivationReported(returned))
1058                    .unwrap_or_else(|error| {
1059                        panic!("the activation fold returns its real Actions: {error}")
1060                    });
1061                let Actions {
1062                    sends,
1063                    creates,
1064                    become_,
1065                } = actions;
1066                let FifoRequests {
1067                    worker_observations,
1068                    worker_initializations,
1069                    worker_activations,
1070                    customer_outcomes,
1071                    worker_assignments,
1072                    worker_preparations,
1073                    restart_schedules,
1074                    worker_shutdowns,
1075                    diagnostics,
1076                } = sends;
1077                assert!(worker_observations.is_empty());
1078                assert!(worker_initializations.is_empty());
1079                assert!(worker_activations.is_empty());
1080                assert!(customer_outcomes.as_slice().is_empty());
1081                assert!(worker_assignments.is_empty());
1082                assert!(worker_preparations.is_empty());
1083                assert!(restart_schedules.is_empty());
1084                assert!(worker_shutdowns.is_empty());
1085                assert!(creates.is_empty());
1086                assert!(matches!(become_, Step::Continue));
1087                let mut diagnostics = diagnostics.into_requests();
1088                assert_eq!(diagnostics.len(), 1);
1089                let DiagnosticAction::Terminal {
1090                    diagnostic:
1091                        FifoDiagnostic {
1092                            cause:
1093                                FifoDiagnosticCause::Unexpected(FifoEvent::WorkerActivationReported(
1094                                    input,
1095                                )),
1096                        },
1097                } = diagnostics.remove(0)
1098                else {
1099                    panic!(
1100                        "a foreign attempt returns the complete original activation through the existing Unexpected policy"
1101                    )
1102                };
1103                assert_eq!(input.worker(), worker);
1104                assert!(input.attempt() == &foreign_attempt);
1105                returned = input;
1106                let PoolState::Operating(operating) = &pool.state else {
1107                    panic!("foreign activation must not retire the original pool")
1108                };
1109                assert_eq!(operating.members.len(), 1);
1110                assert!(operating.backlog.is_empty());
1111                assert_eq!(operating.cursor, 0);
1112                let member = &operating.members[0];
1113                let role = member.role.clone().into_role();
1114                assert!(Arc::ptr_eq(&role, &expected_role));
1115                assert_eq!(*role, 23);
1116                assert_eq!(member.recoveries, RecoveryCount::default());
1117                let MemberState::Worker(Worker {
1118                    current,
1119                    phase: WorkerPhase::Activating { attempt, stopped },
1120                }) = &member.state
1121                else {
1122                    panic!("foreign report preserves the exact original worker phase")
1123                };
1124                assert_current_worker(current, &worker, &expected_actor);
1125                assert!(attempt == &original_attempt);
1126                let Some(stopped) = stopped else {
1127                    panic!("the original stopped child remains owned")
1128                };
1129                assert_eq!(stopped.child, creation);
1130                assert!(matches!(&stopped.outcome, Ok(Exit::Normal)));
1131                assert_eq!(stopped.at, stopped_at);
1132            }
1133            let (returned_worker, returned_attempt, outcome) = returned.into_parts();
1134            assert_eq!(returned_worker, worker);
1135            assert!(returned_attempt == foreign_attempt);
1136            match (report, outcome) {
1137                (WorkerActivationOutcome::Started, WorkerActivationOutcome::Started)
1138                | (WorkerActivationOutcome::Ready(()), WorkerActivationOutcome::Ready(())) => {}
1139                _ => panic!("the original complete activation outcome is preserved"),
1140            }
1141        }
1142        drop(original);
1143    }
1144
1145    #[tokio::test]
1146    async fn fifo_original_activation_started_then_ready_preserves_all_lanes() {
1147        let mut creations = CreationSequence::new();
1148        let creation = creations.issue().expect("one original worker creation");
1149        let worker = WorkerAttempt::issued(creation);
1150        let actor =
1151            EstablishedActor::<StopOnShutdown<ActivationWorker>>::issued(ActivationEndpoint(19));
1152        let expected_actor = actor.clone();
1153        let original = initialized_activation(worker.clone(), actor.recipient());
1154        let original_attempt = original.attempt();
1155        let started = original.started();
1156        let role = RoleName::new(23_u64);
1157        let expected_role = role.clone().into_role();
1158        let member = Member {
1159            role,
1160            recoveries: RecoveryCount::default(),
1161            state: MemberState::Worker(Worker {
1162                current: CurrentWorker::new(worker.clone(), expected_actor.clone()),
1163                phase: WorkerPhase::ActivationDispatched {
1164                    attempt: original_attempt.clone(),
1165                    stopped: None,
1166                },
1167            }),
1168        };
1169        let pool: FifoPool<u64, ActivationWorker, ImmediateActivation, Never, Infallible, u8, u16> =
1170            FifoPool {
1171                state: PoolState::Operating(FifoOperating {
1172                    members: vec![member],
1173                    backlog: BTreeMap::new(),
1174                    cursor: 0,
1175                }),
1176                activation: ActivationPolicy::new(1).expect("one actual activation slot"),
1177                recovery: PoolRecovery::<Never>::temporary(PoolFailureReaction::RetireRole).into(),
1178                restarts: RestartBudget::empty(),
1179                backlog: BacklogCapacity::new(2),
1180                interruption: Interruption::Fail,
1181                actor_drain: ActorDrainPolicy::WaitForActorGraph,
1182                diagnostics: DiagnosticDisposition::terminate(),
1183                creations: CreationSequence::new(),
1184                jobs: AcceptedJobSequence::new(),
1185                assignments: AssignmentSequence::new(),
1186                shutdowns: ShutdownSequence::new(),
1187                next_restart_timer: 1,
1188            };
1189        let mut pool = Active { behavior: pool };
1190
1191        let inputs = [
1192            (
1193                started,
1194                WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Started,
1195            ),
1196            (
1197                original.activate().await,
1198                WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Ready(()),
1199            ),
1200        ];
1201        for (input, report) in inputs {
1202            let actions = pool
1203                .transition(FifoEvent::WorkerActivationReported(input))
1204                .unwrap_or_else(|error| {
1205                    panic!("the original activation returns complete Actions: {error}")
1206                });
1207            let Actions {
1208                sends,
1209                creates,
1210                become_,
1211            } = actions;
1212            let FifoRequests {
1213                worker_observations,
1214                worker_initializations,
1215                worker_activations,
1216                customer_outcomes,
1217                worker_assignments,
1218                worker_preparations,
1219                restart_schedules,
1220                worker_shutdowns,
1221                diagnostics,
1222            } = sends;
1223            assert!(worker_observations.is_empty());
1224            assert!(worker_initializations.is_empty());
1225            assert!(worker_activations.is_empty());
1226            assert!(customer_outcomes.as_slice().is_empty());
1227            assert!(worker_assignments.is_empty());
1228            assert!(worker_preparations.is_empty());
1229            assert!(restart_schedules.is_empty());
1230            assert!(worker_shutdowns.is_empty());
1231            assert!(diagnostics.is_empty());
1232            assert!(creates.is_empty());
1233            assert!(matches!(become_, Step::Continue));
1234            let PoolState::Operating(operating) = &pool.state else {
1235                panic!("original activation leaves the pool operating")
1236            };
1237            assert_eq!(operating.members.len(), 1);
1238            assert!(operating.backlog.is_empty());
1239            assert_eq!(operating.cursor, 0);
1240            let member = &operating.members[0];
1241            let actual_role = member.role.clone().into_role();
1242            assert!(Arc::ptr_eq(&actual_role, &expected_role));
1243            assert_eq!(*actual_role, 23);
1244            assert_eq!(member.recoveries, RecoveryCount::default());
1245            let MemberState::Worker(Worker { current, phase }) = &member.state else {
1246                panic!("the actual original worker stays owned")
1247            };
1248            assert_current_worker(current, &worker, &expected_actor);
1249            match (report, phase) {
1250                (
1251                    WorkerActivationOutcome::Started,
1252                    WorkerPhase::Activating {
1253                        attempt,
1254                        stopped: None,
1255                    },
1256                ) => {
1257                    assert!(attempt == &original_attempt);
1258                }
1259                (WorkerActivationOutcome::Ready(()), WorkerPhase::Idle) => {}
1260                _ => panic!("actual Started enters Activating and actual Ready enters Idle"),
1261            }
1262        }
1263    }
1264}