1use core::cmp::Ordering;
4use core::num::NonZeroU64;
5use core::ops::Bound::{Excluded, Unbounded};
6use std::collections::BTreeMap;
7
8use behavior::{
9 Actions, ActiveTurn, Address, Behavior, BehaviorActed, BehaviorAddr, BehaviorBase, Births,
10 ChildHead, ChildReport, CreateChild, CreationId, CreationSequence, Creations, EndpointAddress,
11 EstablishedActor, InitializationTurn, InterpreterRequests, ItemSettlement, MessageProtocol,
12 Never, SendEffects, SettledItem, SourceActions, Step, User,
13};
14
15use crate::{
16 ActivationPlan, ActorDrainPolicy, ChildStopped, DeliveryRoute, DiagnosticAction,
17 DiagnosticDisposition, InitialWorkerOutcome, ObserveChild, ProxyOutcome, ReplacementOutcome,
18 ReplyRoute, ScheduleAfter, WorkerStartResult, WorkerSubmission,
19};
20
21use super::drain::ShutdownDeadline;
22use super::{ActivationPolicy, EntryCapacity, ProxyInputResult, ProxyOperation, StableProxy};
23use crate::atomic::proxy_creation::{StableProxyCreation, identity, resolve};
24use crate::atomic::stable_proxy::ProxyOperationWitness;
25
26mod diagnostic;
27mod entry;
28mod event;
29mod lifecycle;
30mod protocol;
31mod requests;
32
33pub use diagnostic::DynamicDiagnostic;
34pub use event::DynamicSupervisorEvent;
35pub use lifecycle::{
36 CancellationOutcome, DynamicLifecycle, EntryRetirement, EntryStopFailure,
37 EntryStopFailureReason, InterruptedWorker, ReplacementFailure, WorkerChange,
38 WorkerChangeInterruption,
39};
40pub use protocol::{
41 CancelAuthority, CancellationReceipt, DynamicCommand, DynamicStatus, QueryReply,
42 ReplaceRejection, StartRejection, StopRejection, WorkerChangeReceipt, WorkerChangeRejection,
43};
44pub use requests::DynamicSupervisorRequests;
45
46use entry::{
47 DynamicEntry, DynamicEntryPhase, ProxyInputCustody, ProxyInputPurpose, ProxyRetirement,
48 ProxyShutdown, RetiringWork, ServiceAvailability, WorkerChangeDisposition,
49};
50
51#[derive(Clone, Copy, Debug, Eq, PartialEq)]
53pub enum UnexpectedExit {
54 KeepEmpty,
56 Retire,
58}
59
60enum SupervisorAvailability {
61 Accepting(ActorDrainPolicy),
62 ShuttingDown(ShutdownDeadline),
63}
64
65pub struct DynamicSupervisor<Key, Worker, Plan, LifecycleRoute, DiagnosticRoute>
67where
68 Worker: Behavior,
69 Plan: ActivationPlan,
70 BehaviorAddr<Worker>: EndpointAddress,
71 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
72{
73 services: BTreeMap<Key, DynamicEntry<Worker, Plan>>,
74 entries: EntryCapacity,
75 activation: ActivationPolicy,
76 generations: Option<NonZeroU64>,
77 operations: Option<NonZeroU64>,
78 creations: CreationSequence,
79 unexpected_exit: UnexpectedExit,
80 availability: SupervisorAvailability,
81 lifecycle: LifecycleRoute,
82 diagnostics: DiagnosticDisposition<DiagnosticRoute>,
83}
84
85#[must_use]
87pub fn dynamic<Key, Worker, Plan, LifecycleRoute, DiagnosticRoute>(
88 entries: EntryCapacity,
89 activation: ActivationPolicy,
90 unexpected_exit: UnexpectedExit,
91 actor_drain: ActorDrainPolicy,
92 lifecycle: LifecycleRoute,
93 diagnostics: DiagnosticDisposition<DiagnosticRoute>,
94) -> DynamicSupervisor<Key, Worker, Plan, LifecycleRoute, DiagnosticRoute>
95where
96 Worker: Behavior,
97 Plan: ActivationPlan,
98 BehaviorAddr<Worker>: EndpointAddress,
99 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
100{
101 DynamicSupervisor {
102 services: BTreeMap::new(),
103 entries,
104 activation,
105 generations: Some(NonZeroU64::MIN),
106 operations: Some(NonZeroU64::MIN),
107 creations: CreationSequence::new(),
108 unexpected_exit,
109 availability: SupervisorAvailability::Accepting(actor_drain),
110 lifecycle,
111 diagnostics,
112 }
113}
114
115impl<Key, Worker, Plan, LifecycleRoute, DiagnosticRouteType> BehaviorBase
116 for DynamicSupervisor<Key, Worker, Plan, LifecycleRoute, DiagnosticRouteType>
117where
118 Worker: Behavior,
119 Plan: ActivationPlan,
120 BehaviorAddr<Worker>: EndpointAddress,
121 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
122{
123 type Base = Self;
124
125 fn base(&self) -> &Self::Base {
126 self
127 }
128}
129
130impl<Key, Worker, Plan, LifecycleRoute, DiagnosticRouteType>
131 DynamicSupervisor<Key, Worker, Plan, LifecycleRoute, DiagnosticRouteType>
132where
133 Key: Clone + Ord + Send + Sync,
134 Worker: Behavior + Send,
135 Plan: ActivationPlan,
136 BehaviorAddr<Worker>: EndpointAddress,
137 <BehaviorAddr<Worker> as Address>::Nonce: Send,
138 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
139 EstablishedActor<StableProxy<Worker, Plan>>: Clone + Send,
140 LifecycleRoute: DeliveryRoute<
141 Protocol = MessageProtocol<BehaviorAddr<Worker>, DynamicLifecycle<Key, Worker, Plan>>,
142 > + Clone,
143 LifecycleRoute::Sends: behavior::SendsFor<DynamicSupervisorEvent<Key, Worker, Plan>>,
144 DiagnosticRouteType: crate::DiagnosticRoute<DynamicDiagnostic<Key, Worker, Plan>> + Clone,
145{
146 fn take_proxy_entry(
147 &mut self,
148 creation: CreationId,
149 ) -> Option<(Key, DynamicEntry<Worker, Plan>)> {
150 self.services
151 .extract_if(.., |_, entry| entry.creation == creation)
152 .next()
153 }
154
155 fn reject_input(
156 &mut self,
157 input: DynamicSupervisorEvent<Key, Worker, Plan>,
158 ) -> BehaviorActed<Self> {
159 Err(input)
160 }
161
162 fn authorize_waiting(&mut self) -> SourceActions<ProxyOperation<behavior::Here, Worker, Plan>> {
163 let occupied: usize = self
164 .services
165 .values()
166 .map(|service| match service.phase {
167 DynamicEntryPhase::WaitingForProxyInput { .. }
168 | DynamicEntryPhase::WaitingForProxyOutcome { .. }
169 | DynamicEntryPhase::ReplacementAwaitingReceipt { .. }
170 | DynamicEntryPhase::ReplacementAwaitingOutcome { .. }
171 | DynamicEntryPhase::Retiring(ProxyRetirement {
172 work:
173 RetiringWork::Transferred {
174 input:
175 ProxyInputCustody::AwaitingSettlement(_)
176 | ProxyInputCustody::AwaitingOutcome(_),
177 ..
178 },
179 ..
180 }) => 1,
181 _ => 0,
182 })
183 .sum();
184 match occupied.cmp(&self.activation.maximum()) {
185 Ordering::Equal | Ordering::Greater => return SourceActions::empty(),
186 Ordering::Less => {}
187 }
188 let waiting = self
189 .services
190 .iter()
191 .filter_map(|(key, entry)| match entry.phase {
192 DynamicEntryPhase::WaitingForActivation { .. }
193 | DynamicEntryPhase::ReplacementQueued { .. } => {
194 Some((entry.operation, key.clone()))
195 }
196 _ => None,
197 })
198 .min_by_key(|(operation, _)| *operation);
199 let Some((_, key)) = waiting else {
200 return SourceActions::empty();
201 };
202 let Some(entry) = self.services.remove(&key) else {
203 return SourceActions::empty();
204 };
205 let DynamicEntry {
206 creation,
207 generation,
208 operation,
209 phase,
210 } = entry;
211 match phase {
212 DynamicEntryPhase::WaitingForActivation { proxy, submission } => {
213 let (witness, request) = ProxyOperation::initial(creation, submission);
214 self.services.insert(
215 key,
216 DynamicEntry {
217 creation,
218 generation,
219 operation,
220 phase: DynamicEntryPhase::WaitingForProxyInput { proxy, witness },
221 },
222 );
223 SourceActions::sending(request)
224 }
225 DynamicEntryPhase::ReplacementQueued {
226 proxy,
227 service,
228 submission,
229 } => {
230 let (witness, request) = ProxyOperation::replacement(creation, submission);
231 self.services.insert(
232 key,
233 DynamicEntry {
234 creation,
235 generation,
236 operation,
237 phase: DynamicEntryPhase::ReplacementAwaitingReceipt {
238 proxy,
239 service,
240 witness,
241 },
242 },
243 );
244 SourceActions::sending(request)
245 }
246 phase => {
247 self.services.insert(
248 key,
249 DynamicEntry {
250 creation,
251 generation,
252 operation,
253 phase,
254 },
255 );
256 SourceActions::empty()
257 }
258 }
259 }
260
261 fn start(
262 &mut self,
263 key: Key,
264 submission: WorkerSubmission<Worker, Plan>,
265 reply_to: ReplyRoute<
266 MessageProtocol<
267 BehaviorAddr<Worker>,
268 Result<
269 WorkerChangeReceipt<Key>,
270 WorkerChangeRejection<Key, Worker, Plan, StartRejection>,
271 >,
272 >,
273 >,
274 ) -> BehaviorActed<Self> {
275 let admission = if self.services.contains_key(&key) {
276 Err(StartRejection::AlreadyExists)
277 } else if self.services.len() >= self.entries.maximum() {
278 Err(StartRejection::AtCapacity)
279 } else {
280 match (self.generations, self.operations) {
281 (Some(generation), Some(operation)) => Ok((generation, operation)),
282 (None, _) => Err(StartRejection::EntryGenerationExhausted),
283 (Some(_), None) => Err(StartRejection::OperationExhausted),
284 }
285 };
286 let (generation, operation) = match admission {
287 Ok(admission) => admission,
288 Err(reason) => {
289 let mut requests = DynamicSupervisorRequests::empty();
290 requests.start_replies = reply_to.deliver(Err(WorkerChangeRejection {
291 key,
292 submission,
293 reason,
294 }));
295 return Ok(Actions::send(requests));
296 }
297 };
298 let creation = match self.creations.issue() {
299 Some(creation) => creation,
300 None => {
301 let mut requests = DynamicSupervisorRequests::empty();
302 requests.start_replies = reply_to.deliver(Err(WorkerChangeRejection {
303 key,
304 submission,
305 reason: StartRejection::ProxyCreationExhausted,
306 }));
307 return Ok(Actions::send(requests));
308 }
309 };
310 self.generations = generation.checked_add(1);
311 self.operations = operation.checked_add(1);
312 let authority = CancelAuthority::issued(key.clone(), operation.get());
313 let replaced = self.services.insert(
314 key.clone(),
315 DynamicEntry {
316 creation,
317 generation,
318 operation,
319 phase: DynamicEntryPhase::CreatingProxy { submission },
320 },
321 );
322 debug_assert!(replaced.is_none(), "duplicate admission was rejected");
323 let mut requests = DynamicSupervisorRequests::empty();
324 requests.proxy_observations =
325 InterpreterRequests::one(ObserveChild::<Worker::Protocol, ChildHead>::new(creation));
326 requests.start_replies = reply_to.deliver(Ok(WorkerChangeReceipt {
327 key: key.clone(),
328 cancel: authority,
329 }));
330 let creations = Creations::one(CreateChild::birth(
331 creation,
332 StableProxy::<Worker, Plan>::activated(),
333 ));
334 Ok(Actions::new(requests, creations, Step::Continue))
335 }
336
337 fn replace(
338 &mut self,
339 key: Key,
340 submission: WorkerSubmission<Worker, Plan>,
341 reply_to: ReplyRoute<
342 MessageProtocol<
343 BehaviorAddr<Worker>,
344 Result<
345 WorkerChangeReceipt<Key>,
346 WorkerChangeRejection<Key, Worker, Plan, ReplaceRejection>,
347 >,
348 >,
349 >,
350 ) -> BehaviorActed<Self> {
351 let Some(entry) = self.services.remove(&key) else {
352 let mut requests = DynamicSupervisorRequests::empty();
353 requests.replace_replies = reply_to.deliver(Err(WorkerChangeRejection {
354 key,
355 submission,
356 reason: ReplaceRejection::Unknown,
357 }));
358 return Ok(Actions::send(requests));
359 };
360 let DynamicEntry {
361 creation,
362 generation,
363 operation: previous_operation,
364 phase,
365 } = entry;
366 let (creation, proxy, service, latest) = match phase {
367 DynamicEntryPhase::Available {
368 proxy,
369 service,
370 latest,
371 } => (creation, proxy, service, latest),
372 phase => {
373 self.services.insert(
374 key.clone(),
375 DynamicEntry {
376 creation,
377 generation,
378 operation: previous_operation,
379 phase,
380 },
381 );
382 let mut requests = DynamicSupervisorRequests::empty();
383 requests.replace_replies = reply_to.deliver(Err(WorkerChangeRejection {
384 key,
385 submission,
386 reason: ReplaceRejection::Unavailable,
387 }));
388 return Ok(Actions::send(requests));
389 }
390 };
391 let Some(operation) = self.operations else {
392 self.services.insert(
393 key.clone(),
394 DynamicEntry {
395 creation,
396 generation,
397 operation: previous_operation,
398 phase: DynamicEntryPhase::Available {
399 proxy,
400 service,
401 latest,
402 },
403 },
404 );
405 let mut requests = DynamicSupervisorRequests::empty();
406 requests.replace_replies = reply_to.deliver(Err(WorkerChangeRejection {
407 key,
408 submission,
409 reason: ReplaceRejection::OperationExhausted,
410 }));
411 return Ok(Actions::send(requests));
412 };
413 self.operations = operation.checked_add(1);
414 let authority = CancelAuthority::issued(key.clone(), operation.get());
415 self.services.insert(
416 key.clone(),
417 DynamicEntry {
418 creation,
419 generation,
420 operation,
421 phase: DynamicEntryPhase::ReplacementQueued {
422 proxy,
423 service,
424 submission,
425 },
426 },
427 );
428 let mut requests = DynamicSupervisorRequests::empty();
429 requests.proxy_operations = self.authorize_waiting();
430 requests.replace_replies = reply_to.deliver(Ok(WorkerChangeReceipt {
431 key,
432 cancel: authority,
433 }));
434 Ok(Actions::send(requests))
435 }
436
437 fn stop(
438 &mut self,
439 key: Key,
440 reply_to: ReplyRoute<
441 MessageProtocol<BehaviorAddr<Worker>, Result<Key, StopRejection<Key>>>,
442 >,
443 ) -> BehaviorActed<Self> {
444 let Some(entry) = self.services.remove(&key) else {
445 let mut requests = DynamicSupervisorRequests::empty();
446 requests.stop_replies = reply_to.deliver(Err(StopRejection::Unknown { key }));
447 return Ok(Actions::send(requests));
448 };
449 let DynamicEntry {
450 creation,
451 generation,
452 operation: previous_operation,
453 phase,
454 } = entry;
455 let (creation, proxy, restorable, operation, lifecycle) = match phase {
456 DynamicEntryPhase::Available { proxy, service, .. } => {
457 let Some(operation) = self.operations else {
458 return self.stop_operation_exhausted(
459 key,
460 creation,
461 generation,
462 previous_operation,
463 DynamicEntryPhase::Available {
464 proxy,
465 service,
466 latest: WorkerChangeDisposition::Committed,
467 },
468 reply_to,
469 );
470 };
471 self.operations = operation.checked_add(1);
472 (
473 creation,
474 proxy,
475 Some(service),
476 operation,
477 SendEffects::empty(),
478 )
479 }
480 DynamicEntryPhase::CreatingProxy { submission } => {
481 let Some(operation) = self.operations else {
482 return self.stop_operation_exhausted(
483 key,
484 creation,
485 generation,
486 previous_operation,
487 DynamicEntryPhase::CreatingProxy { submission },
488 reply_to,
489 );
490 };
491 self.operations = operation.checked_add(1);
492 self.services.insert(
493 key.clone(),
494 DynamicEntry {
495 creation,
496 generation,
497 operation,
498 phase: DynamicEntryPhase::StoppingProxyCreation {
499 start_operation: previous_operation,
500 submission,
501 },
502 },
503 );
504 let mut requests = DynamicSupervisorRequests::empty();
505 requests.stop_replies = reply_to.deliver(Ok(key));
506 return Ok(Actions::send(requests));
507 }
508 DynamicEntryPhase::WaitingForActivation { proxy, submission } => {
509 let Some(operation) = self.operations else {
510 return self.stop_operation_exhausted(
511 key,
512 creation,
513 generation,
514 previous_operation,
515 DynamicEntryPhase::WaitingForActivation { proxy, submission },
516 reply_to,
517 );
518 };
519 self.operations = operation.checked_add(1);
520 let lifecycle =
521 self.lifecycle
522 .clone()
523 .deliver(DynamicLifecycle::WorkerChangeInterrupted {
524 key: key.clone(),
525 generation: generation.get(),
526 operation: previous_operation.get(),
527 interruption: WorkerChangeInterruption::ExplicitStop {
528 operation: operation.get(),
529 },
530 worker: InterruptedWorker::Submission(submission),
531 });
532 (creation, proxy, None, operation, lifecycle)
533 }
534 DynamicEntryPhase::WaitingForProxyInput { proxy, witness } => {
535 let Some(stop_operation) = self.operations else {
536 return self.stop_operation_exhausted(
537 key,
538 creation,
539 generation,
540 previous_operation,
541 DynamicEntryPhase::WaitingForProxyInput { proxy, witness },
542 reply_to,
543 );
544 };
545 self.operations = stop_operation.checked_add(1);
546 let (retirement, shutdown) = ProxyRetirement::begin(
547 creation,
548 proxy,
549 RetiringWork::Transferred {
550 purpose: ProxyInputPurpose::InterruptedStart(previous_operation),
551 input: ProxyInputCustody::AwaitingSettlement(witness),
552 },
553 );
554 self.services.insert(
555 key.clone(),
556 DynamicEntry {
557 creation,
558 generation,
559 operation: stop_operation,
560 phase: DynamicEntryPhase::Retiring(retirement),
561 },
562 );
563 let mut requests = DynamicSupervisorRequests::empty();
564 requests.proxy_operations = SourceActions::sending(shutdown);
565 requests.stop_replies = reply_to.deliver(Ok(key));
566 return Ok(Actions::send(requests));
567 }
568 DynamicEntryPhase::WaitingForProxyOutcome {
569 proxy,
570 operation: proxy_operation,
571 } => {
572 let Some(stop_operation) = self.operations else {
573 return self.stop_operation_exhausted(
574 key,
575 creation,
576 generation,
577 previous_operation,
578 DynamicEntryPhase::WaitingForProxyOutcome {
579 proxy,
580 operation: proxy_operation,
581 },
582 reply_to,
583 );
584 };
585 self.operations = stop_operation.checked_add(1);
586 let (retirement, shutdown) = ProxyRetirement::begin(
587 creation,
588 proxy,
589 RetiringWork::Transferred {
590 purpose: ProxyInputPurpose::InterruptedStart(previous_operation),
591 input: ProxyInputCustody::AwaitingOutcome(proxy_operation),
592 },
593 );
594 self.services.insert(
595 key.clone(),
596 DynamicEntry {
597 creation,
598 generation,
599 operation: stop_operation,
600 phase: DynamicEntryPhase::Retiring(retirement),
601 },
602 );
603 let mut requests = DynamicSupervisorRequests::empty();
604 requests.proxy_operations = SourceActions::sending(shutdown);
605 requests.stop_replies = reply_to.deliver(Ok(key));
606 return Ok(Actions::send(requests));
607 }
608 phase @ (DynamicEntryPhase::Stopping { .. }
609 | DynamicEntryPhase::StoppingProxyCreation { .. }) => {
610 self.services.insert(
611 key.clone(),
612 DynamicEntry {
613 creation,
614 generation,
615 operation: previous_operation,
616 phase,
617 },
618 );
619 let mut requests = DynamicSupervisorRequests::empty();
620 requests.stop_replies =
621 reply_to.deliver(Err(StopRejection::AlreadyStopping { key }));
622 return Ok(Actions::send(requests));
623 }
624 phase => {
625 self.services.insert(
626 key.clone(),
627 DynamicEntry {
628 creation,
629 generation,
630 operation: previous_operation,
631 phase,
632 },
633 );
634 let mut requests = DynamicSupervisorRequests::empty();
635 requests.stop_replies = reply_to.deliver(Err(StopRejection::Unavailable { key }));
636 return Ok(Actions::send(requests));
637 }
638 };
639 let (witness, shutdown) = ProxyOperation::shutdown(creation);
640 self.services.insert(
641 key.clone(),
642 DynamicEntry {
643 creation,
644 generation,
645 operation,
646 phase: DynamicEntryPhase::Stopping {
647 proxy,
648 restorable,
649 shutdown: ProxyShutdown::AwaitingSettlement(witness),
650 stopped: None,
651 },
652 },
653 );
654 let mut requests = DynamicSupervisorRequests::empty();
655 requests.proxy_operations = SourceActions::sending(shutdown);
656 requests.stop_replies = reply_to.deliver(Ok(key));
657 requests.lifecycle = lifecycle;
658 Ok(Actions::send(requests))
659 }
660
661 fn stop_operation_exhausted(
662 &mut self,
663 key: Key,
664 creation: CreationId,
665 generation: NonZeroU64,
666 operation: NonZeroU64,
667 phase: DynamicEntryPhase<Worker, Plan>,
668 reply_to: ReplyRoute<
669 MessageProtocol<BehaviorAddr<Worker>, Result<Key, StopRejection<Key>>>,
670 >,
671 ) -> BehaviorActed<Self> {
672 self.services.insert(
673 key.clone(),
674 DynamicEntry {
675 creation,
676 generation,
677 operation,
678 phase,
679 },
680 );
681 let mut requests = DynamicSupervisorRequests::empty();
682 requests.stop_replies = reply_to.deliver(Err(StopRejection::OperationExhausted { key }));
683 Ok(Actions::send(requests))
684 }
685
686 fn next_service_key(&self, previous: Option<&Key>) -> Option<Key> {
687 match previous {
688 Some(previous) => self.services.range((Excluded(previous), Unbounded)).next(),
689 None => self.services.first_key_value(),
690 }
691 .map(|(key, _)| key.clone())
692 }
693
694 fn begin_shutdown(&mut self) -> BehaviorActed<Self> {
695 let actor_drain = match &self.availability {
696 SupervisorAvailability::Accepting(actor_drain) => *actor_drain,
697 SupervisorAvailability::ShuttingDown(_) => return Ok(Actions::cont()),
698 };
699 let (deadline, schedule) = ShutdownDeadline::begin(actor_drain);
700 self.availability = SupervisorAvailability::ShuttingDown(deadline);
701 let mut requests: <Self as Behavior>::Sends = SendEffects::empty();
702 match schedule {
703 Some(schedule) => requests.shutdown_schedules.send(schedule),
704 None => {}
705 }
706
707 let mut next = self.next_service_key(None);
708 while let Some(key) = next {
709 let entry = match self.services.remove(&key) {
710 Some(entry) => entry,
711 None => {
712 next = self.next_service_key(Some(&key));
713 continue;
714 }
715 };
716 next = self.next_service_key(Some(&key));
717 let DynamicEntry {
718 creation,
719 generation,
720 operation,
721 phase,
722 } = entry;
723 let (phase, shutdown, interruption) = phase.shutdown(creation, operation);
724 match shutdown {
725 Some(shutdown) => requests.proxy_operations.send(shutdown),
726 None => {}
727 }
728 match interruption {
729 Some((operation, interruption, submission)) => {
730 requests.lifecycle.append(self.lifecycle.clone().deliver(
731 DynamicLifecycle::WorkerChangeInterrupted {
732 key: key.clone(),
733 generation: generation.get(),
734 operation: operation.get(),
735 interruption,
736 worker: InterruptedWorker::Submission(submission),
737 },
738 ))
739 }
740 None => {}
741 }
742 self.services.insert(
743 key,
744 DynamicEntry {
745 creation,
746 generation,
747 operation,
748 phase,
749 },
750 );
751 }
752 Ok(Actions::send(requests))
753 }
754
755 fn apply_shutdown_step(&self, acted: BehaviorActed<Self>) -> BehaviorActed<Self> {
756 acted.map(|mut actions| {
757 match (&self.availability, self.services.first_key_value()) {
758 (SupervisorAvailability::Accepting(_), _) => {}
759 (SupervisorAvailability::ShuttingDown(_), None) => {
760 actions.become_ = Step::Stop(behavior::Stopped);
761 }
762 (
763 SupervisorAvailability::ShuttingDown(ShutdownDeadline::NotScheduled(returned)),
764 Some(_),
765 ) => {
766 let _ = returned;
767 actions.become_ = Step::Stop(behavior::Stopped);
768 }
769 (
770 SupervisorAvailability::ShuttingDown(ShutdownDeadline::Elapsed(elapsed)),
771 Some(_),
772 ) => {
773 let _ = elapsed;
774 actions.become_ = Step::Stop(behavior::Stopped);
775 }
776 (SupervisorAvailability::ShuttingDown(_), Some(_)) => {}
777 }
778 actions
779 })
780 }
781
782 fn command(&mut self, command: DynamicCommand<Key, Worker, Plan>) -> BehaviorActed<Self> {
783 match &self.availability {
784 SupervisorAvailability::Accepting(_) => {}
785 SupervisorAvailability::ShuttingDown(_) => {
786 return self.command_while_shutting_down(command);
787 }
788 }
789 match command {
790 DynamicCommand::Start {
791 key,
792 submission,
793 reply_to,
794 } => self.start(key, submission, reply_to),
795 DynamicCommand::Replace {
796 key,
797 submission,
798 reply_to,
799 } => self.replace(key, submission, reply_to),
800 DynamicCommand::Query { key, reply_to } => {
801 let reply = match self.services.get(&key) {
802 Some(entry) => QueryReply::Known {
803 key,
804 status: entry.phase.status(),
805 },
806 None => QueryReply::Unknown { key },
807 };
808 let mut requests = DynamicSupervisorRequests::empty();
809 requests.query_replies = reply_to.deliver(reply);
810 Ok(Actions::send(requests))
811 }
812 DynamicCommand::Cancel {
813 authority,
814 reply_to,
815 } => self.cancel(authority, reply_to),
816 DynamicCommand::Stop { key, reply_to } => self.stop(key, reply_to),
817 }
818 }
819
820 fn command_while_shutting_down(
821 &mut self,
822 command: DynamicCommand<Key, Worker, Plan>,
823 ) -> BehaviorActed<Self> {
824 let mut requests = DynamicSupervisorRequests::empty();
825 match command {
826 DynamicCommand::Start {
827 key,
828 submission,
829 reply_to,
830 } => {
831 requests.start_replies = reply_to.deliver(Err(WorkerChangeRejection {
832 key,
833 submission,
834 reason: StartRejection::ShuttingDown,
835 }));
836 }
837 DynamicCommand::Replace {
838 key,
839 submission,
840 reply_to,
841 } => {
842 requests.replace_replies = reply_to.deliver(Err(WorkerChangeRejection {
843 key,
844 submission,
845 reason: ReplaceRejection::ShuttingDown,
846 }));
847 }
848 DynamicCommand::Stop { key, reply_to } => {
849 requests.stop_replies = reply_to.deliver(Err(StopRejection::ShuttingDown { key }));
850 }
851 DynamicCommand::Query { key, reply_to } => {
852 let reply = match self.services.get(&key) {
853 Some(_) => QueryReply::Known {
854 key,
855 status: DynamicStatus::Draining,
856 },
857 None => QueryReply::Unknown { key },
858 };
859 requests.query_replies = reply_to.deliver(reply);
860 }
861 DynamicCommand::Cancel {
862 authority,
863 reply_to,
864 } => {
865 requests.cancel_replies =
866 reply_to.deliver(CancellationReceipt::Draining { authority });
867 }
868 }
869 Ok(Actions::send(requests))
870 }
871
872 fn cancel(
873 &mut self,
874 authority: CancelAuthority<Key>,
875 reply_to: ReplyRoute<
876 MessageProtocol<BehaviorAddr<Worker>, CancellationReceipt<Key, Worker, Plan>>,
877 >,
878 ) -> BehaviorActed<Self> {
879 let key = authority.key().clone();
880 let Some(entry) = self.services.remove(&key) else {
881 let mut requests = DynamicSupervisorRequests::empty();
882 requests.cancel_replies = reply_to.deliver(CancellationReceipt::Stale { authority });
883 return Ok(Actions::send(requests));
884 };
885 if authority.operation() != entry.operation.get() {
886 self.services.insert(key, entry);
887 let mut requests = DynamicSupervisorRequests::empty();
888 requests.cancel_replies = reply_to.deliver(CancellationReceipt::Stale { authority });
889 return Ok(Actions::send(requests));
890 }
891 let DynamicEntry {
892 creation,
893 generation,
894 operation,
895 phase,
896 } = entry;
897 let mut proxy_operations = SourceActions::empty();
898 let (phase, receipt) = match phase {
899 DynamicEntryPhase::CreatingProxy { submission } => (
900 DynamicEntryPhase::CancellingProxyCreation,
901 CancellationReceipt::Returned {
902 authority,
903 submission,
904 },
905 ),
906 DynamicEntryPhase::WaitingForActivation { proxy, submission } => {
907 let (retirement, shutdown_request) = ProxyRetirement::begin(
908 creation,
909 proxy,
910 RetiringWork::CancelledWorkerReturned(WorkerChange::Start),
911 );
912 proxy_operations.send(shutdown_request);
913 (
914 DynamicEntryPhase::Retiring(retirement),
915 CancellationReceipt::Returned {
916 authority,
917 submission,
918 },
919 )
920 }
921 DynamicEntryPhase::WaitingForProxyInput { proxy, witness } => {
922 let (retirement, shutdown_request) = ProxyRetirement::begin(
923 creation,
924 proxy,
925 RetiringWork::Transferred {
926 purpose: ProxyInputPurpose::Cancellation(WorkerChange::Start),
927 input: ProxyInputCustody::AwaitingSettlement(witness),
928 },
929 );
930 proxy_operations.send(shutdown_request);
931 (
932 DynamicEntryPhase::Retiring(retirement),
933 CancellationReceipt::Pending {
934 authority,
935 phase: DynamicStatus::AwaitingProxy,
936 },
937 )
938 }
939 DynamicEntryPhase::WaitingForProxyOutcome {
940 proxy,
941 operation: proxy_operation,
942 } => {
943 let (retirement, shutdown_request) = ProxyRetirement::begin(
944 creation,
945 proxy,
946 RetiringWork::Transferred {
947 purpose: ProxyInputPurpose::Cancellation(WorkerChange::Start),
948 input: ProxyInputCustody::AwaitingOutcome(proxy_operation),
949 },
950 );
951 proxy_operations.send(shutdown_request);
952 (
953 DynamicEntryPhase::Retiring(retirement),
954 CancellationReceipt::Pending {
955 authority,
956 phase: DynamicStatus::AwaitingProxy,
957 },
958 )
959 }
960 phase @ DynamicEntryPhase::CancellingProxyCreation => {
961 (phase, CancellationReceipt::Cancelled { authority })
962 }
963 phase @ DynamicEntryPhase::Retiring(ProxyRetirement {
964 work:
965 RetiringWork::CancelledWorkerReturned(_)
966 | RetiringWork::Transferred {
967 purpose: ProxyInputPurpose::Cancellation(_),
968 ..
969 },
970 ..
971 }) => (phase, CancellationReceipt::Cancelled { authority }),
972 phase @ DynamicEntryPhase::Retiring(_) => (
973 phase,
974 CancellationReceipt::Committed {
975 authority,
976 resulting_phase: DynamicStatus::Retiring,
977 },
978 ),
979 phase @ (DynamicEntryPhase::ShuttingDownProxyCreation
980 | DynamicEntryPhase::DrainingProxyCreationStop) => {
981 (phase, CancellationReceipt::Draining { authority })
982 }
983 DynamicEntryPhase::Available {
984 proxy,
985 service,
986 latest,
987 } => {
988 let resulting_phase = match &service {
989 ServiceAvailability::Ready { .. } => DynamicStatus::Ready {
990 proxy: proxy.recipient(),
991 },
992 ServiceAvailability::Empty { .. } => DynamicStatus::Empty,
993 };
994 match latest {
995 WorkerChangeDisposition::Committed => (
996 DynamicEntryPhase::Available {
997 proxy,
998 service,
999 latest: WorkerChangeDisposition::Committed,
1000 },
1001 CancellationReceipt::Committed {
1002 authority,
1003 resulting_phase,
1004 },
1005 ),
1006 WorkerChangeDisposition::Cancelled => (
1007 DynamicEntryPhase::Available {
1008 proxy,
1009 service,
1010 latest: WorkerChangeDisposition::Cancelled,
1011 },
1012 CancellationReceipt::Cancelled { authority },
1013 ),
1014 }
1015 }
1016 DynamicEntryPhase::ReplacementQueued {
1017 proxy,
1018 service,
1019 submission,
1020 } => (
1021 DynamicEntryPhase::Available {
1022 proxy,
1023 service,
1024 latest: WorkerChangeDisposition::Cancelled,
1025 },
1026 CancellationReceipt::Returned {
1027 authority,
1028 submission,
1029 },
1030 ),
1031 DynamicEntryPhase::ReplacementAwaitingReceipt {
1032 proxy,
1033 service: _,
1034 witness,
1035 } => {
1036 let (retirement, shutdown_request) = ProxyRetirement::begin(
1037 creation,
1038 proxy,
1039 RetiringWork::Transferred {
1040 purpose: ProxyInputPurpose::Cancellation(WorkerChange::Replacement),
1041 input: ProxyInputCustody::AwaitingSettlement(witness),
1042 },
1043 );
1044 proxy_operations.send(shutdown_request);
1045 (
1046 DynamicEntryPhase::Retiring(retirement),
1047 CancellationReceipt::Pending {
1048 authority,
1049 phase: DynamicStatus::Replacing,
1050 },
1051 )
1052 }
1053 DynamicEntryPhase::ReplacementAwaitingOutcome {
1054 proxy,
1055 service: _,
1056 operation: proxy_operation,
1057 } => {
1058 let (retirement, shutdown_request) = ProxyRetirement::begin(
1059 creation,
1060 proxy,
1061 RetiringWork::Transferred {
1062 purpose: ProxyInputPurpose::Cancellation(WorkerChange::Replacement),
1063 input: ProxyInputCustody::AwaitingOutcome(proxy_operation),
1064 },
1065 );
1066 proxy_operations.send(shutdown_request);
1067 (
1068 DynamicEntryPhase::Retiring(retirement),
1069 CancellationReceipt::Pending {
1070 authority,
1071 phase: DynamicStatus::Replacing,
1072 },
1073 )
1074 }
1075 phase @ (DynamicEntryPhase::Stopping { .. }
1076 | DynamicEntryPhase::StoppingProxyCreation { .. }) => (
1077 phase,
1078 CancellationReceipt::Committed {
1079 authority,
1080 resulting_phase: DynamicStatus::Stopping,
1081 },
1082 ),
1083 };
1084 self.services.insert(
1085 key,
1086 DynamicEntry {
1087 creation,
1088 generation,
1089 operation,
1090 phase,
1091 },
1092 );
1093 let mut requests = DynamicSupervisorRequests::empty();
1094 requests.proxy_operations = proxy_operations;
1095 requests.cancel_replies = reply_to.deliver(receipt);
1096 Ok(Actions::send(requests))
1097 }
1098
1099 fn accept_creations(
1100 &mut self,
1101 proxies: behavior::CreationsSettled<BehaviorAddr<Worker>, StableProxy<Worker, Plan>>,
1102 ) -> BehaviorActed<Self> {
1103 let settlements = match proxies.into_settlement() {
1104 behavior::CreationSettlement::Settled(settlements) => settlements,
1105 settlement => {
1106 return self.reject_input(DynamicSupervisorEvent::ProxyCreationsSettled(
1107 behavior::CreationsSettled::new(settlement),
1108 ));
1109 }
1110 };
1111 let settlement = match settlements.into_one() {
1112 Ok(settlement) => settlement,
1113 Err(settlements) => {
1114 return self.reject_input(DynamicSupervisorEvent::ProxyCreationsSettled(
1115 behavior::CreationsSettled::new(behavior::CreationSettlement::Settled(
1116 settlements,
1117 )),
1118 ));
1119 }
1120 };
1121 let (creation, _) = identity(&settlement);
1122 let Some((key, entry)) = self.take_proxy_entry(creation) else {
1123 return self.reject_input(DynamicSupervisorEvent::ProxyCreationsSettled(
1124 behavior::CreationsSettled::new(behavior::CreationSettlement::Settled(
1125 Creations::one(settlement),
1126 )),
1127 ));
1128 };
1129 let DynamicEntry {
1130 creation,
1131 generation,
1132 operation,
1133 phase,
1134 } = entry;
1135 let (creation, submission) = match phase {
1136 DynamicEntryPhase::CreatingProxy { submission } => (creation, submission),
1137 DynamicEntryPhase::ShuttingDownProxyCreation => match resolve(settlement) {
1138 StableProxyCreation::Committed(proxy) => {
1139 let (retirement, shutdown_request) = ProxyRetirement::begin(
1140 creation,
1141 proxy,
1142 RetiringWork::SupervisorShutdown(None),
1143 );
1144 self.services.insert(
1145 key,
1146 DynamicEntry {
1147 creation,
1148 generation,
1149 operation,
1150 phase: DynamicEntryPhase::Retiring(retirement),
1151 },
1152 );
1153 let mut requests = DynamicSupervisorRequests::empty();
1154 requests.proxy_operations = SourceActions::sending(shutdown_request);
1155 return Ok(Actions::send(requests));
1156 }
1157 StableProxyCreation::Rejected(creation_settlement) => {
1158 let mut requests = DynamicSupervisorRequests::empty();
1159 requests.lifecycle = self.entry_retired(
1160 key.clone(),
1161 generation.get(),
1162 EntryRetirement::Shutdown,
1163 );
1164 requests.diagnostics = InterpreterRequests::one(self.diagnostics.action(
1165 DynamicDiagnostic::ProxyCreationRejected {
1166 key,
1167 generation: generation.get(),
1168 creation: creation_settlement,
1169 },
1170 ));
1171 return Ok(Actions::send(requests));
1172 }
1173 },
1174 DynamicEntryPhase::DrainingProxyCreationStop => match resolve(settlement) {
1175 StableProxyCreation::Committed(proxy) => {
1176 let (witness, shutdown_request) = ProxyOperation::shutdown(creation);
1177 self.services.insert(
1178 key,
1179 DynamicEntry {
1180 creation,
1181 generation,
1182 operation,
1183 phase: DynamicEntryPhase::Stopping {
1184 proxy,
1185 restorable: None,
1186 shutdown: ProxyShutdown::AwaitingSettlement(witness),
1187 stopped: None,
1188 },
1189 },
1190 );
1191 let mut requests = DynamicSupervisorRequests::empty();
1192 requests.proxy_operations = SourceActions::sending(shutdown_request);
1193 return Ok(Actions::send(requests));
1194 }
1195 StableProxyCreation::Rejected(creation_settlement) => {
1196 let mut requests = DynamicSupervisorRequests::empty();
1197 requests.lifecycle =
1198 self.entry_retired(key.clone(), generation.get(), EntryRetirement::Stop);
1199 requests.diagnostics = InterpreterRequests::one(self.diagnostics.action(
1200 DynamicDiagnostic::ProxyCreationRejected {
1201 key,
1202 generation: generation.get(),
1203 creation: creation_settlement,
1204 },
1205 ));
1206 return Ok(Actions::send(requests));
1207 }
1208 },
1209 DynamicEntryPhase::CancellingProxyCreation => match resolve(settlement) {
1210 StableProxyCreation::Committed(proxy) => {
1211 let (retirement, shutdown_request) = ProxyRetirement::begin(
1212 creation,
1213 proxy,
1214 RetiringWork::CancelledWorkerReturned(WorkerChange::Start),
1215 );
1216 self.services.insert(
1217 key,
1218 DynamicEntry {
1219 creation,
1220 generation,
1221 operation,
1222 phase: DynamicEntryPhase::Retiring(retirement),
1223 },
1224 );
1225 let mut proxy_operations = SourceActions::empty();
1226 proxy_operations.send(shutdown_request);
1227 let mut requests = DynamicSupervisorRequests::empty();
1228 requests.proxy_operations = proxy_operations;
1229 return Ok(Actions::send(requests));
1230 }
1231 StableProxyCreation::Rejected(creation_settlement) => {
1232 let mut requests = DynamicSupervisorRequests::empty();
1233 requests.lifecycle = self.operation_cancelled(
1234 key.clone(),
1235 generation.get(),
1236 operation.get(),
1237 WorkerChange::Start,
1238 CancellationOutcome::WorkerReturned,
1239 );
1240 requests.diagnostics = InterpreterRequests::one(self.diagnostics.action(
1241 DynamicDiagnostic::ProxyCreationRejected {
1242 key,
1243 generation: generation.get(),
1244 creation: creation_settlement,
1245 },
1246 ));
1247 return Ok(Actions::send(requests));
1248 }
1249 },
1250 DynamicEntryPhase::StoppingProxyCreation {
1251 start_operation,
1252 submission,
1253 } => match resolve(settlement) {
1254 StableProxyCreation::Committed(proxy) => {
1255 let (witness, shutdown) = ProxyOperation::shutdown(creation);
1256 self.services.insert(
1257 key.clone(),
1258 DynamicEntry {
1259 creation,
1260 generation,
1261 operation,
1262 phase: DynamicEntryPhase::Stopping {
1263 proxy,
1264 restorable: None,
1265 shutdown: ProxyShutdown::AwaitingSettlement(witness),
1266 stopped: None,
1267 },
1268 },
1269 );
1270 let mut requests = DynamicSupervisorRequests::empty();
1271 requests.proxy_operations = SourceActions::sending(shutdown);
1272 requests.lifecycle =
1273 self.lifecycle
1274 .clone()
1275 .deliver(DynamicLifecycle::WorkerChangeInterrupted {
1276 key,
1277 generation: generation.get(),
1278 operation: start_operation.get(),
1279 interruption: WorkerChangeInterruption::ExplicitStop {
1280 operation: operation.get(),
1281 },
1282 worker: InterruptedWorker::Submission(submission),
1283 });
1284 return Ok(Actions::send(requests));
1285 }
1286 StableProxyCreation::Rejected(creation_settlement) => {
1287 let mut requests = DynamicSupervisorRequests::empty();
1288 requests.lifecycle = self
1289 .lifecycle
1290 .clone()
1291 .deliver(DynamicLifecycle::WorkerChangeInterrupted {
1292 key: key.clone(),
1293 generation: generation.get(),
1294 operation: start_operation.get(),
1295 interruption: WorkerChangeInterruption::ExplicitStop {
1296 operation: operation.get(),
1297 },
1298 worker: InterruptedWorker::Submission(submission),
1299 })
1300 .combine(self.entry_retired(
1301 key.clone(),
1302 generation.get(),
1303 EntryRetirement::Stop,
1304 ));
1305 requests.diagnostics = InterpreterRequests::one(self.diagnostics.action(
1306 DynamicDiagnostic::ProxyCreationRejected {
1307 key,
1308 generation: generation.get(),
1309 creation: creation_settlement,
1310 },
1311 ));
1312 return Ok(Actions::send(requests));
1313 }
1314 },
1315 phase => {
1316 self.services.insert(
1317 key,
1318 DynamicEntry {
1319 creation,
1320 generation,
1321 operation,
1322 phase,
1323 },
1324 );
1325 return self.reject_input(DynamicSupervisorEvent::ProxyCreationsSettled(
1326 behavior::CreationsSettled::new(behavior::CreationSettlement::Settled(
1327 Creations::one(settlement),
1328 )),
1329 ));
1330 }
1331 };
1332 match resolve(settlement) {
1333 StableProxyCreation::Committed(proxy) => {
1334 self.services.insert(
1335 key,
1336 DynamicEntry {
1337 creation,
1338 generation,
1339 operation,
1340 phase: DynamicEntryPhase::WaitingForActivation { proxy, submission },
1341 },
1342 );
1343 let mut requests = DynamicSupervisorRequests::empty();
1344 requests.proxy_operations = self.authorize_waiting();
1345 Ok(Actions::send(requests))
1346 }
1347 StableProxyCreation::Rejected(creation_settlement) => {
1348 let mut requests = DynamicSupervisorRequests::empty();
1349 requests.lifecycle =
1350 self.lifecycle
1351 .clone()
1352 .deliver(DynamicLifecycle::StartCreationRejected {
1353 key,
1354 generation: generation.get(),
1355 submission,
1356 creation: creation_settlement,
1357 });
1358 Ok(Actions::send(requests))
1359 }
1360 }
1361 }
1362
1363 fn accept_proxy_input(
1364 &mut self,
1365 input: ProxyInputResult<behavior::Here, Worker, Plan>,
1366 ) -> BehaviorActed<Self> {
1367 let creation = proxy_input_creation(&input);
1368 let Some((key, entry)) = self.take_proxy_entry(creation) else {
1369 return self.reject_input(DynamicSupervisorEvent::ProxyInputSettled(input));
1370 };
1371 let DynamicEntry {
1372 creation,
1373 generation,
1374 operation,
1375 phase,
1376 } = entry;
1377 match phase {
1378 DynamicEntryPhase::WaitingForProxyInput { proxy, witness } => {
1379 self.accept_start_input(key, generation, operation, creation, proxy, witness, input)
1380 }
1381 DynamicEntryPhase::ReplacementAwaitingReceipt {
1382 proxy,
1383 service,
1384 witness,
1385 } => self.accept_replacement_input(
1386 key,
1387 DynamicEntry {
1388 creation,
1389 generation,
1390 operation,
1391 phase: DynamicEntryPhase::ReplacementAwaitingReceipt {
1392 proxy,
1393 service,
1394 witness,
1395 },
1396 },
1397 input,
1398 ),
1399 DynamicEntryPhase::Stopping {
1400 proxy,
1401 restorable,
1402 shutdown: ProxyShutdown::AwaitingSettlement(witness),
1403 stopped,
1404 } => self.accept_stopping_input(
1405 key, generation, operation, creation, proxy, restorable, witness, stopped, input,
1406 ),
1407 phase @ DynamicEntryPhase::Retiring(_) => self.accept_retiring_input(
1408 key,
1409 DynamicEntry {
1410 creation,
1411 generation,
1412 operation,
1413 phase,
1414 },
1415 input,
1416 ),
1417 phase => {
1418 self.services.insert(
1419 key,
1420 DynamicEntry {
1421 creation,
1422 generation,
1423 operation,
1424 phase,
1425 },
1426 );
1427 self.reject_input(DynamicSupervisorEvent::ProxyInputSettled(input))
1428 }
1429 }
1430 }
1431
1432 fn accept_start_input(
1433 &mut self,
1434 key: Key,
1435 generation: NonZeroU64,
1436 operation: NonZeroU64,
1437 creation: CreationId,
1438 proxy: EstablishedActor<StableProxy<Worker, Plan>>,
1439 witness: ProxyOperationWitness,
1440 input: ProxyInputResult<behavior::Here, Worker, Plan>,
1441 ) -> BehaviorActed<Self> {
1442 let input = match witness.admit(input) {
1443 Ok(input) => input,
1444 Err((witness, input)) => {
1445 self.services.insert(
1446 key,
1447 DynamicEntry {
1448 creation,
1449 generation,
1450 operation,
1451 phase: DynamicEntryPhase::WaitingForProxyInput { proxy, witness },
1452 },
1453 );
1454 return self.reject_input(DynamicSupervisorEvent::ProxyInputSettled(input));
1455 }
1456 };
1457 match input {
1458 SettledItem::Attempted(ItemSettlement::Accepted(receipt)) => {
1459 let (_, proxy, proxy_operation) = receipt.into_parts();
1460 self.services.insert(
1461 key,
1462 DynamicEntry {
1463 creation,
1464 generation,
1465 operation,
1466 phase: DynamicEntryPhase::WaitingForProxyOutcome {
1467 proxy,
1468 operation: proxy_operation,
1469 },
1470 },
1471 );
1472 Ok(Actions::cont())
1473 }
1474 input => {
1475 let (retirement, operation_request) =
1476 ProxyRetirement::begin(creation, proxy, RetiringWork::StartFailed);
1477 self.services.insert(
1478 key.clone(),
1479 DynamicEntry {
1480 creation,
1481 generation,
1482 operation,
1483 phase: DynamicEntryPhase::Retiring(retirement),
1484 },
1485 );
1486 let mut proxy_operations = SourceActions::empty();
1487 proxy_operations.send(operation_request);
1488 let waiting = self.authorize_waiting();
1489 let mut requests = DynamicSupervisorRequests::empty();
1490 requests.proxy_operations = proxy_operations;
1491 requests.proxy_operations.append(waiting);
1492 requests.lifecycle =
1493 self.lifecycle
1494 .clone()
1495 .deliver(DynamicLifecycle::StartInputRejected {
1496 key: key.clone(),
1497 generation: generation.get(),
1498 });
1499 requests.diagnostics = InterpreterRequests::one(self.diagnostics.action(
1500 DynamicDiagnostic::ProxyInputRejected {
1501 key,
1502 generation: generation.get(),
1503 input,
1504 },
1505 ));
1506 Ok(Actions::send(requests))
1507 }
1508 }
1509 }
1510
1511 fn accept_replacement_input(
1512 &mut self,
1513 key: Key,
1514 entry: DynamicEntry<Worker, Plan>,
1515 input: ProxyInputResult<behavior::Here, Worker, Plan>,
1516 ) -> BehaviorActed<Self> {
1517 let DynamicEntry {
1518 creation,
1519 generation,
1520 operation,
1521 phase:
1522 DynamicEntryPhase::ReplacementAwaitingReceipt {
1523 proxy,
1524 service,
1525 witness,
1526 },
1527 } = entry
1528 else {
1529 return self.reject_input(DynamicSupervisorEvent::ProxyInputSettled(input));
1530 };
1531 let input = match witness.admit(input) {
1532 Ok(input) => input,
1533 Err((witness, input)) => {
1534 self.services.insert(
1535 key,
1536 DynamicEntry {
1537 creation,
1538 generation,
1539 operation,
1540 phase: DynamicEntryPhase::ReplacementAwaitingReceipt {
1541 proxy,
1542 service,
1543 witness,
1544 },
1545 },
1546 );
1547 return self.reject_input(DynamicSupervisorEvent::ProxyInputSettled(input));
1548 }
1549 };
1550 match input {
1551 SettledItem::Attempted(ItemSettlement::Accepted(receipt)) => {
1552 let (_, proxy, proxy_operation) = receipt.into_parts();
1553 self.services.insert(
1554 key,
1555 DynamicEntry {
1556 creation,
1557 generation,
1558 operation,
1559 phase: DynamicEntryPhase::ReplacementAwaitingOutcome {
1560 proxy,
1561 service,
1562 operation: proxy_operation,
1563 },
1564 },
1565 );
1566 Ok(Actions::cont())
1567 }
1568 input => {
1569 self.services.insert(
1570 key.clone(),
1571 DynamicEntry {
1572 creation,
1573 generation,
1574 operation,
1575 phase: DynamicEntryPhase::Available {
1576 proxy,
1577 service,
1578 latest: WorkerChangeDisposition::Committed,
1579 },
1580 },
1581 );
1582 let mut requests = DynamicSupervisorRequests::empty();
1583 requests.proxy_operations = self.authorize_waiting();
1584 requests.lifecycle =
1585 self.lifecycle
1586 .clone()
1587 .deliver(DynamicLifecycle::ReplacementInputRejected {
1588 key: key.clone(),
1589 generation: generation.get(),
1590 operation: operation.get(),
1591 });
1592 requests.diagnostics = InterpreterRequests::one(self.diagnostics.action(
1593 DynamicDiagnostic::ProxyInputRejected {
1594 key,
1595 generation: generation.get(),
1596 input,
1597 },
1598 ));
1599 Ok(Actions::send(requests))
1600 }
1601 }
1602 }
1603
1604 fn accept_stopping_input(
1605 &mut self,
1606 key: Key,
1607 generation: NonZeroU64,
1608 operation: NonZeroU64,
1609 creation: CreationId,
1610 proxy: EstablishedActor<StableProxy<Worker, Plan>>,
1611 restorable: Option<ServiceAvailability<Plan::Ready>>,
1612 witness: ProxyOperationWitness,
1613 stopped: Option<ChildStopped<BehaviorAddr<Worker>>>,
1614 input: ProxyInputResult<behavior::Here, Worker, Plan>,
1615 ) -> BehaviorActed<Self> {
1616 let input = match witness.admit(input) {
1617 Ok(input) => input,
1618 Err((witness, input)) => {
1619 self.services.insert(
1620 key,
1621 DynamicEntry {
1622 creation,
1623 generation,
1624 operation,
1625 phase: DynamicEntryPhase::Stopping {
1626 proxy,
1627 restorable,
1628 shutdown: ProxyShutdown::AwaitingSettlement(witness),
1629 stopped,
1630 },
1631 },
1632 );
1633 return self.reject_input(DynamicSupervisorEvent::ProxyInputSettled(input));
1634 }
1635 };
1636 let (outcome, input) = match entry_stop_result(stopped, input) {
1637 Ok(receipt) => {
1638 let (_, proxy, shutdown) = receipt.into_parts();
1639 let mut requests = DynamicSupervisorRequests::empty();
1640 match stopped {
1641 Some(stopped) => {
1642 requests.lifecycle = self
1643 .lifecycle
1644 .clone()
1645 .deliver(DynamicLifecycle::StopFinished {
1646 key: key.clone(),
1647 generation: generation.get(),
1648 operation: operation.get(),
1649 proxy,
1650 result: Ok(stopped),
1651 })
1652 .combine(self.entry_retired(
1653 key,
1654 generation.get(),
1655 EntryRetirement::Stop,
1656 ));
1657 }
1658 None => {
1659 self.services.insert(
1660 key,
1661 DynamicEntry {
1662 creation,
1663 generation,
1664 operation,
1665 phase: DynamicEntryPhase::Stopping {
1666 proxy,
1667 restorable,
1668 shutdown: ProxyShutdown::Accepted(shutdown),
1669 stopped: None,
1670 },
1671 },
1672 );
1673 }
1674 }
1675 return Ok(Actions::send(requests));
1676 }
1677 Err(rejected) => rejected,
1678 };
1679 let mut requests = DynamicSupervisorRequests::empty();
1680 requests.diagnostics = InterpreterRequests::one(self.diagnostics.action(
1681 DynamicDiagnostic::ProxyShutdownRejected {
1682 key: key.clone(),
1683 generation: generation.get(),
1684 input,
1685 },
1686 ));
1687 requests.lifecycle = match (outcome, restorable) {
1688 (failure @ EntryStopFailure { stopped: None, .. }, Some(service)) => {
1689 self.services.insert(
1690 key.clone(),
1691 DynamicEntry {
1692 creation,
1693 generation,
1694 operation,
1695 phase: DynamicEntryPhase::Available {
1696 proxy: proxy.clone(),
1697 service,
1698 latest: WorkerChangeDisposition::Committed,
1699 },
1700 },
1701 );
1702 self.lifecycle
1703 .clone()
1704 .deliver(DynamicLifecycle::StopFinished {
1705 key,
1706 generation: generation.get(),
1707 operation: operation.get(),
1708 proxy,
1709 result: Err(failure),
1710 })
1711 }
1712 (failure @ EntryStopFailure { stopped: None, .. }, None) => {
1713 self.services.insert(
1714 key.clone(),
1715 DynamicEntry {
1716 creation,
1717 generation,
1718 operation,
1719 phase: DynamicEntryPhase::Retiring(ProxyRetirement {
1720 proxy: proxy.clone(),
1721 shutdown: ProxyShutdown::Rejected,
1722 work: RetiringWork::InterruptedStart,
1723 stopped: None,
1724 }),
1725 },
1726 );
1727 self.lifecycle
1728 .clone()
1729 .deliver(DynamicLifecycle::StopFinished {
1730 key,
1731 generation: generation.get(),
1732 operation: operation.get(),
1733 proxy,
1734 result: Err(failure),
1735 })
1736 }
1737 (
1738 failure @ EntryStopFailure {
1739 stopped: Some(_), ..
1740 },
1741 _,
1742 ) => self
1743 .lifecycle
1744 .clone()
1745 .deliver(DynamicLifecycle::StopFinished {
1746 key: key.clone(),
1747 generation: generation.get(),
1748 operation: operation.get(),
1749 proxy,
1750 result: Err(failure),
1751 })
1752 .combine(self.entry_retired(key, generation.get(), EntryRetirement::Stop)),
1753 };
1754 Ok(Actions::send(requests))
1755 }
1756
1757 fn accept_retiring_input(
1758 &mut self,
1759 key: Key,
1760 entry: DynamicEntry<Worker, Plan>,
1761 input: ProxyInputResult<behavior::Here, Worker, Plan>,
1762 ) -> BehaviorActed<Self> {
1763 let DynamicEntry {
1764 creation,
1765 generation,
1766 operation,
1767 phase:
1768 DynamicEntryPhase::Retiring(ProxyRetirement {
1769 proxy,
1770 shutdown,
1771 work,
1772 stopped,
1773 }),
1774 } = entry
1775 else {
1776 return self.reject_input(DynamicSupervisorEvent::ProxyInputSettled(input));
1777 };
1778 let mut requests = DynamicSupervisorRequests::empty();
1779 let (proxy, shutdown, remaining_input) = match shutdown {
1780 ProxyShutdown::AwaitingSettlement(witness) => match witness.admit(input) {
1781 Ok(SettledItem::Attempted(ItemSettlement::Accepted(receipt))) => {
1782 let (_, proxy, operation) = receipt.into_parts();
1783 (proxy, ProxyShutdown::Accepted(operation), None)
1784 }
1785 Ok(input) => match &work {
1786 RetiringWork::Transferred {
1787 purpose: ProxyInputPurpose::InterruptedStart(_),
1788 ..
1789 } => match entry_stop_result(stopped, input) {
1790 Ok(receipt) => {
1791 let (_, proxy, operation) = receipt.into_parts();
1792 (proxy, ProxyShutdown::Accepted(operation), None)
1793 }
1794 Err((failure, input)) => {
1795 requests.lifecycle =
1796 self.lifecycle
1797 .clone()
1798 .deliver(DynamicLifecycle::StopFinished {
1799 key: key.clone(),
1800 generation: generation.get(),
1801 operation: operation.get(),
1802 proxy: proxy.clone(),
1803 result: Err(failure),
1804 });
1805 requests.diagnostics = InterpreterRequests::one(
1806 self.diagnostics
1807 .action(DynamicDiagnostic::ProxyShutdownRejected {
1808 key: key.clone(),
1809 generation: generation.get(),
1810 input,
1811 }),
1812 );
1813 (proxy, ProxyShutdown::Rejected, None)
1814 }
1815 },
1816 _ => {
1817 requests.diagnostics = InterpreterRequests::one(self.diagnostics.action(
1818 DynamicDiagnostic::ProxyShutdownRejected {
1819 key: key.clone(),
1820 generation: generation.get(),
1821 input,
1822 },
1823 ));
1824 (proxy, ProxyShutdown::Rejected, None)
1825 }
1826 },
1827 Err((witness, input)) => (
1828 proxy,
1829 ProxyShutdown::AwaitingSettlement(witness),
1830 Some(input),
1831 ),
1832 },
1833 shutdown => (proxy, shutdown, Some(input)),
1834 };
1835 let work = match remaining_input {
1836 None => work,
1837 Some(input) => match work.settle_input(input) {
1838 Ok((work, rejected)) => {
1839 if let Some(input) = rejected {
1840 requests.diagnostics = InterpreterRequests::one(self.diagnostics.action(
1841 DynamicDiagnostic::ProxyInputRejected {
1842 key: key.clone(),
1843 generation: generation.get(),
1844 input,
1845 },
1846 ));
1847 }
1848 work
1849 }
1850 Err((work, input)) => {
1851 self.services.insert(
1852 key,
1853 DynamicEntry {
1854 creation,
1855 generation,
1856 operation,
1857 phase: DynamicEntryPhase::Retiring(ProxyRetirement {
1858 proxy,
1859 shutdown,
1860 work,
1861 stopped,
1862 }),
1863 },
1864 );
1865 return self.reject_input(DynamicSupervisorEvent::ProxyInputSettled(input));
1866 }
1867 },
1868 };
1869 let retirement = ProxyRetirement {
1870 proxy,
1871 shutdown,
1872 work,
1873 stopped,
1874 };
1875 requests.lifecycle = requests
1876 .lifecycle
1877 .combine(self.settle_retirement(key, creation, generation, operation, retirement));
1878 requests.proxy_operations = self.authorize_waiting();
1879 Ok(Actions::send(requests))
1880 }
1881
1882 #[cfg(test)]
1883 #[expect(
1884 dead_code,
1885 reason = "compile-only proof of closed retirement lifecycle construction"
1886 )]
1887 fn retirement_lifecycle_is_closed(
1888 &self,
1889 key: Key,
1890 generation: u64,
1891 operation: u64,
1892 cause: EntryRetirement,
1893 change: WorkerChange,
1894 outcome: CancellationOutcome<Worker, Plan>,
1895 ) {
1896 let _ = self.entry_retired(key.clone(), generation, cause);
1897 let _ = self.operation_cancelled(key, generation, operation, change, outcome);
1898 }
1899
1900 fn entry_retired(
1901 &self,
1902 key: Key,
1903 generation: u64,
1904 cause: EntryRetirement,
1905 ) -> LifecycleRoute::Sends {
1906 self.lifecycle
1907 .clone()
1908 .deliver(DynamicLifecycle::EntryRetired {
1909 key,
1910 generation,
1911 cause,
1912 })
1913 }
1914
1915 fn operation_cancelled(
1916 &self,
1917 key: Key,
1918 generation: u64,
1919 operation: u64,
1920 change: WorkerChange,
1921 outcome: CancellationOutcome<Worker, Plan>,
1922 ) -> LifecycleRoute::Sends {
1923 let lifecycle = self.lifecycle.clone();
1924 lifecycle
1925 .clone()
1926 .deliver(DynamicLifecycle::OperationCancelled {
1927 key: key.clone(),
1928 generation,
1929 operation,
1930 change,
1931 outcome,
1932 })
1933 .combine(lifecycle.deliver(DynamicLifecycle::EntryRetired {
1934 key,
1935 generation,
1936 cause: EntryRetirement::Cancellation,
1937 }))
1938 }
1939
1940 fn settle_retirement(
1941 &mut self,
1942 key: Key,
1943 creation: CreationId,
1944 generation: NonZeroU64,
1945 operation: NonZeroU64,
1946 retirement: ProxyRetirement<Worker, Plan>,
1947 ) -> LifecycleRoute::Sends {
1948 match self.close_retirement(key.clone(), generation.get(), operation.get(), retirement) {
1949 Ok(lifecycle) => lifecycle,
1950 Err(retirement) => {
1951 self.services.insert(
1952 key,
1953 DynamicEntry {
1954 creation,
1955 generation,
1956 operation,
1957 phase: DynamicEntryPhase::Retiring(retirement),
1958 },
1959 );
1960 SendEffects::empty()
1961 }
1962 }
1963 }
1964
1965 fn close_retirement(
1966 &self,
1967 key: Key,
1968 generation: u64,
1969 operation: u64,
1970 retirement: ProxyRetirement<Worker, Plan>,
1971 ) -> Result<LifecycleRoute::Sends, ProxyRetirement<Worker, Plan>> {
1972 let ProxyRetirement {
1973 proxy,
1974 shutdown,
1975 work,
1976 stopped,
1977 } = retirement;
1978 let stopped = match (shutdown, stopped) {
1979 (shutdown @ (ProxyShutdown::Accepted(_) | ProxyShutdown::Rejected), Some(stopped)) => {
1980 (shutdown, stopped)
1981 }
1982 (shutdown, stopped) => {
1983 return Err(ProxyRetirement {
1984 proxy,
1985 shutdown,
1986 work,
1987 stopped,
1988 });
1989 }
1990 };
1991 match work {
1992 RetiringWork::StartFailed => {
1993 Ok(self.entry_retired(key, generation, EntryRetirement::StartFailed))
1994 }
1995 RetiringWork::UnexpectedWorkerStopped => {
1996 Ok(self.entry_retired(key, generation, EntryRetirement::UnexpectedWorkerStopped))
1997 }
1998 RetiringWork::InterruptedStart => {
1999 Ok(self.entry_retired(key, generation, EntryRetirement::Stop))
2000 }
2001 RetiringWork::SupervisorShutdown(service) => {
2002 drop(service);
2003 Ok(self.entry_retired(key, generation, EntryRetirement::Shutdown))
2004 }
2005 RetiringWork::CancelledWorkerReturned(change) => Ok(self.operation_cancelled(
2006 key,
2007 generation,
2008 operation,
2009 change,
2010 CancellationOutcome::WorkerReturned,
2011 )),
2012 RetiringWork::Transferred {
2013 purpose: ProxyInputPurpose::Cancellation(change),
2014 input,
2015 } => {
2016 let outcome = match input {
2017 ProxyInputCustody::InputRejected => CancellationOutcome::ProxyInputRejected,
2018 ProxyInputCustody::ProxyReported(outcome) => {
2019 CancellationOutcome::ProxyReported { outcome }
2020 }
2021 input => {
2022 return Err(ProxyRetirement {
2023 proxy,
2024 shutdown: stopped.0,
2025 work: RetiringWork::Transferred {
2026 purpose: ProxyInputPurpose::Cancellation(change),
2027 input,
2028 },
2029 stopped: Some(stopped.1),
2030 });
2031 }
2032 };
2033 Ok(self.operation_cancelled(key, generation, operation, change, outcome))
2034 }
2035 RetiringWork::Transferred {
2036 purpose: ProxyInputPurpose::InterruptedStart(start_operation),
2037 input,
2038 } => {
2039 let shutdown = match stopped.0 {
2040 ProxyShutdown::Accepted(operation) => Ok(operation),
2041 ProxyShutdown::Rejected => Err(()),
2042 ProxyShutdown::AwaitingSettlement(witness) => {
2043 return Err(ProxyRetirement {
2044 proxy,
2045 shutdown: ProxyShutdown::AwaitingSettlement(witness),
2046 work: RetiringWork::Transferred {
2047 purpose: ProxyInputPurpose::InterruptedStart(start_operation),
2048 input,
2049 },
2050 stopped: Some(stopped.1),
2051 });
2052 }
2053 };
2054 let worker = match input.into_interrupted_worker() {
2055 Ok(worker) => worker,
2056 Err(input) => {
2057 return Err(ProxyRetirement {
2058 proxy,
2059 shutdown: match shutdown {
2060 Ok(operation) => ProxyShutdown::Accepted(operation),
2061 Err(()) => ProxyShutdown::Rejected,
2062 },
2063 work: RetiringWork::Transferred {
2064 purpose: ProxyInputPurpose::InterruptedStart(start_operation),
2065 input,
2066 },
2067 stopped: Some(stopped.1),
2068 });
2069 }
2070 };
2071 let interrupted =
2072 self.lifecycle
2073 .clone()
2074 .deliver(DynamicLifecycle::WorkerChangeInterrupted {
2075 key: key.clone(),
2076 generation,
2077 operation: start_operation.get(),
2078 interruption: WorkerChangeInterruption::ExplicitStop { operation },
2079 worker,
2080 });
2081 let retired = self.entry_retired(key.clone(), generation, EntryRetirement::Stop);
2082 match shutdown {
2083 Ok(shutdown_operation) => {
2084 drop(shutdown_operation);
2085 Ok(interrupted
2086 .combine(self.lifecycle.clone().deliver(
2087 DynamicLifecycle::StopFinished {
2088 key,
2089 generation,
2090 operation,
2091 proxy,
2092 result: Ok(stopped.1),
2093 },
2094 ))
2095 .combine(retired))
2096 }
2097 Err(()) => {
2098 drop(proxy);
2099 Ok(interrupted.combine(retired))
2100 }
2101 }
2102 }
2103 RetiringWork::Transferred { purpose, input } => {
2104 let (worker_operation, change) = match &purpose {
2105 ProxyInputPurpose::SupervisorStart(operation) => {
2106 (*operation, WorkerChange::Start)
2107 }
2108 ProxyInputPurpose::SupervisorReplacement { operation, service } => {
2109 let _ = service;
2110 (*operation, WorkerChange::Replacement)
2111 }
2112 ProxyInputPurpose::Cancellation(_) | ProxyInputPurpose::InterruptedStart(_) => {
2113 return Err(ProxyRetirement {
2114 proxy,
2115 shutdown: stopped.0,
2116 work: RetiringWork::Transferred { purpose, input },
2117 stopped: Some(stopped.1),
2118 });
2119 }
2120 };
2121 let worker = match input.into_interrupted_worker() {
2122 Ok(worker) => worker,
2123 Err(input) => {
2124 return Err(ProxyRetirement {
2125 proxy,
2126 shutdown: stopped.0,
2127 work: RetiringWork::Transferred { purpose, input },
2128 stopped: Some(stopped.1),
2129 });
2130 }
2131 };
2132 drop((proxy, stopped));
2133 Ok(self
2134 .lifecycle
2135 .clone()
2136 .deliver(DynamicLifecycle::WorkerChangeInterrupted {
2137 key: key.clone(),
2138 generation,
2139 operation: worker_operation.get(),
2140 interruption: WorkerChangeInterruption::SupervisorShutdown { change },
2141 worker,
2142 })
2143 .combine(self.entry_retired(key, generation, EntryRetirement::Shutdown)))
2144 }
2145 }
2146 }
2147
2148 fn accept_proxy_stop(
2149 &mut self,
2150 stopped: crate::ChildStopped<BehaviorAddr<Worker>>,
2151 ) -> BehaviorActed<Self> {
2152 let Some((key, entry)) = self.take_proxy_entry(stopped.child) else {
2153 let mut requests = DynamicSupervisorRequests::empty();
2154 requests.diagnostics = InterpreterRequests::one(
2155 self.diagnostics
2156 .action(DynamicDiagnostic::RejectedProxyStop { stopped }),
2157 );
2158 return Ok(Actions::send(requests));
2159 };
2160 let DynamicEntry {
2161 creation,
2162 generation,
2163 operation,
2164 phase,
2165 } = entry;
2166 match phase {
2167 DynamicEntryPhase::Stopping {
2168 proxy,
2169 restorable,
2170 shutdown: ProxyShutdown::AwaitingSettlement(witness),
2171 stopped: None,
2172 } => {
2173 self.services.insert(
2174 key,
2175 DynamicEntry {
2176 creation,
2177 generation,
2178 operation,
2179 phase: DynamicEntryPhase::Stopping {
2180 proxy,
2181 restorable,
2182 shutdown: ProxyShutdown::AwaitingSettlement(witness),
2183 stopped: Some(stopped),
2184 },
2185 },
2186 );
2187 Ok(Actions::cont())
2188 }
2189 DynamicEntryPhase::Stopping {
2190 proxy,
2191 restorable: _,
2192 shutdown: ProxyShutdown::Accepted(shutdown),
2193 stopped: None,
2194 } => {
2195 drop(shutdown);
2196 let mut requests = DynamicSupervisorRequests::empty();
2197 requests.lifecycle = self
2198 .lifecycle
2199 .clone()
2200 .deliver(DynamicLifecycle::StopFinished {
2201 key: key.clone(),
2202 generation: generation.get(),
2203 operation: operation.get(),
2204 proxy,
2205 result: Ok(stopped),
2206 })
2207 .combine(self.entry_retired(key, generation.get(), EntryRetirement::Stop));
2208 Ok(Actions::send(requests))
2209 }
2210 DynamicEntryPhase::Retiring(ProxyRetirement {
2211 proxy,
2212 shutdown,
2213 work,
2214 stopped: None,
2215 }) => {
2216 let mut requests = DynamicSupervisorRequests::empty();
2217 let retirement = ProxyRetirement {
2218 proxy,
2219 shutdown,
2220 work,
2221 stopped: Some(stopped),
2222 };
2223 requests.lifecycle =
2224 self.settle_retirement(key, creation, generation, operation, retirement);
2225 Ok(Actions::send(requests))
2226 }
2227 phase => {
2228 self.services.insert(
2229 key,
2230 DynamicEntry {
2231 creation,
2232 generation,
2233 operation,
2234 phase,
2235 },
2236 );
2237 let mut requests = DynamicSupervisorRequests::empty();
2238 requests.diagnostics = InterpreterRequests::one(
2239 self.diagnostics
2240 .action(DynamicDiagnostic::RejectedProxyStop { stopped }),
2241 );
2242 Ok(Actions::send(requests))
2243 }
2244 }
2245 }
2246
2247 fn replacement_rejected(
2248 &mut self,
2249 key: Key,
2250 entry: DynamicEntry<Worker, Plan>,
2251 report: ChildReport<ProxyOutcome<Worker, Plan>>,
2252 ) -> BehaviorActed<Self> {
2253 self.services.insert(key, entry);
2254 let mut requests = DynamicSupervisorRequests::empty();
2255 requests.diagnostics = InterpreterRequests::one(
2256 self.diagnostics
2257 .action(DynamicDiagnostic::RejectedProxyOutcome { report }),
2258 );
2259 Ok(Actions::send(requests))
2260 }
2261
2262 fn accept_proxy_report(
2263 &mut self,
2264 report: ChildReport<ProxyOutcome<Worker, Plan>>,
2265 ) -> BehaviorActed<Self> {
2266 let child = report.child;
2267 let Some((key, entry)) = self.take_proxy_entry(child) else {
2268 let mut requests = DynamicSupervisorRequests::empty();
2269 requests.diagnostics = InterpreterRequests::one(
2270 self.diagnostics
2271 .action(DynamicDiagnostic::RejectedProxyOutcome { report }),
2272 );
2273 return Ok(Actions::send(requests));
2274 };
2275 let DynamicEntry {
2276 creation,
2277 generation,
2278 operation,
2279 phase,
2280 } = entry;
2281 match (phase, report.report) {
2282 (
2283 DynamicEntryPhase::WaitingForProxyOutcome {
2284 proxy,
2285 operation: proxy_operation,
2286 },
2287 ProxyOutcome::Initial {
2288 outcome:
2289 InitialWorkerOutcome::Resolved {
2290 result: WorkerStartResult::Ready { attempt, readiness },
2291 },
2292 },
2293 ) => {
2294 drop(proxy_operation);
2295 self.services.insert(
2296 key.clone(),
2297 DynamicEntry {
2298 creation,
2299 generation,
2300 operation,
2301 phase: DynamicEntryPhase::Available {
2302 proxy: proxy.clone(),
2303 service: ServiceAvailability::Ready {
2304 worker: attempt,
2305 readiness,
2306 },
2307 latest: WorkerChangeDisposition::Committed,
2308 },
2309 },
2310 );
2311 let waiting = self.authorize_waiting();
2312 let mut requests = DynamicSupervisorRequests::empty();
2313 requests.proxy_operations = waiting;
2314 requests.lifecycle = self.lifecycle.clone().deliver(DynamicLifecycle::Started {
2315 key,
2316 generation: generation.get(),
2317 operation: operation.get(),
2318 proxy,
2319 });
2320 Ok(Actions::send(requests))
2321 }
2322 (
2323 DynamicEntryPhase::WaitingForProxyOutcome {
2324 proxy,
2325 operation: proxy_operation,
2326 },
2327 ProxyOutcome::Initial { outcome },
2328 ) => {
2329 drop(proxy_operation);
2330 let (retirement, shutdown_request) =
2331 ProxyRetirement::begin(creation, proxy, RetiringWork::StartFailed);
2332 self.services.insert(
2333 key.clone(),
2334 DynamicEntry {
2335 creation,
2336 generation,
2337 operation,
2338 phase: DynamicEntryPhase::Retiring(retirement),
2339 },
2340 );
2341 let mut proxy_operations = SourceActions::empty();
2342 proxy_operations.send(shutdown_request);
2343 let waiting = self.authorize_waiting();
2344 let mut requests = DynamicSupervisorRequests::empty();
2345 requests.proxy_operations = proxy_operations;
2346 requests.proxy_operations.append(waiting);
2347 requests.lifecycle =
2348 self.lifecycle
2349 .clone()
2350 .deliver(DynamicLifecycle::StartOutcomeRejected {
2351 key,
2352 generation: generation.get(),
2353 operation: operation.get(),
2354 outcome,
2355 });
2356 Ok(Actions::send(requests))
2357 }
2358 (
2359 DynamicEntryPhase::ReplacementAwaitingOutcome {
2360 proxy,
2361 service,
2362 operation: proxy_operation,
2363 },
2364 ProxyOutcome::Replacement { outcome },
2365 ) => {
2366 let current_worker = match &service {
2367 ServiceAvailability::Ready { worker, .. }
2368 | ServiceAvailability::Empty { previous: worker } => worker,
2369 };
2370 match &outcome {
2371 ReplacementOutcome::WorkerAttemptsExhausted { replaces, .. }
2372 | ReplacementOutcome::CancelledBeforeBirth { replaces, .. }
2373 | ReplacementOutcome::Resolved { replaces, .. }
2374 if replaces != current_worker =>
2375 {
2376 return self.replacement_rejected(
2377 key,
2378 DynamicEntry {
2379 creation,
2380 generation,
2381 operation,
2382 phase: DynamicEntryPhase::ReplacementAwaitingOutcome {
2383 proxy,
2384 service,
2385 operation: proxy_operation,
2386 },
2387 },
2388 ChildReport::new(child, ProxyOutcome::Replacement { outcome }),
2389 );
2390 }
2391 _ => {}
2392 }
2393
2394 let replacement = match outcome {
2395 ReplacementOutcome::Resolved {
2396 replaces: _,
2397 result: WorkerStartResult::Ready { attempt, readiness },
2398 } => Ok(ServiceAvailability::Ready {
2399 worker: attempt,
2400 readiness,
2401 }),
2402 ReplacementOutcome::NotReplaceable {
2403 worker,
2404 activation,
2405 phase,
2406 } => Err((
2407 service,
2408 ReplacementFailure::ProxyRefused {
2409 worker,
2410 activation,
2411 phase,
2412 },
2413 )),
2414 ReplacementOutcome::WorkerAttemptsExhausted {
2415 replaces: _,
2416 worker,
2417 activation,
2418 } => Err((
2419 service,
2420 ReplacementFailure::WorkerAttemptsExhausted { worker, activation },
2421 )),
2422 ReplacementOutcome::Resolved {
2423 replaces,
2424 result:
2425 WorkerStartResult::CreationRejected {
2426 rejection,
2427 activation,
2428 stopped,
2429 },
2430 } => Err((
2431 ServiceAvailability::Empty { previous: replaces },
2432 ReplacementFailure::WorkerCreationRejected {
2433 rejection,
2434 activation,
2435 stopped,
2436 },
2437 )),
2438 ReplacementOutcome::Resolved {
2439 replaces: _,
2440 result: WorkerStartResult::Unavailable { attempt, drain },
2441 } => Err((
2442 ServiceAvailability::Empty { previous: attempt },
2443 ReplacementFailure::WorkerUnavailable { drain },
2444 )),
2445 outcome @ ReplacementOutcome::CancelledBeforeBirth { .. } => {
2446 return self.replacement_rejected(
2447 key,
2448 DynamicEntry {
2449 creation,
2450 generation,
2451 operation,
2452 phase: DynamicEntryPhase::ReplacementAwaitingOutcome {
2453 proxy,
2454 service,
2455 operation: proxy_operation,
2456 },
2457 },
2458 ChildReport::new(child, ProxyOutcome::Replacement { outcome }),
2459 );
2460 }
2461 };
2462 let (service, lifecycle) = match replacement {
2463 Ok(service) => (
2464 service,
2465 DynamicLifecycle::Replaced {
2466 key: key.clone(),
2467 generation: generation.get(),
2468 operation: operation.get(),
2469 proxy: proxy.clone(),
2470 },
2471 ),
2472 Err((service, failure)) => (
2473 service,
2474 DynamicLifecycle::ReplacementFailed {
2475 key: key.clone(),
2476 generation: generation.get(),
2477 operation: operation.get(),
2478 failure,
2479 },
2480 ),
2481 };
2482 drop(proxy_operation);
2483 self.services.insert(
2484 key,
2485 DynamicEntry {
2486 creation,
2487 generation,
2488 operation,
2489 phase: DynamicEntryPhase::Available {
2490 proxy,
2491 service,
2492 latest: WorkerChangeDisposition::Committed,
2493 },
2494 },
2495 );
2496 let mut requests = DynamicSupervisorRequests::empty();
2497 requests.proxy_operations = self.authorize_waiting();
2498 requests.lifecycle = self.lifecycle.clone().deliver(lifecycle);
2499 Ok(Actions::send(requests))
2500 }
2501 (
2502 DynamicEntryPhase::Retiring(ProxyRetirement {
2503 proxy,
2504 shutdown,
2505 work:
2506 RetiringWork::Transferred {
2507 purpose,
2508 input: ProxyInputCustody::AwaitingOutcome(proxy_operation),
2509 },
2510 stopped,
2511 }),
2512 outcome @ (ProxyOutcome::Initial { .. } | ProxyOutcome::Replacement { .. }),
2513 ) => {
2514 match (&purpose, &outcome) {
2515 (
2516 ProxyInputPurpose::Cancellation(WorkerChange::Start)
2517 | ProxyInputPurpose::InterruptedStart(_)
2518 | ProxyInputPurpose::SupervisorStart(_),
2519 ProxyOutcome::Initial { .. },
2520 )
2521 | (
2522 ProxyInputPurpose::Cancellation(WorkerChange::Replacement)
2523 | ProxyInputPurpose::SupervisorReplacement { .. },
2524 ProxyOutcome::Replacement { .. },
2525 ) => {}
2526 _ => {
2527 return self.replacement_rejected(
2528 key,
2529 DynamicEntry {
2530 creation,
2531 generation,
2532 operation,
2533 phase: DynamicEntryPhase::Retiring(ProxyRetirement {
2534 proxy,
2535 shutdown,
2536 work: RetiringWork::Transferred {
2537 purpose,
2538 input: ProxyInputCustody::AwaitingOutcome(proxy_operation),
2539 },
2540 stopped,
2541 }),
2542 },
2543 ChildReport::new(child, outcome),
2544 );
2545 }
2546 }
2547 drop(proxy_operation);
2548 let retirement = ProxyRetirement {
2549 proxy,
2550 shutdown,
2551 work: RetiringWork::Transferred {
2552 purpose,
2553 input: ProxyInputCustody::ProxyReported(outcome),
2554 },
2555 stopped,
2556 };
2557 let mut requests = DynamicSupervisorRequests::empty();
2558 requests.lifecycle =
2559 self.settle_retirement(key, creation, generation, operation, retirement);
2560 requests.proxy_operations = self.authorize_waiting();
2561 Ok(Actions::send(requests))
2562 }
2563 (
2564 phase,
2565 ProxyOutcome::Unavailable {
2566 sender,
2567 phase: proxy_phase,
2568 command,
2569 },
2570 ) => {
2571 self.services.insert(
2572 key.clone(),
2573 DynamicEntry {
2574 creation,
2575 generation,
2576 operation,
2577 phase,
2578 },
2579 );
2580 let mut requests = DynamicSupervisorRequests::empty();
2581 requests.lifecycle =
2582 self.lifecycle
2583 .clone()
2584 .deliver(DynamicLifecycle::CommandUnavailable {
2585 key,
2586 generation: generation.get(),
2587 sender,
2588 proxy_phase,
2589 command,
2590 });
2591 Ok(Actions::send(requests))
2592 }
2593 (
2594 DynamicEntryPhase::Available {
2595 proxy,
2596 service: ServiceAvailability::Ready { worker, readiness },
2597 latest,
2598 },
2599 ProxyOutcome::WorkerStopped {
2600 worker: stopped_worker,
2601 stopped,
2602 },
2603 ) if worker == stopped_worker => {
2604 let disposition = self.unexpected_exit;
2605 let mut proxy_operations = SourceActions::empty();
2606 let phase = match disposition {
2607 UnexpectedExit::KeepEmpty => DynamicEntryPhase::Available {
2608 proxy,
2609 service: ServiceAvailability::Empty { previous: worker },
2610 latest,
2611 },
2612 UnexpectedExit::Retire => {
2613 let (retirement, shutdown_request) = ProxyRetirement::begin(
2614 creation,
2615 proxy,
2616 RetiringWork::UnexpectedWorkerStopped,
2617 );
2618 proxy_operations.send(shutdown_request);
2619 DynamicEntryPhase::Retiring(retirement)
2620 }
2621 };
2622 self.services.insert(
2623 key.clone(),
2624 DynamicEntry {
2625 creation,
2626 generation,
2627 operation,
2628 phase,
2629 },
2630 );
2631 let mut requests = DynamicSupervisorRequests::empty();
2632 requests.proxy_operations = proxy_operations;
2633 requests.lifecycle =
2634 self.lifecycle
2635 .clone()
2636 .deliver(DynamicLifecycle::UnexpectedWorkerStopped {
2637 key,
2638 generation: generation.get(),
2639 worker: stopped_worker,
2640 readiness,
2641 stopped,
2642 disposition,
2643 });
2644 Ok(Actions::send(requests))
2645 }
2646 (phase, outcome) => {
2647 self.services.insert(
2648 key,
2649 DynamicEntry {
2650 creation,
2651 generation,
2652 operation,
2653 phase,
2654 },
2655 );
2656 let mut requests = DynamicSupervisorRequests::empty();
2657 requests.diagnostics = InterpreterRequests::one(self.diagnostics.action(
2658 DynamicDiagnostic::RejectedProxyOutcome {
2659 report: ChildReport::new(child, outcome),
2660 },
2661 ));
2662 Ok(Actions::send(requests))
2663 }
2664 }
2665 }
2666}
2667
2668fn entry_stop_result<Worker, Plan>(
2669 stopped: Option<ChildStopped<BehaviorAddr<Worker>>>,
2670 input: ProxyInputResult<behavior::Here, Worker, Plan>,
2671) -> Result<
2672 crate::ProxyInputReceipt<Worker, Plan>,
2673 (
2674 EntryStopFailure<BehaviorAddr<Worker>>,
2675 ProxyInputResult<behavior::Here, Worker, Plan>,
2676 ),
2677>
2678where
2679 Worker: Behavior,
2680 Plan: ActivationPlan,
2681 BehaviorAddr<Worker>: EndpointAddress,
2682 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
2683{
2684 match input {
2685 SettledItem::Attempted(ItemSettlement::Accepted(receipt)) => Ok(receipt),
2686 SettledItem::Attempted(ItemSettlement::Rejected { item, reason }) => Err((
2687 EntryStopFailure {
2688 reason: EntryStopFailureReason::ControlRejected(reason),
2689 stopped,
2690 },
2691 SettledItem::Attempted(ItemSettlement::Rejected { item, reason }),
2692 )),
2693 SettledItem::Attempted(ItemSettlement::Corrupt { item, fault }) => Err((
2694 EntryStopFailure {
2695 reason: EntryStopFailureReason::InterpreterCorrupt(fault),
2696 stopped,
2697 },
2698 SettledItem::Attempted(ItemSettlement::Corrupt { item, fault }),
2699 )),
2700 SettledItem::Unattempted(item) => Err((
2701 EntryStopFailure {
2702 reason: EntryStopFailureReason::InterpretationSkipped,
2703 stopped,
2704 },
2705 SettledItem::Unattempted(item),
2706 )),
2707 SettledItem::Attempted(ItemSettlement::Blocked { prerequisite, .. }) => {
2708 match prerequisite {}
2709 }
2710 }
2711}
2712
2713fn proxy_input_creation<Worker, Plan>(
2714 input: &ProxyInputResult<behavior::Here, Worker, Plan>,
2715) -> CreationId
2716where
2717 Worker: Behavior,
2718 Plan: ActivationPlan,
2719 BehaviorAddr<Worker>: EndpointAddress,
2720 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
2721{
2722 match input {
2723 SettledItem::Attempted(ItemSettlement::Accepted(receipt)) => receipt.creation(),
2724 SettledItem::Attempted(
2725 ItemSettlement::Rejected {
2726 item: operation, ..
2727 }
2728 | ItemSettlement::Corrupt {
2729 item: operation, ..
2730 },
2731 )
2732 | SettledItem::Unattempted(operation) => operation.creation(),
2733 SettledItem::Attempted(ItemSettlement::Blocked { prerequisite, .. }) => {
2734 match *prerequisite {}
2735 }
2736 }
2737}
2738
2739impl<Key, Worker, Plan, LifecycleRoute, DiagnosticRouteType> Behavior
2740 for DynamicSupervisor<Key, Worker, Plan, LifecycleRoute, DiagnosticRouteType>
2741where
2742 Key: Clone + Ord + Send + Sync,
2743 Worker: Behavior + Send,
2744 Plan: ActivationPlan,
2745 BehaviorAddr<Worker>: EndpointAddress,
2746 <BehaviorAddr<Worker> as Address>::Nonce: Send,
2747 StableProxy<Worker, Plan>: Behavior<Protocol = Worker::Protocol>,
2748 EstablishedActor<StableProxy<Worker, Plan>>: Clone + Send,
2749 LifecycleRoute: DeliveryRoute<
2750 Protocol = MessageProtocol<BehaviorAddr<Worker>, DynamicLifecycle<Key, Worker, Plan>>,
2751 > + Clone,
2752 LifecycleRoute::Sends: behavior::SendsFor<DynamicSupervisorEvent<Key, Worker, Plan>>,
2753 DiagnosticRouteType: crate::DiagnosticRoute<DynamicDiagnostic<Key, Worker, Plan>> + Clone,
2754{
2755 type Protocol = MessageProtocol<BehaviorAddr<Worker>, DynamicCommand<Key, Worker, Plan>>;
2756 type Event = DynamicSupervisorEvent<Key, Worker, Plan>;
2757 type Sends = DynamicSupervisorRequests<
2758 InterpreterRequests<ObserveChild<Worker::Protocol, ChildHead>>,
2759 SourceActions<ProxyOperation<behavior::Here, Worker, Plan>>,
2760 SourceActions<ScheduleAfter>,
2761 <ReplyRoute<
2762 MessageProtocol<
2763 BehaviorAddr<Worker>,
2764 Result<
2765 WorkerChangeReceipt<Key>,
2766 WorkerChangeRejection<Key, Worker, Plan, StartRejection>,
2767 >,
2768 >,
2769 > as DeliveryRoute>::Sends,
2770 <ReplyRoute<
2771 MessageProtocol<
2772 BehaviorAddr<Worker>,
2773 Result<
2774 WorkerChangeReceipt<Key>,
2775 WorkerChangeRejection<Key, Worker, Plan, ReplaceRejection>,
2776 >,
2777 >,
2778 > as DeliveryRoute>::Sends,
2779 <ReplyRoute<
2780 MessageProtocol<BehaviorAddr<Worker>, Result<Key, StopRejection<Key>>>,
2781 > as DeliveryRoute>::Sends,
2782 <ReplyRoute<
2783 MessageProtocol<BehaviorAddr<Worker>, QueryReply<Key, Worker::Protocol>>,
2784 > as DeliveryRoute>::Sends,
2785 <ReplyRoute<
2786 MessageProtocol<BehaviorAddr<Worker>, CancellationReceipt<Key, Worker, Plan>>,
2787 > as DeliveryRoute>::Sends,
2788 LifecycleRoute::Sends,
2789 InterpreterRequests<
2790 DiagnosticAction<DiagnosticRouteType, DynamicDiagnostic<Key, Worker, Plan>>,
2791 >,
2792 >;
2793 type Ph = Never;
2794 type Error = DynamicSupervisorEvent<Key, Worker, Plan>;
2795 type Birth = Births<StableProxy<Worker, Plan>>;
2796
2797 fn init(&mut self, _: InitializationTurn) -> BehaviorActed<Self> {
2798 Ok(Actions::cont())
2799 }
2800
2801 fn transition(&mut self, _: ActiveTurn, input: Self::Event) -> BehaviorActed<Self> {
2802 let acted = match input {
2803 DynamicSupervisorEvent::Command(User { message, .. }) => self.command(message),
2804 DynamicSupervisorEvent::ProxyCreationsSettled(proxies) => {
2805 self.accept_creations(proxies)
2806 }
2807 DynamicSupervisorEvent::ProxyInputSettled(input) => self.accept_proxy_input(input),
2808 DynamicSupervisorEvent::ProxyReported(report) => self.accept_proxy_report(report),
2809 DynamicSupervisorEvent::ProxyDiagnosed(report) => {
2810 let mut requests = DynamicSupervisorRequests::empty();
2811 requests.diagnostics = InterpreterRequests::one(
2812 self.diagnostics
2813 .action(DynamicDiagnostic::ProxyDiagnosticReported { report }),
2814 );
2815 Ok(Actions::send(requests))
2816 }
2817 DynamicSupervisorEvent::ProxyStopped(stopped) => self.accept_proxy_stop(stopped),
2818 DynamicSupervisorEvent::Shutdown(_) => self.begin_shutdown(),
2819 DynamicSupervisorEvent::ShutdownScheduleSettled(input) => {
2820 match &mut self.availability {
2821 SupervisorAvailability::ShuttingDown(deadline) => deadline
2822 .accept_schedule(input)
2823 .map(|()| Actions::cont())
2824 .map_err(DynamicSupervisorEvent::ShutdownScheduleSettled),
2825 SupervisorAvailability::Accepting(_) => {
2826 Err(DynamicSupervisorEvent::ShutdownScheduleSettled(input))
2827 }
2828 }
2829 }
2830 DynamicSupervisorEvent::ShutdownElapsed(elapsed) => match &mut self.availability {
2831 SupervisorAvailability::ShuttingDown(deadline) => deadline
2832 .accept_elapsed(elapsed)
2833 .map(|()| Actions::cont())
2834 .map_err(DynamicSupervisorEvent::ShutdownElapsed),
2835 SupervisorAvailability::Accepting(_) => {
2836 Err(DynamicSupervisorEvent::ShutdownElapsed(elapsed))
2837 }
2838 },
2839 };
2840 self.apply_shutdown_step(acted)
2841 }
2842}