1use core::ops::ControlFlow;
4
5use behavior::{
6 Address, Behavior, BehaviorAddr, ChildCreationOutcome, CreateChild, CreationId, CreationKind,
7 CreationSequence, CreationSettlement, CreationsSettled, EndpointAddress, EstablishedActor,
8 Here, Ingress, InjectEvent, ItemSettlement, SettledItem,
9};
10
11use crate::{
12 ChildStopped, EstablishedShutdownResolved, ShutdownEstablished, ShutdownId, ShutdownRequested,
13 StopOnShutdown,
14};
15
16use crate::atomic::worker::{
17 WorkerCreationOutcome, WorkerCreationSettlement, settle_worker_creation,
18};
19use crate::atomic::{
20 ActivationPlan, InitializationAttempt, WorkerAttempt, WorkerCreationRejection,
21 WorkerInitializationReport,
22};
23
24use super::ProxyDrain;
25
26#[cfg(test)]
27mod direct_worker_admission_contract {
28 use behavior::{Behavior, BehaviorAddr, EndpointAddress};
29
30 use crate::ChildStopped;
31
32 use super::{CurrentWorker, PendingWorker, WorkerCreationSettlement};
33 use crate::atomic::{ActivationPlan, WorkerInitializationReport};
34
35 #[expect(dead_code, reason = "compile contract for exact worker input")]
36 fn current_stop<W>(
37 worker: &CurrentWorker<W>,
38 input: ChildStopped<BehaviorAddr<W>>,
39 ) -> Result<ChildStopped<BehaviorAddr<W>>, ChildStopped<BehaviorAddr<W>>>
40 where
41 W: Behavior,
42 BehaviorAddr<W>: EndpointAddress,
43 {
44 worker.admit_stop(input)
45 }
46
47 #[expect(dead_code, reason = "compile contract for exact worker input")]
48 fn current_initialization<W, P>(
49 worker: &CurrentWorker<W>,
50 input: WorkerInitializationReport<W, P>,
51 ) -> Result<WorkerInitializationReport<W, P>, WorkerInitializationReport<W, P>>
52 where
53 W: Behavior,
54 P: ActivationPlan,
55 BehaviorAddr<W>: EndpointAddress,
56 {
57 worker.admit_initialization(input)
58 }
59
60 #[expect(dead_code, reason = "compile contract for exact worker input")]
61 fn pending_stop<W, P>(
62 worker: &PendingWorker<P>,
63 input: ChildStopped<BehaviorAddr<W>>,
64 ) -> Result<ChildStopped<BehaviorAddr<W>>, ChildStopped<BehaviorAddr<W>>>
65 where
66 W: Behavior,
67 BehaviorAddr<W>: EndpointAddress,
68 {
69 worker.admit_stop(input)
70 }
71
72 #[expect(dead_code, reason = "compile contract for exact worker input")]
73 fn pending_creation<W, P>(
74 worker: &PendingWorker<P>,
75 input: WorkerCreationSettlement<W>,
76 ) -> Result<WorkerCreationSettlement<W>, WorkerCreationSettlement<W>>
77 where
78 W: Behavior,
79 BehaviorAddr<W>: EndpointAddress,
80 {
81 worker.admit_creation(input)
82 }
83}
84
85pub enum WorkerStartResult<W, P>
87where
88 W: Behavior,
89 P: ActivationPlan,
90 BehaviorAddr<W>: EndpointAddress,
91 StopOnShutdown<W>: Behavior<Protocol = W::Protocol, Ph = W::Ph, Birth = W::Birth>,
92{
93 CreationRejected {
95 rejection: WorkerCreationRejection<W>,
96 activation: P,
97 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
98 },
99 Ready {
101 attempt: WorkerAttempt,
102 readiness: P::Ready,
103 },
104 Unavailable {
106 attempt: WorkerAttempt,
107 drain: ProxyDrain<W, P>,
108 },
109}
110
111pub(in super::super) struct PendingWorker<P> {
112 attempt: WorkerAttempt,
113 initialization: InitializationAttempt,
114 kind: CreationKind,
115 activation: P,
116}
117
118impl<P> PendingWorker<P> {
119 pub(in super::super) fn birth<W>(
120 creations: &mut CreationSequence,
121 worker: W,
122 activation: P,
123 ) -> Result<
124 (
125 Self,
126 CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
127 ),
128 (W, P),
129 >
130 where
131 W: Behavior,
132 {
133 let Some(creation) = creations.issue() else {
134 return Err((worker, activation));
135 };
136 let request = CreateChild::birth(creation, StopOnShutdown::new(worker));
137 Ok(Self::await_creation(activation, request))
138 }
139
140 pub(in super::super) fn replacement<W>(
141 creations: &mut CreationSequence,
142 previous: CreationId,
143 worker: W,
144 activation: P,
145 ) -> Result<
146 (
147 Self,
148 CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
149 ),
150 (W, P),
151 >
152 where
153 W: Behavior,
154 {
155 let Some(creation) = creations.issue() else {
156 return Err((worker, activation));
157 };
158 let request = CreateChild::replacement(creation, previous, StopOnShutdown::new(worker));
159 Ok(Self::await_creation(activation, request))
160 }
161
162 fn await_creation<W>(
163 activation: P,
164 request: CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
165 ) -> (
166 Self,
167 CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
168 )
169 where
170 W: Behavior,
171 {
172 let creation = request.id();
173 let kind = request.kind();
174 let attempt = WorkerAttempt::issued(creation);
175 let initialization = InitializationAttempt::issued(&attempt);
176 (
177 Self {
178 attempt,
179 initialization,
180 kind,
181 activation,
182 },
183 request,
184 )
185 }
186
187 pub(in super::super) const fn creation(&self) -> CreationId {
188 self.attempt.creation()
189 }
190
191 pub(in super::super) fn cancel<W>(
192 self,
193 creation: CreateChild<BehaviorAddr<StopOnShutdown<W>>, StopOnShutdown<W>>,
194 ) -> (WorkerAttempt, W, P)
195 where
196 W: Behavior,
197 {
198 let (_, worker, _) = creation.into_parts();
199 (self.attempt, worker.into_inner(), self.activation)
200 }
201
202 pub(in super::super) fn admit_stop<A>(
203 &self,
204 stopped: ChildStopped<A>,
205 ) -> Result<ChildStopped<A>, ChildStopped<A>>
206 where
207 A: Address,
208 {
209 if stopped.child == self.attempt.creation() {
210 Ok(stopped)
211 } else {
212 Err(stopped)
213 }
214 }
215
216 fn admit_creation<W>(
217 &self,
218 settlement: WorkerCreationSettlement<W>,
219 ) -> Result<WorkerCreationSettlement<W>, WorkerCreationSettlement<W>>
220 where
221 W: Behavior,
222 BehaviorAddr<W>: EndpointAddress,
223 {
224 let received = match &settlement {
225 SettledItem::Unattempted(creation) => (creation.id(), creation.kind()),
226 SettledItem::Attempted(ItemSettlement::Rejected { item, .. })
227 | SettledItem::Attempted(ItemSettlement::Blocked { item, .. })
228 | SettledItem::Attempted(ItemSettlement::Corrupt { item, .. }) => {
229 (item.id(), item.kind())
230 }
231 SettledItem::Attempted(ItemSettlement::Accepted(creation)) => match creation {
232 ChildCreationOutcome::Established(child) => (child.id(), child.kind()),
233 ChildCreationOutcome::InitializationRejected { creation, .. }
234 | ChildCreationOutcome::InitializationPanicked { creation }
235 | ChildCreationOutcome::HostRejected { creation, .. } => {
236 (creation.id(), creation.kind())
237 }
238 },
239 };
240 if received == (self.attempt.creation(), self.kind) {
241 Ok(settlement)
242 } else {
243 Err(settlement)
244 }
245 }
246}
247
248enum WorkerShutdown<P>
249where
250 P: behavior::Protocol,
251{
252 Requested(ShutdownId),
253 Settled(EstablishedShutdownResolved<P>),
254}
255
256pub(in super::super) struct WorkerStopping<W>
257where
258 W: Behavior,
259 BehaviorAddr<W>: EndpointAddress,
260{
261 worker: CurrentWorker<W>,
262 shutdown: WorkerShutdown<W::Protocol>,
263 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
264}
265
266pub(in super::super) struct StoppedWorker<W>
267where
268 W: Behavior,
269 BehaviorAddr<W>: EndpointAddress,
270{
271 pub(in super::super) worker: CurrentWorker<W>,
272 pub(in super::super) shutdown: EstablishedShutdownResolved<W::Protocol>,
273 pub(in super::super) stopped: ChildStopped<BehaviorAddr<W>>,
274}
275
276impl<W> WorkerStopping<W>
277where
278 W: Behavior,
279 BehaviorAddr<W>: EndpointAddress,
280{
281 pub(in super::super) const fn worker(&self) -> &CurrentWorker<W> {
282 &self.worker
283 }
284
285 pub(in super::super) const fn stopped(&self) -> Option<&ChildStopped<BehaviorAddr<W>>> {
286 self.stopped.as_ref()
287 }
288
289 pub(in super::super) fn begin(
290 worker: CurrentWorker<W>,
291 ) -> (Self, ShutdownEstablished<StopOnShutdown<W>, Here>)
292 where
293 StopOnShutdown<W>: Behavior<Protocol = W::Protocol>,
294 <StopOnShutdown<W> as Behavior>::Event: InjectEvent<ShutdownRequested, Here>,
295 {
296 let shutdown = ShutdownEstablished::new(
297 ShutdownId(worker.attempt.creation().get()),
298 worker.actor.clone(),
299 Ingress::new(),
300 );
301 (
302 Self {
303 worker,
304 shutdown: WorkerShutdown::Requested(shutdown.id),
305 stopped: None,
306 },
307 shutdown,
308 )
309 }
310
311 pub(in super::super) fn shutdown_resolved(
312 self,
313 input: EstablishedShutdownResolved<W::Protocol>,
314 ) -> Result<ControlFlow<StoppedWorker<W>, Self>, (Self, EstablishedShutdownResolved<W::Protocol>)>
315 {
316 let Self {
317 worker,
318 shutdown,
319 stopped,
320 } = self;
321 match shutdown {
322 WorkerShutdown::Requested(expected) if input.id() == expected => Ok(Self {
323 worker,
324 shutdown: WorkerShutdown::Settled(input),
325 stopped,
326 }
327 .complete()),
328 shutdown => Err((
329 Self {
330 worker,
331 shutdown,
332 stopped,
333 },
334 input,
335 )),
336 }
337 }
338
339 pub(in super::super) fn worker_stopped(
340 self,
341 input: ChildStopped<BehaviorAddr<W>>,
342 ) -> Result<ControlFlow<StoppedWorker<W>, Self>, (Self, ChildStopped<BehaviorAddr<W>>)> {
343 let Self {
344 worker,
345 shutdown,
346 stopped,
347 } = self;
348 match stopped {
349 None => match worker.admit_stop(input) {
350 Ok(stopped) => Ok(Self {
351 worker,
352 shutdown,
353 stopped: Some(stopped),
354 }
355 .complete()),
356 Err(input) => Err((
357 Self {
358 worker,
359 shutdown,
360 stopped: None,
361 },
362 input,
363 )),
364 },
365 Some(stopped) => Err((
366 Self {
367 worker,
368 shutdown,
369 stopped: Some(stopped),
370 },
371 input,
372 )),
373 }
374 }
375
376 pub(in super::super) fn into_stopped_before_shutdown(
377 self,
378 ) -> Result<(CurrentWorker<W>, ShutdownId, ChildStopped<BehaviorAddr<W>>), Self> {
379 match (self.shutdown, self.stopped) {
380 (WorkerShutdown::Requested(shutdown), Some(stopped)) => {
381 Ok((self.worker, shutdown, stopped))
382 }
383 (shutdown, stopped) => Err(Self {
384 worker: self.worker,
385 shutdown,
386 stopped,
387 }),
388 }
389 }
390
391 fn complete(self) -> ControlFlow<StoppedWorker<W>, Self> {
392 match (self.shutdown, self.stopped) {
393 (WorkerShutdown::Settled(shutdown), Some(stopped)) => {
394 ControlFlow::Break(StoppedWorker {
395 worker: self.worker,
396 shutdown,
397 stopped,
398 })
399 }
400 (shutdown, stopped) => ControlFlow::Continue(Self {
401 worker: self.worker,
402 shutdown,
403 stopped,
404 }),
405 }
406 }
407}
408
409pub(in super::super) struct CurrentWorker<W>
410where
411 W: Behavior,
412 BehaviorAddr<W>: EndpointAddress,
413{
414 pub(in super::super) attempt: WorkerAttempt,
415 pub(in super::super) initialization: InitializationAttempt,
416 pub(in super::super) actor: EstablishedActor<StopOnShutdown<W>>,
417}
418
419impl<W> CurrentWorker<W>
420where
421 W: Behavior,
422 BehaviorAddr<W>: EndpointAddress,
423{
424 pub(in super::super) fn admit_stop(
425 &self,
426 stopped: ChildStopped<BehaviorAddr<W>>,
427 ) -> Result<ChildStopped<BehaviorAddr<W>>, ChildStopped<BehaviorAddr<W>>> {
428 if stopped.child == self.attempt.creation() {
429 Ok(stopped)
430 } else {
431 Err(stopped)
432 }
433 }
434
435 pub(in super::super) fn admit_initialization<P>(
436 &self,
437 input: WorkerInitializationReport<W, P>,
438 ) -> Result<WorkerInitializationReport<W, P>, WorkerInitializationReport<W, P>>
439 where
440 P: ActivationPlan,
441 {
442 match (input.worker(), input.initialization()) {
443 (worker, initialization)
444 if worker == &self.attempt && initialization == &self.initialization =>
445 {
446 Ok(input)
447 }
448 _ => Err(input),
449 }
450 }
451}
452
453pub(in super::super) enum WorkerCreation<W, P>
454where
455 W: Behavior,
456 BehaviorAddr<W>: EndpointAddress,
457{
458 Initializing {
459 worker: CurrentWorker<W>,
460 activation: P,
461 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
462 },
463 Rejected {
464 rejection: WorkerCreationRejection<W>,
465 activation: P,
466 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
467 },
468 Unexpected {
469 worker: PendingWorker<P>,
470 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
471 workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
472 },
473}
474
475impl<P> PendingWorker<P> {
476 pub(in super::super) fn created<W>(
477 self,
478 workers: CreationsSettled<BehaviorAddr<W>, StopOnShutdown<W>>,
479 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
480 ) -> WorkerCreation<W, P>
481 where
482 W: Behavior,
483 BehaviorAddr<W>: EndpointAddress,
484 StopOnShutdown<W>:
485 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
486 {
487 match workers.into_settlement() {
488 CreationSettlement::Rejected { creations, reason } => {
489 let creation = match creations.into_one() {
490 Ok(creation) => creation,
491 Err(creations) => {
492 return WorkerCreation::Unexpected {
493 worker: self,
494 stopped,
495 workers: CreationsSettled::new(CreationSettlement::Rejected {
496 creations,
497 reason,
498 }),
499 };
500 }
501 };
502 if (creation.id(), creation.kind()) != (self.creation(), self.kind) {
503 return WorkerCreation::Unexpected {
504 worker: self,
505 stopped,
506 workers: CreationsSettled::new(CreationSettlement::Rejected {
507 creations: behavior::Creations::one(creation),
508 reason,
509 }),
510 };
511 }
512 let (_, worker, _) = creation.into_parts();
513 WorkerCreation::Rejected {
514 rejection: WorkerCreationRejection::NamespaceExhausted {
515 worker: worker.into_inner(),
516 },
517 activation: self.activation,
518 stopped,
519 }
520 }
521 CreationSettlement::Settled(settlements) => {
522 let settlement = match settlements.into_one() {
523 Ok(settlement) => settlement,
524 Err(settlements) => {
525 return WorkerCreation::Unexpected {
526 worker: self,
527 stopped,
528 workers: CreationsSettled::new(CreationSettlement::Settled(
529 settlements,
530 )),
531 };
532 }
533 };
534 self.created_one(settlement, stopped)
535 }
536 settlement @ CreationSettlement::Corrupt { .. } => WorkerCreation::Unexpected {
537 worker: self,
538 stopped,
539 workers: CreationsSettled::new(settlement),
540 },
541 }
542 }
543
544 fn created_one<W>(
545 self,
546 worker: WorkerCreationSettlement<W>,
547 stopped: Option<ChildStopped<BehaviorAddr<W>>>,
548 ) -> WorkerCreation<W, P>
549 where
550 W: Behavior,
551 BehaviorAddr<W>: EndpointAddress,
552 StopOnShutdown<W>:
553 Behavior<Protocol = W::Protocol, Error = W::Error, Ph = W::Ph, Birth = W::Birth>,
554 {
555 match self.admit_creation(worker) {
556 Ok(settlement) => match settle_worker_creation(settlement) {
557 WorkerCreationOutcome::Established(actor) => WorkerCreation::Initializing {
558 worker: CurrentWorker {
559 attempt: self.attempt,
560 initialization: self.initialization,
561 actor,
562 },
563 activation: self.activation,
564 stopped,
565 },
566 WorkerCreationOutcome::Rejected(rejection) => WorkerCreation::Rejected {
567 rejection,
568 activation: self.activation,
569 stopped,
570 },
571 },
572 Err(settlement) => WorkerCreation::Unexpected {
573 worker: self,
574 stopped,
575 workers: CreationsSettled::new(CreationSettlement::Settled(
576 behavior::Creations::one(settlement),
577 )),
578 },
579 }
580 }
581}