1use core::{num::NonZeroUsize, ops::ControlFlow};
4use std::collections::BTreeMap;
5use std::mem;
6
7use behavior::{
8 Actions, ActiveTurn, Address, Behavior, BehaviorActed, BehaviorAddr, BehaviorBase, Births,
9 ChildCreationSettled, ChildHead, ChildNamespaceExhausted, ChildReport, CreateChild,
10 CreationKind, CreationSequence, CreationSettlement, Creations, CreationsSettled,
11 EndpointAddress, Here, InitializationTurn, InjectEvent, InterpreterFault, InterpreterRequests,
12 ItemSettlement, MessageProtocol, Never, Protocol, SendEffects, SettledItem, SourceActions,
13 User,
14};
15use thiserror::Error;
16
17use crate::atomic::drain::{ForcedRetirementCause, ShutdownDeadline};
18use crate::{
19 DeliveryRoute, DiagnosticAction, DiagnosticDisposition, DiagnosticRoute,
20 EstablishedShutdownResolved, ObserveChild, ReplyRoute, ScheduleAfter, ShutdownEstablished,
21 ShutdownRequested, StopOnShutdown,
22};
23
24use super::RoleName;
25use super::worker::{InitialWorkerRejection, prepare_initial_workers};
26use super::{
27 ActivationPlan, ActivationPolicy, ActorDrainPolicy, BeginActivation, InitializeWorker,
28 OrderedRoles, PrepareWorkers, PreparedWorker, WorkerCreationRejection, WorkerPreparation,
29 WorkerSource, WorkerSubmission,
30};
31
32mod event;
33mod job;
34mod protocol;
35mod requests;
36
37pub use event::FifoEvent;
38pub use protocol::{
39 AdmissionRejection, AssignedReturnReason, FifoCommand, FifoDiagnostic, FifoOutcome,
40 FifoOutcomeKind, QueuedReturnReason,
41};
42pub use requests::FifoRequests;
43
44use super::pool::assignment::{
45 AcceptedJobSequence, AdmissionOrdinal, AssignedJob, AssignmentReceiptOutcome,
46 AssignmentRejectionOutcome, AssignmentSequence, AssignmentShutdown, CorrelationMatch,
47 CustomerJob, WorkerCompletionOutcome, WorkerExitOutcome,
48};
49use super::pool::worker as direct_worker;
50use super::pool::worker::{
51 Member, MemberState, PreparedReplacement, RetiringWorker, RetiringWorkerActivation,
52 RetiringWorkerInitialization, ShutdownJoin, Worker, WorkerCreationAdmission, WorkerCustody,
53 WorkerDeparture, WorkerPhase, WorkerPreparationError, WorkerPreparationStartFault,
54 WorkerRecoveryPreparation, WorkerReplacementError, WorkerReplacementRelease, creation_identity,
55};
56use super::pool::{
57 AssignWorker, Assignment, AssignmentReceipt, BacklogCapacity, CompletesAssignments, Completion,
58 Interruption, PoolFailureReaction, PoolRecovery, ShutdownSequence, SubmissionId,
59};
60use super::pool::{PoolRecoveryState, WorkerRecoveryDecision};
61use super::restart::{RecoveryRelease, RestartAdmission, RestartBudget, admit_restart};
62use super::schedule::ScheduleKey;
63use super::worker::{WorkerActivationOutcome, WorkerPreparationExpectation, stop_kind};
64use job::QueuedJob;
65use protocol::FifoDiagnosticCause;
66
67type FifoActions<Role, W, P, Source, DiagnosticRoute, Job, WorkerResult> = Actions<
68 BehaviorAddr<W>,
69 Never,
70 FifoRequests<
71 InterpreterRequests<ObserveChild<<W as Behavior>::Protocol, ChildHead>>,
72 InterpreterRequests<InitializeWorker<W, P>>,
73 InterpreterRequests<BeginActivation<W, P>>,
74 <CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult> as DeliveryRoute>::Sends,
75 SourceActions<AssignWorker<<W as Behavior>::Protocol, Job>>,
76 SourceActions<PrepareWorkers<Source, Role, W, P>>,
77 SourceActions<ScheduleAfter>,
78 InterpreterRequests<ShutdownEstablished<StopOnShutdown<W>, Here>>,
79 InterpreterRequests<
80 DiagnosticAction<
81 DiagnosticRoute,
82 FifoDiagnostic<Role, W, P, Source, Job, WorkerResult>,
83 >,
84 >,
85 >,
86 Births<StopOnShutdown<W>>,
87>;
88
89type CustomerRoute<A, Role, Job, WorkerResult> =
90 ReplyRoute<MessageProtocol<A, FifoOutcome<Role, Job, WorkerResult>>>;
91
92struct FifoOperating<Role, W, P, Job, WorkerResult>
93where
94 W: Behavior,
95 BehaviorAddr<W>: EndpointAddress,
96{
97 members: Vec<
98 Member<
99 Role,
100 W,
101 P,
102 Job,
103 WorkerResult,
104 CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult>,
105 >,
106 >,
107 backlog: BTreeMap<
108 AdmissionOrdinal,
109 QueuedJob<RoleName<Role>, Job, CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult>>,
110 >,
111 cursor: usize,
112}
113
114enum FifoDispatch<Role, W, Job, WorkerResult>
115where
116 W: Behavior,
117 W::Protocol: Protocol<Msg = Assignment<Job>>,
118 BehaviorAddr<W>: EndpointAddress,
119{
120 Assigned(AssignWorker<W::Protocol, Job>),
121 WorkerUnavailable(
122 QueuedJob<RoleName<Role>, Job, CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult>>,
123 ),
124 CorrelationUnavailable(
125 QueuedJob<RoleName<Role>, Job, CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult>>,
126 ),
127}
128
129#[derive(Clone, Copy)]
130enum CreationBatchRejection {
131 NamespaceExhausted,
132 InterpreterCorrupt(InterpreterFault),
133}
134
135impl CreationBatchRejection {
136 fn worker<W>(self, worker: W) -> WorkerCreationRejection<W>
137 where
138 W: Behavior,
139 {
140 match self {
141 Self::NamespaceExhausted => WorkerCreationRejection::NamespaceExhausted { worker },
142 Self::InterpreterCorrupt(fault) => {
143 WorkerCreationRejection::InterpreterCorrupt { worker, fault }
144 }
145 }
146 }
147
148 fn settlement<W>(
149 self,
150 creations: Creations<CreateChild<BehaviorAddr<W>, StopOnShutdown<W>>>,
151 ) -> CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>
152 where
153 W: Behavior + BehaviorBase,
154 BehaviorAddr<W>: EndpointAddress,
155 StopOnShutdown<W>:
156 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
157 {
158 match self {
159 Self::NamespaceExhausted => CreationsSettled::new(CreationSettlement::Rejected {
160 creations,
161 reason: ChildNamespaceExhausted,
162 }),
163 Self::InterpreterCorrupt(fault) => {
164 CreationsSettled::new(CreationSettlement::Corrupt { creations, fault })
165 }
166 }
167 }
168}
169
170impl<Role, W, P, Job, WorkerResult> FifoOperating<Role, W, P, Job, WorkerResult>
171where
172 W: Behavior + BehaviorBase,
173 W::Protocol: Protocol<Msg = Assignment<Job>>,
174 P: ActivationPlan,
175 BehaviorAddr<W>: EndpointAddress,
176 StopOnShutdown<W>:
177 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
178 Job: Clone,
179{
180 fn available_position(&self) -> Option<usize> {
181 if self.members.is_empty() {
182 return None;
183 }
184 (0..self.members.len()).find_map(|offset| {
185 let position = (self.cursor + offset) % self.members.len();
186 match &self.members[position].state {
187 MemberState::Worker(Worker {
188 phase: WorkerPhase::Idle,
189 ..
190 }) => Some(position),
191 MemberState::Creating(_)
192 | MemberState::Worker(_)
193 | MemberState::Recovering(_)
194 | MemberState::Retired => None,
195 }
196 })
197 }
198
199 fn creation_positions(
200 &self,
201 identities: impl IntoIterator<Item = (behavior::CreationId, CreationKind)>,
202 ) -> Option<Vec<usize>> {
203 direct_worker::ordered_creation_positions(
204 self.members.iter().map(Member::expected_creation),
205 identities,
206 )
207 }
208
209 fn recoverable_position<Source>(&self, recovery: &PoolRecoveryState<Source>) -> Option<usize> {
210 self.members.iter().position(|member| match &member.state {
211 MemberState::Retired => false,
212 MemberState::Worker(Worker {
213 phase: WorkerPhase::Stopping(_),
214 ..
215 })
216 | MemberState::Recovering(_) => match recovery {
217 PoolRecoveryState::Permanent { .. } | PoolRecoveryState::Transient { .. } => true,
218 PoolRecoveryState::Temporary { .. } => false,
219 },
220 MemberState::Creating(_) | MemberState::Worker(_) => true,
221 })
222 }
223
224 fn dispatch(
225 &mut self,
226 assignments: &mut AssignmentSequence,
227 queued: QueuedJob<
228 RoleName<Role>,
229 Job,
230 CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult>,
231 >,
232 ) -> FifoDispatch<Role, W, Job, WorkerResult> {
233 let Some(position) = self.available_position() else {
234 return FifoDispatch::WorkerUnavailable(queued);
235 };
236 let member = &mut self.members[position];
237 let MemberState::Worker(Worker { current, phase }) = &mut member.state else {
238 return FifoDispatch::WorkerUnavailable(queued);
239 };
240 let WorkerPhase::Idle = phase else {
241 return FifoDispatch::WorkerUnavailable(queued);
242 };
243 let Some((correlation, execution)) =
244 assignments.assign(¤t.attempt, queued.customer.payload.clone())
245 else {
246 return FifoDispatch::CorrelationUnavailable(queued);
247 };
248 let action = AssignWorker::new(current.recipient(), &correlation, execution);
249 *phase = WorkerPhase::Busy(AssignedJob::new(queued.customer, correlation));
250 self.cursor = position + 1;
251 FifoDispatch::Assigned(action)
252 }
253
254 fn assignment_receipt_position(&self, receipt: &AssignmentReceipt) -> Option<usize> {
255 self.members.iter().position(|member| match &member.state {
256 MemberState::Worker(Worker {
257 phase: WorkerPhase::Busy(assignment),
258 ..
259 }) => matches!(assignment.compare_receipt(receipt), CorrelationMatch::Exact),
260 MemberState::Creating(_)
261 | MemberState::Worker(_)
262 | MemberState::Recovering(_)
263 | MemberState::Retired => false,
264 })
265 }
266
267 fn completion_position(
268 &self,
269 child: behavior::CreationId,
270 completion: &Completion<WorkerResult>,
271 ) -> Option<usize> {
272 self.members.iter().position(|member| match &member.state {
273 MemberState::Worker(Worker {
274 current,
275 phase: WorkerPhase::Busy(assignment),
276 }) if current.attempt.creation() == child => matches!(
277 assignment.compare_completion(completion),
278 CorrelationMatch::Exact
279 ),
280 MemberState::Creating(_)
281 | MemberState::Worker(_)
282 | MemberState::Recovering(_)
283 | MemberState::Retired => false,
284 })
285 }
286
287 fn worker_position(&self, child: behavior::CreationId) -> Option<usize> {
288 self.members.iter().position(|member| match &member.state {
289 MemberState::Creating(worker) => worker.attempt.creation() == child,
290 MemberState::Worker(Worker { current, phase: _ }) => {
291 current.attempt.creation() == child
292 }
293 MemberState::Recovering(_) | MemberState::Retired => false,
294 })
295 }
296
297 fn worker_shutdown_position(&self, request: crate::ShutdownId) -> Option<usize> {
298 self.members.iter().position(|member| match &member.state {
299 MemberState::Worker(Worker {
300 phase: WorkerPhase::Stopping(join),
301 ..
302 }) => join.request() == request,
303 MemberState::Creating(_)
304 | MemberState::Worker(_)
305 | MemberState::Recovering(_)
306 | MemberState::Retired => false,
307 })
308 }
309
310 fn preparation_start_position<Source>(
311 &self,
312 input: &behavior::ActionItemResult<PrepareWorkers<Source, Role, W, P>>,
313 ) -> Option<usize>
314 where
315 Role: Send + Sync,
316 W: Send,
317 P: ActivationPlan,
318 Source: WorkerSource<Role, W, P>,
319 {
320 self.members.iter().position(|member| match &member.state {
321 MemberState::Recovering(direct_worker::RecoveringWorker::Preparing {
322 preparation: expected,
323 ..
324 }) => expected.accepts_issued(input),
325 MemberState::Creating(_)
326 | MemberState::Worker(_)
327 | MemberState::Recovering(_)
328 | MemberState::Retired => false,
329 })
330 }
331
332 fn preparation_return_position<Source>(
333 &self,
334 input: &WorkerPreparation<Source, Role, W, P>,
335 ) -> Option<usize>
336 where
337 Role: Send + Sync,
338 W: Send,
339 P: ActivationPlan,
340 Source: WorkerSource<Role, W, P>,
341 {
342 self.members.iter().position(|member| match &member.state {
343 MemberState::Recovering(direct_worker::RecoveringWorker::Preparing {
344 preparation: expected,
345 ..
346 }) => expected.accepts_return(input),
347 MemberState::Creating(_)
348 | MemberState::Worker(_)
349 | MemberState::Recovering(_)
350 | MemberState::Retired => false,
351 })
352 }
353
354 fn waiting_recovery_position(&self) -> Option<usize> {
355 self.members.iter().position(|member| {
356 matches!(
357 &member.state,
358 MemberState::Recovering(direct_worker::RecoveringWorker::WaitingForSource { .. })
359 )
360 })
361 }
362
363 fn restart_schedule_position(
364 &self,
365 input: &behavior::ActionItemResult<ScheduleAfter>,
366 ) -> Option<usize> {
367 self.members.iter().position(|member| match &member.state {
368 MemberState::Recovering(recovering) => recovering.accepts_restart_schedule(input),
369 MemberState::Creating(_) | MemberState::Worker(_) | MemberState::Retired => false,
370 })
371 }
372
373 fn restart_timer_position(&self, elapsed: &crate::TimerElapsed) -> Option<usize> {
374 self.members.iter().position(|member| match &member.state {
375 MemberState::Recovering(recovering) => recovering.accepts_restart_timer(elapsed),
376 MemberState::Creating(_) | MemberState::Worker(_) | MemberState::Retired => false,
377 })
378 }
379}
380
381impl<Role, W, P, Source, Diagnostics, Job, WorkerResult>
382 FifoPool<Role, W, P, Source, Diagnostics, Job, WorkerResult>
383where
384 Role: Send + Sync,
385 W: Behavior + BehaviorBase + Send,
386 W::Protocol: Protocol<Msg = Assignment<Job>>,
387 W::Sends: CompletesAssignments<WorkerResult = WorkerResult>,
388 P: ActivationPlan,
389 Source: WorkerSource<Role, W, P>,
390 Diagnostics: DiagnosticRoute<FifoDiagnostic<Role, W, P, Source, Job, WorkerResult>> + Clone,
391 BehaviorAddr<W>: EndpointAddress,
392 <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
393 StopOnShutdown<W>:
394 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
395 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
396 Job: Clone + Send,
397 WorkerResult: Send,
398{
399 fn accept_submission(
400 &mut self,
401 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
402 submission: SubmissionId,
403 payload: Job,
404 customer: CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult>,
405 ) -> (
406 FifoOperating<Role, W, P, Job, WorkerResult>,
407 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
408 ) {
409 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
410 Actions::cont();
411 match operating.recoverable_position(&self.recovery) {
412 Some(_) => {}
413 None => {
414 actions.sends.customer_outcomes = customer.deliver(FifoOutcome::rejected(
415 submission,
416 payload,
417 AdmissionRejection::NoRecoverableWorkers,
418 ));
419 return (operating, actions);
420 }
421 }
422
423 let ready = match operating.available_position() {
424 Some(position) => Some(position),
425 None => match operating.backlog.len().cmp(&self.backlog.maximum()) {
426 core::cmp::Ordering::Less => None,
427 core::cmp::Ordering::Equal | core::cmp::Ordering::Greater => {
428 actions.sends.customer_outcomes = customer.deliver(FifoOutcome::rejected(
429 submission,
430 payload,
431 AdmissionRejection::BacklogFull,
432 ));
433 return (operating, actions);
434 }
435 },
436 };
437
438 let Some((job, admitted)) = self.jobs.issue() else {
439 actions.sends.customer_outcomes = customer.deliver(FifoOutcome::rejected(
440 submission,
441 payload,
442 AdmissionRejection::JobCorrelationUnavailable,
443 ));
444 return (operating, actions);
445 };
446 let admission_customer = customer.clone();
447 let queued = QueuedJob {
448 customer: CustomerJob {
449 id: job,
450 admitted,
451 payload,
452 customer,
453 },
454 assigned_role: None,
455 };
456
457 match ready {
458 None => {
459 actions.sends.customer_outcomes =
460 admission_customer.deliver(FifoOutcome::accepted(submission, job));
461 operating.backlog.insert(admitted, queued);
462 }
463 Some(_) => match operating.dispatch(&mut self.assignments, queued) {
464 FifoDispatch::Assigned(assignment) => {
465 actions.sends.customer_outcomes =
466 admission_customer.deliver(FifoOutcome::accepted(submission, job));
467 actions.sends.worker_assignments.send(assignment);
468 }
469 FifoDispatch::WorkerUnavailable(queued) => {
470 match operating.backlog.len().cmp(&self.backlog.maximum()) {
471 core::cmp::Ordering::Less => {
472 actions.sends.customer_outcomes =
473 admission_customer.deliver(FifoOutcome::accepted(submission, job));
474 operating.backlog.insert(admitted, queued);
475 }
476 core::cmp::Ordering::Equal | core::cmp::Ordering::Greater => {
477 actions.sends.customer_outcomes =
478 queued.customer.customer.deliver(FifoOutcome::rejected(
479 submission,
480 queued.customer.payload,
481 AdmissionRejection::BacklogFull,
482 ));
483 }
484 }
485 }
486 FifoDispatch::CorrelationUnavailable(queued) => {
487 actions.sends.customer_outcomes =
488 queued.customer.customer.deliver(FifoOutcome::rejected(
489 submission,
490 queued.customer.payload,
491 AdmissionRejection::AssignmentCorrelationUnavailable,
492 ));
493 }
494 },
495 }
496 (operating, actions)
497 }
498
499 fn fill_fifo(
500 &mut self,
501 operating: &mut FifoOperating<Role, W, P, Job, WorkerResult>,
502 actions: &mut FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
503 ) {
504 loop {
505 match operating.available_position() {
506 Some(_) => {}
507 None => return,
508 }
509 let Some((admitted, queued)) = operating.backlog.pop_first() else {
510 return;
511 };
512 match operating.dispatch(&mut self.assignments, queued) {
513 FifoDispatch::Assigned(assignment) => {
514 actions.sends.worker_assignments.send(assignment);
515 }
516 FifoDispatch::WorkerUnavailable(queued) => {
517 operating.backlog.insert(admitted, queued);
518 return;
519 }
520 FifoDispatch::CorrelationUnavailable(queued) => {
521 let CustomerJob {
522 id,
523 admitted: _,
524 payload,
525 customer,
526 } = queued.customer;
527 let outcome = match queued.assigned_role {
528 None => FifoOutcome::returned_queued(
529 id,
530 payload,
531 QueuedReturnReason::AssignmentCorrelationUnavailable,
532 ),
533 Some(role) => FifoOutcome::returned_assigned(
534 id,
535 role,
536 payload,
537 AssignedReturnReason::RetryPreparationRejected,
538 ),
539 };
540 actions
541 .sends
542 .customer_outcomes
543 .append(customer.deliver(outcome));
544 }
545 }
546 }
547 }
548
549 fn return_unrecoverable_jobs(
550 &self,
551 operating: &mut FifoOperating<Role, W, P, Job, WorkerResult>,
552 actions: &mut FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
553 ) {
554 if operating.recoverable_position(&self.recovery).is_some() {
555 return;
556 }
557 for (_, queued) in mem::take(&mut operating.backlog) {
558 let CustomerJob {
559 id,
560 admitted: _,
561 payload,
562 customer,
563 } = queued.customer;
564 actions
565 .sends
566 .customer_outcomes
567 .append(customer.deliver(FifoOutcome::returned_queued(
568 id,
569 payload,
570 QueuedReturnReason::NoRecoverableWorkers,
571 )));
572 }
573 }
574
575 fn request_preparation(
576 role: RoleName<Role>,
577 recoveries: super::restart::RecoveryCount,
578 previous: super::WorkerAttempt,
579 stopped: crate::ChildStopped<BehaviorAddr<W>>,
580 source: Source,
581 mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
582 ) -> (
583 Member<
584 Role,
585 W,
586 P,
587 Job,
588 WorkerResult,
589 CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult>,
590 >,
591 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
592 ) {
593 let (ticket, request) = PrepareWorkers::new(source, role.clone(), Vec::new());
594 actions.sends.worker_preparations.send(request);
595 (
596 Member {
597 role,
598 recoveries,
599 state: MemberState::Recovering(direct_worker::RecoveringWorker::Preparing {
600 previous,
601 stopped,
602 preparation: WorkerPreparationExpectation::issued(ticket),
603 }),
604 },
605 actions,
606 )
607 }
608
609 fn prepare_next_waiting(
610 &mut self,
611 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
612 actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
613 ) -> (
614 FifoOperating<Role, W, P, Job, WorkerResult>,
615 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
616 ) {
617 let Some(position) = operating.waiting_recovery_position() else {
618 return (operating, actions);
619 };
620 let member = operating.members.remove(position);
621 let Member {
622 role,
623 recoveries,
624 state:
625 MemberState::Recovering(direct_worker::RecoveringWorker::WaitingForSource {
626 previous,
627 stopped,
628 }),
629 } = member
630 else {
631 operating.members.insert(position, member);
632 return (operating, actions);
633 };
634 let Some(source) = self.recovery.claim_waiting_source() else {
635 operating.members.insert(
636 position,
637 Member {
638 role,
639 recoveries,
640 state: MemberState::Recovering(
641 direct_worker::RecoveringWorker::WaitingForSource { previous, stopped },
642 ),
643 },
644 );
645 return (operating, actions);
646 };
647 let (member, actions) =
648 Self::request_preparation(role, recoveries, previous, stopped, source, actions);
649 operating.members.insert(position, member);
650 (operating, actions)
651 }
652
653 fn reject_replacement(
654 &mut self,
655 operating: FifoOperating<Role, W, P, Job, WorkerResult>,
656 position: usize,
657 role: RoleName<Role>,
658 recoveries: super::restart::RecoveryCount,
659 previous: super::WorkerAttempt,
660 stopped: crate::ChildStopped<BehaviorAddr<W>>,
661 submission: WorkerSubmission<W, P>,
662 returned_source: Option<Source>,
663 error: WorkerReplacementError,
664 mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
665 ) -> (
666 PoolState<Role, W, P, Job, WorkerResult>,
667 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
668 ) {
669 actions
670 .sends
671 .diagnostics
672 .append(InterpreterRequests::one(self.diagnostics.action(
673 FifoDiagnostic::new(FifoDiagnosticCause::WorkerReplacementFailed {
674 role: role.clone(),
675 previous,
676 stopped,
677 submission,
678 returned_source,
679 error,
680 }),
681 )));
682 self.retire_role_after_error(operating, position, role, recoveries, actions)
683 }
684
685 fn reject_preparation(
686 &mut self,
687 operating: FifoOperating<Role, W, P, Job, WorkerResult>,
688 position: usize,
689 role: RoleName<Role>,
690 recoveries: super::restart::RecoveryCount,
691 previous: super::WorkerAttempt,
692 stopped: crate::ChildStopped<BehaviorAddr<W>>,
693 source: Source,
694 error: WorkerPreparationError<Source::WorkerRejection, Source::SourceRejection>,
695 ) -> (
696 PoolState<Role, W, P, Job, WorkerResult>,
697 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
698 ) {
699 let (returned_source, error) = match self.recovery.restore_source(source) {
700 Ok(()) => (None, error),
701 Err(source) => (Some(source), WorkerPreparationError::SourceStateCorrupt),
702 };
703 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
704 Actions::cont();
705 actions
706 .sends
707 .diagnostics
708 .append(InterpreterRequests::one(self.diagnostics.action(
709 FifoDiagnostic::new(FifoDiagnosticCause::WorkerPreparationFailed {
710 role: role.clone(),
711 previous,
712 stopped,
713 returned_source,
714 error,
715 }),
716 )));
717 self.retire_role_after_error(operating, position, role, recoveries, actions)
718 }
719
720 fn retire_role_after_error(
721 &mut self,
722 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
723 position: usize,
724 role: RoleName<Role>,
725 recoveries: super::restart::RecoveryCount,
726 actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
727 ) -> (
728 PoolState<Role, W, P, Job, WorkerResult>,
729 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
730 ) {
731 operating.members.insert(
732 position,
733 Member {
734 role,
735 recoveries,
736 state: MemberState::Retired,
737 },
738 );
739 match self.recovery.failure() {
740 PoolFailureReaction::RetireRole => {
741 let (mut operating, mut actions) = self.prepare_next_waiting(operating, actions);
742 self.fill_fifo(&mut operating, &mut actions);
743 self.return_unrecoverable_jobs(&mut operating, &mut actions);
744 (PoolState::Operating(operating), actions)
745 }
746 PoolFailureReaction::StopPool => self.begin_shutdown(operating, actions),
747 }
748 }
749
750 fn return_source_after_failed_replacement(
751 &mut self,
752 operating: FifoOperating<Role, W, P, Job, WorkerResult>,
753 position: usize,
754 role: RoleName<Role>,
755 recoveries: super::restart::RecoveryCount,
756 previous: super::WorkerAttempt,
757 stopped: crate::ChildStopped<BehaviorAddr<W>>,
758 source: Source,
759 submission: WorkerSubmission<W, P>,
760 error: WorkerReplacementError,
761 ) -> (
762 PoolState<Role, W, P, Job, WorkerResult>,
763 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
764 ) {
765 let (returned_source, error) = match self.recovery.restore_source(source) {
766 Ok(()) => (None, error),
767 Err(source) => (Some(source), WorkerReplacementError::SourceStateCorrupt),
768 };
769 self.reject_replacement(
770 operating,
771 position,
772 role,
773 recoveries,
774 previous,
775 stopped,
776 submission,
777 returned_source,
778 error,
779 Actions::cont(),
780 )
781 }
782
783 fn recover_worker(
784 &mut self,
785 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
786 position: usize,
787 role: RoleName<Role>,
788 recoveries: super::restart::RecoveryCount,
789 previous: super::WorkerAttempt,
790 stopped: crate::ChildStopped<BehaviorAddr<W>>,
791 mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
792 ) -> (
793 PoolState<Role, W, P, Job, WorkerResult>,
794 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
795 ) {
796 let stop = stop_kind(&stopped.outcome);
797 match self.recovery.decide(stop) {
798 WorkerRecoveryDecision::PrepareWorker(source) => {
799 let (member, next_actions) =
800 Self::request_preparation(role, recoveries, previous, stopped, source, actions);
801 actions = next_actions;
802 operating.members.insert(position, member);
803 self.fill_fifo(&mut operating, &mut actions);
804 (PoolState::Operating(operating), actions)
805 }
806 WorkerRecoveryDecision::WaitForSource => {
807 operating.members.insert(
808 position,
809 Member {
810 role,
811 recoveries,
812 state: MemberState::Recovering(
813 direct_worker::RecoveringWorker::WaitingForSource { previous, stopped },
814 ),
815 },
816 );
817 self.fill_fifo(&mut operating, &mut actions);
818 (PoolState::Operating(operating), actions)
819 }
820 WorkerRecoveryDecision::RetireRole => {
821 operating.members.insert(
822 position,
823 Member {
824 role,
825 recoveries,
826 state: MemberState::Retired,
827 },
828 );
829 self.fill_fifo(&mut operating, &mut actions);
830 self.return_unrecoverable_jobs(&mut operating, &mut actions);
831 (PoolState::Operating(operating), actions)
832 }
833 WorkerRecoveryDecision::StopPool => {
834 operating.members.insert(
835 position,
836 Member {
837 role,
838 recoveries,
839 state: MemberState::Retired,
840 },
841 );
842 self.begin_shutdown(operating, actions)
843 }
844 }
845 }
846
847 fn continue_after_pre_ready_stop(
848 &mut self,
849 operating: FifoOperating<Role, W, P, Job, WorkerResult>,
850 position: usize,
851 role: RoleName<Role>,
852 recoveries: super::restart::RecoveryCount,
853 previous: super::WorkerAttempt,
854 stopped: crate::ChildStopped<BehaviorAddr<W>>,
855 actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
856 ) -> (
857 PoolState<Role, W, P, Job, WorkerResult>,
858 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
859 ) {
860 let (state, mut actions) = self.recover_worker(
861 operating, position, role, recoveries, previous, stopped, actions,
862 );
863 match state {
864 PoolState::Operating(mut operating) => {
865 self.authorize_waiting(&mut operating, &mut actions);
866 (PoolState::Operating(operating), actions)
867 }
868 state => (state, actions),
869 }
870 }
871
872 fn continue_after_returned_activation(
873 &mut self,
874 operating: FifoOperating<Role, W, P, Job, WorkerResult>,
875 position: usize,
876 role: RoleName<Role>,
877 recoveries: super::restart::RecoveryCount,
878 previous: super::WorkerAttempt,
879 stopped: crate::ChildStopped<BehaviorAddr<W>>,
880 returned_worker: super::WorkerAttempt,
881 returned_attempt: super::worker::WorkerActivationGrant<W, P>,
882 outcome: WorkerActivationOutcome<W, P>,
883 ) -> (
884 PoolState<Role, W, P, Job, WorkerResult>,
885 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
886 ) {
887 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
888 Actions::cont();
889 actions
890 .sends
891 .diagnostics
892 .append(InterpreterRequests::one(self.diagnostics.action(
893 FifoDiagnostic::new(FifoDiagnosticCause::WorkerActivationReturned {
894 role: role.clone(),
895 worker: returned_worker,
896 activation: returned_attempt,
897 outcome,
898 }),
899 )));
900 self.continue_after_pre_ready_stop(
901 operating, position, role, recoveries, previous, stopped, actions,
902 )
903 }
904
905 fn accept_worker_preparation_start(
906 &mut self,
907 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
908 input: behavior::ActionItemResult<PrepareWorkers<Source, Role, W, P>>,
909 ) -> Result<
910 (
911 PoolState<Role, W, P, Job, WorkerResult>,
912 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
913 ),
914 (
915 FifoOperating<Role, W, P, Job, WorkerResult>,
916 behavior::ActionItemResult<PrepareWorkers<Source, Role, W, P>>,
917 ),
918 >
919 where
920 Role: Eq,
921 {
922 let Some(position) = operating.preparation_start_position(&input) else {
923 return Err((operating, input));
924 };
925 let member = operating.members.remove(position);
926 match member.accept_preparation_start(input) {
927 Ok(ControlFlow::Continue(member)) => {
928 operating.members.insert(position, member);
929 Ok((PoolState::Operating(operating), Actions::cont()))
930 }
931 Ok(ControlFlow::Break(failed)) => {
932 let error = match failed.fault {
933 WorkerPreparationStartFault::InterpreterCorrupt(fault) => {
934 WorkerPreparationError::InterpreterCorrupt(fault)
935 }
936 WorkerPreparationStartFault::InterpretationSkipped => {
937 WorkerPreparationError::InterpretationSkipped
938 }
939 };
940 Ok(self.reject_preparation(
941 operating,
942 position,
943 failed.role,
944 failed.recoveries,
945 failed.previous,
946 failed.stopped,
947 failed.source,
948 error,
949 ))
950 }
951 Err((member, input)) => {
952 operating.members.insert(position, member);
953 Err((operating, input))
954 }
955 }
956 }
957
958 fn accept_worker_preparation_return(
959 &mut self,
960 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
961 input: WorkerPreparation<Source, Role, W, P>,
962 ) -> Result<
963 (
964 PoolState<Role, W, P, Job, WorkerResult>,
965 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
966 ),
967 (
968 FifoOperating<Role, W, P, Job, WorkerResult>,
969 WorkerPreparation<Source, Role, W, P>,
970 ),
971 >
972 where
973 Role: Eq,
974 {
975 let (limit, release) = match &self.recovery {
976 PoolRecoveryState::Permanent { limit, release, .. }
977 | PoolRecoveryState::Transient { limit, release, .. } => (*limit, *release),
978 PoolRecoveryState::Temporary { .. } => return Err((operating, input)),
979 };
980 let Some(position) = operating.preparation_return_position(&input) else {
981 return Err((operating, input));
982 };
983 let member = operating.members.remove(position);
984 let prepared = match member.accept_preparation_return(input) {
985 Ok(WorkerRecoveryPreparation::Ready(prepared)) => prepared,
986 Ok(WorkerRecoveryPreparation::Failed {
987 role,
988 recoveries,
989 previous,
990 stopped,
991 source,
992 error,
993 }) => {
994 return Ok(self.reject_preparation(
995 operating, position, role, recoveries, previous, stopped, source, error,
996 ));
997 }
998 Err((member, input)) => {
999 operating.members.insert(position, member);
1000 return Err((operating, input));
1001 }
1002 };
1003 let PreparedReplacement {
1004 role,
1005 mut recoveries,
1006 previous,
1007 stopped,
1008 source,
1009 submission,
1010 } = prepared;
1011 let budget = mem::replace(&mut self.restarts, RestartBudget::empty());
1012 match admit_restart(
1013 &mut recoveries,
1014 budget,
1015 limit,
1016 release,
1017 stopped.at,
1018 NonZeroUsize::MIN,
1019 ) {
1020 RestartAdmission::Proposed(proposal) => {
1021 let release = proposal.release();
1022 let release = match release {
1023 RecoveryRelease::Immediate => match self.creations.issue() {
1024 Some(creation) => Ok(WorkerReplacementRelease::Now(creation)),
1025 None => Err(WorkerReplacementError::WorkerCreationsExhausted),
1026 },
1027 RecoveryRelease::Delayed(delay) => match self.issue_restart_timer() {
1028 Some(timer) => Ok(WorkerReplacementRelease::After { timer, delay }),
1029 None => Err(WorkerReplacementError::RestartTimersExhausted),
1030 },
1031 };
1032 let release = match release {
1033 Ok(release) => release,
1034 Err(error) => {
1035 self.restarts = proposal.decline();
1036 return Ok(self.return_source_after_failed_replacement(
1037 operating, position, role, recoveries, previous, stopped, source,
1038 submission, error,
1039 ));
1040 }
1041 };
1042 if let Err(source) = self.recovery.restore_source(source) {
1043 self.restarts = proposal.decline();
1044 return Ok(self.reject_replacement(
1045 operating,
1046 position,
1047 role,
1048 recoveries,
1049 previous,
1050 stopped,
1051 submission,
1052 Some(source),
1053 WorkerReplacementError::SourceStateCorrupt,
1054 Actions::cont(),
1055 ));
1056 }
1057 self.restarts = proposal.accept();
1058 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
1059 Actions::cont();
1060 match release {
1061 WorkerReplacementRelease::Now(creation) => {
1062 let prepared = PreparedWorker { role, submission };
1063 let (member, worker, observation) = Member::begin_replacement(
1064 prepared,
1065 recoveries,
1066 previous.creation(),
1067 creation,
1068 );
1069 operating.members.insert(position, member);
1070 actions.creates.extend([worker]);
1071 actions
1072 .sends
1073 .worker_observations
1074 .append(InterpreterRequests::one(observation));
1075 }
1076 WorkerReplacementRelease::After { timer, delay } => {
1077 operating.members.insert(
1078 position,
1079 Member {
1080 role,
1081 recoveries,
1082 state: MemberState::Recovering(
1083 direct_worker::RecoveringWorker::Scheduling {
1084 previous,
1085 stopped,
1086 submission,
1087 timer,
1088 },
1089 ),
1090 },
1091 );
1092 actions.sends.restart_schedules.send(timer.after(delay));
1093 }
1094 }
1095 let (operating, actions) = self.prepare_next_waiting(operating, actions);
1096 Ok((PoolState::Operating(operating), actions))
1097 }
1098 RestartAdmission::Denied { budget, reason } => {
1099 self.restarts = budget;
1100 Ok(self.return_source_after_failed_replacement(
1101 operating,
1102 position,
1103 role,
1104 recoveries,
1105 previous,
1106 stopped,
1107 source,
1108 submission,
1109 WorkerReplacementError::RestartDenied(reason),
1110 ))
1111 }
1112 }
1113 }
1114
1115 fn accept_restart_schedule(
1116 &mut self,
1117 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
1118 input: behavior::ActionItemResult<ScheduleAfter>,
1119 ) -> Result<
1120 (
1121 PoolState<Role, W, P, Job, WorkerResult>,
1122 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
1123 ),
1124 (
1125 FifoOperating<Role, W, P, Job, WorkerResult>,
1126 behavior::ActionItemResult<ScheduleAfter>,
1127 ),
1128 > {
1129 let Some(position) = operating.restart_schedule_position(&input) else {
1130 return Err((operating, input));
1131 };
1132 let member = operating.members.remove(position);
1133 let Member {
1134 role,
1135 recoveries,
1136 state,
1137 } = member;
1138 let MemberState::Recovering(worker) = state else {
1139 operating.members.insert(
1140 position,
1141 Member {
1142 role,
1143 recoveries,
1144 state,
1145 },
1146 );
1147 return Err((operating, input));
1148 };
1149 match worker.admit_restart_schedule(input) {
1150 Ok(ControlFlow::Continue(worker)) => {
1151 operating.members.insert(
1152 position,
1153 Member {
1154 role,
1155 recoveries,
1156 state: MemberState::Recovering(worker),
1157 },
1158 );
1159 Ok((PoolState::Operating(operating), Actions::cont()))
1160 }
1161 Ok(ControlFlow::Break((replacement, settlement))) => {
1162 let direct_worker::PendingWorkerReplacement {
1163 previous,
1164 stopped,
1165 submission,
1166 } = replacement;
1167 Ok(self.reject_replacement(
1168 operating,
1169 position,
1170 role,
1171 recoveries,
1172 previous,
1173 stopped,
1174 submission,
1175 None,
1176 WorkerReplacementError::RestartScheduleReturned(settlement),
1177 Actions::cont(),
1178 ))
1179 }
1180 Err((worker, settlement)) => {
1181 operating.members.insert(
1182 position,
1183 Member {
1184 role,
1185 recoveries,
1186 state: MemberState::Recovering(worker),
1187 },
1188 );
1189 Err((operating, settlement))
1190 }
1191 }
1192 }
1193
1194 fn accept_restart_timer(
1195 &mut self,
1196 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
1197 elapsed: crate::TimerElapsed,
1198 ) -> Result<
1199 (
1200 PoolState<Role, W, P, Job, WorkerResult>,
1201 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
1202 ),
1203 (
1204 FifoOperating<Role, W, P, Job, WorkerResult>,
1205 crate::TimerElapsed,
1206 ),
1207 > {
1208 let Some(position) = operating.restart_timer_position(&elapsed) else {
1209 return Err((operating, elapsed));
1210 };
1211 let member = operating.members.remove(position);
1212 let Member {
1213 role,
1214 recoveries,
1215 state,
1216 } = member;
1217 let MemberState::Recovering(worker) = state else {
1218 operating.members.insert(
1219 position,
1220 Member {
1221 role,
1222 recoveries,
1223 state,
1224 },
1225 );
1226 return Err((operating, elapsed));
1227 };
1228 let replacement = match worker.admit_restart_timer(elapsed) {
1229 Ok(replacement) => replacement,
1230 Err((worker, elapsed)) => {
1231 operating.members.insert(
1232 position,
1233 Member {
1234 role,
1235 recoveries,
1236 state: MemberState::Recovering(worker),
1237 },
1238 );
1239 return Err((operating, elapsed));
1240 }
1241 };
1242 let direct_worker::PendingWorkerReplacement {
1243 previous,
1244 stopped,
1245 submission,
1246 } = replacement;
1247 let Some(creation) = self.creations.issue() else {
1248 return Ok(self.reject_replacement(
1249 operating,
1250 position,
1251 role,
1252 recoveries,
1253 previous,
1254 stopped,
1255 submission,
1256 None,
1257 WorkerReplacementError::WorkerCreationsExhausted,
1258 Actions::cont(),
1259 ));
1260 };
1261 let prepared = PreparedWorker { role, submission };
1262 let (member, worker, observation) =
1263 Member::begin_replacement(prepared, recoveries, previous.creation(), creation);
1264 operating.members.insert(position, member);
1265 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
1266 Actions::cont();
1267 actions.creates.extend([worker]);
1268 actions
1269 .sends
1270 .worker_observations
1271 .append(InterpreterRequests::one(observation));
1272 Ok((PoolState::Operating(operating), actions))
1273 }
1274
1275 fn accept_assignment_receipt(
1276 &mut self,
1277 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
1278 receipt: AssignmentReceipt,
1279 ) -> Result<
1280 (
1281 PoolState<Role, W, P, Job, WorkerResult>,
1282 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
1283 ),
1284 (
1285 FifoOperating<Role, W, P, Job, WorkerResult>,
1286 AssignmentReceipt,
1287 ),
1288 > {
1289 let Some(position) = operating.assignment_receipt_position(&receipt) else {
1290 return Err((operating, receipt));
1291 };
1292 let Member {
1293 role,
1294 recoveries,
1295 state,
1296 } = operating.members.remove(position);
1297 let MemberState::Worker(Worker {
1298 current,
1299 phase: WorkerPhase::Busy(assignment),
1300 }) = state
1301 else {
1302 operating.members.insert(
1303 position,
1304 Member {
1305 role,
1306 recoveries,
1307 state,
1308 },
1309 );
1310 return Err((operating, receipt));
1311 };
1312 let acceptance = match assignment.accept_receipt(receipt) {
1313 Ok(acceptance) => acceptance,
1314 Err((assignment, receipt)) => {
1315 operating.members.insert(
1316 position,
1317 Member {
1318 role,
1319 recoveries,
1320 state: MemberState::Worker(Worker {
1321 current,
1322 phase: WorkerPhase::Busy(assignment),
1323 }),
1324 },
1325 );
1326 return Err((operating, receipt));
1327 }
1328 };
1329 match acceptance {
1330 AssignmentReceiptOutcome::AwaitingCompletion(assignment) => {
1331 operating.members.insert(
1332 position,
1333 Member {
1334 role,
1335 recoveries,
1336 state: MemberState::Worker(Worker {
1337 current,
1338 phase: WorkerPhase::Busy(assignment),
1339 }),
1340 },
1341 );
1342 Ok((PoolState::Operating(operating), Actions::cont()))
1343 }
1344 AssignmentReceiptOutcome::JobCompleted {
1345 customer,
1346 result,
1347 stopped,
1348 } => Ok(self.complete_assignment(
1349 operating, position, role, recoveries, current, customer, result, stopped,
1350 )),
1351 AssignmentReceiptOutcome::JobInterrupted {
1352 customer,
1353 stopped,
1354 late_completion,
1355 } => Ok(self.interrupt_assignment(
1356 operating,
1357 position,
1358 role,
1359 recoveries,
1360 current,
1361 customer,
1362 stopped,
1363 late_completion,
1364 )),
1365 }
1366 }
1367
1368 fn accept_completion(
1369 &mut self,
1370 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
1371 child: behavior::CreationId,
1372 completion: Completion<WorkerResult>,
1373 ) -> Result<
1374 (
1375 PoolState<Role, W, P, Job, WorkerResult>,
1376 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
1377 ),
1378 (
1379 FifoOperating<Role, W, P, Job, WorkerResult>,
1380 ChildReport<Completion<WorkerResult>>,
1381 ),
1382 > {
1383 let Some(position) = operating.completion_position(child, &completion) else {
1384 return Err((operating, ChildReport::new(child, completion)));
1385 };
1386 let Member {
1387 role,
1388 recoveries,
1389 state,
1390 } = operating.members.remove(position);
1391 let MemberState::Worker(Worker {
1392 current,
1393 phase: WorkerPhase::Busy(assignment),
1394 }) = state
1395 else {
1396 operating.members.insert(
1397 position,
1398 Member {
1399 role,
1400 recoveries,
1401 state,
1402 },
1403 );
1404 return Err((operating, ChildReport::new(child, completion)));
1405 };
1406 let admission = match assignment.accept_completion(completion) {
1407 Ok(admission) => admission,
1408 Err((assignment, completion)) => {
1409 operating.members.insert(
1410 position,
1411 Member {
1412 role,
1413 recoveries,
1414 state: MemberState::Worker(Worker {
1415 current,
1416 phase: WorkerPhase::Busy(assignment),
1417 }),
1418 },
1419 );
1420 return Err((operating, ChildReport::new(child, completion)));
1421 }
1422 };
1423 match admission {
1424 WorkerCompletionOutcome::AwaitingReceipt(assignment) => {
1425 operating.members.insert(
1426 position,
1427 Member {
1428 role,
1429 recoveries,
1430 state: MemberState::Worker(Worker {
1431 current,
1432 phase: WorkerPhase::Busy(assignment),
1433 }),
1434 },
1435 );
1436 Ok((PoolState::Operating(operating), Actions::cont()))
1437 }
1438 WorkerCompletionOutcome::JobCompleted {
1439 customer,
1440 result,
1441 stopped,
1442 } => Ok(self.complete_assignment(
1443 operating, position, role, recoveries, current, customer, result, stopped,
1444 )),
1445 }
1446 }
1447
1448 fn request_worker_shutdown(
1449 &self,
1450 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
1451 position: usize,
1452 role: RoleName<Role>,
1453 recoveries: super::restart::RecoveryCount,
1454 current: direct_worker::CurrentWorker<W>,
1455 shutdown: crate::ShutdownId,
1456 mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
1457 ) -> (
1458 FifoOperating<Role, W, P, Job, WorkerResult>,
1459 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
1460 ) {
1461 actions
1462 .sends
1463 .worker_shutdowns
1464 .append(InterpreterRequests::one(current.shutdown_request(shutdown)));
1465 operating.members.insert(
1466 position,
1467 Member {
1468 role,
1469 recoveries,
1470 state: MemberState::Worker(Worker {
1471 current,
1472 phase: WorkerPhase::Stopping(ShutdownJoin::AwaitingBoth(shutdown)),
1473 }),
1474 },
1475 );
1476 (operating, actions)
1477 }
1478
1479 fn reject_assignment_delivery(
1480 &mut self,
1481 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
1482 item: AssignWorker<W::Protocol, Job>,
1483 reason: behavior::ExactDeliveryReason,
1484 ) -> (
1485 PoolState<Role, W, P, Job, WorkerResult>,
1486 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
1487 ) {
1488 let receipt = item.receipt();
1489 let Some(position) = operating.assignment_receipt_position(&receipt) else {
1490 return self.diagnose(
1491 PoolState::Operating(operating),
1492 FifoEvent::AssignmentSettled(SettledItem::Attempted(ItemSettlement::Rejected {
1493 item,
1494 reason,
1495 })),
1496 );
1497 };
1498 let Some(shutdown) = self
1499 .shutdowns
1500 .reserve(1)
1501 .and_then(|reserved| reserved.into_iter().next())
1502 else {
1503 let (state, actions) = self.diagnose(
1504 PoolState::Operating(operating),
1505 FifoEvent::AssignmentSettled(SettledItem::Attempted(ItemSettlement::Rejected {
1506 item,
1507 reason,
1508 })),
1509 );
1510 return match state {
1511 PoolState::Operating(operating) => self.begin_shutdown(operating, actions),
1512 state => (state, actions),
1513 };
1514 };
1515 let (target, returned, receipt) = item.into_parts();
1516 let Member {
1517 role,
1518 recoveries,
1519 state,
1520 } = operating.members.remove(position);
1521 let MemberState::Worker(Worker {
1522 current,
1523 phase: WorkerPhase::Busy(assignment),
1524 }) = state
1525 else {
1526 operating.members.insert(
1527 position,
1528 Member {
1529 role,
1530 recoveries,
1531 state,
1532 },
1533 );
1534 return self.diagnose(
1535 PoolState::Operating(operating),
1536 FifoEvent::AssignmentSettled(SettledItem::Attempted(ItemSettlement::Rejected {
1537 item: AssignWorker::returned(target, returned, receipt),
1538 reason,
1539 })),
1540 );
1541 };
1542 let rejected = match assignment.accept_rejection(returned) {
1543 Ok(rejected) => rejected,
1544 Err((assignment, returned)) => {
1545 operating.members.insert(
1546 position,
1547 Member {
1548 role,
1549 recoveries,
1550 state: MemberState::Worker(Worker {
1551 current,
1552 phase: WorkerPhase::Busy(assignment),
1553 }),
1554 },
1555 );
1556 return self.diagnose(
1557 PoolState::Operating(operating),
1558 FifoEvent::AssignmentSettled(SettledItem::Attempted(
1559 ItemSettlement::Rejected {
1560 item: AssignWorker::returned(target, returned, receipt),
1561 reason,
1562 },
1563 )),
1564 );
1565 }
1566 };
1567 match rejected {
1568 AssignmentRejectionOutcome::JobReturned { customer, stopped } => {
1569 operating.backlog.insert(
1570 customer.admitted,
1571 QueuedJob {
1572 customer,
1573 assigned_role: Some(role.clone()),
1574 },
1575 );
1576 let actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
1577 Actions::cont();
1578 match stopped {
1579 None => {
1580 let (mut operating, mut actions) = self.request_worker_shutdown(
1581 operating, position, role, recoveries, current, shutdown, actions,
1582 );
1583 self.fill_fifo(&mut operating, &mut actions);
1584 (PoolState::Operating(operating), actions)
1585 }
1586 Some(stopped) => self.recover_worker(
1587 operating,
1588 position,
1589 role,
1590 recoveries,
1591 current.attempt,
1592 stopped,
1593 actions,
1594 ),
1595 }
1596 }
1597 AssignmentRejectionOutcome::ConflictingCompletion {
1598 customer,
1599 assignment,
1600 completion,
1601 stopped,
1602 } => {
1603 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
1604 Actions::cont();
1605 actions.sends.customer_outcomes =
1606 customer.customer.deliver(FifoOutcome::returned_assigned(
1607 customer.id,
1608 role.clone(),
1609 customer.payload,
1610 AssignedReturnReason::ContradictoryAssignmentSettlement,
1611 ));
1612 actions.sends.diagnostics.append(InterpreterRequests::one(
1613 self.diagnostics
1614 .action(FifoDiagnostic::new(FifoDiagnosticCause::Unexpected(
1615 FifoEvent::AssignmentSettled(SettledItem::Attempted(
1616 ItemSettlement::Rejected {
1617 item: AssignWorker::returned(target, assignment, receipt),
1618 reason,
1619 },
1620 )),
1621 ))),
1622 ));
1623 actions.sends.diagnostics.append(InterpreterRequests::one(
1624 self.diagnostics
1625 .action(FifoDiagnostic::new(FifoDiagnosticCause::Unexpected(
1626 FifoEvent::WorkerCompleted(ChildReport::new(
1627 current.attempt.creation(),
1628 completion,
1629 )),
1630 ))),
1631 ));
1632 match stopped {
1633 None => {
1634 let (operating, actions) = self.request_worker_shutdown(
1635 operating, position, role, recoveries, current, shutdown, actions,
1636 );
1637 self.begin_shutdown(operating, actions)
1638 }
1639 Some(stopped) => {
1640 let (state, actions) = self.recover_worker(
1641 operating,
1642 position,
1643 role,
1644 recoveries,
1645 current.attempt,
1646 stopped,
1647 actions,
1648 );
1649 match state {
1650 PoolState::Operating(operating) => {
1651 self.begin_shutdown(operating, actions)
1652 }
1653 state => (state, actions),
1654 }
1655 }
1656 }
1657 }
1658 }
1659 }
1660
1661 fn accept_worker_stop(
1662 &mut self,
1663 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
1664 stopped: crate::ChildStopped<BehaviorAddr<W>>,
1665 ) -> Result<
1666 (
1667 PoolState<Role, W, P, Job, WorkerResult>,
1668 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
1669 ),
1670 (
1671 FifoOperating<Role, W, P, Job, WorkerResult>,
1672 crate::ChildStopped<BehaviorAddr<W>>,
1673 ),
1674 > {
1675 let Some(position) = operating.worker_position(stopped.child) else {
1676 return Err((operating, stopped));
1677 };
1678 let Member {
1679 role,
1680 recoveries,
1681 state,
1682 } = operating.members.remove(position);
1683 let (current, phase) = match state {
1684 MemberState::Creating(mut worker) => {
1685 if worker.stopped.is_none() {
1686 worker.stopped = Some(stopped);
1687 operating.members.insert(
1688 position,
1689 Member {
1690 role,
1691 recoveries,
1692 state: MemberState::Creating(worker),
1693 },
1694 );
1695 return Ok((PoolState::Operating(operating), Actions::cont()));
1696 }
1697 operating.members.insert(
1698 position,
1699 Member {
1700 role,
1701 recoveries,
1702 state: MemberState::Creating(worker),
1703 },
1704 );
1705 return Err((operating, stopped));
1706 }
1707 MemberState::Worker(Worker { current, phase }) => (current, phase),
1708 state @ (MemberState::Recovering(_) | MemberState::Retired) => {
1709 operating.members.insert(
1710 position,
1711 Member {
1712 role,
1713 recoveries,
1714 state,
1715 },
1716 );
1717 return Err((operating, stopped));
1718 }
1719 };
1720 match phase {
1721 WorkerPhase::Initializing {
1722 initialization,
1723 stopped: None,
1724 } => {
1725 operating.members.insert(
1726 position,
1727 Member {
1728 role,
1729 recoveries,
1730 state: MemberState::Worker(Worker {
1731 current,
1732 phase: WorkerPhase::Initializing {
1733 initialization,
1734 stopped: Some(stopped),
1735 },
1736 }),
1737 },
1738 );
1739 Ok((PoolState::Operating(operating), Actions::cont()))
1740 }
1741 WorkerPhase::WaitingForActivation { permit, activation } => {
1742 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
1743 Actions::cont();
1744 actions.sends.diagnostics.append(InterpreterRequests::one(
1745 self.diagnostics.action(FifoDiagnostic::new(
1746 FifoDiagnosticCause::UnusedActivation {
1747 role: role.clone(),
1748 permit: Some(permit),
1749 activation,
1750 },
1751 )),
1752 ));
1753 Ok(self.continue_after_pre_ready_stop(
1754 operating,
1755 position,
1756 role,
1757 recoveries,
1758 current.attempt,
1759 stopped,
1760 actions,
1761 ))
1762 }
1763 WorkerPhase::ActivationDispatched {
1764 attempt,
1765 stopped: None,
1766 } => {
1767 operating.members.insert(
1768 position,
1769 Member {
1770 role,
1771 recoveries,
1772 state: MemberState::Worker(Worker {
1773 current,
1774 phase: WorkerPhase::ActivationDispatched {
1775 attempt,
1776 stopped: Some(stopped),
1777 },
1778 }),
1779 },
1780 );
1781 Ok((PoolState::Operating(operating), Actions::cont()))
1782 }
1783 WorkerPhase::Activating {
1784 attempt,
1785 stopped: None,
1786 } => {
1787 operating.members.insert(
1788 position,
1789 Member {
1790 role,
1791 recoveries,
1792 state: MemberState::Worker(Worker {
1793 current,
1794 phase: WorkerPhase::Activating {
1795 attempt,
1796 stopped: Some(stopped),
1797 },
1798 }),
1799 },
1800 );
1801 Ok((PoolState::Operating(operating), Actions::cont()))
1802 }
1803 WorkerPhase::Idle => Ok(self.recover_worker(
1804 operating,
1805 position,
1806 role,
1807 recoveries,
1808 current.attempt,
1809 stopped,
1810 Actions::cont(),
1811 )),
1812 WorkerPhase::Busy(assignment) => {
1813 let admission = match assignment.accept_worker_exit(stopped) {
1814 Ok(admission) => admission,
1815 Err((assignment, stopped)) => {
1816 operating.members.insert(
1817 position,
1818 Member {
1819 role,
1820 recoveries,
1821 state: MemberState::Worker(Worker {
1822 current,
1823 phase: WorkerPhase::Busy(assignment),
1824 }),
1825 },
1826 );
1827 return Err((operating, stopped));
1828 }
1829 };
1830 match admission {
1831 WorkerExitOutcome::AwaitingReceipt(assignment) => {
1832 operating.members.insert(
1833 position,
1834 Member {
1835 role,
1836 recoveries,
1837 state: MemberState::Worker(Worker {
1838 current,
1839 phase: WorkerPhase::Busy(assignment),
1840 }),
1841 },
1842 );
1843 Ok((PoolState::Operating(operating), Actions::cont()))
1844 }
1845 WorkerExitOutcome::JobInterrupted {
1846 customer,
1847 stopped,
1848 late_completion,
1849 } => Ok(self.interrupt_assignment(
1850 operating,
1851 position,
1852 role,
1853 recoveries,
1854 current,
1855 customer,
1856 stopped,
1857 late_completion,
1858 )),
1859 }
1860 }
1861 WorkerPhase::Stopping(join) => match join.stopped(stopped, ¤t.attempt) {
1862 Ok(core::ops::ControlFlow::Continue(join)) => {
1863 operating.members.insert(
1864 position,
1865 Member {
1866 role,
1867 recoveries,
1868 state: MemberState::Worker(Worker {
1869 current,
1870 phase: WorkerPhase::Stopping(join),
1871 }),
1872 },
1873 );
1874 Ok((PoolState::Operating(operating), Actions::cont()))
1875 }
1876 Ok(core::ops::ControlFlow::Break((shutdown, stopped))) => Ok(self
1877 .finish_worker_shutdown(
1878 operating, position, role, recoveries, current, shutdown, stopped,
1879 )),
1880 Err((join, stopped)) => {
1881 operating.members.insert(
1882 position,
1883 Member {
1884 role,
1885 recoveries,
1886 state: MemberState::Worker(Worker {
1887 current,
1888 phase: WorkerPhase::Stopping(join),
1889 }),
1890 },
1891 );
1892 Err((operating, stopped))
1893 }
1894 },
1895 phase @ (WorkerPhase::Initializing {
1896 stopped: Some(_), ..
1897 }
1898 | WorkerPhase::ActivationDispatched {
1899 stopped: Some(_), ..
1900 }
1901 | WorkerPhase::Activating {
1902 stopped: Some(_), ..
1903 }) => {
1904 operating.members.insert(
1905 position,
1906 Member {
1907 role,
1908 recoveries,
1909 state: MemberState::Worker(Worker { current, phase }),
1910 },
1911 );
1912 Err((operating, stopped))
1913 }
1914 }
1915 }
1916
1917 fn finish_worker_shutdown(
1918 &mut self,
1919 operating: FifoOperating<Role, W, P, Job, WorkerResult>,
1920 position: usize,
1921 role: RoleName<Role>,
1922 recoveries: super::restart::RecoveryCount,
1923 current: direct_worker::CurrentWorker<W>,
1924 shutdown: EstablishedShutdownResolved<W::Protocol>,
1925 stopped: crate::ChildStopped<BehaviorAddr<W>>,
1926 ) -> (
1927 PoolState<Role, W, P, Job, WorkerResult>,
1928 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
1929 ) {
1930 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
1931 Actions::cont();
1932 match shutdown {
1933 EstablishedShutdownResolved::Accepted { .. } => {}
1934 shutdown @ EstablishedShutdownResolved::Rejected { .. } => {
1935 actions.sends.diagnostics.append(InterpreterRequests::one(
1936 self.diagnostics.action(FifoDiagnostic::new(
1937 FifoDiagnosticCause::WorkerShutdownRejected {
1938 role: role.clone(),
1939 shutdown,
1940 },
1941 )),
1942 ));
1943 }
1944 }
1945 self.recover_worker(
1946 operating,
1947 position,
1948 role,
1949 recoveries,
1950 current.attempt,
1951 stopped,
1952 actions,
1953 )
1954 }
1955
1956 fn accept_worker_shutdown(
1957 &mut self,
1958 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
1959 shutdown: EstablishedShutdownResolved<W::Protocol>,
1960 ) -> Result<
1961 (
1962 PoolState<Role, W, P, Job, WorkerResult>,
1963 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
1964 ),
1965 (
1966 FifoOperating<Role, W, P, Job, WorkerResult>,
1967 EstablishedShutdownResolved<W::Protocol>,
1968 ),
1969 > {
1970 let Some(position) = operating.worker_shutdown_position(shutdown.id()) else {
1971 return Err((operating, shutdown));
1972 };
1973 let Member {
1974 role,
1975 recoveries,
1976 state,
1977 } = operating.members.remove(position);
1978 let MemberState::Worker(Worker {
1979 current,
1980 phase: WorkerPhase::Stopping(join),
1981 }) = state
1982 else {
1983 operating.members.insert(
1984 position,
1985 Member {
1986 role,
1987 recoveries,
1988 state,
1989 },
1990 );
1991 return Err((operating, shutdown));
1992 };
1993 match join.settled(shutdown) {
1994 Ok(core::ops::ControlFlow::Continue(join)) => {
1995 operating.members.insert(
1996 position,
1997 Member {
1998 role,
1999 recoveries,
2000 state: MemberState::Worker(Worker {
2001 current,
2002 phase: WorkerPhase::Stopping(join),
2003 }),
2004 },
2005 );
2006 Ok((PoolState::Operating(operating), Actions::cont()))
2007 }
2008 Ok(core::ops::ControlFlow::Break((shutdown, stopped))) => Ok(self
2009 .finish_worker_shutdown(
2010 operating, position, role, recoveries, current, shutdown, stopped,
2011 )),
2012 Err((join, shutdown)) => {
2013 operating.members.insert(
2014 position,
2015 Member {
2016 role,
2017 recoveries,
2018 state: MemberState::Worker(Worker {
2019 current,
2020 phase: WorkerPhase::Stopping(join),
2021 }),
2022 },
2023 );
2024 Err((operating, shutdown))
2025 }
2026 }
2027 }
2028
2029 fn complete_assignment(
2030 &mut self,
2031 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
2032 position: usize,
2033 role: RoleName<Role>,
2034 recoveries: super::restart::RecoveryCount,
2035 current: direct_worker::CurrentWorker<W>,
2036 customer: CustomerJob<Job, CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult>>,
2037 result: WorkerResult,
2038 stopped: Option<crate::ChildStopped<BehaviorAddr<W>>>,
2039 ) -> (
2040 PoolState<Role, W, P, Job, WorkerResult>,
2041 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
2042 ) {
2043 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
2044 Actions::cont();
2045 actions.sends.customer_outcomes =
2046 customer
2047 .customer
2048 .deliver(FifoOutcome::completed(customer.id, role.clone(), result));
2049 match stopped {
2050 None => {
2051 operating.members.insert(
2052 position,
2053 Member {
2054 role,
2055 recoveries,
2056 state: MemberState::Worker(Worker {
2057 current,
2058 phase: WorkerPhase::Idle,
2059 }),
2060 },
2061 );
2062 self.fill_fifo(&mut operating, &mut actions);
2063 (PoolState::Operating(operating), actions)
2064 }
2065 Some(stopped) => self.recover_worker(
2066 operating,
2067 position,
2068 role,
2069 recoveries,
2070 current.attempt,
2071 stopped,
2072 actions,
2073 ),
2074 }
2075 }
2076
2077 fn interrupt_assignment(
2078 &mut self,
2079 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
2080 position: usize,
2081 role: RoleName<Role>,
2082 recoveries: super::restart::RecoveryCount,
2083 current: direct_worker::CurrentWorker<W>,
2084 customer: CustomerJob<Job, CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult>>,
2085 stopped: crate::ChildStopped<BehaviorAddr<W>>,
2086 late_completion: Option<Completion<WorkerResult>>,
2087 ) -> (
2088 PoolState<Role, W, P, Job, WorkerResult>,
2089 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
2090 ) {
2091 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
2092 Actions::cont();
2093 if let Some(completion) = late_completion {
2094 let input = FifoEvent::WorkerCompleted(ChildReport::new(
2095 current.attempt.creation(),
2096 completion,
2097 ));
2098 actions
2099 .sends
2100 .diagnostics
2101 .append(InterpreterRequests::one(self.diagnostics.action(
2102 FifoDiagnostic::new(FifoDiagnosticCause::Unexpected(input)),
2103 )));
2104 }
2105 match self.interruption {
2106 Interruption::Fail => {
2107 actions.sends.customer_outcomes =
2108 customer.customer.deliver(FifoOutcome::returned_assigned(
2109 customer.id,
2110 role.clone(),
2111 customer.payload,
2112 AssignedReturnReason::WorkerStopped,
2113 ));
2114 }
2115 Interruption::Retry => {
2116 operating.backlog.insert(
2117 customer.admitted,
2118 QueuedJob {
2119 customer,
2120 assigned_role: Some(role.clone()),
2121 },
2122 );
2123 }
2124 }
2125 self.recover_worker(
2126 operating,
2127 position,
2128 role,
2129 recoveries,
2130 current.attempt,
2131 stopped,
2132 actions,
2133 )
2134 }
2135
2136 fn authorize_waiting(
2137 &mut self,
2138 operating: &mut FifoOperating<Role, W, P, Job, WorkerResult>,
2139 actions: &mut FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
2140 ) {
2141 let occupied = operating
2142 .members
2143 .iter()
2144 .filter_map(|member| match &member.state {
2145 MemberState::Worker(Worker {
2146 phase:
2147 WorkerPhase::ActivationDispatched { .. }
2148 | WorkerPhase::Activating { .. },
2149 ..
2150 }) => Some(()),
2151 MemberState::Creating(_)
2152 | MemberState::Worker(_)
2153 | MemberState::Recovering(_)
2154 | MemberState::Retired => None,
2155 })
2156 .count();
2157 let available = self.activation.maximum().saturating_sub(occupied);
2158 for _ in 0..available {
2159 let position = operating
2160 .members
2161 .iter()
2162 .enumerate()
2163 .find_map(|(position, member)| match &member.state {
2164 MemberState::Worker(Worker {
2165 phase: WorkerPhase::WaitingForActivation { .. },
2166 ..
2167 }) => Some(position),
2168 MemberState::Creating(_)
2169 | MemberState::Worker(_)
2170 | MemberState::Recovering(_)
2171 | MemberState::Retired => None,
2172 });
2173 let Some(position) = position else {
2174 return;
2175 };
2176 let member = operating.members.remove(position);
2177 match member.authorize_activation() {
2178 Ok((member, request)) => {
2179 operating.members.insert(position, member);
2180 actions
2181 .sends
2182 .worker_activations
2183 .append(InterpreterRequests::one(request));
2184 }
2185 Err(member) => {
2186 operating.members.insert(position, member);
2187 return;
2188 }
2189 }
2190 }
2191 }
2192
2193 fn accept_initialization(
2194 &mut self,
2195 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
2196 mut input: super::WorkerInitializationReport<W, P>,
2197 ) -> Result<
2198 (
2199 PoolState<Role, W, P, Job, WorkerResult>,
2200 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
2201 ),
2202 (
2203 FifoOperating<Role, W, P, Job, WorkerResult>,
2204 super::WorkerInitializationReport<W, P>,
2205 ),
2206 > {
2207 if let Some(position) = operating
2208 .members
2209 .iter()
2210 .position(|member| member.worker_attempt() == Some(input.worker()))
2211 {
2212 let member = operating.members.remove(position);
2213 match member.admit_initialization(input) {
2214 Ok(member) => {
2215 operating.members.insert(position, member);
2216 let mut actions = Actions::cont();
2217 self.authorize_waiting(&mut operating, &mut actions);
2218 return Ok((PoolState::Operating(operating), actions));
2219 }
2220 Err((member, returned)) => {
2221 operating.members.insert(position, member);
2222 input = returned;
2223 }
2224 }
2225 }
2226 for position in 0..operating.members.len() {
2227 let Member {
2228 role,
2229 recoveries,
2230 state,
2231 } = operating.members.remove(position);
2232 let MemberState::Worker(Worker {
2233 current,
2234 phase:
2235 WorkerPhase::Initializing {
2236 initialization,
2237 stopped,
2238 },
2239 }) = state
2240 else {
2241 operating.members.insert(
2242 position,
2243 Member {
2244 role,
2245 recoveries,
2246 state,
2247 },
2248 );
2249 continue;
2250 };
2251 if &initialization != input.initialization() {
2252 operating.members.insert(
2253 position,
2254 Member {
2255 role,
2256 recoveries,
2257 state: MemberState::Worker(Worker {
2258 current,
2259 phase: WorkerPhase::Initializing {
2260 initialization,
2261 stopped,
2262 },
2263 }),
2264 },
2265 );
2266 continue;
2267 }
2268 match (stopped, input) {
2269 (
2270 None,
2271 super::WorkerInitializationReport::Stopped {
2272 activation,
2273 stopped,
2274 ..
2275 },
2276 ) if stopped.child == current.attempt.creation() => {
2277 let mut actions: FifoActions<
2278 Role,
2279 W,
2280 P,
2281 Source,
2282 Diagnostics,
2283 Job,
2284 WorkerResult,
2285 > = Actions::cont();
2286 actions.sends.diagnostics.append(InterpreterRequests::one(
2287 self.diagnostics.action(FifoDiagnostic::new(
2288 FifoDiagnosticCause::UnusedActivation {
2289 role: role.clone(),
2290 permit: None,
2291 activation,
2292 },
2293 )),
2294 ));
2295 return Ok(self.continue_after_pre_ready_stop(
2296 operating,
2297 position,
2298 role,
2299 recoveries,
2300 current.attempt,
2301 stopped,
2302 actions,
2303 ));
2304 }
2305 (Some(stopped), returned) => {
2306 let mut actions: FifoActions<
2307 Role,
2308 W,
2309 P,
2310 Source,
2311 Diagnostics,
2312 Job,
2313 WorkerResult,
2314 > = Actions::cont();
2315 actions.sends.diagnostics.append(InterpreterRequests::one(
2316 self.diagnostics.action(FifoDiagnostic::new(
2317 FifoDiagnosticCause::Unexpected(FifoEvent::WorkerInitialization(
2318 returned,
2319 )),
2320 )),
2321 ));
2322 return Ok(self.continue_after_pre_ready_stop(
2323 operating,
2324 position,
2325 role,
2326 recoveries,
2327 current.attempt,
2328 stopped,
2329 actions,
2330 ));
2331 }
2332 (None, returned @ super::WorkerInitializationReport::EffectsRejected { .. }) => {
2333 let mut actions: FifoActions<
2334 Role,
2335 W,
2336 P,
2337 Source,
2338 Diagnostics,
2339 Job,
2340 WorkerResult,
2341 > = Actions::cont();
2342 actions.sends.diagnostics.append(InterpreterRequests::one(
2343 self.diagnostics.action(FifoDiagnostic::new(
2344 FifoDiagnosticCause::Unexpected(FifoEvent::WorkerInitialization(
2345 returned,
2346 )),
2347 )),
2348 ));
2349 let shutdown = self
2350 .shutdowns
2351 .reserve(1)
2352 .and_then(|reserved| reserved.into_iter().next());
2353 let Some(shutdown) = shutdown else {
2354 operating.members.insert(
2355 position,
2356 Member {
2357 role,
2358 recoveries,
2359 state: MemberState::Worker(Worker {
2360 current,
2361 phase: WorkerPhase::Initializing {
2362 initialization,
2363 stopped: None,
2364 },
2365 }),
2366 },
2367 );
2368 return Ok(self.begin_shutdown(operating, actions));
2369 };
2370 let (operating, actions) = self.request_worker_shutdown(
2371 operating, position, role, recoveries, current, shutdown, actions,
2372 );
2373 return Ok((PoolState::Operating(operating), actions));
2374 }
2375 (None, returned) => {
2376 input = returned;
2377 operating.members.insert(
2378 position,
2379 Member {
2380 role,
2381 recoveries,
2382 state: MemberState::Worker(Worker {
2383 current,
2384 phase: WorkerPhase::Initializing {
2385 initialization,
2386 stopped: None,
2387 },
2388 }),
2389 },
2390 );
2391 }
2392 }
2393 }
2394 Err((operating, input))
2395 }
2396
2397 fn accept_activation(
2398 &mut self,
2399 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
2400 mut input: super::WorkerActivation<W, P>,
2401 ) -> Result<
2402 (
2403 PoolState<Role, W, P, Job, WorkerResult>,
2404 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
2405 ),
2406 (
2407 FifoOperating<Role, W, P, Job, WorkerResult>,
2408 super::WorkerActivation<W, P>,
2409 ),
2410 > {
2411 if let Some(position) = operating
2412 .members
2413 .iter()
2414 .position(|member| member.worker_attempt() == Some(&input.worker()))
2415 {
2416 let member = operating.members.remove(position);
2417 match member.admit_activation(input) {
2418 Ok(member) => {
2419 let mut actions = Actions::cont();
2420 match &member.state {
2421 MemberState::Worker(Worker {
2422 phase: WorkerPhase::Idle,
2423 ..
2424 }) => {
2425 operating.members.insert(position, member);
2426 self.fill_fifo(&mut operating, &mut actions);
2427 }
2428 MemberState::Creating(_)
2429 | MemberState::Worker(_)
2430 | MemberState::Recovering(_)
2431 | MemberState::Retired => {
2432 operating.members.insert(position, member);
2433 }
2434 }
2435 self.authorize_waiting(&mut operating, &mut actions);
2436 return Ok((PoolState::Operating(operating), actions));
2437 }
2438 Err((member, returned)) => {
2439 operating.members.insert(position, member);
2440 input = returned;
2441 }
2442 }
2443 }
2444 for position in 0..operating.members.len() {
2445 let Member {
2446 role,
2447 recoveries,
2448 state,
2449 } = operating.members.remove(position);
2450 let MemberState::Worker(Worker { current, phase }) = state else {
2451 operating.members.insert(
2452 position,
2453 Member {
2454 role,
2455 recoveries,
2456 state,
2457 },
2458 );
2459 continue;
2460 };
2461 if current.attempt != input.worker() {
2462 operating.members.insert(
2463 position,
2464 Member {
2465 role,
2466 recoveries,
2467 state: MemberState::Worker(Worker { current, phase }),
2468 },
2469 );
2470 continue;
2471 }
2472 match phase {
2473 WorkerPhase::ActivationDispatched {
2474 attempt,
2475 stopped: None,
2476 } if &attempt == input.attempt() => {
2477 let (returned_worker, returned_attempt, outcome) = input.into_parts();
2478 let mut actions: FifoActions<
2479 Role,
2480 W,
2481 P,
2482 Source,
2483 Diagnostics,
2484 Job,
2485 WorkerResult,
2486 > = Actions::cont();
2487 actions.sends.diagnostics.append(InterpreterRequests::one(
2488 self.diagnostics.action(FifoDiagnostic::new(
2489 FifoDiagnosticCause::WorkerActivationReturned {
2490 role: role.clone(),
2491 worker: returned_worker,
2492 activation: returned_attempt,
2493 outcome,
2494 },
2495 )),
2496 ));
2497 let shutdown = self
2498 .shutdowns
2499 .reserve(1)
2500 .and_then(|reserved| reserved.into_iter().next());
2501 let Some(shutdown) = shutdown else {
2502 operating.members.insert(
2503 position,
2504 Member {
2505 role,
2506 recoveries,
2507 state: MemberState::Worker(Worker {
2508 current,
2509 phase: WorkerPhase::ActivationDispatched {
2510 attempt,
2511 stopped: None,
2512 },
2513 }),
2514 },
2515 );
2516 return Ok(self.begin_shutdown(operating, actions));
2517 };
2518 let (mut operating, mut actions) = self.request_worker_shutdown(
2519 operating, position, role, recoveries, current, shutdown, actions,
2520 );
2521 self.authorize_waiting(&mut operating, &mut actions);
2522 return Ok((PoolState::Operating(operating), actions));
2523 }
2524 WorkerPhase::Activating {
2525 attempt,
2526 stopped: None,
2527 } if &attempt == input.attempt() => {
2528 let (returned_worker, returned_attempt, outcome) = input.into_parts();
2529 match outcome {
2530 WorkerActivationOutcome::Started => {
2531 let mut actions: FifoActions<
2532 Role,
2533 W,
2534 P,
2535 Source,
2536 Diagnostics,
2537 Job,
2538 WorkerResult,
2539 > = Actions::cont();
2540 actions.sends.diagnostics.append(InterpreterRequests::one(
2541 self.diagnostics.action(FifoDiagnostic::new(
2542 FifoDiagnosticCause::WorkerActivationReturned {
2543 role: role.clone(),
2544 worker: returned_worker,
2545 activation: returned_attempt,
2546 outcome: WorkerActivationOutcome::Started,
2547 },
2548 )),
2549 ));
2550 operating.members.insert(
2551 position,
2552 Member {
2553 role,
2554 recoveries,
2555 state: MemberState::Worker(Worker {
2556 current,
2557 phase: WorkerPhase::Activating {
2558 attempt,
2559 stopped: None,
2560 },
2561 }),
2562 },
2563 );
2564 return Ok((PoolState::Operating(operating), actions));
2565 }
2566 outcome @ (WorkerActivationOutcome::StartRejected { .. }
2567 | WorkerActivationOutcome::Ready(_)
2568 | WorkerActivationOutcome::Rejected(_)) => {
2569 let mut actions: FifoActions<
2570 Role,
2571 W,
2572 P,
2573 Source,
2574 Diagnostics,
2575 Job,
2576 WorkerResult,
2577 > = Actions::cont();
2578 actions.sends.diagnostics.append(InterpreterRequests::one(
2579 self.diagnostics.action(FifoDiagnostic::new(
2580 FifoDiagnosticCause::WorkerActivationReturned {
2581 role: role.clone(),
2582 worker: returned_worker,
2583 activation: returned_attempt,
2584 outcome,
2585 },
2586 )),
2587 ));
2588 let shutdown = self
2589 .shutdowns
2590 .reserve(1)
2591 .and_then(|reserved| reserved.into_iter().next());
2592 let Some(shutdown) = shutdown else {
2593 operating.members.insert(
2594 position,
2595 Member {
2596 role,
2597 recoveries,
2598 state: MemberState::Worker(Worker {
2599 current,
2600 phase: WorkerPhase::Activating {
2601 attempt,
2602 stopped: None,
2603 },
2604 }),
2605 },
2606 );
2607 return Ok(self.begin_shutdown(operating, actions));
2608 };
2609 let (mut operating, mut actions) = self.request_worker_shutdown(
2610 operating, position, role, recoveries, current, shutdown, actions,
2611 );
2612 self.authorize_waiting(&mut operating, &mut actions);
2613 return Ok((PoolState::Operating(operating), actions));
2614 }
2615 }
2616 }
2617 WorkerPhase::ActivationDispatched {
2618 attempt,
2619 stopped: Some(stopped),
2620 } if &attempt == input.attempt() => {
2621 let (returned_worker, returned_attempt, outcome) = input.into_parts();
2622 match outcome {
2623 outcome @ (WorkerActivationOutcome::Started
2624 | WorkerActivationOutcome::StartRejected { .. }
2625 | WorkerActivationOutcome::Ready(_)
2626 | WorkerActivationOutcome::Rejected(_)) => {
2627 return Ok(self.continue_after_returned_activation(
2628 operating,
2629 position,
2630 role,
2631 recoveries,
2632 current.attempt,
2633 stopped,
2634 returned_worker,
2635 returned_attempt,
2636 outcome,
2637 ));
2638 }
2639 }
2640 }
2641 WorkerPhase::Activating {
2642 attempt,
2643 stopped: Some(stopped),
2644 } if &attempt == input.attempt() => {
2645 let (returned_worker, returned_attempt, outcome) = input.into_parts();
2646 match outcome {
2647 WorkerActivationOutcome::Started => {
2648 let mut actions: FifoActions<
2649 Role,
2650 W,
2651 P,
2652 Source,
2653 Diagnostics,
2654 Job,
2655 WorkerResult,
2656 > = Actions::cont();
2657 actions.sends.diagnostics.append(InterpreterRequests::one(
2658 self.diagnostics.action(FifoDiagnostic::new(
2659 FifoDiagnosticCause::WorkerActivationReturned {
2660 role: role.clone(),
2661 worker: returned_worker,
2662 activation: returned_attempt,
2663 outcome: WorkerActivationOutcome::Started,
2664 },
2665 )),
2666 ));
2667 operating.members.insert(
2668 position,
2669 Member {
2670 role,
2671 recoveries,
2672 state: MemberState::Worker(Worker {
2673 current,
2674 phase: WorkerPhase::Activating {
2675 attempt,
2676 stopped: Some(stopped),
2677 },
2678 }),
2679 },
2680 );
2681 return Ok((PoolState::Operating(operating), actions));
2682 }
2683 outcome @ (WorkerActivationOutcome::StartRejected { .. }
2684 | WorkerActivationOutcome::Ready(_)
2685 | WorkerActivationOutcome::Rejected(_)) => {
2686 return Ok(self.continue_after_returned_activation(
2687 operating,
2688 position,
2689 role,
2690 recoveries,
2691 current.attempt,
2692 stopped,
2693 returned_worker,
2694 returned_attempt,
2695 outcome,
2696 ));
2697 }
2698 }
2699 }
2700 phase => {
2701 operating.members.insert(
2702 position,
2703 Member {
2704 role,
2705 recoveries,
2706 state: MemberState::Worker(Worker { current, phase }),
2707 },
2708 );
2709 }
2710 }
2711 }
2712 Err((operating, input))
2713 }
2714
2715 fn reject_creation_batch(
2716 &self,
2717 operating: &mut FifoOperating<Role, W, P, Job, WorkerResult>,
2718 creations: Creations<CreateChild<BehaviorAddr<W>, StopOnShutdown<W>>>,
2719 rejection: CreationBatchRejection,
2720 ) -> Result<
2721 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
2722 CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
2723 > {
2724 let Some(positions) = operating.creation_positions(
2725 creations
2726 .iter()
2727 .map(|creation| (creation.id(), creation.kind())),
2728 ) else {
2729 return Err(rejection.settlement(creations));
2730 };
2731 let mut returned: BTreeMap<_, _> = positions.into_iter().zip(creations).collect();
2732 let mut members = Vec::with_capacity(operating.members.len());
2733 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
2734 Actions::cont();
2735 for (position, member) in operating.members.drain(..).enumerate() {
2736 let Some(creation) = returned.remove(&position) else {
2737 members.push(member);
2738 continue;
2739 };
2740 let Member {
2741 role,
2742 recoveries,
2743 state,
2744 } = member;
2745 match state {
2746 MemberState::Creating(worker) => {
2747 let (_, returned_worker, _) = creation.into_parts();
2748 actions.sends.diagnostics.append(InterpreterRequests::one(
2749 self.diagnostics.action(FifoDiagnostic::new(
2750 FifoDiagnosticCause::WorkerReturned {
2751 role: role.clone(),
2752 rejection: rejection.worker(returned_worker.into_inner()),
2753 activation: worker.activation,
2754 stopped: worker.stopped,
2755 },
2756 )),
2757 ));
2758 members.push(Member {
2759 role,
2760 recoveries,
2761 state: MemberState::Retired,
2762 });
2763 }
2764 state => {
2765 members.push(Member {
2766 role,
2767 recoveries,
2768 state,
2769 });
2770 let workers = rejection.settlement(Creations::one(creation));
2771 actions.sends.diagnostics.append(InterpreterRequests::one(
2772 self.diagnostics.action(FifoDiagnostic::new(
2773 FifoDiagnosticCause::WorkersReturned(workers),
2774 )),
2775 ));
2776 }
2777 }
2778 }
2779 operating.members = members;
2780 Ok(actions)
2781 }
2782
2783 fn accept_creations(
2784 &mut self,
2785 mut operating: FifoOperating<Role, W, P, Job, WorkerResult>,
2786 workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
2787 ) -> Result<
2788 (
2789 PoolState<Role, W, P, Job, WorkerResult>,
2790 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
2791 ),
2792 (
2793 FifoOperating<Role, W, P, Job, WorkerResult>,
2794 CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
2795 ),
2796 > {
2797 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
2798 Actions::cont();
2799 let mut failure = None;
2800 match workers.into_settlement() {
2801 CreationSettlement::Corrupt { creations, fault } => {
2802 let rejection = CreationBatchRejection::InterpreterCorrupt(fault);
2803 match self.reject_creation_batch(&mut operating, creations, rejection) {
2804 Ok(returned) => {
2805 actions = returned;
2806 failure = Some(self.recovery.failure());
2807 }
2808 Err(workers) => return Err((operating, workers)),
2809 }
2810 }
2811 CreationSettlement::Settled(settlements) => {
2812 let identities = settlements
2813 .iter()
2814 .map(creation_identity)
2815 .collect::<Option<Vec<_>>>();
2816 let Some(positions) =
2817 identities.and_then(|identities| operating.creation_positions(identities))
2818 else {
2819 return Err((
2820 operating,
2821 CreationsSettled::new(CreationSettlement::Settled(settlements)),
2822 ));
2823 };
2824 let mut returned: BTreeMap<_, _> = positions.into_iter().zip(settlements).collect();
2825 let mut members = Vec::with_capacity(operating.members.len());
2826 for (position, member) in operating.members.into_iter().enumerate() {
2827 let Some(settlement) = returned.remove(&position) else {
2828 members.push(member);
2829 continue;
2830 };
2831 match member.admit_creation(ChildCreationSettled::new(settlement)) {
2832 WorkerCreationAdmission::Initializing { member, request } => {
2833 members.push(member);
2834 actions
2835 .sends
2836 .worker_initializations
2837 .append(InterpreterRequests::one(request));
2838 }
2839 WorkerCreationAdmission::Returned {
2840 member,
2841 role,
2842 rejection,
2843 activation,
2844 stopped,
2845 } => {
2846 members.push(member);
2847 actions.sends.diagnostics.append(InterpreterRequests::one(
2848 self.diagnostics.action(FifoDiagnostic::new(
2849 FifoDiagnosticCause::WorkerReturned {
2850 role,
2851 rejection,
2852 activation,
2853 stopped,
2854 },
2855 )),
2856 ));
2857 failure = Some(self.recovery.failure());
2858 }
2859 WorkerCreationAdmission::Unrelated { member, creation } => {
2860 members.push(member);
2861 let workers = CreationsSettled::new(CreationSettlement::Settled(
2862 Creations::one(creation.into_settlement()),
2863 ));
2864 actions.sends.diagnostics.append(InterpreterRequests::one(
2865 self.diagnostics.action(FifoDiagnostic::new(
2866 FifoDiagnosticCause::WorkersReturned(workers),
2867 )),
2868 ));
2869 }
2870 }
2871 }
2872 operating.members = members;
2873 }
2874 CreationSettlement::Rejected {
2875 creations,
2876 reason: _,
2877 } => {
2878 let rejection = CreationBatchRejection::NamespaceExhausted;
2879 match self.reject_creation_batch(&mut operating, creations, rejection) {
2880 Ok(returned) => {
2881 actions = returned;
2882 failure = Some(self.recovery.failure());
2883 }
2884 Err(workers) => return Err((operating, workers)),
2885 }
2886 }
2887 }
2888
2889 match failure {
2890 None => Ok((PoolState::Operating(operating), actions)),
2891 Some(PoolFailureReaction::RetireRole) => {
2892 self.fill_fifo(&mut operating, &mut actions);
2893 self.return_unrecoverable_jobs(&mut operating, &mut actions);
2894 Ok((PoolState::Operating(operating), actions))
2895 }
2896 Some(PoolFailureReaction::StopPool) => Ok(self.begin_shutdown(operating, actions)),
2897 }
2898 }
2899
2900 fn issue_restart_timer(&mut self) -> Option<ScheduleKey> {
2901 ScheduleKey::issue(&mut self.next_restart_timer)
2902 }
2903
2904 fn begin_shutdown(
2905 &mut self,
2906 operating: FifoOperating<Role, W, P, Job, WorkerResult>,
2907 mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
2908 ) -> (
2909 PoolState<Role, W, P, Job, WorkerResult>,
2910 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
2911 ) {
2912 let FifoOperating {
2913 members,
2914 backlog,
2915 cursor: _,
2916 } = operating;
2917 for (_, queued) in backlog {
2918 let CustomerJob {
2919 id,
2920 admitted: _,
2921 payload,
2922 customer,
2923 } = queued.customer;
2924 actions
2925 .sends
2926 .customer_outcomes
2927 .append(customer.deliver(FifoOutcome::returned_queued(
2928 id,
2929 payload,
2930 QueuedReturnReason::PoolShutdown,
2931 )));
2932 }
2933
2934 let required = members.iter().filter_map(Member::shutdown_target).count();
2935 let reserved = self.shutdowns.reserve(required);
2936 let mut ids = match reserved {
2937 Some(ids) => Some(ids.into_iter()),
2938 None => None,
2939 };
2940 let mut draining = Vec::with_capacity(members.len());
2941 for member in members {
2942 let role = member.role_name();
2943 let id = match (&mut ids, member.shutdown_target()) {
2944 (Some(ids), Some(_)) => ids.next(),
2945 _ => None,
2946 };
2947 let direct_worker::MemberRetirement {
2948 member,
2949 request,
2950 custody,
2951 } = member.shutdown(id);
2952 if let Some(request) = request {
2953 actions
2954 .sends
2955 .worker_shutdowns
2956 .append(InterpreterRequests::one(request));
2957 }
2958 match custody {
2959 direct_worker::RetirementReturns::NoTransfer => {}
2960 direct_worker::RetirementReturns::Assignment(AssignmentShutdown {
2961 customer,
2962 worker,
2963 stopped: _,
2964 completion,
2965 }) => {
2966 actions
2967 .sends
2968 .customer_outcomes
2969 .append(customer.customer.deliver(FifoOutcome::returned_assigned(
2970 customer.id,
2971 role,
2972 customer.payload,
2973 AssignedReturnReason::PoolShutdown,
2974 )));
2975 if let Some(completion) = completion {
2976 let input = FifoEvent::WorkerCompleted(ChildReport::new(
2977 worker.creation(),
2978 completion,
2979 ));
2980 actions.sends.diagnostics.append(InterpreterRequests::one(
2981 self.diagnostics.action(FifoDiagnostic::new(
2982 FifoDiagnosticCause::Unexpected(input),
2983 )),
2984 ));
2985 }
2986 }
2987 direct_worker::RetirementReturns::Replacement(replacement) => {
2988 self.retain_cancelled_replacement(
2989 role,
2990 replacement,
2991 WorkerReplacementError::PoolShutdown,
2992 &mut actions,
2993 );
2994 }
2995 }
2996 draining.push(member);
2997 }
2998
2999 match ids {
3000 None => {
3001 actions.become_ = behavior::Step::Stop(behavior::Stopped);
3002 (
3003 PoolState::ForcedRetirement {
3004 members: draining,
3005 cause: ForcedRetirementCause::WorkerShutdownIdsExhausted,
3006 },
3007 actions,
3008 )
3009 }
3010 Some(_) => match draining.iter().find_map(RetiringWorker::waiting) {
3011 None => {
3012 actions.become_ = behavior::Step::Stop(behavior::Stopped);
3013 (PoolState::Stopped, actions)
3014 }
3015 Some(_) => {
3016 let (deadline, schedule) = ShutdownDeadline::begin(self.actor_drain);
3017 if let Some(schedule) = schedule {
3018 actions.sends.restart_schedules.send(schedule);
3019 }
3020 (
3021 PoolState::Draining {
3022 members: draining,
3023 deadline,
3024 },
3025 actions,
3026 )
3027 }
3028 },
3029 }
3030 }
3031
3032 fn retain_drained_worker(
3033 &self,
3034 worker: WorkerDeparture<Role, W, P>,
3035 actions: &mut FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3036 ) {
3037 let WorkerDeparture {
3038 role,
3039 startup,
3040 shutdown,
3041 } = worker;
3042 match startup {
3043 Some(direct_worker::WorkerStartupCustody::ActivationPlan(activation)) => {
3044 actions.sends.diagnostics.append(InterpreterRequests::one(
3045 self.diagnostics.action(FifoDiagnostic::new(
3046 FifoDiagnosticCause::UnusedActivation {
3047 role: role.clone(),
3048 permit: None,
3049 activation,
3050 },
3051 )),
3052 ));
3053 }
3054 Some(direct_worker::WorkerStartupCustody::ActivationPermit { permit, activation }) => {
3055 actions.sends.diagnostics.append(InterpreterRequests::one(
3056 self.diagnostics.action(FifoDiagnostic::new(
3057 FifoDiagnosticCause::UnusedActivation {
3058 role: role.clone(),
3059 permit: Some(permit),
3060 activation,
3061 },
3062 )),
3063 ));
3064 }
3065 Some(
3066 direct_worker::WorkerStartupCustody::Initializing(_)
3067 | direct_worker::WorkerStartupCustody::ActivationStart(_)
3068 | direct_worker::WorkerStartupCustody::Activation(_),
3069 )
3070 | None => {}
3071 }
3072 match shutdown {
3073 EstablishedShutdownResolved::Accepted { .. } => {}
3074 shutdown @ EstablishedShutdownResolved::Rejected { .. } => {
3075 actions.sends.diagnostics.append(InterpreterRequests::one(
3076 self.diagnostics.action(FifoDiagnostic::new(
3077 FifoDiagnosticCause::WorkerShutdownRejected { role, shutdown },
3078 )),
3079 ));
3080 }
3081 }
3082 }
3083
3084 fn retain_worker(
3085 &mut self,
3086 custody: WorkerCustody<Role, W, P>,
3087 workers: &mut Vec<RetiringWorker<Role, W, P>>,
3088 actions: &mut FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3089 ) -> Option<ForcedRetirementCause> {
3090 match custody {
3091 WorkerCustody::Running {
3092 role,
3093 worker,
3094 activation,
3095 } => {
3096 let reserved = self.shutdowns.reserve(1);
3097 let shutdown = reserved.and_then(|mut shutdowns| shutdowns.pop());
3098 let exhausted = match shutdown {
3099 Some(_) => None,
3100 None => Some(ForcedRetirementCause::WorkerShutdownIdsExhausted),
3101 };
3102 let (worker, request) = direct_worker::retire_worker(
3103 role,
3104 worker,
3105 Some(direct_worker::WorkerStartupCustody::ActivationPlan(
3106 activation,
3107 )),
3108 shutdown,
3109 );
3110 workers.push(worker);
3111 if let Some(request) = request {
3112 actions
3113 .sends
3114 .worker_shutdowns
3115 .append(InterpreterRequests::one(request));
3116 }
3117 exhausted
3118 }
3119 WorkerCustody::Returned {
3120 role,
3121 rejection,
3122 activation,
3123 stopped,
3124 } => {
3125 actions.sends.diagnostics.append(InterpreterRequests::one(
3126 self.diagnostics.action(FifoDiagnostic::new(
3127 FifoDiagnosticCause::WorkerReturned {
3128 role,
3129 rejection,
3130 activation,
3131 stopped,
3132 },
3133 )),
3134 ));
3135 None
3136 }
3137 WorkerCustody::Stopped {
3138 role,
3139 worker,
3140 activation,
3141 stopped,
3142 } => {
3143 actions.sends.diagnostics.append(InterpreterRequests::one(
3144 self.diagnostics.action(FifoDiagnostic::new(
3145 FifoDiagnosticCause::StoppedWorkerEstablished {
3146 role,
3147 worker,
3148 activation,
3149 stopped,
3150 },
3151 )),
3152 ));
3153 None
3154 }
3155 }
3156 }
3157
3158 fn retain_cancelled_replacement(
3159 &self,
3160 role: RoleName<Role>,
3161 replacement: direct_worker::PendingWorkerReplacement<W, P>,
3162 error: WorkerReplacementError,
3163 actions: &mut FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3164 ) {
3165 let direct_worker::PendingWorkerReplacement {
3166 previous,
3167 stopped,
3168 submission,
3169 } = replacement;
3170 actions
3171 .sends
3172 .diagnostics
3173 .append(InterpreterRequests::one(self.diagnostics.action(
3174 FifoDiagnostic::new(FifoDiagnosticCause::WorkerReplacementFailed {
3175 role,
3176 previous,
3177 stopped,
3178 submission,
3179 returned_source: None,
3180 error,
3181 }),
3182 )));
3183 }
3184
3185 fn reject_drain_creations(
3186 &mut self,
3187 mut members: Vec<RetiringWorker<Role, W, P>>,
3188 deadline: ShutdownDeadline,
3189 creations: Creations<CreateChild<BehaviorAddr<W>, StopOnShutdown<W>>>,
3190 rejection: CreationBatchRejection,
3191 ) -> (
3192 PoolState<Role, W, P, Job, WorkerResult>,
3193 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3194 ) {
3195 let positions = direct_worker::ordered_creation_positions(
3196 members.iter().map(RetiringWorker::expected_creation),
3197 creations
3198 .iter()
3199 .map(|creation| (creation.id(), creation.kind())),
3200 );
3201 let Some(positions) = positions else {
3202 return self.diagnose(
3203 PoolState::Draining { members, deadline },
3204 FifoEvent::WorkerCreationsSettled(rejection.settlement(creations)),
3205 );
3206 };
3207 let mut returned: BTreeMap<_, _> = positions.into_iter().zip(creations).collect();
3208 let mut draining = Vec::with_capacity(members.len());
3209 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
3210 Actions::cont();
3211 let mut forced = None;
3212 for (position, member) in members.drain(..).enumerate() {
3213 let Some(creation) = returned.remove(&position) else {
3214 draining.push(member);
3215 continue;
3216 };
3217 let (id, returned_worker, kind) = creation.into_parts();
3218 let returned_worker = rejection.worker(returned_worker.into_inner());
3219 match member.accept_creation_rejection(id, kind, returned_worker) {
3220 Ok(custody) => {
3221 let exhausted = self.retain_worker(custody, &mut draining, &mut actions);
3222 forced = forced.or(exhausted);
3223 }
3224 Err((member, id, kind, rejection)) => {
3225 draining.push(member);
3226 actions.sends.diagnostics.append(InterpreterRequests::one(
3227 self.diagnostics.action(FifoDiagnostic::new(
3228 FifoDiagnosticCause::UnmatchedWorkerReturn {
3229 id,
3230 kind,
3231 rejection,
3232 },
3233 )),
3234 ));
3235 }
3236 }
3237 }
3238 match forced {
3239 Some(cause) => {
3240 actions.become_ = behavior::Step::Stop(behavior::Stopped);
3241 (
3242 PoolState::ForcedRetirement {
3243 members: draining,
3244 cause,
3245 },
3246 actions,
3247 )
3248 }
3249 None => self.continue_draining(draining, deadline, actions),
3250 }
3251 }
3252
3253 fn accept_drain_creations(
3254 &mut self,
3255 members: Vec<RetiringWorker<Role, W, P>>,
3256 deadline: ShutdownDeadline,
3257 workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
3258 ) -> (
3259 PoolState<Role, W, P, Job, WorkerResult>,
3260 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3261 ) {
3262 let settlements = match workers.into_settlement() {
3263 CreationSettlement::Rejected {
3264 creations,
3265 reason: _,
3266 } => {
3267 return self.reject_drain_creations(
3268 members,
3269 deadline,
3270 creations,
3271 CreationBatchRejection::NamespaceExhausted,
3272 );
3273 }
3274 CreationSettlement::Corrupt { creations, fault } => {
3275 return self.reject_drain_creations(
3276 members,
3277 deadline,
3278 creations,
3279 CreationBatchRejection::InterpreterCorrupt(fault),
3280 );
3281 }
3282 CreationSettlement::Settled(settlements) => settlements,
3283 };
3284 let identities = settlements
3285 .iter()
3286 .map(creation_identity)
3287 .collect::<Option<Vec<_>>>();
3288 let positions = identities.and_then(|identities| {
3289 direct_worker::ordered_creation_positions(
3290 members.iter().map(RetiringWorker::expected_creation),
3291 identities,
3292 )
3293 });
3294 let Some(positions) = positions else {
3295 return self.diagnose(
3296 PoolState::Draining { members, deadline },
3297 FifoEvent::WorkerCreationsSettled(CreationsSettled::new(
3298 CreationSettlement::Settled(settlements),
3299 )),
3300 );
3301 };
3302 let mut returned: BTreeMap<_, _> = positions.into_iter().zip(settlements).collect();
3303 let mut draining = Vec::with_capacity(members.len());
3304 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
3305 Actions::cont();
3306 let mut forced = None;
3307 for (position, member) in members.into_iter().enumerate() {
3308 let Some(settlement) = returned.remove(&position) else {
3309 draining.push(member);
3310 continue;
3311 };
3312 match member.accept_creation(ChildCreationSettled::new(settlement)) {
3313 Ok(custody) => {
3314 let exhausted = self.retain_worker(custody, &mut draining, &mut actions);
3315 forced = forced.or(exhausted);
3316 }
3317 Err((member, creation)) => {
3318 actions.sends.diagnostics.append(InterpreterRequests::one(
3319 self.diagnostics.action(FifoDiagnostic::new(
3320 FifoDiagnosticCause::WorkersReturned(CreationsSettled::new(
3321 CreationSettlement::Settled(Creations::one(
3322 creation.into_settlement(),
3323 )),
3324 )),
3325 )),
3326 ));
3327 draining.push(member);
3328 }
3329 }
3330 }
3331 match forced {
3332 Some(cause) => {
3333 actions.become_ = behavior::Step::Stop(behavior::Stopped);
3334 (
3335 PoolState::ForcedRetirement {
3336 members: draining,
3337 cause,
3338 },
3339 actions,
3340 )
3341 }
3342 None => self.continue_draining(draining, deadline, actions),
3343 }
3344 }
3345
3346 fn accept_drain_initialization(
3347 &self,
3348 members: Vec<RetiringWorker<Role, W, P>>,
3349 deadline: ShutdownDeadline,
3350 input: super::WorkerInitializationReport<W, P>,
3351 ) -> (
3352 PoolState<Role, W, P, Job, WorkerResult>,
3353 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3354 ) {
3355 match direct_worker::accept_worker_initialization(members, input) {
3356 RetiringWorkerInitialization::Returned { workers, report } => self.diagnose(
3357 PoolState::Draining {
3358 members: workers,
3359 deadline,
3360 },
3361 FifoEvent::WorkerInitialization(report),
3362 ),
3363 RetiringWorkerInitialization::WorkerStopped {
3364 workers,
3365 departure,
3366 role,
3367 activation,
3368 } => {
3369 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
3370 Actions::cont();
3371 actions.sends.diagnostics.append(InterpreterRequests::one(
3372 self.diagnostics.action(FifoDiagnostic::new(
3373 FifoDiagnosticCause::UnusedActivation {
3374 role,
3375 permit: None,
3376 activation,
3377 },
3378 )),
3379 ));
3380 if let Some(departure) = departure {
3381 self.retain_drained_worker(departure, &mut actions);
3382 }
3383 self.continue_draining(workers, deadline, actions)
3384 }
3385 }
3386 }
3387
3388 fn accept_drain_activation(
3389 &self,
3390 members: Vec<RetiringWorker<Role, W, P>>,
3391 deadline: ShutdownDeadline,
3392 input: super::WorkerActivation<W, P>,
3393 ) -> (
3394 PoolState<Role, W, P, Job, WorkerResult>,
3395 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3396 ) {
3397 match direct_worker::accept_worker_activation(members, input) {
3398 RetiringWorkerActivation::Started { workers } => {
3399 self.continue_draining(workers, deadline, Actions::cont())
3400 }
3401 RetiringWorkerActivation::Returned {
3402 workers,
3403 role,
3404 worker,
3405 activation,
3406 outcome,
3407 } => {
3408 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
3409 Actions::cont();
3410 actions.sends.diagnostics.append(InterpreterRequests::one(
3411 self.diagnostics.action(FifoDiagnostic::new(
3412 FifoDiagnosticCause::WorkerActivationReturned {
3413 role,
3414 worker,
3415 activation,
3416 outcome,
3417 },
3418 )),
3419 ));
3420 self.continue_draining(workers, deadline, actions)
3421 }
3422 RetiringWorkerActivation::Unrelated {
3423 workers,
3424 activation,
3425 } => self.diagnose(
3426 PoolState::Draining {
3427 members: workers,
3428 deadline,
3429 },
3430 FifoEvent::WorkerActivationReported(activation),
3431 ),
3432 }
3433 }
3434
3435 fn accept_drain_preparation_start(
3436 &mut self,
3437 members: Vec<RetiringWorker<Role, W, P>>,
3438 deadline: ShutdownDeadline,
3439 input: behavior::ActionItemResult<PrepareWorkers<Source, Role, W, P>>,
3440 ) -> (
3441 PoolState<Role, W, P, Job, WorkerResult>,
3442 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3443 )
3444 where
3445 Role: Eq,
3446 {
3447 let accepted = direct_worker::accept_retiring_worker_preparation_start(members, input);
3448 match accepted {
3449 Ok(ControlFlow::Continue(members)) => {
3450 self.continue_draining(members, deadline, Actions::cont())
3451 }
3452 Ok(ControlFlow::Break((members, failed))) => {
3453 let error = match failed.fault {
3454 WorkerPreparationStartFault::InterpreterCorrupt(fault) => {
3455 WorkerPreparationError::InterpreterCorrupt(fault)
3456 }
3457 WorkerPreparationStartFault::InterpretationSkipped => {
3458 WorkerPreparationError::InterpretationSkipped
3459 }
3460 };
3461 let (returned_source, error) = match self.recovery.restore_source(failed.source) {
3462 Ok(()) => (None, error),
3463 Err(source) => (Some(source), WorkerPreparationError::SourceStateCorrupt),
3464 };
3465 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
3466 Actions::cont();
3467 actions.sends.diagnostics.append(InterpreterRequests::one(
3468 self.diagnostics.action(FifoDiagnostic::new(
3469 FifoDiagnosticCause::WorkerPreparationFailed {
3470 role: failed.role,
3471 previous: failed.previous,
3472 stopped: failed.stopped,
3473 returned_source,
3474 error,
3475 },
3476 )),
3477 ));
3478 self.continue_draining(members, deadline, actions)
3479 }
3480 Err((workers, input)) => self.diagnose(
3481 PoolState::Draining {
3482 members: workers,
3483 deadline,
3484 },
3485 FifoEvent::WorkerPreparationStarted(input),
3486 ),
3487 }
3488 }
3489
3490 fn accept_drain_preparation_return(
3491 &mut self,
3492 members: Vec<RetiringWorker<Role, W, P>>,
3493 deadline: ShutdownDeadline,
3494 input: WorkerPreparation<Source, Role, W, P>,
3495 ) -> (
3496 PoolState<Role, W, P, Job, WorkerResult>,
3497 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3498 )
3499 where
3500 Role: Eq,
3501 {
3502 let (members, accepted) =
3503 match direct_worker::accept_retiring_worker_preparation_return(members, input) {
3504 Ok(accepted) => accepted,
3505 Err((workers, input)) => {
3506 return self.diagnose(
3507 PoolState::Draining {
3508 members: workers,
3509 deadline,
3510 },
3511 FifoEvent::WorkerPreparationReturned(input),
3512 );
3513 }
3514 };
3515 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
3516 Actions::cont();
3517 let diagnostic = match accepted {
3518 WorkerRecoveryPreparation::Ready(direct_worker::PreparedReplacement {
3519 role,
3520 recoveries: _,
3521 previous,
3522 stopped,
3523 source,
3524 submission,
3525 }) => {
3526 let (returned_source, error) = match self.recovery.restore_source(source) {
3527 Ok(()) => (None, WorkerReplacementError::PoolShutdown),
3528 Err(source) => (Some(source), WorkerReplacementError::SourceStateCorrupt),
3529 };
3530 FifoDiagnosticCause::WorkerReplacementFailed {
3531 role,
3532 previous,
3533 stopped,
3534 submission,
3535 returned_source,
3536 error,
3537 }
3538 }
3539 WorkerRecoveryPreparation::Failed {
3540 role,
3541 recoveries: _,
3542 previous,
3543 stopped,
3544 source,
3545 error,
3546 } => {
3547 let (returned_source, error) = match self.recovery.restore_source(source) {
3548 Ok(()) => (None, error),
3549 Err(source) => (Some(source), WorkerPreparationError::SourceStateCorrupt),
3550 };
3551 FifoDiagnosticCause::WorkerPreparationFailed {
3552 role,
3553 previous,
3554 stopped,
3555 returned_source,
3556 error,
3557 }
3558 }
3559 };
3560 actions.sends.diagnostics.append(InterpreterRequests::one(
3561 self.diagnostics.action(FifoDiagnostic::new(diagnostic)),
3562 ));
3563 self.continue_draining(members, deadline, actions)
3564 }
3565
3566 fn continue_draining(
3567 &self,
3568 members: Vec<RetiringWorker<Role, W, P>>,
3569 deadline: ShutdownDeadline,
3570 mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3571 ) -> (
3572 PoolState<Role, W, P, Job, WorkerResult>,
3573 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3574 ) {
3575 match deadline {
3576 ShutdownDeadline::NotScheduled(settlement) => {
3577 actions.become_ = behavior::Step::Stop(behavior::Stopped);
3578 (
3579 PoolState::ForcedRetirement {
3580 members,
3581 cause: ForcedRetirementCause::DeadlineNotScheduled(settlement),
3582 },
3583 actions,
3584 )
3585 }
3586 ShutdownDeadline::Elapsed(elapsed) => {
3587 actions.become_ = behavior::Step::Stop(behavior::Stopped);
3588 (
3589 PoolState::ForcedRetirement {
3590 members,
3591 cause: ForcedRetirementCause::DeadlineElapsed(elapsed),
3592 },
3593 actions,
3594 )
3595 }
3596 deadline @ (ShutdownDeadline::Unlimited
3597 | ShutdownDeadline::Scheduling(_)
3598 | ShutdownDeadline::Waiting(_)) => {
3599 match members.iter().find_map(RetiringWorker::waiting) {
3600 Some(_) => (PoolState::Draining { members, deadline }, actions),
3601 None => {
3602 actions.become_ = behavior::Step::Stop(behavior::Stopped);
3603 (PoolState::Stopped, actions)
3604 }
3605 }
3606 }
3607 }
3608 }
3609
3610 fn accept_drain_stop(
3611 &self,
3612 members: Vec<RetiringWorker<Role, W, P>>,
3613 deadline: ShutdownDeadline,
3614 stopped: crate::ChildStopped<BehaviorAddr<W>>,
3615 ) -> (
3616 PoolState<Role, W, P, Job, WorkerResult>,
3617 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3618 ) {
3619 match direct_worker::accept_worker_exit(members, stopped) {
3620 Ok(ControlFlow::Continue(members)) => {
3621 self.continue_draining(members, deadline, Actions::cont())
3622 }
3623 Ok(ControlFlow::Break((members, worker))) => {
3624 let mut actions = Actions::cont();
3625 self.retain_drained_worker(worker, &mut actions);
3626 self.continue_draining(members, deadline, actions)
3627 }
3628 Err((members, stopped)) => self.diagnose(
3629 PoolState::Draining { members, deadline },
3630 FifoEvent::WorkerStopped(stopped),
3631 ),
3632 }
3633 }
3634
3635 fn accept_drain_shutdown(
3636 &self,
3637 members: Vec<RetiringWorker<Role, W, P>>,
3638 deadline: ShutdownDeadline,
3639 settled: EstablishedShutdownResolved<W::Protocol>,
3640 ) -> (
3641 PoolState<Role, W, P, Job, WorkerResult>,
3642 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3643 ) {
3644 match direct_worker::accept_worker_shutdown(members, settled) {
3645 Ok(ControlFlow::Continue(members)) => {
3646 self.continue_draining(members, deadline, Actions::cont())
3647 }
3648 Ok(ControlFlow::Break((members, worker))) => {
3649 let mut actions = Actions::cont();
3650 self.retain_drained_worker(worker, &mut actions);
3651 self.continue_draining(members, deadline, actions)
3652 }
3653 Err((members, settled)) => self.diagnose(
3654 PoolState::Draining { members, deadline },
3655 FifoEvent::WorkerShutdownSettled(settled),
3656 ),
3657 }
3658 }
3659
3660 fn accept_drain_schedule(
3661 &self,
3662 mut members: Vec<RetiringWorker<Role, W, P>>,
3663 mut deadline: ShutdownDeadline,
3664 settled: behavior::ActionItemResult<ScheduleAfter>,
3665 ) -> (
3666 PoolState<Role, W, P, Job, WorkerResult>,
3667 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3668 ) {
3669 let replacement = members.iter().position(|member| match &member.state {
3670 direct_worker::RetirementStatus::AwaitingRestartSchedule { timer, .. } => {
3671 timer.accepts_result(&settled)
3672 }
3673 direct_worker::RetirementStatus::AwaitingCreation(_)
3674 | direct_worker::RetirementStatus::Established { .. }
3675 | direct_worker::RetirementStatus::AwaitingPreparation { .. }
3676 | direct_worker::RetirementStatus::Drained => false,
3677 });
3678 if let Some(position) = replacement {
3679 let RetiringWorker { role, state } = members.remove(position);
3680 let direct_worker::RetirementStatus::AwaitingRestartSchedule {
3681 replacement,
3682 timer: _,
3683 } = state
3684 else {
3685 members.insert(position, RetiringWorker { role, state });
3686 return self.diagnose(
3687 PoolState::Draining { members, deadline },
3688 FifoEvent::RestartScheduleSettled(settled),
3689 );
3690 };
3691 members.insert(
3692 position,
3693 RetiringWorker {
3694 role: role.clone(),
3695 state: direct_worker::RetirementStatus::Drained,
3696 },
3697 );
3698 let mut actions = Actions::cont();
3699 self.retain_cancelled_replacement(
3700 role,
3701 replacement,
3702 WorkerReplacementError::RestartScheduleReturned(settled),
3703 &mut actions,
3704 );
3705 return self.continue_draining(members, deadline, actions);
3706 }
3707 match deadline.accept_schedule(settled) {
3708 Ok(()) => self.continue_draining(members, deadline, Actions::cont()),
3709 Err(settled) => self.diagnose(
3710 PoolState::Draining { members, deadline },
3711 FifoEvent::RestartScheduleSettled(settled),
3712 ),
3713 }
3714 }
3715
3716 fn accept_drain_deadline(
3717 &self,
3718 members: Vec<RetiringWorker<Role, W, P>>,
3719 mut deadline: ShutdownDeadline,
3720 elapsed: crate::TimerElapsed,
3721 ) -> (
3722 PoolState<Role, W, P, Job, WorkerResult>,
3723 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3724 ) {
3725 match deadline.accept_elapsed(elapsed) {
3726 Ok(()) => self.continue_draining(members, deadline, Actions::cont()),
3727 Err(elapsed) => self.diagnose(
3728 PoolState::Draining { members, deadline },
3729 FifoEvent::RestartElapsed(elapsed),
3730 ),
3731 }
3732 }
3733
3734 fn diagnose(
3735 &self,
3736 state: PoolState<Role, W, P, Job, WorkerResult>,
3737 input: FifoEvent<
3738 Role,
3739 W,
3740 P,
3741 Job,
3742 WorkerResult,
3743 behavior::ActionItemResult<PrepareWorkers<Source, Role, W, P>>,
3744 WorkerPreparation<Source, Role, W, P>,
3745 >,
3746 ) -> (
3747 PoolState<Role, W, P, Job, WorkerResult>,
3748 FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult>,
3749 ) {
3750 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
3751 Actions::cont();
3752 actions.sends.diagnostics = InterpreterRequests::one(
3753 self.diagnostics
3754 .action(FifoDiagnostic::new(FifoDiagnosticCause::Unexpected(input))),
3755 );
3756 (state, actions)
3757 }
3758}
3759
3760impl<Role, W, P, Source, DiagnosticRoute, Job, WorkerResult> BehaviorBase
3761 for FifoPool<Role, W, P, Source, DiagnosticRoute, Job, WorkerResult>
3762where
3763 Role: Send + Sync,
3764 W: Behavior + Send,
3765 W::Protocol: behavior::Protocol<Msg = Assignment<Job>>,
3766 W::Sends: CompletesAssignments<WorkerResult = WorkerResult>,
3767 P: ActivationPlan,
3768 Source: WorkerSource<Role, W, P>,
3769 BehaviorAddr<W>: EndpointAddress,
3770 Job: Send,
3771{
3772 type Base = Self;
3773
3774 fn base(&self) -> &Self::Base {
3775 self
3776 }
3777}
3778
3779enum PoolState<Role, W, P, Job, WorkerResult>
3780where
3781 W: Behavior,
3782 BehaviorAddr<W>: EndpointAddress,
3783{
3784 Constructed(Vec<PreparedWorker<Role, W, P>>),
3785 Operating(FifoOperating<Role, W, P, Job, WorkerResult>),
3786 Draining {
3787 members: Vec<RetiringWorker<Role, W, P>>,
3788 deadline: ShutdownDeadline,
3789 },
3790 Stopped,
3791 ForcedRetirement {
3792 #[expect(
3793 dead_code,
3794 reason = "Bombay's retirement custodian receives every unresolved member"
3795 )]
3796 members: Vec<RetiringWorker<Role, W, P>>,
3797 #[expect(
3798 dead_code,
3799 reason = "Bombay's retirement custodian receives the exact forced-retirement cause"
3800 )]
3801 cause: ForcedRetirementCause,
3802 },
3803}
3804
3805#[derive(Debug, Error)]
3807pub enum FifoError {
3808 #[error("FIFO pool initialization is no longer available")]
3810 InitializationUnavailable,
3811 #[error("FIFO pool worker creation identifiers are exhausted")]
3813 WorkerCreationsExhausted,
3814}
3815
3816impl<Role, W, P, Job, WorkerResult> PoolState<Role, W, P, Job, WorkerResult>
3817where
3818 W: Behavior + BehaviorBase,
3819 P: ActivationPlan,
3820 BehaviorAddr<W>: EndpointAddress,
3821 StopOnShutdown<W>:
3822 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
3823{
3824 fn start(
3825 self,
3826 creations: &mut CreationSequence,
3827 ) -> Result<
3828 (
3829 Self,
3830 Vec<ObserveChild<W::Protocol, behavior::ChildHead>>,
3831 Creations<CreateChild<BehaviorAddr<W>, StopOnShutdown<W>>>,
3832 ),
3833 (Self, FifoError),
3834 > {
3835 let prepared = match self {
3836 Self::Constructed(prepared) => prepared,
3837 state => return Err((state, FifoError::InitializationUnavailable)),
3838 };
3839 let mut ids = Vec::with_capacity(prepared.len());
3840 for _ in 0..prepared.len() {
3841 let Some(id) = creations.issue() else {
3842 return Err((
3843 Self::Constructed(prepared),
3844 FifoError::WorkerCreationsExhausted,
3845 ));
3846 };
3847 ids.push(id);
3848 }
3849
3850 let mut observations = Vec::with_capacity(prepared.len());
3851 let mut workers = Creations::empty();
3852 let members = prepared
3853 .into_iter()
3854 .zip(ids)
3855 .map(|(prepared, creation)| {
3856 let (member, worker, observation) = Member::begin(prepared, creation);
3857 workers.extend([worker]);
3858 observations.push(observation);
3859 member
3860 })
3861 .collect();
3862
3863 Ok((
3864 Self::Operating(FifoOperating {
3865 members,
3866 backlog: BTreeMap::new(),
3867 cursor: 0,
3868 }),
3869 observations,
3870 workers,
3871 ))
3872 }
3873}
3874
3875pub struct FifoPool<Role, W, P, Source, DiagnosticRoute, Job, WorkerResult>
3877where
3878 W: Behavior,
3879 BehaviorAddr<W>: EndpointAddress,
3880{
3881 state: PoolState<Role, W, P, Job, WorkerResult>,
3882 activation: ActivationPolicy,
3883 recovery: PoolRecoveryState<Source>,
3884 restarts: super::restart::RestartBudget,
3885 backlog: BacklogCapacity,
3886 interruption: Interruption,
3887 actor_drain: ActorDrainPolicy,
3888 diagnostics: DiagnosticDisposition<DiagnosticRoute>,
3889 creations: behavior::CreationSequence,
3890 jobs: AcceptedJobSequence,
3891 assignments: AssignmentSequence,
3892 shutdowns: ShutdownSequence,
3893 next_restart_timer: u64,
3894}
3895
3896pub struct FifoConstructionRejected<Factory, Role, W, P, Rejection, Source, DiagnosticRoute> {
3898 pub workers: InitialWorkerRejection<Factory, Role, W, P, Rejection>,
3900 pub activation: ActivationPolicy,
3902 pub recovery: PoolRecovery<Source>,
3904 pub backlog: BacklogCapacity,
3906 pub interruption: Interruption,
3908 pub actor_drain: ActorDrainPolicy,
3910 pub diagnostics: DiagnosticDisposition<DiagnosticRoute>,
3912}
3913
3914pub fn fifo<Factory, Role, W, P, Rejection, Source, DiagnosticRoute, Job, WorkerResult>(
3916 factory: Factory,
3917 roles: OrderedRoles<Role>,
3918 activation: ActivationPolicy,
3919 recovery: PoolRecovery<Source>,
3920 backlog: BacklogCapacity,
3921 interruption: Interruption,
3922 actor_drain: ActorDrainPolicy,
3923 diagnostics: DiagnosticDisposition<DiagnosticRoute>,
3924) -> Result<
3925 FifoPool<Role, W, P, Source, DiagnosticRoute, Job, WorkerResult>,
3926 FifoConstructionRejected<Factory, Role, W, P, Rejection, Source, DiagnosticRoute>,
3927>
3928where
3929 Factory: FnMut(&Role) -> core::result::Result<WorkerSubmission<W, P>, Rejection>,
3930 Role: Eq,
3931 W: Behavior,
3932 W::Protocol: behavior::Protocol<Msg = Assignment<Job>>,
3933 W::Sends: CompletesAssignments<WorkerResult = WorkerResult>,
3934 BehaviorAddr<W>: EndpointAddress,
3935{
3936 let prepared = match prepare_initial_workers(factory, roles) {
3937 Ok(prepared) => prepared,
3938 Err(workers) => {
3939 return Err(FifoConstructionRejected {
3940 workers,
3941 activation,
3942 recovery,
3943 backlog,
3944 interruption,
3945 actor_drain,
3946 diagnostics,
3947 });
3948 }
3949 };
3950 Ok(FifoPool {
3951 state: PoolState::Constructed(prepared),
3952 activation,
3953 recovery: recovery.into(),
3954 restarts: super::restart::RestartBudget::empty(),
3955 backlog,
3956 interruption,
3957 actor_drain,
3958 diagnostics,
3959 creations: behavior::CreationSequence::new(),
3960 jobs: AcceptedJobSequence::new(),
3961 assignments: AssignmentSequence::new(),
3962 shutdowns: ShutdownSequence::new(),
3963 next_restart_timer: 1,
3964 })
3965}
3966
3967impl<Role, W, P, Source, Diagnostics, Job, WorkerResult> Behavior
3968 for FifoPool<Role, W, P, Source, Diagnostics, Job, WorkerResult>
3969where
3970 Role: Eq + Send + Sync,
3971 W: Behavior + BehaviorBase + Send,
3972 W::Protocol: Protocol<Msg = Assignment<Job>>,
3973 W::Sends: CompletesAssignments<WorkerResult = WorkerResult>,
3974 P: ActivationPlan,
3975 Source: WorkerSource<Role, W, P>,
3976 Diagnostics: DiagnosticRoute<FifoDiagnostic<Role, W, P, Source, Job, WorkerResult>> + Clone,
3977 BehaviorAddr<W>: EndpointAddress,
3978 <BehaviorAddr<W> as Address>::Nonce: Send,
3979 <BehaviorAddr<W> as EndpointAddress>::Established<W::Protocol>: Send,
3980 StopOnShutdown<W>:
3981 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
3982 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
3983 Job: Clone + Send,
3984 WorkerResult: Send,
3985{
3986 type Protocol =
3987 MessageProtocol<BehaviorAddr<W>, FifoCommand<BehaviorAddr<W>, Role, Job, WorkerResult>>;
3988 type Event = FifoEvent<
3989 Role,
3990 W,
3991 P,
3992 Job,
3993 WorkerResult,
3994 behavior::ActionItemResult<PrepareWorkers<Source, Role, W, P>>,
3995 WorkerPreparation<Source, Role, W, P>,
3996 >;
3997 type Sends = FifoRequests<
3998 InterpreterRequests<ObserveChild<W::Protocol, ChildHead>>,
3999 InterpreterRequests<InitializeWorker<W, P>>,
4000 InterpreterRequests<BeginActivation<W, P>>,
4001 <CustomerRoute<BehaviorAddr<W>, Role, Job, WorkerResult> as DeliveryRoute>::Sends,
4002 SourceActions<AssignWorker<W::Protocol, Job>>,
4003 SourceActions<PrepareWorkers<Source, Role, W, P>>,
4004 SourceActions<ScheduleAfter>,
4005 InterpreterRequests<ShutdownEstablished<StopOnShutdown<W>, Here>>,
4006 InterpreterRequests<
4007 DiagnosticAction<Diagnostics, FifoDiagnostic<Role, W, P, Source, Job, WorkerResult>>,
4008 >,
4009 >;
4010 type Ph = Never;
4011 type Error = FifoError;
4012 type Birth = Births<StopOnShutdown<W>>;
4013
4014 fn init(&mut self, _: InitializationTurn) -> BehaviorActed<Self> {
4015 let state = mem::replace(&mut self.state, PoolState::Stopped);
4016 match state.start(&mut self.creations) {
4017 Ok((state, observations, workers)) => {
4018 self.state = state;
4019 let mut sends = FifoRequests::empty();
4020 sends.worker_observations = InterpreterRequests::new(observations);
4021 Ok(Actions::new(sends, workers, behavior::Step::Continue))
4022 }
4023 Err((state, error)) => {
4024 self.state = state;
4025 Err(error)
4026 }
4027 }
4028 }
4029
4030 fn transition(&mut self, _: ActiveTurn, input: Self::Event) -> BehaviorActed<Self> {
4031 let state = mem::replace(&mut self.state, PoolState::Stopped);
4032 let (state, actions) = match (state, input) {
4033 (
4034 PoolState::Operating(operating),
4035 FifoEvent::Command(User {
4036 message:
4037 FifoCommand::Submit {
4038 submission,
4039 payload,
4040 customer,
4041 },
4042 ..
4043 }),
4044 ) => {
4045 let (operating, actions) =
4046 self.accept_submission(operating, submission, payload, customer);
4047 (PoolState::Operating(operating), actions)
4048 }
4049 (
4050 PoolState::Operating(operating),
4051 FifoEvent::Command(User {
4052 message: FifoCommand::Shutdown,
4053 ..
4054 })
4055 | FifoEvent::Shutdown(_),
4056 ) => self.begin_shutdown(operating, Actions::cont()),
4057 (PoolState::Operating(operating), FifoEvent::WorkerCreationsSettled(workers)) => {
4058 match self.accept_creations(operating, workers) {
4059 Ok(result) => result,
4060 Err((operating, workers)) => self.diagnose(
4061 PoolState::Operating(operating),
4062 FifoEvent::WorkerCreationsSettled(workers),
4063 ),
4064 }
4065 }
4066 (PoolState::Operating(operating), FifoEvent::WorkerInitialization(input)) => {
4067 match self.accept_initialization(operating, input) {
4068 Ok(result) => result,
4069 Err((operating, input)) => self.diagnose(
4070 PoolState::Operating(operating),
4071 FifoEvent::WorkerInitialization(input),
4072 ),
4073 }
4074 }
4075 (PoolState::Operating(operating), FifoEvent::WorkerActivationReported(input)) => {
4076 match self.accept_activation(operating, input) {
4077 Ok(result) => result,
4078 Err((operating, input)) => self.diagnose(
4079 PoolState::Operating(operating),
4080 FifoEvent::WorkerActivationReported(input),
4081 ),
4082 }
4083 }
4084 (PoolState::Operating(operating), FifoEvent::WorkerPreparationStarted(input)) => {
4085 match self.accept_worker_preparation_start(operating, input) {
4086 Ok(result) => result,
4087 Err((operating, input)) => self.diagnose(
4088 PoolState::Operating(operating),
4089 FifoEvent::WorkerPreparationStarted(input),
4090 ),
4091 }
4092 }
4093 (PoolState::Operating(operating), FifoEvent::WorkerPreparationReturned(input)) => {
4094 match self.accept_worker_preparation_return(operating, input) {
4095 Ok(result) => result,
4096 Err((operating, input)) => self.diagnose(
4097 PoolState::Operating(operating),
4098 FifoEvent::WorkerPreparationReturned(input),
4099 ),
4100 }
4101 }
4102 (PoolState::Operating(operating), FifoEvent::RestartScheduleSettled(input)) => {
4103 match self.accept_restart_schedule(operating, input) {
4104 Ok(result) => result,
4105 Err((operating, input)) => self.diagnose(
4106 PoolState::Operating(operating),
4107 FifoEvent::RestartScheduleSettled(input),
4108 ),
4109 }
4110 }
4111 (PoolState::Operating(operating), FifoEvent::RestartElapsed(elapsed)) => {
4112 match self.accept_restart_timer(operating, elapsed) {
4113 Ok(result) => result,
4114 Err((operating, elapsed)) => self.diagnose(
4115 PoolState::Operating(operating),
4116 FifoEvent::RestartElapsed(elapsed),
4117 ),
4118 }
4119 }
4120 (
4121 PoolState::Operating(operating),
4122 FifoEvent::AssignmentSettled(SettledItem::Attempted(ItemSettlement::Accepted(
4123 receipt,
4124 ))),
4125 ) => match self.accept_assignment_receipt(operating, receipt) {
4126 Ok(result) => result,
4127 Err((operating, receipt)) => self.diagnose(
4128 PoolState::Operating(operating),
4129 FifoEvent::AssignmentSettled(SettledItem::Attempted(ItemSettlement::Accepted(
4130 receipt,
4131 ))),
4132 ),
4133 },
4134 (
4135 PoolState::Operating(operating),
4136 FifoEvent::AssignmentSettled(SettledItem::Attempted(ItemSettlement::Rejected {
4137 item,
4138 reason,
4139 })),
4140 ) => self.reject_assignment_delivery(operating, item, reason),
4141 (PoolState::Operating(operating), FifoEvent::WorkerCompleted(completion)) => {
4142 let ChildReport { child, report } = completion;
4143 match self.accept_completion(operating, child, report) {
4144 Ok(result) => result,
4145 Err((operating, completion)) => self.diagnose(
4146 PoolState::Operating(operating),
4147 FifoEvent::WorkerCompleted(completion),
4148 ),
4149 }
4150 }
4151 (PoolState::Operating(operating), FifoEvent::WorkerStopped(stopped)) => {
4152 match self.accept_worker_stop(operating, stopped) {
4153 Ok(result) => result,
4154 Err((operating, stopped)) => self.diagnose(
4155 PoolState::Operating(operating),
4156 FifoEvent::WorkerStopped(stopped),
4157 ),
4158 }
4159 }
4160 (PoolState::Operating(operating), FifoEvent::WorkerShutdownSettled(shutdown)) => {
4161 match self.accept_worker_shutdown(operating, shutdown) {
4162 Ok(result) => result,
4163 Err((operating, shutdown)) => self.diagnose(
4164 PoolState::Operating(operating),
4165 FifoEvent::WorkerShutdownSettled(shutdown),
4166 ),
4167 }
4168 }
4169 (
4170 PoolState::Draining { members, deadline },
4171 FifoEvent::Command(User {
4172 message:
4173 FifoCommand::Submit {
4174 submission,
4175 payload,
4176 customer,
4177 },
4178 ..
4179 }),
4180 ) => {
4181 let mut actions: FifoActions<Role, W, P, Source, Diagnostics, Job, WorkerResult> =
4182 Actions::cont();
4183 actions.sends.customer_outcomes = customer.deliver(FifoOutcome::rejected(
4184 submission,
4185 payload,
4186 AdmissionRejection::ShuttingDown,
4187 ));
4188 (PoolState::Draining { members, deadline }, actions)
4189 }
4190 (
4191 state @ PoolState::Draining { .. },
4192 FifoEvent::Command(User {
4193 message: FifoCommand::Shutdown,
4194 ..
4195 })
4196 | FifoEvent::Shutdown(_),
4197 ) => (state, Actions::cont()),
4198 (
4199 PoolState::Draining { members, deadline },
4200 FifoEvent::WorkerCreationsSettled(workers),
4201 ) => self.accept_drain_creations(members, deadline, workers),
4202 (PoolState::Draining { members, deadline }, FifoEvent::WorkerStopped(stopped)) => {
4203 self.accept_drain_stop(members, deadline, stopped)
4204 }
4205 (
4206 PoolState::Draining { members, deadline },
4207 FifoEvent::WorkerInitialization(initialization),
4208 ) => self.accept_drain_initialization(members, deadline, initialization),
4209 (
4210 PoolState::Draining { members, deadline },
4211 FifoEvent::WorkerActivationReported(activation),
4212 ) => self.accept_drain_activation(members, deadline, activation),
4213 (
4214 PoolState::Draining { members, deadline },
4215 FifoEvent::WorkerShutdownSettled(settled),
4216 ) => self.accept_drain_shutdown(members, deadline, settled),
4217 (
4218 PoolState::Draining { members, deadline },
4219 FifoEvent::WorkerPreparationStarted(started),
4220 ) => self.accept_drain_preparation_start(members, deadline, started),
4221 (
4222 PoolState::Draining { members, deadline },
4223 FifoEvent::WorkerPreparationReturned(returned),
4224 ) => self.accept_drain_preparation_return(members, deadline, returned),
4225 (
4226 PoolState::Draining { members, deadline },
4227 FifoEvent::RestartScheduleSettled(settled),
4228 ) => self.accept_drain_schedule(members, deadline, settled),
4229 (PoolState::Draining { members, deadline }, FifoEvent::RestartElapsed(elapsed)) => {
4230 self.accept_drain_deadline(members, deadline, elapsed)
4231 }
4232 (state @ (PoolState::Stopped | PoolState::ForcedRetirement { .. }), input) => {
4233 let (state, mut actions) = self.diagnose(state, input);
4234 actions.become_ = behavior::Step::Stop(behavior::Stopped);
4235 (state, actions)
4236 }
4237 (state, input) => self.diagnose(state, input),
4238 };
4239 self.state = state;
4240 Ok(actions)
4241 }
4242}