1use behavior::{Address, Behavior, BehaviorAddr, EndpointAddress, MessageProtocol, Protocol};
4
5use crate::{ChildStopped, EstablishedShutdownResolved, ReplyRoute};
6
7use super::super::pool::worker::{
8 CurrentWorker, PendingWorkerReplacement, WorkerPreparationError, WorkerReplacementError,
9};
10use super::super::worker::{WorkerActivationGrant, WorkerActivationOutcome};
11use super::super::{
12 ActivationPermit, ActivationPlan, JobId, PrepareWorkers, RoleName, SubmissionId, WorkerSource,
13};
14use super::super::{WorkerAttempt, WorkerCreationRejection};
15use super::{BindingEvidence, BindingExpectation, BindingRequestId};
16
17pub(super) enum BindingChange<Role> {
18 Rebalance(Role),
19 Unbind,
20}
21
22pub struct BindingCommand<A, Key, Role>
27where
28 A: Address + EndpointAddress,
29{
30 request: BindingRequestId,
31 key: Key,
32 expected: BindingExpectation,
33 change: BindingChange<Role>,
34 reply: ReplyRoute<MessageProtocol<A, BindingReply<A, Key, Role>>>,
35}
36
37impl<A, Key, Role> BindingCommand<A, Key, Role>
38where
39 A: Address + EndpointAddress,
40{
41 pub(super) fn rebalance<Route>(
42 request: BindingRequestId,
43 key: Key,
44 expected: BindingExpectation,
45 target: Role,
46 reply: Route,
47 ) -> Self
48 where
49 Route: Into<ReplyRoute<MessageProtocol<A, BindingReply<A, Key, Role>>>>,
50 {
51 Self {
52 request,
53 key,
54 expected,
55 change: BindingChange::Rebalance(target),
56 reply: reply.into(),
57 }
58 }
59
60 pub(super) fn unbind<Route>(
61 request: BindingRequestId,
62 key: Key,
63 expected: BindingExpectation,
64 reply: Route,
65 ) -> Self
66 where
67 Route: Into<ReplyRoute<MessageProtocol<A, BindingReply<A, Key, Role>>>>,
68 {
69 Self {
70 request,
71 key,
72 expected,
73 change: BindingChange::Unbind,
74 reply: reply.into(),
75 }
76 }
77
78 #[must_use]
80 pub const fn request(&self) -> BindingRequestId {
81 self.request
82 }
83
84 #[must_use]
86 pub const fn key(&self) -> &Key {
87 &self.key
88 }
89
90 #[must_use]
92 pub const fn expectation(&self) -> &BindingExpectation {
93 &self.expected
94 }
95
96 pub(super) fn reply(&self) -> &ReplyRoute<MessageProtocol<A, BindingReply<A, Key, Role>>> {
97 &self.reply
98 }
99
100 pub(super) const fn change(&self) -> &BindingChange<Role> {
101 &self.change
102 }
103
104 pub(super) fn from_parts(
105 request: BindingRequestId,
106 key: Key,
107 expected: BindingExpectation,
108 change: BindingChange<Role>,
109 reply: ReplyRoute<MessageProtocol<A, BindingReply<A, Key, Role>>>,
110 ) -> Self {
111 Self {
112 request,
113 key,
114 expected,
115 change,
116 reply,
117 }
118 }
119
120 pub(super) fn into_parts(
121 self,
122 ) -> (
123 BindingRequestId,
124 Key,
125 BindingExpectation,
126 BindingChange<Role>,
127 ReplyRoute<MessageProtocol<A, BindingReply<A, Key, Role>>>,
128 ) {
129 (
130 self.request,
131 self.key,
132 self.expected,
133 self.change,
134 self.reply,
135 )
136 }
137}
138
139#[derive(Clone, Debug, Eq, PartialEq)]
141pub enum BindingRejection {
142 StaleExpectation {
144 actual: BindingExpectation,
146 },
147 UnknownTarget,
149 TargetUnavailable,
151 BindingCapacityExhausted,
153 GenerationExhausted,
155 ShuttingDown,
157}
158
159pub enum BindingReply<A, Key, Role>
161where
162 A: Address + EndpointAddress,
163{
164 Bound {
166 request: BindingRequestId,
168 current: BindingEvidence<Role>,
170 },
171 Rebalanced {
173 request: BindingRequestId,
175 prior: BindingEvidence<Role>,
177 current: BindingEvidence<Role>,
179 },
180 Unchanged {
182 request: BindingRequestId,
184 current: BindingEvidence<Role>,
186 },
187 Unbound {
189 request: BindingRequestId,
191 key: Key,
193 removed: BindingEvidence<Role>,
195 },
196 AlreadyUnbound {
198 command: BindingCommand<A, Key, Role>,
200 },
201 Rejected {
203 command: BindingCommand<A, Key, Role>,
205 reason: BindingRejection,
207 },
208}
209
210#[derive(Clone, Copy, Debug, Eq, PartialEq)]
212pub enum KeyedAdmissionRejection {
213 ShuttingDown,
215 BindingCapacityExhausted,
217 UnknownSelectedRole,
219 RoleUnavailable,
221 RoleBacklogFull,
223 NoRecoverableWorkers,
225 BindingGenerationExhausted,
227 JobCorrelationExhausted,
229 AssignmentCorrelationExhausted,
231}
232
233#[derive(Clone, Copy, Debug, Eq, PartialEq)]
235pub enum KeyedQueuedReturnReason {
236 PoolShutdown,
238 RolePermanentlyUnavailable,
240 NoRecoverableWorkers,
242 AssignmentCorrelationExhausted,
244}
245
246#[derive(Clone, Copy, Debug, Eq, PartialEq)]
248pub enum KeyedAssignedReturnReason {
249 PoolShutdown,
251 RolePermanentlyUnavailable,
253 WorkerStopped,
255 AssignmentReturned,
257 RetryPreparationRejected,
259 ContradictoryAssignmentSettlement,
261}
262
263pub enum KeyedOutcome<Key, Role, Job, WorkerResult> {
265 Accepted {
267 submission: SubmissionId,
269 job: JobId,
271 binding: BindingEvidence<Role>,
273 },
274 Rejected {
276 submission: SubmissionId,
278 key: Key,
280 payload: Job,
282 reason: KeyedAdmissionRejection,
284 },
285 Completed {
287 job: JobId,
289 binding: BindingEvidence<Role>,
291 worker_result: WorkerResult,
293 },
294 ReturnedQueued {
296 job: JobId,
298 binding: BindingEvidence<Role>,
300 payload: Job,
302 reason: KeyedQueuedReturnReason,
304 },
305 ReturnedAssigned {
307 job: JobId,
309 binding: BindingEvidence<Role>,
311 payload: Job,
313 reason: KeyedAssignedReturnReason,
315 },
316}
317
318pub(super) enum KeyedRequest<A, Key, Role, Job, WorkerResult>
319where
320 A: Address + EndpointAddress,
321{
322 Submit {
323 submission: SubmissionId,
324 key: Key,
325 payload: Job,
326 customer: ReplyRoute<MessageProtocol<A, KeyedOutcome<Key, Role, Job, WorkerResult>>>,
327 },
328 Binding(BindingCommand<A, Key, Role>),
329 Shutdown,
330}
331
332pub struct KeyedCommand<A, Key, Role, Job, WorkerResult>
334where
335 A: Address + EndpointAddress,
336{
337 request: KeyedRequest<A, Key, Role, Job, WorkerResult>,
338}
339
340impl<A, Key, Role, Job, WorkerResult> KeyedCommand<A, Key, Role, Job, WorkerResult>
341where
342 A: Address + EndpointAddress,
343{
344 #[must_use]
346 pub fn submit<Route>(submission: SubmissionId, key: Key, payload: Job, customer: Route) -> Self
347 where
348 Route: Into<ReplyRoute<MessageProtocol<A, KeyedOutcome<Key, Role, Job, WorkerResult>>>>,
349 {
350 Self {
351 request: KeyedRequest::Submit {
352 submission,
353 key,
354 payload,
355 customer: customer.into(),
356 },
357 }
358 }
359
360 #[must_use]
362 pub fn rebalance<Route>(
363 request: BindingRequestId,
364 key: Key,
365 expected: BindingExpectation,
366 target: Role,
367 reply: Route,
368 ) -> Self
369 where
370 Route: Into<ReplyRoute<MessageProtocol<A, BindingReply<A, Key, Role>>>>,
371 {
372 Self {
373 request: KeyedRequest::Binding(BindingCommand::rebalance(
374 request, key, expected, target, reply,
375 )),
376 }
377 }
378
379 #[must_use]
381 pub fn unbind<Route>(
382 request: BindingRequestId,
383 key: Key,
384 expected: BindingExpectation,
385 reply: Route,
386 ) -> Self
387 where
388 Route: Into<ReplyRoute<MessageProtocol<A, BindingReply<A, Key, Role>>>>,
389 {
390 Self {
391 request: KeyedRequest::Binding(BindingCommand::unbind(request, key, expected, reply)),
392 }
393 }
394
395 #[must_use]
397 pub const fn binding(command: BindingCommand<A, Key, Role>) -> Self {
398 Self {
399 request: KeyedRequest::Binding(command),
400 }
401 }
402
403 #[must_use]
405 pub const fn shutdown() -> Self {
406 Self {
407 request: KeyedRequest::Shutdown,
408 }
409 }
410
411 pub(super) fn into_request(self) -> KeyedRequest<A, Key, Role, Job, WorkerResult> {
412 self.request
413 }
414
415 pub(super) const fn from_request(
416 request: KeyedRequest<A, Key, Role, Job, WorkerResult>,
417 ) -> Self {
418 Self { request }
419 }
420}
421
422#[expect(dead_code, reason = "transferred intact to diagnostic custody")]
423pub(super) enum KeyedDiagnosticCause<Role, W, P, Source, Key, Job, WorkerResult>
424where
425 Role: Send + Sync,
426 W: Behavior + Send,
427 W::Protocol: Protocol<Msg = super::Assignment<Job>>,
428 P: ActivationPlan,
429 Source: WorkerSource<Role, W, P>,
430 BehaviorAddr<W>: EndpointAddress,
431 <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
432 Job: Send,
433{
434 BindingRemoved {
435 key: Key,
436 binding: BindingEvidence<Role>,
437 },
438 WorkerReturned {
439 role: RoleName<Role>,
440 rejection: WorkerCreationRejection<W>,
441 activation: P,
442 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
443 },
444 StoppedWorkerEstablished {
445 role: RoleName<Role>,
446 worker: CurrentWorker<W>,
447 activation: P,
448 stopped: ChildStopped<BehaviorAddr<W>>,
449 },
450 UnmatchedWorkerReturn {
451 id: behavior::CreationId,
452 kind: behavior::CreationKind,
453 rejection: WorkerCreationRejection<W>,
454 },
455 UnusedActivation {
456 role: RoleName<Role>,
457 permit: Option<ActivationPermit<W>>,
458 activation: P,
459 },
460 WorkerActivationReturned {
461 role: RoleName<Role>,
462 worker: WorkerAttempt,
463 activation: WorkerActivationGrant<W, P>,
464 outcome: WorkerActivationOutcome<W, P>,
465 },
466 WorkerShutdownRejected {
467 role: RoleName<Role>,
468 shutdown: EstablishedShutdownResolved<W::Protocol>,
469 },
470 WorkerReplacementCancelled {
471 role: RoleName<Role>,
472 replacement: PendingWorkerReplacement<W, P>,
473 },
474 WorkerPreparationFailed {
475 role: RoleName<Role>,
476 previous: WorkerAttempt,
477 stopped: ChildStopped<BehaviorAddr<W>>,
478 returned_source: Option<Source>,
479 error: WorkerPreparationError<Source::WorkerRejection, Source::SourceRejection>,
480 },
481 WorkerReplacementFailed {
482 role: RoleName<Role>,
483 previous: WorkerAttempt,
484 stopped: ChildStopped<BehaviorAddr<W>>,
485 submission: super::super::WorkerSubmission<W, P>,
486 returned_source: Option<Source>,
487 error: WorkerReplacementError,
488 },
489 Unexpected(
490 super::KeyedEvent<
491 Role,
492 W,
493 P,
494 Key,
495 Job,
496 WorkerResult,
497 behavior::ActionItemResult<PrepareWorkers<Source, Role, W, P>>,
498 super::super::WorkerPreparation<Source, Role, W, P>,
499 >,
500 ),
501}
502
503pub struct KeyedDiagnostic<Role, W, P, Source, Key, Job, WorkerResult>
505where
506 Role: Send + Sync,
507 W: Behavior + Send,
508 W::Protocol: Protocol<Msg = super::Assignment<Job>>,
509 P: ActivationPlan,
510 Source: WorkerSource<Role, W, P>,
511 BehaviorAddr<W>: EndpointAddress,
512 <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
513 Job: Send,
514{
515 #[allow(
516 dead_code,
517 reason = "Bombay's diagnostic custodian receives the complete opaque cause"
518 )]
519 cause: KeyedDiagnosticCause<Role, W, P, Source, Key, Job, WorkerResult>,
520}
521
522impl<Role, W, P, Source, Key, Job, WorkerResult>
523 From<KeyedDiagnosticCause<Role, W, P, Source, Key, Job, WorkerResult>>
524 for KeyedDiagnostic<Role, W, P, Source, Key, Job, WorkerResult>
525where
526 Role: Send + Sync,
527 W: Behavior + Send,
528 W::Protocol: Protocol<Msg = super::Assignment<Job>>,
529 P: ActivationPlan,
530 Source: WorkerSource<Role, W, P>,
531 BehaviorAddr<W>: EndpointAddress,
532 <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
533 Job: Send,
534{
535 fn from(cause: KeyedDiagnosticCause<Role, W, P, Source, Key, Job, WorkerResult>) -> Self {
536 Self { cause }
537 }
538}
539
540impl<Role, W, P, Source, Key, Job, WorkerResult> core::fmt::Debug
541 for KeyedDiagnostic<Role, W, P, Source, Key, Job, WorkerResult>
542where
543 Role: Send + Sync,
544 W: Behavior + Send,
545 W::Protocol: Protocol<Msg = super::Assignment<Job>>,
546 P: ActivationPlan,
547 Source: WorkerSource<Role, W, P>,
548 BehaviorAddr<W>: EndpointAddress,
549 <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
550 Job: Send,
551{
552 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
553 formatter
554 .debug_struct("KeyedDiagnostic")
555 .finish_non_exhaustive()
556 }
557}
558
559#[cfg(test)]
560mod live_activation_correlation {
561 use super::super::binding::{BindingCapacity, BindingTable};
562 use super::super::requests::KeyedRequests;
563 use super::super::role::RoleCell;
564 use super::super::{KeyedEvent, KeyedOperating, KeyedPool, KeyedPoolState};
565 use super::{KeyedDiagnostic, KeyedDiagnosticCause};
566 use crate::atomic::pool::ShutdownSequence;
567 use crate::atomic::pool::assignment::{AcceptedJobSequence, AssignmentSequence};
568 use crate::atomic::pool::worker::activation_pool_correlation::{
569 ActivationEndpoint, ActivationWorker, assert_current_worker, initialized_activation,
570 };
571 use crate::atomic::pool::worker::{CurrentWorker, Member, MemberState, Worker, WorkerPhase};
572 use crate::atomic::restart::RecoveryCount;
573 use crate::atomic::restart::RestartBudget;
574 use crate::atomic::worker::WorkerActivationOutcome;
575 use crate::atomic::{
576 ActivationPolicy, ActorDrainPolicy, BacklogCapacity, ImmediateActivation, Interruption,
577 PoolFailureReaction, PoolRecovery, RoleName, WorkerAttempt,
578 };
579 use crate::{
580 Active, ChildStopped, DiagnosticAction, DiagnosticDisposition, Exit, StopOnShutdown,
581 };
582 use behavior::{Actions, CreationSequence, EstablishedActor, Never, Step};
583 use core::convert::Infallible;
584 use std::sync::Arc;
585 use std::time::Instant;
586 fn role_for_key(_: &u64) -> u64 {
587 23
588 }
589
590 #[tokio::test]
591 async fn keyed_dispatched_live_returns_the_original_foreign_activation_twice() {
592 let mut creations = CreationSequence::new();
593 let creation = creations.issue().expect("one original worker creation");
594 let worker = WorkerAttempt::issued(creation);
595 let actor =
596 EstablishedActor::<StopOnShutdown<ActivationWorker>>::issued(ActivationEndpoint(19));
597 let expected_actor = actor.clone();
598 let original = initialized_activation(worker.clone(), actor.recipient());
599 let foreign = initialized_activation(worker.clone(), actor.recipient());
600 let original_attempt = original.attempt();
601 let foreign_attempt = foreign.attempt();
602 assert!(original_attempt != foreign_attempt);
603 let inputs = [
604 (
605 foreign.started(),
606 WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Started,
607 ),
608 (
609 foreign.activate().await,
610 WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Ready(()),
611 ),
612 ];
613 for (input, report) in inputs {
614 let role = RoleName::new(23_u64);
615 let expected_role = role.clone().into_role();
616 let member = Member {
617 role,
618 recoveries: RecoveryCount::default(),
619 state: MemberState::Worker(Worker {
620 current: CurrentWorker::new(worker.clone(), expected_actor.clone()),
621 phase: WorkerPhase::ActivationDispatched {
622 attempt: original_attempt.clone(),
623 stopped: None,
624 },
625 }),
626 };
627 let binding_capacity = BindingCapacity::new(2).expect("two real binding slots");
628 let pool: KeyedPool<
629 u64,
630 ActivationWorker,
631 ImmediateActivation,
632 Never,
633 fn(&u64) -> u64,
634 Infallible,
635 u64,
636 u8,
637 u16,
638 > = KeyedPool {
639 state: KeyedPoolState::Operating(KeyedOperating {
640 roles: vec![RoleCell::new(member, BacklogCapacity::new(2))],
641 bindings: BindingTable::new(binding_capacity),
642 }),
643 selector: role_for_key,
644 activation: ActivationPolicy::new(1).expect("one actual activation slot"),
645 recovery: PoolRecovery::<Never>::temporary(PoolFailureReaction::RetireRole).into(),
646 restarts: RestartBudget::empty(),
647 backlog: BacklogCapacity::new(2),
648 binding_capacity,
649 interruption: Interruption::Fail,
650 actor_drain: ActorDrainPolicy::WaitForActorGraph,
651 diagnostics: DiagnosticDisposition::terminate(),
652 creations: CreationSequence::new(),
653 jobs: AcceptedJobSequence::new(),
654 assignments: AssignmentSequence::new(),
655 shutdowns: ShutdownSequence::new(),
656 next_restart_timer: 1,
657 };
658 let mut pool = Active { behavior: pool };
659 let mut returned = input;
660 for _ in 0..2 {
661 let actions = pool
662 .transition(KeyedEvent::WorkerActivationReported(returned))
663 .unwrap_or_else(|error| {
664 panic!("the activation fold returns its real Actions: {error}")
665 });
666 let Actions {
667 sends,
668 creates,
669 become_,
670 } = actions;
671 let KeyedRequests {
672 worker_observations,
673 worker_initializations,
674 worker_activations,
675 customer_outcomes,
676 binding_replies,
677 worker_assignments,
678 worker_preparations,
679 restart_schedules,
680 worker_shutdowns,
681 diagnostics,
682 } = sends;
683 assert!(worker_observations.is_empty());
684 assert!(worker_initializations.is_empty());
685 assert!(worker_activations.is_empty());
686 assert!(customer_outcomes.is_empty());
687 assert!(binding_replies.as_slice().is_empty());
688 assert!(worker_assignments.is_empty());
689 assert!(worker_preparations.is_empty());
690 assert!(restart_schedules.is_empty());
691 assert!(worker_shutdowns.is_empty());
692 assert!(creates.is_empty());
693 assert!(matches!(become_, Step::Continue));
694 let mut diagnostics = diagnostics.into_requests();
695 assert_eq!(diagnostics.len(), 1);
696 let DiagnosticAction::Terminal {
697 diagnostic:
698 KeyedDiagnostic {
699 cause:
700 KeyedDiagnosticCause::Unexpected(KeyedEvent::WorkerActivationReported(
701 input,
702 )),
703 },
704 } = diagnostics.remove(0)
705 else {
706 panic!(
707 "a foreign attempt returns the complete original activation through the existing Unexpected policy"
708 )
709 };
710 assert_eq!(input.worker(), worker);
711 assert!(input.attempt() == &foreign_attempt);
712 returned = input;
713 let KeyedPoolState::Operating(operating) = &pool.state else {
714 panic!("foreign activation must not retire the original pool")
715 };
716 assert_eq!(operating.roles.len(), 1);
717 assert!(operating.roles[0].queue.is_empty());
718 assert!(operating.bindings.binding(&7_u64).is_none());
719 let member = &operating.roles[0].member;
720 let role = member.role.clone().into_role();
721 assert!(Arc::ptr_eq(&role, &expected_role));
722 assert_eq!(*role, 23);
723 assert_eq!(member.recoveries, RecoveryCount::default());
724 let MemberState::Worker(Worker {
725 current,
726 phase: WorkerPhase::ActivationDispatched { attempt, stopped },
727 }) = &member.state
728 else {
729 panic!("foreign report preserves the exact original worker phase")
730 };
731 assert_current_worker(current, &worker, &expected_actor);
732 assert!(attempt == &original_attempt);
733 assert!(stopped.is_none());
734 }
735 let (returned_worker, returned_attempt, outcome) = returned.into_parts();
736 assert_eq!(returned_worker, worker);
737 assert!(returned_attempt == foreign_attempt);
738 match (report, outcome) {
739 (WorkerActivationOutcome::Started, WorkerActivationOutcome::Started)
740 | (WorkerActivationOutcome::Ready(()), WorkerActivationOutcome::Ready(())) => {}
741 _ => panic!("the original complete activation outcome is preserved"),
742 }
743 }
744 drop(original);
745 }
746
747 #[tokio::test]
748 async fn keyed_dispatched_stopped_returns_the_original_foreign_activation_twice() {
749 let mut creations = CreationSequence::new();
750 let creation = creations.issue().expect("one original worker creation");
751 let worker = WorkerAttempt::issued(creation);
752 let actor =
753 EstablishedActor::<StopOnShutdown<ActivationWorker>>::issued(ActivationEndpoint(19));
754 let expected_actor = actor.clone();
755 let original = initialized_activation(worker.clone(), actor.recipient());
756 let foreign = initialized_activation(worker.clone(), actor.recipient());
757 let original_attempt = original.attempt();
758 let foreign_attempt = foreign.attempt();
759 assert!(original_attempt != foreign_attempt);
760 let stopped_at = Instant::now();
761 let inputs = [
762 (
763 foreign.started(),
764 WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Started,
765 ),
766 (
767 foreign.activate().await,
768 WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Ready(()),
769 ),
770 ];
771 for (input, report) in inputs {
772 let role = RoleName::new(23_u64);
773 let expected_role = role.clone().into_role();
774 let member = Member {
775 role,
776 recoveries: RecoveryCount::default(),
777 state: MemberState::Worker(Worker {
778 current: CurrentWorker::new(worker.clone(), expected_actor.clone()),
779 phase: WorkerPhase::ActivationDispatched {
780 attempt: original_attempt.clone(),
781 stopped: Some(ChildStopped::new(creation, Ok(Exit::Normal), stopped_at)),
782 },
783 }),
784 };
785 let binding_capacity = BindingCapacity::new(2).expect("two real binding slots");
786 let pool: KeyedPool<
787 u64,
788 ActivationWorker,
789 ImmediateActivation,
790 Never,
791 fn(&u64) -> u64,
792 Infallible,
793 u64,
794 u8,
795 u16,
796 > = KeyedPool {
797 state: KeyedPoolState::Operating(KeyedOperating {
798 roles: vec![RoleCell::new(member, BacklogCapacity::new(2))],
799 bindings: BindingTable::new(binding_capacity),
800 }),
801 selector: role_for_key,
802 activation: ActivationPolicy::new(1).expect("one actual activation slot"),
803 recovery: PoolRecovery::<Never>::temporary(PoolFailureReaction::RetireRole).into(),
804 restarts: RestartBudget::empty(),
805 backlog: BacklogCapacity::new(2),
806 binding_capacity,
807 interruption: Interruption::Fail,
808 actor_drain: ActorDrainPolicy::WaitForActorGraph,
809 diagnostics: DiagnosticDisposition::terminate(),
810 creations: CreationSequence::new(),
811 jobs: AcceptedJobSequence::new(),
812 assignments: AssignmentSequence::new(),
813 shutdowns: ShutdownSequence::new(),
814 next_restart_timer: 1,
815 };
816 let mut pool = Active { behavior: pool };
817 let mut returned = input;
818 for _ in 0..2 {
819 let actions = pool
820 .transition(KeyedEvent::WorkerActivationReported(returned))
821 .unwrap_or_else(|error| {
822 panic!("the activation fold returns its real Actions: {error}")
823 });
824 let Actions {
825 sends,
826 creates,
827 become_,
828 } = actions;
829 let KeyedRequests {
830 worker_observations,
831 worker_initializations,
832 worker_activations,
833 customer_outcomes,
834 binding_replies,
835 worker_assignments,
836 worker_preparations,
837 restart_schedules,
838 worker_shutdowns,
839 diagnostics,
840 } = sends;
841 assert!(worker_observations.is_empty());
842 assert!(worker_initializations.is_empty());
843 assert!(worker_activations.is_empty());
844 assert!(customer_outcomes.is_empty());
845 assert!(binding_replies.as_slice().is_empty());
846 assert!(worker_assignments.is_empty());
847 assert!(worker_preparations.is_empty());
848 assert!(restart_schedules.is_empty());
849 assert!(worker_shutdowns.is_empty());
850 assert!(creates.is_empty());
851 assert!(matches!(become_, Step::Continue));
852 let mut diagnostics = diagnostics.into_requests();
853 assert_eq!(diagnostics.len(), 1);
854 let DiagnosticAction::Terminal {
855 diagnostic:
856 KeyedDiagnostic {
857 cause:
858 KeyedDiagnosticCause::Unexpected(KeyedEvent::WorkerActivationReported(
859 input,
860 )),
861 },
862 } = diagnostics.remove(0)
863 else {
864 panic!(
865 "a foreign attempt returns the complete original activation through the existing Unexpected policy"
866 )
867 };
868 assert_eq!(input.worker(), worker);
869 assert!(input.attempt() == &foreign_attempt);
870 returned = input;
871 let KeyedPoolState::Operating(operating) = &pool.state else {
872 panic!("foreign activation must not retire the original pool")
873 };
874 assert_eq!(operating.roles.len(), 1);
875 assert!(operating.roles[0].queue.is_empty());
876 assert!(operating.bindings.binding(&7_u64).is_none());
877 let member = &operating.roles[0].member;
878 let role = member.role.clone().into_role();
879 assert!(Arc::ptr_eq(&role, &expected_role));
880 assert_eq!(*role, 23);
881 assert_eq!(member.recoveries, RecoveryCount::default());
882 let MemberState::Worker(Worker {
883 current,
884 phase: WorkerPhase::ActivationDispatched { attempt, stopped },
885 }) = &member.state
886 else {
887 panic!("foreign report preserves the exact original worker phase")
888 };
889 assert_current_worker(current, &worker, &expected_actor);
890 assert!(attempt == &original_attempt);
891 let Some(stopped) = stopped else {
892 panic!("the original stopped child remains owned")
893 };
894 assert_eq!(stopped.child, creation);
895 assert!(matches!(&stopped.outcome, Ok(Exit::Normal)));
896 assert_eq!(stopped.at, stopped_at);
897 }
898 let (returned_worker, returned_attempt, outcome) = returned.into_parts();
899 assert_eq!(returned_worker, worker);
900 assert!(returned_attempt == foreign_attempt);
901 match (report, outcome) {
902 (WorkerActivationOutcome::Started, WorkerActivationOutcome::Started)
903 | (WorkerActivationOutcome::Ready(()), WorkerActivationOutcome::Ready(())) => {}
904 _ => panic!("the original complete activation outcome is preserved"),
905 }
906 }
907 drop(original);
908 }
909
910 #[tokio::test]
911 async fn keyed_activating_live_returns_the_original_foreign_activation_twice() {
912 let mut creations = CreationSequence::new();
913 let creation = creations.issue().expect("one original worker creation");
914 let worker = WorkerAttempt::issued(creation);
915 let actor =
916 EstablishedActor::<StopOnShutdown<ActivationWorker>>::issued(ActivationEndpoint(19));
917 let expected_actor = actor.clone();
918 let original = initialized_activation(worker.clone(), actor.recipient());
919 let foreign = initialized_activation(worker.clone(), actor.recipient());
920 let original_attempt = original.attempt();
921 let foreign_attempt = foreign.attempt();
922 assert!(original_attempt != foreign_attempt);
923 let inputs = [
924 (
925 foreign.started(),
926 WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Started,
927 ),
928 (
929 foreign.activate().await,
930 WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Ready(()),
931 ),
932 ];
933 for (input, report) in inputs {
934 let role = RoleName::new(23_u64);
935 let expected_role = role.clone().into_role();
936 let member = Member {
937 role,
938 recoveries: RecoveryCount::default(),
939 state: MemberState::Worker(Worker {
940 current: CurrentWorker::new(worker.clone(), expected_actor.clone()),
941 phase: WorkerPhase::Activating {
942 attempt: original_attempt.clone(),
943 stopped: None,
944 },
945 }),
946 };
947 let binding_capacity = BindingCapacity::new(2).expect("two real binding slots");
948 let pool: KeyedPool<
949 u64,
950 ActivationWorker,
951 ImmediateActivation,
952 Never,
953 fn(&u64) -> u64,
954 Infallible,
955 u64,
956 u8,
957 u16,
958 > = KeyedPool {
959 state: KeyedPoolState::Operating(KeyedOperating {
960 roles: vec![RoleCell::new(member, BacklogCapacity::new(2))],
961 bindings: BindingTable::new(binding_capacity),
962 }),
963 selector: role_for_key,
964 activation: ActivationPolicy::new(1).expect("one actual activation slot"),
965 recovery: PoolRecovery::<Never>::temporary(PoolFailureReaction::RetireRole).into(),
966 restarts: RestartBudget::empty(),
967 backlog: BacklogCapacity::new(2),
968 binding_capacity,
969 interruption: Interruption::Fail,
970 actor_drain: ActorDrainPolicy::WaitForActorGraph,
971 diagnostics: DiagnosticDisposition::terminate(),
972 creations: CreationSequence::new(),
973 jobs: AcceptedJobSequence::new(),
974 assignments: AssignmentSequence::new(),
975 shutdowns: ShutdownSequence::new(),
976 next_restart_timer: 1,
977 };
978 let mut pool = Active { behavior: pool };
979 let mut returned = input;
980 for _ in 0..2 {
981 let actions = pool
982 .transition(KeyedEvent::WorkerActivationReported(returned))
983 .unwrap_or_else(|error| {
984 panic!("the activation fold returns its real Actions: {error}")
985 });
986 let Actions {
987 sends,
988 creates,
989 become_,
990 } = actions;
991 let KeyedRequests {
992 worker_observations,
993 worker_initializations,
994 worker_activations,
995 customer_outcomes,
996 binding_replies,
997 worker_assignments,
998 worker_preparations,
999 restart_schedules,
1000 worker_shutdowns,
1001 diagnostics,
1002 } = sends;
1003 assert!(worker_observations.is_empty());
1004 assert!(worker_initializations.is_empty());
1005 assert!(worker_activations.is_empty());
1006 assert!(customer_outcomes.is_empty());
1007 assert!(binding_replies.as_slice().is_empty());
1008 assert!(worker_assignments.is_empty());
1009 assert!(worker_preparations.is_empty());
1010 assert!(restart_schedules.is_empty());
1011 assert!(worker_shutdowns.is_empty());
1012 assert!(creates.is_empty());
1013 assert!(matches!(become_, Step::Continue));
1014 let mut diagnostics = diagnostics.into_requests();
1015 assert_eq!(diagnostics.len(), 1);
1016 let DiagnosticAction::Terminal {
1017 diagnostic:
1018 KeyedDiagnostic {
1019 cause:
1020 KeyedDiagnosticCause::Unexpected(KeyedEvent::WorkerActivationReported(
1021 input,
1022 )),
1023 },
1024 } = diagnostics.remove(0)
1025 else {
1026 panic!(
1027 "a foreign attempt returns the complete original activation through the existing Unexpected policy"
1028 )
1029 };
1030 assert_eq!(input.worker(), worker);
1031 assert!(input.attempt() == &foreign_attempt);
1032 returned = input;
1033 let KeyedPoolState::Operating(operating) = &pool.state else {
1034 panic!("foreign activation must not retire the original pool")
1035 };
1036 assert_eq!(operating.roles.len(), 1);
1037 assert!(operating.roles[0].queue.is_empty());
1038 assert!(operating.bindings.binding(&7_u64).is_none());
1039 let member = &operating.roles[0].member;
1040 let role = member.role.clone().into_role();
1041 assert!(Arc::ptr_eq(&role, &expected_role));
1042 assert_eq!(*role, 23);
1043 assert_eq!(member.recoveries, RecoveryCount::default());
1044 let MemberState::Worker(Worker {
1045 current,
1046 phase: WorkerPhase::Activating { attempt, stopped },
1047 }) = &member.state
1048 else {
1049 panic!("foreign report preserves the exact original worker phase")
1050 };
1051 assert_current_worker(current, &worker, &expected_actor);
1052 assert!(attempt == &original_attempt);
1053 assert!(stopped.is_none());
1054 }
1055 let (returned_worker, returned_attempt, outcome) = returned.into_parts();
1056 assert_eq!(returned_worker, worker);
1057 assert!(returned_attempt == foreign_attempt);
1058 match (report, outcome) {
1059 (WorkerActivationOutcome::Started, WorkerActivationOutcome::Started)
1060 | (WorkerActivationOutcome::Ready(()), WorkerActivationOutcome::Ready(())) => {}
1061 _ => panic!("the original complete activation outcome is preserved"),
1062 }
1063 }
1064 drop(original);
1065 }
1066
1067 #[tokio::test]
1068 async fn keyed_activating_stopped_returns_the_original_foreign_activation_twice() {
1069 let mut creations = CreationSequence::new();
1070 let creation = creations.issue().expect("one original worker creation");
1071 let worker = WorkerAttempt::issued(creation);
1072 let actor =
1073 EstablishedActor::<StopOnShutdown<ActivationWorker>>::issued(ActivationEndpoint(19));
1074 let expected_actor = actor.clone();
1075 let original = initialized_activation(worker.clone(), actor.recipient());
1076 let foreign = initialized_activation(worker.clone(), actor.recipient());
1077 let original_attempt = original.attempt();
1078 let foreign_attempt = foreign.attempt();
1079 assert!(original_attempt != foreign_attempt);
1080 let stopped_at = Instant::now();
1081 let inputs = [
1082 (
1083 foreign.started(),
1084 WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Started,
1085 ),
1086 (
1087 foreign.activate().await,
1088 WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Ready(()),
1089 ),
1090 ];
1091 for (input, report) in inputs {
1092 let role = RoleName::new(23_u64);
1093 let expected_role = role.clone().into_role();
1094 let member = Member {
1095 role,
1096 recoveries: RecoveryCount::default(),
1097 state: MemberState::Worker(Worker {
1098 current: CurrentWorker::new(worker.clone(), expected_actor.clone()),
1099 phase: WorkerPhase::Activating {
1100 attempt: original_attempt.clone(),
1101 stopped: Some(ChildStopped::new(creation, Ok(Exit::Normal), stopped_at)),
1102 },
1103 }),
1104 };
1105 let binding_capacity = BindingCapacity::new(2).expect("two real binding slots");
1106 let pool: KeyedPool<
1107 u64,
1108 ActivationWorker,
1109 ImmediateActivation,
1110 Never,
1111 fn(&u64) -> u64,
1112 Infallible,
1113 u64,
1114 u8,
1115 u16,
1116 > = KeyedPool {
1117 state: KeyedPoolState::Operating(KeyedOperating {
1118 roles: vec![RoleCell::new(member, BacklogCapacity::new(2))],
1119 bindings: BindingTable::new(binding_capacity),
1120 }),
1121 selector: role_for_key,
1122 activation: ActivationPolicy::new(1).expect("one actual activation slot"),
1123 recovery: PoolRecovery::<Never>::temporary(PoolFailureReaction::RetireRole).into(),
1124 restarts: RestartBudget::empty(),
1125 backlog: BacklogCapacity::new(2),
1126 binding_capacity,
1127 interruption: Interruption::Fail,
1128 actor_drain: ActorDrainPolicy::WaitForActorGraph,
1129 diagnostics: DiagnosticDisposition::terminate(),
1130 creations: CreationSequence::new(),
1131 jobs: AcceptedJobSequence::new(),
1132 assignments: AssignmentSequence::new(),
1133 shutdowns: ShutdownSequence::new(),
1134 next_restart_timer: 1,
1135 };
1136 let mut pool = Active { behavior: pool };
1137 let mut returned = input;
1138 for _ in 0..2 {
1139 let actions = pool
1140 .transition(KeyedEvent::WorkerActivationReported(returned))
1141 .unwrap_or_else(|error| {
1142 panic!("the activation fold returns its real Actions: {error}")
1143 });
1144 let Actions {
1145 sends,
1146 creates,
1147 become_,
1148 } = actions;
1149 let KeyedRequests {
1150 worker_observations,
1151 worker_initializations,
1152 worker_activations,
1153 customer_outcomes,
1154 binding_replies,
1155 worker_assignments,
1156 worker_preparations,
1157 restart_schedules,
1158 worker_shutdowns,
1159 diagnostics,
1160 } = sends;
1161 assert!(worker_observations.is_empty());
1162 assert!(worker_initializations.is_empty());
1163 assert!(worker_activations.is_empty());
1164 assert!(customer_outcomes.is_empty());
1165 assert!(binding_replies.as_slice().is_empty());
1166 assert!(worker_assignments.is_empty());
1167 assert!(worker_preparations.is_empty());
1168 assert!(restart_schedules.is_empty());
1169 assert!(worker_shutdowns.is_empty());
1170 assert!(creates.is_empty());
1171 assert!(matches!(become_, Step::Continue));
1172 let mut diagnostics = diagnostics.into_requests();
1173 assert_eq!(diagnostics.len(), 1);
1174 let DiagnosticAction::Terminal {
1175 diagnostic:
1176 KeyedDiagnostic {
1177 cause:
1178 KeyedDiagnosticCause::Unexpected(KeyedEvent::WorkerActivationReported(
1179 input,
1180 )),
1181 },
1182 } = diagnostics.remove(0)
1183 else {
1184 panic!(
1185 "a foreign attempt returns the complete original activation through the existing Unexpected policy"
1186 )
1187 };
1188 assert_eq!(input.worker(), worker);
1189 assert!(input.attempt() == &foreign_attempt);
1190 returned = input;
1191 let KeyedPoolState::Operating(operating) = &pool.state else {
1192 panic!("foreign activation must not retire the original pool")
1193 };
1194 assert_eq!(operating.roles.len(), 1);
1195 assert!(operating.roles[0].queue.is_empty());
1196 assert!(operating.bindings.binding(&7_u64).is_none());
1197 let member = &operating.roles[0].member;
1198 let role = member.role.clone().into_role();
1199 assert!(Arc::ptr_eq(&role, &expected_role));
1200 assert_eq!(*role, 23);
1201 assert_eq!(member.recoveries, RecoveryCount::default());
1202 let MemberState::Worker(Worker {
1203 current,
1204 phase: WorkerPhase::Activating { attempt, stopped },
1205 }) = &member.state
1206 else {
1207 panic!("foreign report preserves the exact original worker phase")
1208 };
1209 assert_current_worker(current, &worker, &expected_actor);
1210 assert!(attempt == &original_attempt);
1211 let Some(stopped) = stopped else {
1212 panic!("the original stopped child remains owned")
1213 };
1214 assert_eq!(stopped.child, creation);
1215 assert!(matches!(&stopped.outcome, Ok(Exit::Normal)));
1216 assert_eq!(stopped.at, stopped_at);
1217 }
1218 let (returned_worker, returned_attempt, outcome) = returned.into_parts();
1219 assert_eq!(returned_worker, worker);
1220 assert!(returned_attempt == foreign_attempt);
1221 match (report, outcome) {
1222 (WorkerActivationOutcome::Started, WorkerActivationOutcome::Started)
1223 | (WorkerActivationOutcome::Ready(()), WorkerActivationOutcome::Ready(())) => {}
1224 _ => panic!("the original complete activation outcome is preserved"),
1225 }
1226 }
1227 drop(original);
1228 }
1229
1230 #[tokio::test]
1231 async fn keyed_original_activation_started_then_ready_preserves_all_lanes() {
1232 let mut creations = CreationSequence::new();
1233 let creation = creations.issue().expect("one original worker creation");
1234 let worker = WorkerAttempt::issued(creation);
1235 let actor =
1236 EstablishedActor::<StopOnShutdown<ActivationWorker>>::issued(ActivationEndpoint(19));
1237 let expected_actor = actor.clone();
1238 let original = initialized_activation(worker.clone(), actor.recipient());
1239 let original_attempt = original.attempt();
1240 let started = original.started();
1241 let role = RoleName::new(23_u64);
1242 let expected_role = role.clone().into_role();
1243 let member = Member {
1244 role,
1245 recoveries: RecoveryCount::default(),
1246 state: MemberState::Worker(Worker {
1247 current: CurrentWorker::new(worker.clone(), expected_actor.clone()),
1248 phase: WorkerPhase::ActivationDispatched {
1249 attempt: original_attempt.clone(),
1250 stopped: None,
1251 },
1252 }),
1253 };
1254 let binding_capacity = BindingCapacity::new(2).expect("two real binding slots");
1255 let pool: KeyedPool<
1256 u64,
1257 ActivationWorker,
1258 ImmediateActivation,
1259 Never,
1260 fn(&u64) -> u64,
1261 Infallible,
1262 u64,
1263 u8,
1264 u16,
1265 > = KeyedPool {
1266 state: KeyedPoolState::Operating(KeyedOperating {
1267 roles: vec![RoleCell::new(member, BacklogCapacity::new(2))],
1268 bindings: BindingTable::new(binding_capacity),
1269 }),
1270 selector: role_for_key,
1271 activation: ActivationPolicy::new(1).expect("one actual activation slot"),
1272 recovery: PoolRecovery::<Never>::temporary(PoolFailureReaction::RetireRole).into(),
1273 restarts: RestartBudget::empty(),
1274 backlog: BacklogCapacity::new(2),
1275 binding_capacity,
1276 interruption: Interruption::Fail,
1277 actor_drain: ActorDrainPolicy::WaitForActorGraph,
1278 diagnostics: DiagnosticDisposition::terminate(),
1279 creations: CreationSequence::new(),
1280 jobs: AcceptedJobSequence::new(),
1281 assignments: AssignmentSequence::new(),
1282 shutdowns: ShutdownSequence::new(),
1283 next_restart_timer: 1,
1284 };
1285 let mut pool = Active { behavior: pool };
1286
1287 let inputs = [
1288 (
1289 started,
1290 WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Started,
1291 ),
1292 (
1293 original.activate().await,
1294 WorkerActivationOutcome::<ActivationWorker, ImmediateActivation>::Ready(()),
1295 ),
1296 ];
1297 for (input, report) in inputs {
1298 let actions = pool
1299 .transition(KeyedEvent::WorkerActivationReported(input))
1300 .unwrap_or_else(|error| {
1301 panic!("the original activation returns complete Actions: {error}")
1302 });
1303 let Actions {
1304 sends,
1305 creates,
1306 become_,
1307 } = actions;
1308 let KeyedRequests {
1309 worker_observations,
1310 worker_initializations,
1311 worker_activations,
1312 customer_outcomes,
1313 binding_replies,
1314 worker_assignments,
1315 worker_preparations,
1316 restart_schedules,
1317 worker_shutdowns,
1318 diagnostics,
1319 } = sends;
1320 assert!(worker_observations.is_empty());
1321 assert!(worker_initializations.is_empty());
1322 assert!(worker_activations.is_empty());
1323 assert!(customer_outcomes.is_empty());
1324 assert!(binding_replies.as_slice().is_empty());
1325 assert!(worker_assignments.is_empty());
1326 assert!(worker_preparations.is_empty());
1327 assert!(restart_schedules.is_empty());
1328 assert!(worker_shutdowns.is_empty());
1329 assert!(diagnostics.is_empty());
1330 assert!(creates.is_empty());
1331 assert!(matches!(become_, Step::Continue));
1332 let KeyedPoolState::Operating(operating) = &pool.state else {
1333 panic!("original activation leaves the pool operating")
1334 };
1335 assert_eq!(operating.roles.len(), 1);
1336 assert!(operating.roles[0].queue.is_empty());
1337 assert!(operating.bindings.binding(&7_u64).is_none());
1338 let member = &operating.roles[0].member;
1339 let actual_role = member.role.clone().into_role();
1340 assert!(Arc::ptr_eq(&actual_role, &expected_role));
1341 assert_eq!(*actual_role, 23);
1342 assert_eq!(member.recoveries, RecoveryCount::default());
1343 let MemberState::Worker(Worker { current, phase }) = &member.state else {
1344 panic!("the actual original worker stays owned")
1345 };
1346 assert_current_worker(current, &worker, &expected_actor);
1347 match (report, phase) {
1348 (
1349 WorkerActivationOutcome::Started,
1350 WorkerPhase::Activating {
1351 attempt,
1352 stopped: None,
1353 },
1354 ) => {
1355 assert!(attempt == &original_attempt);
1356 }
1357 (WorkerActivationOutcome::Ready(()), WorkerPhase::Idle) => {}
1358 _ => panic!("actual Started enters Activating and actual Ready enters Idle"),
1359 }
1360 }
1361 }
1362}