Skip to main content

behavior_actors/atomic/worker/
preparation.rs

1//! Typed worker preparation outside the pure supervisor transition.
2
3use core::marker::PhantomData;
4use core::ops::ControlFlow;
5use std::sync::Arc;
6
7use behavior::{
8    ActionItem, ActionItemResult, Behavior, InterpretationProgress, ItemSettlement, Never,
9    SettledItem, SourceAction, finish_item, prepare_item,
10};
11
12use crate::atomic::RoleName;
13
14use super::{ActivationPlan, PreparedWorker, WorkerSubmission};
15
16pub(in super::super) struct PreparationTicket {
17    token: Arc<()>,
18}
19
20impl PreparationTicket {
21    fn reserve() -> (Self, Self) {
22        let token = Arc::new(());
23        (
24            Self {
25                token: Arc::clone(&token),
26            },
27            Self { token },
28        )
29    }
30
31    pub(in super::super) fn matches(&self, other: &Self) -> bool {
32        Arc::ptr_eq(&self.token, &other.token)
33    }
34}
35
36/// The one preparation correlation retained by an actor across a source task.
37///
38/// An issued request can accept its start settlement. Only after that exact
39/// receipt may the actor accept the late source return. Keeping the phase
40/// beside the ticket denies duplicate starts and early completion without an
41/// arrival-history flag.
42pub(in super::super) enum WorkerPreparationExpectation {
43    Issued(PreparationTicket),
44    Started(PreparationTicket),
45}
46
47impl WorkerPreparationExpectation {
48    pub(in super::super) fn issued(ticket: PreparationTicket) -> Self {
49        Self::Issued(ticket)
50    }
51
52    pub(in super::super) fn accept_start(
53        self,
54        receipt: &WorkerPreparationStarted,
55    ) -> Result<Self, Self> {
56        match self {
57            Self::Issued(ticket) if receipt.accepts(&ticket) => Ok(Self::Started(ticket)),
58            expectation => Err(expectation),
59        }
60    }
61
62    pub(in super::super) fn accepts_issued<Source, Role, Worker, Plan>(
63        &self,
64        input: &ActionItemResult<PrepareWorkers<Source, Role, Worker, Plan>>,
65    ) -> bool
66    where
67        Source: WorkerSource<Role, Worker, Plan>,
68        Role: Send + Sync,
69        Worker: Behavior + Send,
70        Plan: ActivationPlan,
71    {
72        match self {
73            Self::Issued(ticket) => preparation_start_accepts(input, ticket),
74            Self::Started(_) => false,
75        }
76    }
77
78    pub(in super::super) fn accepts_return<Source, Role, Worker, Plan>(
79        &self,
80        returned: &WorkerPreparation<Source, Role, Worker, Plan>,
81    ) -> bool
82    where
83        Source: WorkerSource<Role, Worker, Plan>,
84        Worker: Behavior + Send,
85        Plan: ActivationPlan,
86    {
87        match self {
88            Self::Issued(_) => false,
89            Self::Started(ticket) => returned.accepts(ticket),
90        }
91    }
92}
93
94/// Static declaration implemented by one concrete Bombay worker-source
95/// capability.
96///
97/// This trait declares types only. It deliberately has no method: the source is
98/// an affine capability interpreted by Bombay, not an application callback
99/// invoked by `FixedSupervisor`.
100///
101/// A source declared for another worker cannot prepare this supervisor:
102///
103/// ```compile_fail,E0277
104/// struct DeclaredWorker;
105/// #[behavior::behavior(addr = behavior::MailAddr, message = behavior::Never)]
106/// impl DeclaredWorker {
107///     fn receive(&mut self, _: behavior::MailAddr, message: behavior::Never)
108///         -> behavior::BehaviorActed<Self> {
109///         match message {}
110///     }
111/// }
112/// struct OtherWorker;
113/// #[behavior::behavior(addr = behavior::MailAddr, message = behavior::Never)]
114/// impl OtherWorker {
115///     fn receive(&mut self, _: behavior::MailAddr, message: behavior::Never)
116///         -> behavior::BehaviorActed<Self> {
117///         match message {}
118///     }
119/// }
120/// struct Source;
121/// impl behavior_actors::atomic::WorkerSource<(), DeclaredWorker,
122///     behavior_actors::atomic::ImmediateActivation> for Source
123/// {
124///     type WorkerRejection = behavior::Never;
125///     type SourceRejection = behavior::Never;
126/// }
127/// fn requires_declared<Source>()
128/// where
129///     Source: behavior_actors::atomic::WorkerSource<
130///         (), OtherWorker, behavior_actors::atomic::ImmediateActivation,
131///     >,
132/// {}
133/// requires_declared::<Source>();
134/// ```
135pub trait WorkerSource<Role, Worker, Plan>: Send
136where
137    Worker: Behavior + Send,
138    Plan: ActivationPlan,
139{
140    /// Exact rejection while preparing one selected role.
141    type WorkerRejection: Send;
142    /// Exact rejection before the first worker submission is prepared.
143    type SourceRejection: Send;
144}
145
146impl<Role, Worker, Plan> WorkerSource<Role, Worker, Plan> for Never
147where
148    Worker: Behavior + Send,
149    Plan: ActivationPlan,
150{
151    type WorkerRejection = Never;
152    type SourceRejection = Never;
153}
154
155/// One non-empty ordered request for replacement worker submissions.
156///
157/// The request cannot report a prepared worker before its exact start is
158/// committed:
159///
160/// ```compile_fail,E0599
161/// fn premature_submission<Source, Role, Worker, Plan>(
162///     request: behavior_actors::atomic::PrepareWorkers<Source, Role, Worker, Plan>,
163///     submission: behavior_actors::atomic::WorkerSubmission<Worker, Plan>,
164/// )
165/// where
166///     Source: behavior_actors::atomic::WorkerSource<Role, Worker, Plan>,
167///     Worker: behavior::Behavior + Send,
168///     Plan: behavior_actors::atomic::ActivationPlan,
169/// {
170///     let _ = request.accept(submission);
171/// }
172/// ```
173#[must_use = "worker preparation must settle or transfer outward"]
174pub struct PrepareWorkers<Source, Role, Worker, Plan>
175where
176    Source: WorkerSource<Role, Worker, Plan>,
177    Worker: Behavior + Send,
178    Plan: ActivationPlan,
179{
180    ticket: PreparationTicket,
181    source: Source,
182    first: RoleName<Role>,
183    remaining: Vec<RoleName<Role>>,
184    worker: PhantomData<fn() -> Worker>,
185    plan: PhantomData<fn() -> Plan>,
186}
187
188/// Receipt that the runtime committed one exact worker-preparation start.
189///
190/// This proves neither a completed source nor permission to create a worker.
191/// Its private ticket names the issued preparation without exposing a runtime
192/// address or allowing the actor to forge a receipt for another request.
193#[must_use = "worker preparation start must return to its source"]
194pub struct WorkerPreparationStarted {
195    ticket: PreparationTicket,
196}
197
198impl WorkerPreparationStarted {
199    pub(in super::super) fn accepts(&self, expected: &PreparationTicket) -> bool {
200        self.ticket.matches(expected)
201    }
202}
203
204/// Source work after one accepted start and before its first submission.
205///
206/// Only this phase can return a source rejection. Once it accepts a first
207/// submission, the returned pending request has a nonempty prepared prefix
208/// and can report only worker-specific rejection for later roles.
209#[must_use = "started worker preparation must settle or transfer outward"]
210pub struct StartingWorkerPreparation<Source, Role, Worker, Plan>
211where
212    Source: WorkerSource<Role, Worker, Plan>,
213    Worker: Behavior + Send,
214    Plan: ActivationPlan,
215{
216    ticket: PreparationTicket,
217    source: Source,
218    first: RoleName<Role>,
219    remaining: Vec<RoleName<Role>>,
220    worker: PhantomData<fn() -> Worker>,
221    plan: PhantomData<fn() -> Plan>,
222}
223
224impl<Source, Role, Worker, Plan> StartingWorkerPreparation<Source, Role, Worker, Plan>
225where
226    Source: WorkerSource<Role, Worker, Plan>,
227    Worker: Behavior + Send,
228    Plan: ActivationPlan,
229{
230    /// Borrow the source and first selected application role.
231    #[must_use]
232    pub fn source_and_role(&mut self) -> (&mut Source, &Role) {
233        (&mut self.source, self.first.role())
234    }
235
236    /// Accept the first submission and advance or complete the ordered group.
237    #[must_use]
238    pub fn accept(
239        self,
240        submission: WorkerSubmission<Worker, Plan>,
241    ) -> ControlFlow<
242        WorkerPreparation<Source, Role, Worker, Plan>,
243        PendingWorkerPreparation<Source, Role, Worker, Plan>,
244    > {
245        advance_preparation(
246            self.ticket,
247            self.source,
248            Vec::new(),
249            self.first,
250            self.remaining,
251            submission,
252        )
253    }
254
255    /// Return an exact worker rejection for the first selected role.
256    #[must_use]
257    pub fn reject(
258        self,
259        reason: Source::WorkerRejection,
260    ) -> WorkerPreparation<Source, Role, Worker, Plan> {
261        rejected_preparation(
262            self.ticket,
263            self.source,
264            Vec::new(),
265            self.first,
266            reason,
267            self.remaining,
268        )
269    }
270
271    /// Return the source and its exact rejection before any worker submission.
272    #[must_use]
273    pub fn reject_source(
274        self,
275        reason: Source::SourceRejection,
276    ) -> WorkerPreparation<Source, Role, Worker, Plan> {
277        WorkerPreparation {
278            ticket: self.ticket,
279            outcome: WorkerPreparationOutcome::SourceRejected {
280                source: self.source,
281                failed_role: self.first,
282                reason,
283                remaining: self.remaining,
284            },
285        }
286    }
287}
288
289impl<Source, Role, Worker, Plan> PrepareWorkers<Source, Role, Worker, Plan>
290where
291    Source: WorkerSource<Role, Worker, Plan>,
292    Worker: Behavior + Send,
293    Plan: ActivationPlan,
294{
295    pub(in super::super) fn new(
296        source: Source,
297        first: RoleName<Role>,
298        remaining: Vec<RoleName<Role>>,
299    ) -> (PreparationTicket, Self) {
300        let (expected, ticket) = PreparationTicket::reserve();
301        (
302            expected,
303            Self {
304                ticket,
305                source,
306                first,
307                remaining,
308                worker: PhantomData,
309                plan: PhantomData,
310            },
311        )
312    }
313
314    pub(in super::super) fn accepts(&self, expected: &PreparationTicket) -> bool {
315        expected.matches(&self.ticket)
316    }
317
318    /// Borrow the source and first selected role before committing the start.
319    #[must_use]
320    pub fn source_and_role(&mut self) -> (&mut Source, &Role) {
321        (&mut self.source, self.first.role())
322    }
323
324    pub(in super::super) fn into_parts(self) -> (Source, RoleName<Role>, Vec<RoleName<Role>>) {
325        (self.source, self.first, self.remaining)
326    }
327
328    /// Commit one start and transfer the affine source to its first attempt.
329    ///
330    /// The issued request is consumed, so another start receipt cannot be
331    /// produced from the same source action.
332    #[must_use]
333    pub fn start(
334        self,
335    ) -> (
336        WorkerPreparationStarted,
337        StartingWorkerPreparation<Source, Role, Worker, Plan>,
338    ) {
339        let started = WorkerPreparationStarted {
340            ticket: PreparationTicket {
341                token: Arc::clone(&self.ticket.token),
342            },
343        };
344        let starting = StartingWorkerPreparation {
345            ticket: self.ticket,
346            source: self.source,
347            first: self.first,
348            remaining: self.remaining,
349            worker: PhantomData,
350            plan: PhantomData,
351        };
352        (started, starting)
353    }
354}
355
356pub(in super::super) fn preparation_start_accepts<Source, Role, Worker, Plan>(
357    result: &ActionItemResult<PrepareWorkers<Source, Role, Worker, Plan>>,
358    expected: &PreparationTicket,
359) -> bool
360where
361    Source: WorkerSource<Role, Worker, Plan>,
362    Role: Send + Sync,
363    Worker: Behavior + Send,
364    Plan: ActivationPlan,
365{
366    match result {
367        SettledItem::Attempted(ItemSettlement::Accepted(started)) => started.accepts(expected),
368        SettledItem::Attempted(ItemSettlement::Rejected { item, .. })
369        | SettledItem::Attempted(ItemSettlement::Blocked { item, .. })
370        | SettledItem::Attempted(ItemSettlement::Corrupt { item, .. })
371        | SettledItem::Unattempted(item) => item.accepts(expected),
372    }
373}
374
375pub(in super::super) enum WorkerPreparationOutcome<
376    Source,
377    Role,
378    Worker,
379    Plan,
380    WorkerRejection,
381    SourceRejection,
382> {
383    Prepared {
384        source: Source,
385        members: Vec<PreparedWorker<RoleName<Role>, Worker, Plan>>,
386    },
387    WorkerRejected {
388        source: Source,
389        prepared: Vec<PreparedWorker<RoleName<Role>, Worker, Plan>>,
390        failed_role: RoleName<Role>,
391        reason: WorkerRejection,
392        remaining: Vec<RoleName<Role>>,
393    },
394    SourceRejected {
395        source: Source,
396        failed_role: RoleName<Role>,
397        reason: SourceRejection,
398        remaining: Vec<RoleName<Role>>,
399    },
400}
401
402/// Non-empty remainder of one exact worker-preparation request.
403///
404/// Once the first worker submission has been accepted, a later failure is a
405/// worker rejection; the source cannot be reclassified as rejected:
406///
407/// ```compile_fail,E0599
408/// fn late_source_rejection<Source, Role, Worker, Plan>(
409///     pending: behavior_actors::atomic::PendingWorkerPreparation<Source, Role, Worker, Plan>,
410///     reason: Source::SourceRejection,
411/// )
412/// where
413///     Source: behavior_actors::atomic::WorkerSource<Role, Worker, Plan>,
414///     Worker: behavior::Behavior + Send,
415///     Plan: behavior_actors::atomic::ActivationPlan,
416/// {
417///     let _ = pending.reject_source(reason);
418/// }
419/// ```
420#[must_use = "worker preparation must advance, reject, or transfer outward"]
421pub struct PendingWorkerPreparation<Source, Role, Worker, Plan>
422where
423    Source: WorkerSource<Role, Worker, Plan>,
424    Worker: Behavior + Send,
425    Plan: ActivationPlan,
426{
427    ticket: PreparationTicket,
428    source: Source,
429    prepared: Vec<PreparedWorker<RoleName<Role>, Worker, Plan>>,
430    current: RoleName<Role>,
431    remaining: Vec<RoleName<Role>>,
432}
433
434impl<Source, Role, Worker, Plan> PendingWorkerPreparation<Source, Role, Worker, Plan>
435where
436    Source: WorkerSource<Role, Worker, Plan>,
437    Worker: Behavior + Send,
438    Plan: ActivationPlan,
439{
440    /// Borrow the source and current application role for one Bombay-owned
441    /// preparation attempt.
442    #[must_use]
443    pub fn source_and_role(&mut self) -> (&mut Source, &Role) {
444        (&mut self.source, self.current.role())
445    }
446
447    /// Borrow each accepted role, original worker, and activation plan in order.
448    ///
449    /// The complete pending request, including its original preparation ticket,
450    /// remains owned by the caller. No submission is cloned or transferred.
451    #[must_use]
452    pub fn prepared_workers(&self) -> impl Iterator<Item = (&Role, &Worker, &Plan)> {
453        self.prepared.iter().map(|prepared| {
454            (
455                prepared.role.role(),
456                &prepared.submission.worker,
457                &prepared.submission.activation,
458            )
459        })
460    }
461
462    /// Borrow the untouched roles after the current role in their original order.
463    ///
464    /// The existing `source_and_role` method supplies the source and current role.
465    #[must_use]
466    pub fn remaining_roles(&self) -> impl Iterator<Item = &Role> {
467        self.remaining.iter().map(RoleName::role)
468    }
469
470    /// Accept one submission for the current role.
471    #[must_use]
472    pub fn accept(
473        self,
474        submission: WorkerSubmission<Worker, Plan>,
475    ) -> ControlFlow<
476        WorkerPreparation<Source, Role, Worker, Plan>,
477        PendingWorkerPreparation<Source, Role, Worker, Plan>,
478    > {
479        advance_preparation(
480            self.ticket,
481            self.source,
482            self.prepared,
483            self.current,
484            self.remaining,
485            submission,
486        )
487    }
488
489    /// Return an exact rejection for the current role.
490    #[must_use]
491    pub fn reject(
492        self,
493        reason: Source::WorkerRejection,
494    ) -> WorkerPreparation<Source, Role, Worker, Plan> {
495        rejected_preparation(
496            self.ticket,
497            self.source,
498            self.prepared,
499            self.current,
500            reason,
501            self.remaining,
502        )
503    }
504}
505
506/// Complete result of source work after one accepted preparation start.
507#[must_use = "worker preparation must return to its supervisor or retire outward"]
508pub struct WorkerPreparation<Source, Role, Worker, Plan>
509where
510    Source: WorkerSource<Role, Worker, Plan>,
511    Worker: Behavior + Send,
512    Plan: ActivationPlan,
513{
514    ticket: PreparationTicket,
515    outcome: WorkerPreparationOutcome<
516        Source,
517        Role,
518        Worker,
519        Plan,
520        Source::WorkerRejection,
521        Source::SourceRejection,
522    >,
523}
524
525impl<Source, Role, Worker, Plan> WorkerPreparation<Source, Role, Worker, Plan>
526where
527    Source: WorkerSource<Role, Worker, Plan>,
528    Worker: Behavior + Send,
529    Plan: ActivationPlan,
530{
531    pub(in super::super) fn from_parts(
532        ticket: PreparationTicket,
533        outcome: WorkerPreparationOutcome<
534            Source,
535            Role,
536            Worker,
537            Plan,
538            Source::WorkerRejection,
539            Source::SourceRejection,
540        >,
541    ) -> Self {
542        Self { ticket, outcome }
543    }
544
545    pub(in super::super) fn into_parts(
546        self,
547    ) -> (
548        PreparationTicket,
549        WorkerPreparationOutcome<
550            Source,
551            Role,
552            Worker,
553            Plan,
554            Source::WorkerRejection,
555            Source::SourceRejection,
556        >,
557    ) {
558        (self.ticket, self.outcome)
559    }
560
561    pub(in super::super) fn accepts(&self, expected: &PreparationTicket) -> bool {
562        expected.matches(&self.ticket)
563    }
564}
565
566fn advance_preparation<Source, Role, Worker, Plan>(
567    ticket: PreparationTicket,
568    source: Source,
569    mut prepared: Vec<PreparedWorker<RoleName<Role>, Worker, Plan>>,
570    current: RoleName<Role>,
571    remaining: Vec<RoleName<Role>>,
572    submission: WorkerSubmission<Worker, Plan>,
573) -> ControlFlow<
574    WorkerPreparation<Source, Role, Worker, Plan>,
575    PendingWorkerPreparation<Source, Role, Worker, Plan>,
576>
577where
578    Source: WorkerSource<Role, Worker, Plan>,
579    Worker: Behavior + Send,
580    Plan: ActivationPlan,
581{
582    prepared.push(PreparedWorker {
583        role: current,
584        submission,
585    });
586    let mut remaining = remaining.into_iter();
587    match remaining.next() {
588        None => ControlFlow::Break(WorkerPreparation {
589            ticket,
590            outcome: WorkerPreparationOutcome::Prepared {
591                source,
592                members: prepared,
593            },
594        }),
595        Some(current) => ControlFlow::Continue(PendingWorkerPreparation {
596            ticket,
597            source,
598            prepared,
599            current,
600            remaining: remaining.collect(),
601        }),
602    }
603}
604
605fn rejected_preparation<Source, Role, Worker, Plan>(
606    ticket: PreparationTicket,
607    source: Source,
608    prepared: Vec<PreparedWorker<RoleName<Role>, Worker, Plan>>,
609    failed_role: RoleName<Role>,
610    reason: Source::WorkerRejection,
611    remaining: Vec<RoleName<Role>>,
612) -> WorkerPreparation<Source, Role, Worker, Plan>
613where
614    Source: WorkerSource<Role, Worker, Plan>,
615    Worker: Behavior + Send,
616    Plan: ActivationPlan,
617{
618    WorkerPreparation {
619        ticket,
620        outcome: WorkerPreparationOutcome::WorkerRejected {
621            source,
622            prepared,
623            failed_role,
624            reason,
625            remaining,
626        },
627    }
628}
629
630impl<Source, Role, Worker, Plan> ActionItem for PrepareWorkers<Source, Role, Worker, Plan>
631where
632    Source: WorkerSource<Role, Worker, Plan>,
633    Role: Send + Sync,
634    Worker: Behavior + Send,
635    Plan: ActivationPlan,
636{
637    type Custody = (Option<Self>, Option<Self::Reply>);
638    type Input<'a>
639        = &'a mut Option<Self>
640    where
641        Self: 'a;
642    type Reply = ItemSettlement<Self, Self::Accepted, Self::Rejection, Self::Prerequisite>;
643
644    fn prepare_interpretation(
645        progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
646    ) {
647        prepare_item::<Self>(progress);
648    }
649
650    fn interpretation_input<'a>(
651        custody: &'a mut Self::Custody,
652    ) -> Option<(Self::Input<'a>, &'a mut Option<Self::Reply>)>
653    where
654        Self: 'a,
655    {
656        match custody {
657            (input @ Some(_), received @ None) => Some((input, received)),
658            _ => None,
659        }
660    }
661
662    fn finish_interpretation(
663        progress: &mut Option<InterpretationProgress<Self, Self::Custody, Self::Reply>>,
664    ) {
665        finish_item::<Self>(progress);
666    }
667
668    type Accepted = WorkerPreparationStarted;
669    type Rejection = Never;
670    type Prerequisite = Never;
671}
672
673impl<Source, Role, Worker, Plan> SourceAction for PrepareWorkers<Source, Role, Worker, Plan>
674where
675    Source: WorkerSource<Role, Worker, Plan>,
676    Role: Send + Sync,
677    Worker: Behavior + Send,
678    Plan: ActivationPlan,
679{
680    type Source = Self;
681}
682
683#[cfg(test)]
684mod tests {
685    use core::ops::ControlFlow;
686    use core::ptr;
687    use std::sync::Arc;
688
689    use behavior::{
690        ActionItemResult, ActiveTurn, Behavior, BehaviorActed, MailAddr, MessageProtocol, Never,
691        NoBirths, NoSends, SettledItem, User,
692    };
693
694    use super::{PrepareWorkers, WorkerPreparationOutcome, WorkerSource};
695    use crate::ActivationPlan;
696    use crate::WorkerSubmission;
697    use crate::atomic::RoleName;
698
699    #[derive(Debug, Eq, PartialEq)]
700    enum Role {
701        Api,
702        Storage,
703        Search,
704        Queue,
705        Index,
706    }
707
708    #[derive(Debug, Eq, PartialEq)]
709    struct Worker(u8);
710
711    impl Behavior for Worker {
712        type Protocol = MessageProtocol<MailAddr, Never>;
713        type Event = User<MailAddr, Never>;
714        type Sends = NoSends;
715        type Ph = Never;
716        type Error = Never;
717        type Birth = NoBirths;
718
719        fn transition(&mut self, _: ActiveTurn, event: Self::Event) -> BehaviorActed<Self> {
720            match event.message {}
721        }
722    }
723
724    #[derive(Debug, Eq, PartialEq)]
725    struct Plan(u8);
726
727    impl ActivationPlan for Plan {
728        type Ready = ();
729        type Rejection = Never;
730
731        fn activate(
732            self,
733        ) -> impl core::future::Future<Output = Result<Self::Ready, Self::Rejection>> + Send
734        {
735            core::future::ready(Ok(()))
736        }
737    }
738
739    #[derive(Debug, Eq, PartialEq)]
740    struct Source(u8);
741
742    #[derive(Debug, Eq, PartialEq)]
743    enum WorkerRejection {
744        Unsupported,
745    }
746
747    #[derive(Debug, Eq, PartialEq)]
748    enum SourceRejection {
749        Unavailable,
750    }
751
752    impl WorkerSource<Role, Worker, Plan> for Source {
753        type WorkerRejection = WorkerRejection;
754        type SourceRejection = SourceRejection;
755    }
756
757    type Request = PrepareWorkers<Source, Role, Worker, Plan>;
758
759    fn submission(worker: u8) -> WorkerSubmission<Worker, Plan> {
760        WorkerSubmission::activated(Worker(worker), Plan(worker))
761    }
762
763    #[test]
764    fn two_role_request_advances_once_then_completes() {
765        let api = RoleName::new(Role::Api);
766        let storage = RoleName::new(Role::Storage);
767        let api_name = api.clone();
768        let storage_name = storage.clone();
769        let (ticket, request) = Request::new(Source(3), api_name, vec![storage_name]);
770        let (_started, mut starting) = request.start();
771        let (source, role) = starting.source_and_role();
772        assert_eq!(source, &mut Source(3));
773        assert!(core::ptr::eq(role, api.role()));
774
775        let ControlFlow::Continue(mut pending) = starting.accept(submission(11)) else {
776            panic!("one remaining role cannot complete preparation");
777        };
778        let (source, role) = pending.source_and_role();
779        assert_eq!(source, &mut Source(3));
780        assert!(core::ptr::eq(role, storage.role()));
781        let ControlFlow::Break(preparation) = pending.accept(submission(13)) else {
782            panic!("the final role completes preparation");
783        };
784        let (returned_ticket, outcome) = preparation.into_parts();
785        assert!(ticket.matches(&returned_ticket));
786        let WorkerPreparationOutcome::Prepared { source, members } = outcome else {
787            panic!("accepted roles produce the prepared outcome");
788        };
789
790        assert_eq!(source, Source(3));
791        assert_eq!(members.len(), 2);
792        assert!(core::ptr::eq(members[0].role.role(), api.role()));
793        assert!(core::ptr::eq(members[1].role.role(), storage.role()));
794    }
795
796    #[test]
797    fn rejection_derives_failed_name_and_untouched_suffix() {
798        let api = RoleName::new(Role::Api);
799        let storage = RoleName::new(Role::Storage);
800        let later = RoleName::new(Role::Api);
801        let (ticket, request) =
802            Request::new(Source(7), api.clone(), vec![storage.clone(), later.clone()]);
803        let (_started, starting) = request.start();
804        let ControlFlow::Continue(progress) = starting.accept(submission(17)) else {
805            panic!("two names remain after the first submission");
806        };
807        let preparation = progress.reject(WorkerRejection::Unsupported);
808        let (returned_ticket, outcome) = preparation.into_parts();
809        assert!(ticket.matches(&returned_ticket));
810        let WorkerPreparationOutcome::WorkerRejected {
811            source,
812            prepared,
813            failed_role,
814            reason,
815            remaining,
816        } = outcome
817        else {
818            panic!("rejected construction must select the worker rejection outcome");
819        };
820
821        assert_eq!(source, Source(7));
822        assert_eq!(prepared.len(), 1);
823        assert!(core::ptr::eq(prepared[0].role.role(), api.role()));
824        assert!(core::ptr::eq(failed_role.role(), storage.role()));
825        assert_eq!(reason, WorkerRejection::Unsupported);
826        assert_eq!(remaining.len(), 1);
827        assert!(core::ptr::eq(remaining[0].role(), later.role()));
828    }
829
830    #[test]
831    fn source_rejection_after_start_and_no_attempt_return_the_complete_source() {
832        let api = RoleName::new(Role::Api);
833        let storage = RoleName::new(Role::Storage);
834        let (rejected_ticket, rejected_request) =
835            Request::new(Source(19), api.clone(), vec![storage.clone()]);
836        let (_started, mut starting) = rejected_request.start();
837        let (source, role) = starting.source_and_role();
838        assert_eq!(source, &mut Source(19));
839        assert!(core::ptr::eq(role, api.role()));
840        let rejected = starting.reject_source(SourceRejection::Unavailable);
841        assert!(rejected.accepts(&rejected_ticket));
842        let (_, outcome) = rejected.into_parts();
843        let WorkerPreparationOutcome::SourceRejected { source, reason, .. } = outcome else {
844            panic!("late source rejection returns the exact source");
845        };
846        assert_eq!(source, Source(19));
847        assert_eq!(reason, SourceRejection::Unavailable);
848
849        let (unattempted_ticket, unattempted_request) =
850            Request::new(Source(23), storage.clone(), Vec::new());
851        let unattempted: ActionItemResult<Request> = SettledItem::Unattempted(unattempted_request);
852        let SettledItem::Unattempted(item) = unattempted else {
853            panic!("no-attempt must retain the complete request");
854        };
855        assert!(item.accepts(&unattempted_ticket));
856        let (foreign_ticket, _) = Request::new(Source(29), api, Vec::new());
857        assert!(!item.accepts(&foreign_ticket));
858        let (source, role, remaining) = item.into_parts();
859        assert_eq!(source, Source(23));
860        assert!(core::ptr::eq(role.role(), storage.role()));
861        assert!(remaining.is_empty());
862    }
863
864    #[test]
865    fn pending_observation_preserves_original_workers_roles_and_ticket() {
866        for roles in [
867            [
868                Role::Api,
869                Role::Storage,
870                Role::Search,
871                Role::Queue,
872                Role::Index,
873            ],
874            [
875                Role::Storage,
876                Role::Api,
877                Role::Index,
878                Role::Search,
879                Role::Queue,
880            ],
881        ] {
882            let [first, second, current, later, last] = roles;
883            let first = RoleName::new(first);
884            let second = RoleName::new(second);
885            let current = RoleName::new(current);
886            let later = RoleName::new(later);
887            let last = RoleName::new(last);
888            let first_role = first.clone().into_role();
889            let second_role = second.clone().into_role();
890            let current_role = current.clone().into_role();
891            let later_role = later.clone().into_role();
892            let last_role = last.clone().into_role();
893            let (issued, request) =
894                Request::new(Source(31), first, vec![second, current, later, last]);
895            let ticket_allocation = Arc::downgrade(&issued.token);
896            let (started, starting) = request.start();
897            let ControlFlow::Continue(pending) = starting.accept(submission(11)) else {
898                panic!("four remaining roles cannot complete preparation");
899            };
900            let ControlFlow::Continue(mut pending) = pending.accept(submission(13)) else {
901                panic!("three remaining roles cannot complete preparation");
902            };
903            let prepared: Vec<_> = pending
904                .prepared_workers()
905                .map(|(role, worker, plan)| (ptr::from_ref(role), worker.0, plan.0))
906                .collect();
907            let expected_prepared = [
908                (Arc::as_ptr(&first_role), 11, 11),
909                (Arc::as_ptr(&second_role), 13, 13),
910            ];
911            assert_eq!(prepared, expected_prepared);
912            let remaining: Vec<_> = pending.remaining_roles().map(ptr::from_ref).collect();
913            let expected_remaining = [Arc::as_ptr(&later_role), Arc::as_ptr(&last_role)];
914            assert_eq!(remaining, expected_remaining);
915            let source_current = {
916                let (source, role) = pending.source_and_role();
917                (source.0, ptr::from_ref(role))
918            };
919            assert_eq!(source_current, (31, Arc::as_ptr(&current_role)));
920            assert!(started.accepts(&pending.ticket));
921            assert!(issued.matches(&pending.ticket));
922            assert_eq!(ticket_allocation.strong_count(), 3);
923            drop((pending, started, issued));
924            assert_eq!(ticket_allocation.strong_count(), 0);
925            let role_counts = [
926                Arc::strong_count(&first_role),
927                Arc::strong_count(&second_role),
928                Arc::strong_count(&current_role),
929                Arc::strong_count(&later_role),
930                Arc::strong_count(&last_role),
931            ];
932            assert_eq!(role_counts, [1, 1, 1, 1, 1]);
933        }
934    }
935}