1use std::collections::BTreeMap;
4use std::mem;
5
6use behavior::{
7 Actions, ActiveTurn, Address, Behavior, BehaviorActed, BehaviorAddr, BehaviorBase, Births,
8 ChildCreationOutcome, ChildCreationSettled, ChildReport, CreateChild, CreationKind,
9 CreationSequence, CreationSettlement, Creations, CreationsSettled, EndpointAddress, Here,
10 InitializationTurn, InjectEvent, InterpreterRequests, ItemSettlement, MessageProtocol, Never,
11 Protocol, SendEffects, SettledItem, SourceActions, User,
12};
13use thiserror::Error;
14
15use crate::{
16 DeliveryRoute, DiagnosticAction, DiagnosticDisposition, DiagnosticRoute, ObserveChild,
17 ReplyRoute, ScheduleAfter, ShutdownEstablished, ShutdownRequested, StopOnShutdown,
18};
19
20use super::pool::assignment::{AcceptedJobSequence, AssignedJob, AssignmentSequence, CustomerJob};
21use super::pool::worker as direct_worker;
22use super::pool::worker::{Member, MemberState, RetiringWorker, Worker, WorkerPhase};
23use super::pool::{
24 AssignWorker, Assignment, BacklogCapacity, CompletesAssignments, CustomerDelivery,
25 Interruption, PoolRecovery, PoolRecoveryState, ShutdownSequence, SubmissionId,
26};
27use super::worker::{InitialWorkerRejection, prepare_initial_workers};
28use super::{
29 ActivationPlan, ActivationPolicy, ActorDrainPolicy, BeginActivation, InitializeWorker,
30 OrderedRoles, PrepareWorkers, PreparedWorker, WorkerActivation, WorkerAttempt,
31 WorkerInitializationReport, WorkerPreparation, WorkerSource, WorkerSubmission,
32};
33
34mod assignment;
35mod binding;
36mod event;
37mod job;
38mod protocol;
39mod recovery;
40mod requests;
41mod role;
42mod shutdown;
43
44pub use binding::{
45 BindingCapacity, BindingEvidence, BindingExpectation, BindingGeneration, BindingRequestId,
46};
47pub use event::KeyedEvent;
48pub use protocol::{
49 BindingCommand, BindingRejection, BindingReply, KeyedAdmissionRejection,
50 KeyedAssignedReturnReason, KeyedCommand, KeyedDiagnostic, KeyedOutcome,
51 KeyedQueuedReturnReason,
52};
53pub use requests::KeyedRequests;
54
55use binding::{BindingReservationRejected, BindingTable, ReservedBinding};
56use job::KeyedCustomer;
57use protocol::{BindingChange, KeyedRequest};
58use role::{AdmissionTarget, ManagementTarget, RoleCell};
59
60type CustomerRoute<A, Key, Role, Job, WorkerResult> =
61 ReplyRoute<MessageProtocol<A, KeyedOutcome<Key, Role, Job, WorkerResult>>>;
62
63type KeyedQueue<Role, W, Key, Job, WorkerResult> = BTreeMap<
64 super::pool::assignment::AdmissionOrdinal,
65 CustomerJob<
66 Job,
67 KeyedCustomer<Role, CustomerRoute<BehaviorAddr<W>, Key, Role, Job, WorkerResult>>,
68 >,
69>;
70
71type ManagementRoute<A, Key, Role> = ReplyRoute<MessageProtocol<A, BindingReply<A, Key, Role>>>;
72
73type KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult> = Actions<
74 BehaviorAddr<W>,
75 Never,
76 KeyedRequests<
77 InterpreterRequests<ObserveChild<<W as Behavior>::Protocol, behavior::ChildHead>>,
78 InterpreterRequests<InitializeWorker<W, P>>,
79 InterpreterRequests<BeginActivation<W, P>>,
80 InterpreterRequests<
81 CustomerDelivery<
82 MessageProtocol<BehaviorAddr<W>, KeyedOutcome<Key, Role, Job, WorkerResult>>,
83 >,
84 >,
85 <ManagementRoute<BehaviorAddr<W>, Key, Role> as DeliveryRoute>::Sends,
86 SourceActions<AssignWorker<<W as Behavior>::Protocol, Job>>,
87 SourceActions<PrepareWorkers<Source, Role, W, P>>,
88 SourceActions<ScheduleAfter>,
89 InterpreterRequests<ShutdownEstablished<StopOnShutdown<W>, Here>>,
90 InterpreterRequests<
91 DiagnosticAction<
92 Diagnostics,
93 KeyedDiagnostic<Role, W, P, Source, Key, Job, WorkerResult>,
94 >,
95 >,
96 >,
97 Births<StopOnShutdown<W>>,
98>;
99
100type KeyedRoleCell<Role, W, P, Key, Job, WorkerResult> = RoleCell<
101 Role,
102 W,
103 P,
104 Job,
105 WorkerResult,
106 CustomerRoute<BehaviorAddr<W>, Key, Role, Job, WorkerResult>,
107>;
108
109struct KeyedOperating<Role, W, P, Key, Job, WorkerResult>
110where
111 W: Behavior,
112 BehaviorAddr<W>: EndpointAddress,
113{
114 roles: Vec<KeyedRoleCell<Role, W, P, Key, Job, WorkerResult>>,
115 bindings: BindingTable<Key, Role>,
116}
117
118impl<Role, W, P, Key, Job, WorkerResult> KeyedOperating<Role, W, P, Key, Job, WorkerResult>
119where
120 Role: Eq,
121 W: Behavior + BehaviorBase,
122 P: ActivationPlan,
123 BehaviorAddr<W>: EndpointAddress,
124{
125 fn role_position(&self, role: &Role) -> Option<usize> {
126 self.roles.iter().position(|cell| cell.role() == role)
127 }
128
129 fn binding_expectation(&self, key: &Key) -> BindingExpectation
130 where
131 Key: Ord,
132 {
133 match self.bindings.binding(key) {
134 Some(evidence) => BindingExpectation::Exact(evidence.generation().clone()),
135 None => BindingExpectation::Absent,
136 }
137 }
138
139 fn creation_positions(
140 &self,
141 identities: impl IntoIterator<Item = (behavior::CreationId, CreationKind)>,
142 ) -> Option<Vec<usize>> {
143 direct_worker::ordered_creation_positions(
144 self.roles
145 .iter()
146 .map(|cell| cell.member.expected_creation()),
147 identities,
148 )
149 }
150
151 fn worker_position(&self, worker: &WorkerAttempt) -> Option<usize> {
152 self.roles
153 .iter()
154 .position(|cell| cell.member.worker_attempt() == Some(worker))
155 }
156}
157
158enum AdmissionBinding<'table, Key, Role> {
159 Retained {
160 submitted_key: Key,
161 evidence: BindingEvidence<Role>,
162 },
163 Reserved(ReservedBinding<'table, Key, Role>),
164}
165
166impl<Key, Role> AdmissionBinding<'_, Key, Role> {
167 fn commit(self) -> BindingEvidence<Role> {
168 match self {
169 Self::Retained {
170 submitted_key: _,
171 evidence,
172 } => evidence,
173 Self::Reserved(binding) => binding.commit(),
174 }
175 }
176
177 fn reject(self) -> Key {
178 match self {
179 Self::Retained { submitted_key, .. } => submitted_key,
180 Self::Reserved(binding) => binding.reject(),
181 }
182 }
183}
184
185enum KeyedPoolState<Role, W, P, Key, Job, WorkerResult>
186where
187 W: Behavior,
188 BehaviorAddr<W>: EndpointAddress,
189{
190 Constructed(Vec<PreparedWorker<Role, W, P>>),
191 Operating(KeyedOperating<Role, W, P, Key, Job, WorkerResult>),
192 Retiring {
193 workers: Vec<RetiringWorker<Role, W, P>>,
194 deadline: super::drain::ShutdownDeadline,
195 },
196 Stopped,
197 ForcedRetirement {
198 #[expect(
199 dead_code,
200 reason = "Bombay's retirement custodian receives every unresolved worker"
201 )]
202 workers: Vec<RetiringWorker<Role, W, P>>,
203 #[expect(
204 dead_code,
205 reason = "Bombay's retirement custodian receives the exact retirement reason"
206 )]
207 cause: super::drain::ForcedRetirementCause,
208 },
209}
210
211impl<Role, W, P, Key, Job, WorkerResult> KeyedPoolState<Role, W, P, Key, Job, WorkerResult>
212where
213 Key: Ord,
214 W: Behavior + BehaviorBase,
215 P: ActivationPlan,
216 BehaviorAddr<W>: EndpointAddress,
217 StopOnShutdown<W>:
218 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
219{
220 fn start(
221 self,
222 creations: &mut CreationSequence,
223 backlog: BacklogCapacity,
224 binding_capacity: BindingCapacity,
225 ) -> Result<
226 (
227 Self,
228 Vec<ObserveChild<W::Protocol, behavior::ChildHead>>,
229 Creations<CreateChild<BehaviorAddr<W>, StopOnShutdown<W>>>,
230 ),
231 (Self, KeyedError),
232 > {
233 let prepared = match self {
234 Self::Constructed(prepared) => prepared,
235 state => return Err((state, KeyedError::InitializationUnavailable)),
236 };
237 let mut ids = Vec::with_capacity(prepared.len());
238 for _ in 0..prepared.len() {
239 let Some(id) = creations.issue() else {
240 return Err((
241 Self::Constructed(prepared),
242 KeyedError::WorkerCreationsExhausted,
243 ));
244 };
245 ids.push(id);
246 }
247
248 let mut observations = Vec::with_capacity(prepared.len());
249 let mut workers = Creations::empty();
250 let roles = prepared
251 .into_iter()
252 .zip(ids)
253 .map(|(prepared, creation)| {
254 let (member, worker, observation) = Member::begin(prepared, creation);
255 workers.extend([worker]);
256 observations.push(observation);
257 RoleCell::new(member, backlog)
258 })
259 .collect();
260
261 Ok((
262 Self::Operating(KeyedOperating {
263 roles,
264 bindings: BindingTable::new(binding_capacity),
265 }),
266 observations,
267 workers,
268 ))
269 }
270}
271
272#[derive(Debug, Error)]
274pub enum KeyedError {
275 #[error("keyed pool initialization is no longer available")]
277 InitializationUnavailable,
278 #[error("keyed pool worker creation identifiers are exhausted")]
280 WorkerCreationsExhausted,
281}
282
283pub struct KeyedPool<Role, W, P, Source, Selector, Diagnostics, Key, Job, WorkerResult>
285where
286 W: Behavior,
287 BehaviorAddr<W>: EndpointAddress,
288{
289 state: KeyedPoolState<Role, W, P, Key, Job, WorkerResult>,
290 selector: Selector,
291 activation: ActivationPolicy,
292 recovery: PoolRecoveryState<Source>,
293 restarts: super::restart::RestartBudget,
294 backlog: BacklogCapacity,
295 binding_capacity: BindingCapacity,
296 interruption: Interruption,
297 actor_drain: ActorDrainPolicy,
298 diagnostics: DiagnosticDisposition<Diagnostics>,
299 creations: CreationSequence,
300 jobs: AcceptedJobSequence,
301 assignments: AssignmentSequence,
302 shutdowns: ShutdownSequence,
303 next_restart_timer: u64,
304}
305
306pub struct KeyedConstructionRejected<Factory, Selector, Role, W, P, Rejection, Source, Diagnostics>
308{
309 pub workers: InitialWorkerRejection<Factory, Role, W, P, Rejection>,
311 pub selector: Selector,
313 pub activation: ActivationPolicy,
315 pub recovery: PoolRecovery<Source>,
317 pub backlog: BacklogCapacity,
319 pub binding_capacity: BindingCapacity,
321 pub interruption: Interruption,
323 pub actor_drain: ActorDrainPolicy,
325 pub diagnostics: DiagnosticDisposition<Diagnostics>,
327}
328
329#[expect(
331 clippy::too_many_arguments,
332 reason = "every argument is one explicit KeyedPool policy pending the recorded builder comparison"
333)]
334pub fn keyed<
335 Factory,
336 Selector,
337 Role,
338 W,
339 P,
340 Rejection,
341 Source,
342 Diagnostics,
343 Key,
344 Job,
345 WorkerResult,
346>(
347 factory: Factory,
348 roles: OrderedRoles<Role>,
349 selector: Selector,
350 activation: ActivationPolicy,
351 recovery: PoolRecovery<Source>,
352 backlog: BacklogCapacity,
353 binding_capacity: BindingCapacity,
354 interruption: Interruption,
355 actor_drain: ActorDrainPolicy,
356 diagnostics: DiagnosticDisposition<Diagnostics>,
357) -> Result<
358 KeyedPool<Role, W, P, Source, Selector, Diagnostics, Key, Job, WorkerResult>,
359 KeyedConstructionRejected<Factory, Selector, Role, W, P, Rejection, Source, Diagnostics>,
360>
361where
362 Factory: FnMut(&Role) -> core::result::Result<WorkerSubmission<W, P>, Rejection>,
363 Selector: Fn(&Key) -> Role,
364 Role: Eq,
365 W: Behavior,
366 W::Protocol: Protocol<Msg = Assignment<Job>>,
367 W::Sends: CompletesAssignments<WorkerResult = WorkerResult>,
368 BehaviorAddr<W>: EndpointAddress,
369{
370 let prepared = match prepare_initial_workers(factory, roles) {
371 Ok(prepared) => prepared,
372 Err(workers) => {
373 return Err(KeyedConstructionRejected {
374 workers,
375 selector,
376 activation,
377 recovery,
378 backlog,
379 binding_capacity,
380 interruption,
381 actor_drain,
382 diagnostics,
383 });
384 }
385 };
386 Ok(KeyedPool {
387 state: KeyedPoolState::Constructed(prepared),
388 selector,
389 activation,
390 recovery: recovery.into(),
391 restarts: super::restart::RestartBudget::empty(),
392 backlog,
393 binding_capacity,
394 interruption,
395 actor_drain,
396 diagnostics,
397 creations: CreationSequence::new(),
398 jobs: AcceptedJobSequence::new(),
399 assignments: AssignmentSequence::new(),
400 shutdowns: ShutdownSequence::new(),
401 next_restart_timer: 1,
402 })
403}
404
405impl<Role, W, P, Source, Selector, Diagnostics, Key, Job, WorkerResult> BehaviorBase
406 for KeyedPool<Role, W, P, Source, Selector, Diagnostics, Key, Job, WorkerResult>
407where
408 Role: Send + Sync,
409 W: Behavior + Send,
410 W::Protocol: Protocol<Msg = Assignment<Job>>,
411 W::Sends: CompletesAssignments<WorkerResult = WorkerResult>,
412 P: ActivationPlan,
413 Source: WorkerSource<Role, W, P>,
414 BehaviorAddr<W>: EndpointAddress,
415 Job: Send,
416{
417 type Base = Self;
418
419 fn base(&self) -> &Self::Base {
420 self
421 }
422}
423
424impl<Role, W, P, Source, Selector, Diagnostics, Key, Job, WorkerResult>
425 KeyedPool<Role, W, P, Source, Selector, Diagnostics, Key, Job, WorkerResult>
426where
427 Role: Eq + Send + Sync,
428 W: Behavior + BehaviorBase + Send,
429 W::Protocol: Protocol<Msg = Assignment<Job>>,
430 W::Sends: CompletesAssignments<WorkerResult = WorkerResult>,
431 P: ActivationPlan,
432 Source: WorkerSource<Role, W, P>,
433 Selector: Fn(&Key) -> Role,
434 Diagnostics:
435 DiagnosticRoute<KeyedDiagnostic<Role, W, P, Source, Key, Job, WorkerResult>> + Clone,
436 BehaviorAddr<W>: EndpointAddress,
437 <BehaviorAddr<W> as Address>::Nonce: Send,
438 <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
439 StopOnShutdown<W>:
440 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
441 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
442 Key: Ord + Send,
443 Job: Clone + Send,
444 WorkerResult: Send,
445{
446 fn reject_submission(
447 &self,
448 operating: KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
449 submission: SubmissionId,
450 key: Key,
451 payload: Job,
452 customer: CustomerRoute<BehaviorAddr<W>, Key, Role, Job, WorkerResult>,
453 reason: KeyedAdmissionRejection,
454 ) -> (
455 KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
456 KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult>,
457 ) {
458 let mut actions: KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult> =
459 Actions::cont();
460 actions.sends.customer_outcomes = InterpreterRequests::one(CustomerDelivery::rejected(
461 customer,
462 KeyedOutcome::Rejected {
463 submission,
464 key,
465 payload,
466 reason,
467 },
468 ));
469 (operating, actions)
470 }
471
472 fn accept_submission(
473 &mut self,
474 mut operating: KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
475 submission: SubmissionId,
476 key: Key,
477 payload: Job,
478 customer: CustomerRoute<BehaviorAddr<W>, Key, Role, Job, WorkerResult>,
479 ) -> (
480 KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
481 KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult>,
482 ) {
483 let (position, proposed) = match operating.bindings.binding(&key) {
484 Some(evidence) => {
485 let evidence = evidence.clone();
486 let Some(position) = operating.role_position(evidence.role()) else {
487 return self.reject_submission(
488 operating,
489 submission,
490 key,
491 payload,
492 customer,
493 KeyedAdmissionRejection::RoleUnavailable,
494 );
495 };
496 (
497 position,
498 AdmissionBinding::Retained {
499 submitted_key: key,
500 evidence,
501 },
502 )
503 }
504 None => {
505 let selected = (self.selector)(&key);
506 let Some(position) = operating.role_position(&selected) else {
507 return self.reject_submission(
508 operating,
509 submission,
510 key,
511 payload,
512 customer,
513 KeyedAdmissionRejection::UnknownSelectedRole,
514 );
515 };
516 let role = operating.roles[position].member.role_name();
517 let proposed = match operating.bindings.reserve(key, role.clone()) {
518 Ok(binding) => AdmissionBinding::Reserved(binding),
519 Err(BindingReservationRejected::CapacityExhausted(key)) => {
520 return self.reject_submission(
521 operating,
522 submission,
523 key,
524 payload,
525 customer,
526 KeyedAdmissionRejection::BindingCapacityExhausted,
527 );
528 }
529 Err(BindingReservationRejected::GenerationsExhausted(key)) => {
530 return self.reject_submission(
531 operating,
532 submission,
533 key,
534 payload,
535 customer,
536 KeyedAdmissionRejection::BindingGenerationExhausted,
537 );
538 }
539 Err(BindingReservationRejected::AlreadyBound(key)) => {
540 return self.reject_submission(
541 operating,
542 submission,
543 key,
544 payload,
545 customer,
546 KeyedAdmissionRejection::RoleUnavailable,
547 );
548 }
549 };
550 (position, proposed)
551 }
552 };
553
554 let target = operating.roles[position].admission_target();
555 let rejection = match target {
556 AdmissionTarget::QueueFull => Some(KeyedAdmissionRejection::RoleBacklogFull),
557 AdmissionTarget::Unavailable => Some(KeyedAdmissionRejection::RoleUnavailable),
558 AdmissionTarget::Assign | AdmissionTarget::Queue => None,
559 };
560 if let Some(reason) = rejection {
561 let key = proposed.reject();
562 return self.reject_submission(operating, submission, key, payload, customer, reason);
563 }
564
565 let Some((job, admitted)) = self.jobs.issue() else {
566 let key = proposed.reject();
567 return self.reject_submission(
568 operating,
569 submission,
570 key,
571 payload,
572 customer,
573 KeyedAdmissionRejection::JobCorrelationExhausted,
574 );
575 };
576 let customer_target = customer.clone();
577
578 match target {
579 AdmissionTarget::Assign => {
580 let member = &mut operating.roles[position].member;
581 let super::pool::worker::MemberState::Worker(Worker { current, phase }) =
582 &mut member.state
583 else {
584 let key = proposed.reject();
585 return self.reject_submission(
586 operating,
587 submission,
588 key,
589 payload,
590 customer,
591 KeyedAdmissionRejection::RoleUnavailable,
592 );
593 };
594 let WorkerPhase::Idle = phase else {
595 let key = proposed.reject();
596 return self.reject_submission(
597 operating,
598 submission,
599 key,
600 payload,
601 customer,
602 KeyedAdmissionRejection::RoleUnavailable,
603 );
604 };
605 let Some((correlation, assignment)) =
606 self.assignments.assign(¤t.attempt, payload.clone())
607 else {
608 let key = proposed.reject();
609 return self.reject_submission(
610 operating,
611 submission,
612 key,
613 payload,
614 customer,
615 KeyedAdmissionRejection::AssignmentCorrelationExhausted,
616 );
617 };
618 let request = AssignWorker::new(current.recipient(), &correlation, assignment);
619 let binding = proposed.commit();
620 let obligation = CustomerJob {
621 id: job,
622 admitted,
623 payload,
624 customer: KeyedCustomer {
625 binding: binding.clone(),
626 route: customer,
627 },
628 };
629 *phase = WorkerPhase::Busy(AssignedJob::new(obligation, correlation));
630 let mut actions: KeyedActions<
631 Role,
632 W,
633 P,
634 Source,
635 Diagnostics,
636 Key,
637 Job,
638 WorkerResult,
639 > = Actions::cont();
640 actions.sends.customer_outcomes =
641 InterpreterRequests::one(CustomerDelivery::outcome(
642 customer_target,
643 KeyedOutcome::Accepted {
644 submission,
645 job,
646 binding,
647 },
648 ));
649 actions.sends.worker_assignments.send(request);
650 (operating, actions)
651 }
652 AdmissionTarget::Queue => {
653 let binding = proposed.commit();
654 let obligation = CustomerJob {
655 id: job,
656 admitted,
657 payload,
658 customer: KeyedCustomer {
659 binding: binding.clone(),
660 route: customer,
661 },
662 };
663 operating.roles[position].queue.insert(admitted, obligation);
664 let mut actions: KeyedActions<
665 Role,
666 W,
667 P,
668 Source,
669 Diagnostics,
670 Key,
671 Job,
672 WorkerResult,
673 > = Actions::cont();
674 actions.sends.customer_outcomes =
675 InterpreterRequests::one(CustomerDelivery::outcome(
676 customer_target,
677 KeyedOutcome::Accepted {
678 submission,
679 job,
680 binding,
681 },
682 ));
683 (operating, actions)
684 }
685 AdmissionTarget::QueueFull | AdmissionTarget::Unavailable => {
686 let key = proposed.reject();
687 self.reject_submission(
688 operating,
689 submission,
690 key,
691 payload,
692 customer,
693 KeyedAdmissionRejection::RoleUnavailable,
694 )
695 }
696 }
697 }
698
699 fn fill_role(
700 &mut self,
701 operating: &mut KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
702 position: usize,
703 actions: &mut KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult>,
704 ) {
705 let Some(cell) = operating.roles.get_mut(position) else {
706 return;
707 };
708 let MemberState::Worker(Worker {
709 current,
710 phase: phase @ WorkerPhase::Idle,
711 }) = &mut cell.member.state
712 else {
713 return;
714 };
715 let Some((_, customer)) = cell.queue.pop_first() else {
716 return;
717 };
718 let Some((correlation, assignment)) = self
719 .assignments
720 .assign(¤t.attempt, customer.payload.clone())
721 else {
722 let CustomerJob {
723 id,
724 admitted: _,
725 payload,
726 customer: KeyedCustomer { binding, route },
727 } = customer;
728 actions
729 .sends
730 .customer_outcomes
731 .append(InterpreterRequests::one(CustomerDelivery::outcome(
732 route,
733 KeyedOutcome::ReturnedQueued {
734 job: id,
735 binding,
736 payload,
737 reason: KeyedQueuedReturnReason::AssignmentCorrelationExhausted,
738 },
739 )));
740 return;
741 };
742 let request = AssignWorker::new(current.recipient(), &correlation, assignment);
743 *phase = WorkerPhase::Busy(AssignedJob::new(customer, correlation));
744 actions.sends.worker_assignments.send(request);
745 }
746
747 fn authorize_waiting(
748 &mut self,
749 operating: &mut KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
750 actions: &mut KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult>,
751 ) {
752 let occupied = operating
753 .roles
754 .iter()
755 .filter_map(|cell| match &cell.member.state {
756 MemberState::Worker(Worker {
757 phase:
758 WorkerPhase::ActivationDispatched { .. }
759 | WorkerPhase::Activating { .. },
760 ..
761 }) => Some(()),
762 MemberState::Creating(_)
763 | MemberState::Worker(_)
764 | MemberState::Recovering(_)
765 | MemberState::Retired => None,
766 })
767 .count();
768 let available = self.activation.maximum().saturating_sub(occupied);
769 for _ in 0..available {
770 let position = operating
771 .roles
772 .iter()
773 .enumerate()
774 .find_map(|(position, cell)| match &cell.member.state {
775 MemberState::Worker(Worker {
776 phase: WorkerPhase::WaitingForActivation { .. },
777 ..
778 }) => Some(position),
779 MemberState::Creating(_)
780 | MemberState::Worker(_)
781 | MemberState::Recovering(_)
782 | MemberState::Retired => None,
783 });
784 let Some(position) = position else {
785 return;
786 };
787 let RoleCell {
788 member,
789 queue,
790 capacity,
791 } = operating.roles.remove(position);
792 match member.authorize_activation() {
793 Ok((member, request)) => {
794 operating.roles.insert(
795 position,
796 RoleCell {
797 member,
798 queue,
799 capacity,
800 },
801 );
802 actions
803 .sends
804 .worker_activations
805 .append(InterpreterRequests::one(request));
806 }
807 Err(member) => {
808 operating.roles.insert(
809 position,
810 RoleCell {
811 member,
812 queue,
813 capacity,
814 },
815 );
816 return;
817 }
818 }
819 }
820 }
821
822 fn accept_creations(
823 &mut self,
824 mut operating: KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
825 workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
826 ) -> Result<
827 (
828 KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
829 KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult>,
830 ),
831 (
832 KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
833 CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
834 ),
835 > {
836 let settlements = match workers.into_settlement() {
837 CreationSettlement::Settled(settlements) => settlements,
838 settlement => {
839 return Err((operating, CreationsSettled::new(settlement)));
840 }
841 };
842 let identities = settlements
843 .iter()
844 .map(|settlement| match settlement {
845 SettledItem::Attempted(ItemSettlement::Accepted(
846 ChildCreationOutcome::Established(child),
847 )) => Some((child.id(), child.kind())),
848 SettledItem::Attempted(
849 ItemSettlement::Accepted(
850 ChildCreationOutcome::InitializationRejected { .. }
851 | ChildCreationOutcome::InitializationPanicked { .. }
852 | ChildCreationOutcome::HostRejected { .. },
853 )
854 | ItemSettlement::Rejected { .. }
855 | ItemSettlement::Corrupt { .. }
856 | ItemSettlement::Blocked { .. },
857 )
858 | SettledItem::Unattempted(_) => None,
859 })
860 .collect::<Option<Vec<_>>>();
861 let Some(positions) =
862 identities.and_then(|identities| operating.creation_positions(identities))
863 else {
864 return Err((
865 operating,
866 CreationsSettled::new(CreationSettlement::Settled(settlements)),
867 ));
868 };
869 let mut returned: BTreeMap<_, _> = positions.into_iter().zip(settlements).collect();
870 let mut roles = Vec::with_capacity(operating.roles.len());
871 let mut actions: KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult> =
872 Actions::cont();
873 for (position, cell) in operating.roles.into_iter().enumerate() {
874 let Some(settlement) = returned.remove(&position) else {
875 roles.push(cell);
876 continue;
877 };
878 let RoleCell {
879 member,
880 queue,
881 capacity,
882 } = cell;
883 match member.admit_successful_creation(ChildCreationSettled::new(settlement)) {
884 Ok((member, request)) => {
885 roles.push(RoleCell {
886 member,
887 queue,
888 capacity,
889 });
890 actions
891 .sends
892 .worker_initializations
893 .append(InterpreterRequests::one(request));
894 }
895 Err((member, creation)) => {
896 roles.push(RoleCell {
897 member,
898 queue,
899 capacity,
900 });
901 let returned = CreationsSettled::new(CreationSettlement::Settled(
902 Creations::one(creation.into_settlement()),
903 ));
904 actions.sends.diagnostics.append(InterpreterRequests::one(
905 self.diagnostics.action(KeyedDiagnostic::from(
906 protocol::KeyedDiagnosticCause::Unexpected(
907 KeyedEvent::WorkerCreationsSettled(returned),
908 ),
909 )),
910 ));
911 }
912 }
913 }
914 operating.roles = roles;
915 Ok((operating, actions))
916 }
917
918 fn accept_initialization(
919 &mut self,
920 mut operating: KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
921 input: WorkerInitializationReport<W, P>,
922 ) -> Result<
923 (
924 KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
925 KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult>,
926 ),
927 (
928 KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
929 WorkerInitializationReport<W, P>,
930 ),
931 > {
932 let Some(position) = operating.worker_position(input.worker()) else {
933 return Err((operating, input));
934 };
935 let RoleCell {
936 member,
937 queue,
938 capacity,
939 } = operating.roles.remove(position);
940 match member.admit_initialization(input) {
941 Ok(member) => {
942 operating.roles.insert(
943 position,
944 RoleCell {
945 member,
946 queue,
947 capacity,
948 },
949 );
950 let mut actions = Actions::cont();
951 self.authorize_waiting(&mut operating, &mut actions);
952 Ok((operating, actions))
953 }
954 Err((member, input)) => {
955 operating.roles.insert(
956 position,
957 RoleCell {
958 member,
959 queue,
960 capacity,
961 },
962 );
963 Err((operating, input))
964 }
965 }
966 }
967
968 fn accept_activation(
969 &mut self,
970 mut operating: KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
971 input: WorkerActivation<W, P>,
972 ) -> Result<
973 (
974 KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
975 KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult>,
976 ),
977 (
978 KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
979 WorkerActivation<W, P>,
980 ),
981 > {
982 let Some(position) = operating.worker_position(&input.worker()) else {
983 return Err((operating, input));
984 };
985 let RoleCell {
986 member,
987 queue,
988 capacity,
989 } = operating.roles.remove(position);
990 match member.admit_activation(input) {
991 Ok(member) => {
992 let mut actions = Actions::cont();
993 match &member.state {
994 MemberState::Worker(Worker {
995 phase: WorkerPhase::Idle,
996 ..
997 }) => {
998 operating.roles.insert(
999 position,
1000 RoleCell {
1001 member,
1002 queue,
1003 capacity,
1004 },
1005 );
1006 self.fill_role(&mut operating, position, &mut actions);
1007 }
1008 MemberState::Creating(_)
1009 | MemberState::Worker(_)
1010 | MemberState::Recovering(_)
1011 | MemberState::Retired => {
1012 operating.roles.insert(
1013 position,
1014 RoleCell {
1015 member,
1016 queue,
1017 capacity,
1018 },
1019 );
1020 }
1021 }
1022 self.authorize_waiting(&mut operating, &mut actions);
1023 Ok((operating, actions))
1024 }
1025 Err((member, input)) => {
1026 operating.roles.insert(
1027 position,
1028 RoleCell {
1029 member,
1030 queue,
1031 capacity,
1032 },
1033 );
1034 Err((operating, input))
1035 }
1036 }
1037 }
1038
1039 fn reject_binding(
1040 &self,
1041 operating: KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
1042 command: BindingCommand<BehaviorAddr<W>, Key, Role>,
1043 reason: BindingRejection,
1044 ) -> (
1045 KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
1046 KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult>,
1047 ) {
1048 let target = command.reply().clone();
1049 let mut actions: KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult> =
1050 Actions::cont();
1051 actions.sends.binding_replies = target.deliver(BindingReply::Rejected { command, reason });
1052 (operating, actions)
1053 }
1054
1055 fn accept_binding(
1056 &self,
1057 mut operating: KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
1058 command: BindingCommand<BehaviorAddr<W>, Key, Role>,
1059 ) -> (
1060 KeyedOperating<Role, W, P, Key, Job, WorkerResult>,
1061 KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult>,
1062 ) {
1063 let actual = operating.binding_expectation(command.key());
1064 if command.expectation() != &actual {
1065 return self.reject_binding(
1066 operating,
1067 command,
1068 BindingRejection::StaleExpectation { actual },
1069 );
1070 }
1071
1072 match command.change() {
1073 BindingChange::Rebalance(target) => {
1074 let Some(position) = operating.role_position(target) else {
1075 return self.reject_binding(
1076 operating,
1077 command,
1078 BindingRejection::UnknownTarget,
1079 );
1080 };
1081 if let ManagementTarget::Unavailable = operating.roles[position].management_target()
1082 {
1083 return self.reject_binding(
1084 operating,
1085 command,
1086 BindingRejection::TargetUnavailable,
1087 );
1088 }
1089 let target = operating.roles[position].member.role_name();
1090 match actual {
1091 BindingExpectation::Absent => {
1092 let (request, key, expected, change, reply) = command.into_parts();
1093 match operating.bindings.reserve(key, target) {
1094 Ok(binding) => {
1095 let current = binding.commit();
1096 let mut actions: KeyedActions<
1097 Role,
1098 W,
1099 P,
1100 Source,
1101 Diagnostics,
1102 Key,
1103 Job,
1104 WorkerResult,
1105 > = Actions::cont();
1106 actions.sends.binding_replies =
1107 reply.deliver(BindingReply::Bound { request, current });
1108 (operating, actions)
1109 }
1110 Err(rejection) => {
1111 let (key, reason) = match rejection {
1112 BindingReservationRejected::AlreadyBound(key) => {
1113 let actual = operating.binding_expectation(&key);
1114 (key, BindingRejection::StaleExpectation { actual })
1115 }
1116 BindingReservationRejected::CapacityExhausted(key) => {
1117 (key, BindingRejection::BindingCapacityExhausted)
1118 }
1119 BindingReservationRejected::GenerationsExhausted(key) => {
1120 (key, BindingRejection::GenerationExhausted)
1121 }
1122 };
1123 let command = BindingCommand::from_parts(
1124 request, key, expected, change, reply,
1125 );
1126 self.reject_binding(operating, command, reason)
1127 }
1128 }
1129 }
1130 BindingExpectation::Exact(_) => {
1131 let Some(current) = operating.bindings.binding(command.key()).cloned()
1132 else {
1133 return self.reject_binding(
1134 operating,
1135 command,
1136 BindingRejection::StaleExpectation {
1137 actual: BindingExpectation::Absent,
1138 },
1139 );
1140 };
1141 if current.role() == target.role() {
1142 let (request, _, _, _, reply) = command.into_parts();
1143 let mut actions: KeyedActions<
1144 Role,
1145 W,
1146 P,
1147 Source,
1148 Diagnostics,
1149 Key,
1150 Job,
1151 WorkerResult,
1152 > = Actions::cont();
1153 actions.sends.binding_replies =
1154 reply.deliver(BindingReply::Unchanged { request, current });
1155 return (operating, actions);
1156 }
1157 let Some((prior, current)) = operating
1158 .bindings
1159 .occupied(command.key())
1160 .and_then(|binding| binding.rebind(target))
1161 else {
1162 return self.reject_binding(
1163 operating,
1164 command,
1165 BindingRejection::GenerationExhausted,
1166 );
1167 };
1168 let (request, _, _, _, reply) = command.into_parts();
1169 let mut actions: KeyedActions<
1170 Role,
1171 W,
1172 P,
1173 Source,
1174 Diagnostics,
1175 Key,
1176 Job,
1177 WorkerResult,
1178 > = Actions::cont();
1179 actions.sends.binding_replies = reply.deliver(BindingReply::Rebalanced {
1180 request,
1181 prior,
1182 current,
1183 });
1184 (operating, actions)
1185 }
1186 }
1187 }
1188 BindingChange::Unbind => match actual {
1189 BindingExpectation::Absent => {
1190 let target = command.reply().clone();
1191 let mut actions: KeyedActions<
1192 Role,
1193 W,
1194 P,
1195 Source,
1196 Diagnostics,
1197 Key,
1198 Job,
1199 WorkerResult,
1200 > = Actions::cont();
1201 actions.sends.binding_replies =
1202 target.deliver(BindingReply::AlreadyUnbound { command });
1203 (operating, actions)
1204 }
1205 BindingExpectation::Exact(_) => {
1206 let Some((key, removed)) = operating.bindings.remove(command.key()) else {
1207 return self.reject_binding(
1208 operating,
1209 command,
1210 BindingRejection::StaleExpectation {
1211 actual: BindingExpectation::Absent,
1212 },
1213 );
1214 };
1215 let (request, _, _, _, reply) = command.into_parts();
1216 let mut actions: KeyedActions<
1217 Role,
1218 W,
1219 P,
1220 Source,
1221 Diagnostics,
1222 Key,
1223 Job,
1224 WorkerResult,
1225 > = Actions::cont();
1226 actions.sends.binding_replies = reply.deliver(BindingReply::Unbound {
1227 request,
1228 key,
1229 removed,
1230 });
1231 (operating, actions)
1232 }
1233 },
1234 }
1235 }
1236
1237 fn diagnose(
1238 &self,
1239 state: KeyedPoolState<Role, W, P, Key, Job, WorkerResult>,
1240 input: KeyedEvent<
1241 Role,
1242 W,
1243 P,
1244 Key,
1245 Job,
1246 WorkerResult,
1247 behavior::ActionItemResult<PrepareWorkers<Source, Role, W, P>>,
1248 WorkerPreparation<Source, Role, W, P>,
1249 >,
1250 ) -> (
1251 KeyedPoolState<Role, W, P, Key, Job, WorkerResult>,
1252 KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult>,
1253 ) {
1254 let mut actions: KeyedActions<Role, W, P, Source, Diagnostics, Key, Job, WorkerResult> =
1255 Actions::cont();
1256 actions.sends.diagnostics = InterpreterRequests::one(self.diagnostics.action(
1257 KeyedDiagnostic::from(protocol::KeyedDiagnosticCause::Unexpected(input)),
1258 ));
1259 (state, actions)
1260 }
1261}
1262
1263impl<Role, W, P, Source, Selector, Diagnostics, Key, Job, WorkerResult> Behavior
1264 for KeyedPool<Role, W, P, Source, Selector, Diagnostics, Key, Job, WorkerResult>
1265where
1266 Role: Eq + Send + Sync,
1267 W: Behavior + BehaviorBase + Send,
1268 W::Protocol: Protocol<Msg = Assignment<Job>>,
1269 W::Sends: CompletesAssignments<WorkerResult = WorkerResult>,
1270 P: ActivationPlan,
1271 Source: WorkerSource<Role, W, P>,
1272 Selector: Fn(&Key) -> Role + Send,
1273 Diagnostics:
1274 DiagnosticRoute<KeyedDiagnostic<Role, W, P, Source, Key, Job, WorkerResult>> + Clone,
1275 BehaviorAddr<W>: EndpointAddress,
1276 <BehaviorAddr<W> as Address>::Nonce: Send,
1277 <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
1278 StopOnShutdown<W>:
1279 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
1280 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
1281 Key: Ord + Send,
1282 Job: Clone + Send,
1283 WorkerResult: Send,
1284{
1285 type Protocol = MessageProtocol<
1286 BehaviorAddr<W>,
1287 KeyedCommand<BehaviorAddr<W>, Key, Role, Job, WorkerResult>,
1288 >;
1289 type Event = KeyedEvent<
1290 Role,
1291 W,
1292 P,
1293 Key,
1294 Job,
1295 WorkerResult,
1296 behavior::ActionItemResult<PrepareWorkers<Source, Role, W, P>>,
1297 WorkerPreparation<Source, Role, W, P>,
1298 >;
1299 type Sends = KeyedRequests<
1300 InterpreterRequests<ObserveChild<W::Protocol, behavior::ChildHead>>,
1301 InterpreterRequests<InitializeWorker<W, P>>,
1302 InterpreterRequests<BeginActivation<W, P>>,
1303 InterpreterRequests<
1304 CustomerDelivery<
1305 MessageProtocol<BehaviorAddr<W>, KeyedOutcome<Key, Role, Job, WorkerResult>>,
1306 >,
1307 >,
1308 <ManagementRoute<BehaviorAddr<W>, Key, Role> as DeliveryRoute>::Sends,
1309 SourceActions<AssignWorker<W::Protocol, Job>>,
1310 SourceActions<PrepareWorkers<Source, Role, W, P>>,
1311 SourceActions<ScheduleAfter>,
1312 InterpreterRequests<ShutdownEstablished<StopOnShutdown<W>, Here>>,
1313 InterpreterRequests<
1314 DiagnosticAction<
1315 Diagnostics,
1316 KeyedDiagnostic<Role, W, P, Source, Key, Job, WorkerResult>,
1317 >,
1318 >,
1319 >;
1320 type Ph = Never;
1321 type Error = KeyedError;
1322 type Birth = Births<StopOnShutdown<W>>;
1323
1324 fn init(&mut self, _: InitializationTurn) -> BehaviorActed<Self> {
1325 let state = mem::replace(&mut self.state, KeyedPoolState::Stopped);
1326 match state.start(&mut self.creations, self.backlog, self.binding_capacity) {
1327 Ok((state, observations, workers)) => {
1328 self.state = state;
1329 let mut sends = KeyedRequests::empty();
1330 sends.worker_observations = InterpreterRequests::new(observations);
1331 Ok(Actions::new(sends, workers, behavior::Step::Continue))
1332 }
1333 Err((state, error)) => {
1334 self.state = state;
1335 Err(error)
1336 }
1337 }
1338 }
1339
1340 fn transition(&mut self, _: ActiveTurn, input: Self::Event) -> BehaviorActed<Self> {
1341 let state = mem::replace(&mut self.state, KeyedPoolState::Stopped);
1342 let (state, actions) = match (state, input) {
1343 (state, KeyedEvent::Command(User { from, message })) => {
1344 match (state, message.into_request()) {
1345 (
1346 KeyedPoolState::Operating(operating),
1347 KeyedRequest::Submit {
1348 submission,
1349 key,
1350 payload,
1351 customer,
1352 },
1353 ) => {
1354 let (operating, actions) =
1355 self.accept_submission(operating, submission, key, payload, customer);
1356 (KeyedPoolState::Operating(operating), actions)
1357 }
1358 (KeyedPoolState::Operating(operating), KeyedRequest::Binding(command)) => {
1359 let (operating, actions) = self.accept_binding(operating, command);
1360 (KeyedPoolState::Operating(operating), actions)
1361 }
1362 (KeyedPoolState::Operating(operating), KeyedRequest::Shutdown) => {
1363 self.begin_retirement(operating, Actions::cont())
1364 }
1365 (state, request) => self.diagnose(
1366 state,
1367 KeyedEvent::Command(User::new(from, KeyedCommand::from_request(request))),
1368 ),
1369 }
1370 }
1371 (KeyedPoolState::Operating(operating), KeyedEvent::WorkerCreationsSettled(workers)) => {
1372 match self.accept_creations(operating, workers) {
1373 Ok((operating, actions)) => (KeyedPoolState::Operating(operating), actions),
1374 Err((operating, workers)) => self.diagnose(
1375 KeyedPoolState::Operating(operating),
1376 KeyedEvent::WorkerCreationsSettled(workers),
1377 ),
1378 }
1379 }
1380 (KeyedPoolState::Operating(operating), KeyedEvent::WorkerInitialization(input)) => {
1381 match self.accept_initialization(operating, input) {
1382 Ok((operating, actions)) => (KeyedPoolState::Operating(operating), actions),
1383 Err((operating, input)) => self.diagnose(
1384 KeyedPoolState::Operating(operating),
1385 KeyedEvent::WorkerInitialization(input),
1386 ),
1387 }
1388 }
1389 (KeyedPoolState::Operating(operating), KeyedEvent::WorkerActivationReported(input)) => {
1390 match self.accept_activation(operating, input) {
1391 Ok((operating, actions)) => (KeyedPoolState::Operating(operating), actions),
1392 Err((operating, input)) => self.diagnose(
1393 KeyedPoolState::Operating(operating),
1394 KeyedEvent::WorkerActivationReported(input),
1395 ),
1396 }
1397 }
1398 (
1399 KeyedPoolState::Operating(operating),
1400 KeyedEvent::AssignmentSettled(SettledItem::Attempted(ItemSettlement::Accepted(
1401 receipt,
1402 ))),
1403 ) => match self.accept_assignment_receipt(operating, receipt) {
1404 Ok(result) => result,
1405 Err((operating, receipt)) => self.diagnose(
1406 KeyedPoolState::Operating(operating),
1407 KeyedEvent::AssignmentSettled(SettledItem::Attempted(
1408 ItemSettlement::Accepted(receipt),
1409 )),
1410 ),
1411 },
1412 (
1413 KeyedPoolState::Operating(operating),
1414 KeyedEvent::AssignmentSettled(SettledItem::Attempted(ItemSettlement::Rejected {
1415 item,
1416 reason,
1417 })),
1418 ) => self.reject_assignment_delivery(operating, item, reason),
1419 (
1420 KeyedPoolState::Operating(operating),
1421 KeyedEvent::WorkerCompleted(ChildReport { child, report }),
1422 ) => match self.accept_completion(operating, child, report) {
1423 Ok(result) => result,
1424 Err((operating, completion)) => self.diagnose(
1425 KeyedPoolState::Operating(operating),
1426 KeyedEvent::WorkerCompleted(completion),
1427 ),
1428 },
1429 (KeyedPoolState::Operating(operating), KeyedEvent::WorkerStopped(stopped)) => {
1430 match self.accept_worker_stop(operating, stopped) {
1431 Ok(result) => result,
1432 Err((operating, stopped)) => self.diagnose(
1433 KeyedPoolState::Operating(operating),
1434 KeyedEvent::WorkerStopped(stopped),
1435 ),
1436 }
1437 }
1438 (
1439 KeyedPoolState::Operating(operating),
1440 KeyedEvent::WorkerShutdownSettled(settlement),
1441 ) => match self.accept_quarantine_shutdown(operating, settlement) {
1442 Ok(result) => result,
1443 Err((operating, settlement)) => self.diagnose(
1444 KeyedPoolState::Operating(operating),
1445 KeyedEvent::WorkerShutdownSettled(settlement),
1446 ),
1447 },
1448 (KeyedPoolState::Operating(operating), KeyedEvent::WorkerPreparationStarted(input)) => {
1449 match self.accept_worker_preparation_start(operating, input) {
1450 Ok(result) => result,
1451 Err((operating, input)) => self.diagnose(
1452 KeyedPoolState::Operating(operating),
1453 KeyedEvent::WorkerPreparationStarted(input),
1454 ),
1455 }
1456 }
1457 (
1458 KeyedPoolState::Operating(operating),
1459 KeyedEvent::WorkerPreparationReturned(input),
1460 ) => match self.accept_worker_preparation_return(operating, input) {
1461 Ok(result) => result,
1462 Err((operating, input)) => self.diagnose(
1463 KeyedPoolState::Operating(operating),
1464 KeyedEvent::WorkerPreparationReturned(input),
1465 ),
1466 },
1467 (KeyedPoolState::Operating(operating), KeyedEvent::RestartScheduleSettled(input)) => {
1468 match self.accept_restart_schedule(operating, input) {
1469 Ok(result) => result,
1470 Err((operating, input)) => self.diagnose(
1471 KeyedPoolState::Operating(operating),
1472 KeyedEvent::RestartScheduleSettled(input),
1473 ),
1474 }
1475 }
1476 (KeyedPoolState::Operating(operating), KeyedEvent::RestartElapsed(elapsed)) => {
1477 match self.accept_restart_timer(operating, elapsed) {
1478 Ok(result) => result,
1479 Err((operating, elapsed)) => self.diagnose(
1480 KeyedPoolState::Operating(operating),
1481 KeyedEvent::RestartElapsed(elapsed),
1482 ),
1483 }
1484 }
1485 (KeyedPoolState::Operating(operating), KeyedEvent::Shutdown(_)) => {
1486 self.begin_retirement(operating, Actions::cont())
1487 }
1488 (
1489 KeyedPoolState::Retiring { workers, deadline },
1490 KeyedEvent::WorkerCreationsSettled(creations),
1491 ) => self.accept_worker_creations(workers, deadline, creations),
1492 (
1493 KeyedPoolState::Retiring { workers, deadline },
1494 KeyedEvent::WorkerStopped(stopped),
1495 ) => self.accept_worker_exit(workers, deadline, stopped),
1496 (
1497 KeyedPoolState::Retiring { workers, deadline },
1498 KeyedEvent::WorkerShutdownSettled(settlement),
1499 ) => self.accept_worker_shutdown(workers, deadline, settlement),
1500 (
1501 KeyedPoolState::Retiring { workers, deadline },
1502 KeyedEvent::WorkerInitialization(input),
1503 ) => self.accept_worker_initialization(workers, deadline, input),
1504 (
1505 KeyedPoolState::Retiring { workers, deadline },
1506 KeyedEvent::WorkerActivationReported(input),
1507 ) => self.accept_worker_activation(workers, deadline, input),
1508 (
1509 KeyedPoolState::Retiring { workers, deadline },
1510 KeyedEvent::WorkerPreparationStarted(input),
1511 ) => self.accept_retired_preparation_start(workers, deadline, input),
1512 (
1513 KeyedPoolState::Retiring { workers, deadline },
1514 KeyedEvent::WorkerPreparationReturned(input),
1515 ) => self.accept_retired_preparation_return(workers, deadline, input),
1516 (
1517 KeyedPoolState::Retiring { workers, deadline },
1518 KeyedEvent::RestartScheduleSettled(settlement),
1519 ) => self.accept_retirement_schedule(workers, deadline, settlement),
1520 (
1521 KeyedPoolState::Retiring { workers, deadline },
1522 KeyedEvent::RestartElapsed(elapsed),
1523 ) => self.accept_deadline(workers, deadline, elapsed),
1524 (state, input) => self.diagnose(state, input),
1525 };
1526 self.state = state;
1527 Ok(actions)
1528 }
1529}