1use 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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
21pub enum AdmissionRejection {
22 NoRecoverableWorkers,
24 BacklogFull,
26 JobCorrelationUnavailable,
28 AssignmentCorrelationUnavailable,
30 ShuttingDown,
32}
33
34#[derive(Clone, Copy, Debug, Eq, PartialEq)]
36pub enum QueuedReturnReason {
37 NoRecoverableWorkers,
39 AssignmentCorrelationUnavailable,
41 PoolShutdown,
43}
44
45#[derive(Clone, Copy, Debug, Eq, PartialEq)]
47pub enum AssignedReturnReason {
48 WorkerStopped,
50 RetryPreparationRejected,
52 PoolShutdown,
54 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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
88pub enum FifoOutcomeKind {
89 Accepted,
90 Rejected,
91 Completed,
92 ReturnedQueued,
93 ReturnedAssigned,
94}
95
96pub 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 #[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 #[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 #[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 #[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 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 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 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 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 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
280pub enum FifoCommand<A, Role, Job, WorkerResult>
282where
283 A: Address + EndpointAddress,
284{
285 Submit {
287 submission: SubmissionId,
289 payload: Job,
291 customer: ReplyRoute<MessageProtocol<A, FifoOutcome<Role, Job, WorkerResult>>>,
293 },
294 Shutdown,
296}
297
298impl<A, Role, Job, WorkerResult> FifoCommand<A, Role, Job, WorkerResult>
299where
300 A: Address + EndpointAddress,
301{
302 #[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 #[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
398pub 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 #[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 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}