1use core::num::NonZeroU64;
4use std::sync::Arc;
5
6use behavior::{
7 ActionItem, EstablishedDelivery, EstablishedRecipient, ExactDeliveryReason, InterpretItem,
8 Interpretation, InterpretationProgress, ItemSettlement, Never, Protocol, RecipientAddress,
9 ReportToParent, SourceAction,
10};
11
12use super::super::WorkerAttempt;
13
14#[derive(Clone, Copy, Debug, Eq, Hash, PartialEq)]
16pub struct SubmissionId(u64);
17
18impl SubmissionId {
19 #[must_use]
21 pub const fn new(value: u64) -> Self {
22 Self(value)
23 }
24
25 #[must_use]
27 pub const fn get(self) -> u64 {
28 self.0
29 }
30}
31
32#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
40pub struct JobId(NonZeroU64);
41
42impl JobId {
43 pub(in crate::atomic) const fn issued(value: NonZeroU64) -> Self {
44 Self(value)
45 }
46
47 #[must_use]
49 pub const fn get(self) -> u64 {
50 self.0.get()
51 }
52}
53
54#[must_use = "an assignment must be completed, returned by delivery settlement, or transferred"]
72pub struct Assignment<Job> {
73 payload: Job,
74 authority: CompletionAuthority,
75}
76
77impl<Job> Assignment<Job> {
78 pub(in crate::atomic) const fn issued(payload: Job, authority: CompletionAuthority) -> Self {
79 Self { payload, authority }
80 }
81
82 #[must_use]
84 pub const fn payload(&self) -> &Job {
85 &self.payload
86 }
87
88 #[must_use]
90 pub fn complete<WorkerResult>(
91 self,
92 worker_result: WorkerResult,
93 ) -> ReportToParent<Completion<WorkerResult>> {
94 ReportToParent::new(Completion {
95 worker_result,
96 authority: self.authority,
97 })
98 }
99
100 pub(in crate::atomic) fn into_parts(self) -> (Job, CompletionAuthority) {
101 (self.payload, self.authority)
102 }
103}
104
105impl<Job> core::fmt::Debug for Assignment<Job>
106where
107 Job: core::fmt::Debug,
108{
109 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
110 formatter
111 .debug_struct("Assignment")
112 .field("payload", &self.payload)
113 .finish_non_exhaustive()
114 }
115}
116
117#[must_use = "a completion must return to its pool or remain in terminal custody"]
125pub struct Completion<WorkerResult> {
126 worker_result: WorkerResult,
127 authority: CompletionAuthority,
128}
129
130impl<WorkerResult> Completion<WorkerResult> {
131 pub(in crate::atomic) const fn authority(&self) -> &CompletionAuthority {
132 &self.authority
133 }
134
135 pub(in crate::atomic) fn into_parts(self) -> (WorkerResult, CompletionAuthority) {
136 (self.worker_result, self.authority)
137 }
138}
139
140impl<WorkerResult> core::fmt::Debug for Completion<WorkerResult>
141where
142 WorkerResult: core::fmt::Debug,
143{
144 fn fmt(&self, formatter: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
145 formatter
146 .debug_struct("Completion")
147 .field("worker_result", &self.worker_result)
148 .finish_non_exhaustive()
149 }
150}
151
152#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)]
153pub(in crate::atomic) struct AdmissionOrdinal(NonZeroU64);
154
155#[derive(Clone, Copy, Debug, Eq, PartialEq)]
156struct AssignmentId(NonZeroU64);
157
158#[derive(Debug, Eq, PartialEq)]
159pub(in crate::atomic) struct AcceptedJobSequence {
160 next: Option<NonZeroU64>,
161}
162
163impl AcceptedJobSequence {
164 pub(in crate::atomic) const fn new() -> Self {
165 Self {
166 next: Some(NonZeroU64::MIN),
167 }
168 }
169
170 pub(in crate::atomic) fn issue(&mut self) -> Option<(JobId, AdmissionOrdinal)> {
171 let value = self.next?;
172 self.next = value.checked_add(1);
173 Some((JobId::issued(value), AdmissionOrdinal(value)))
174 }
175}
176
177#[derive(Debug, Eq, PartialEq)]
178pub(in crate::atomic) struct AssignmentSequence {
179 next: Option<NonZeroU64>,
180}
181
182impl AssignmentSequence {
183 pub(in crate::atomic) const fn new() -> Self {
184 Self {
185 next: Some(NonZeroU64::MIN),
186 }
187 }
188
189 fn issue(&mut self) -> Option<AssignmentId> {
190 let value = self.next?;
191 self.next = value.checked_add(1);
192 Some(AssignmentId(value))
193 }
194
195 pub(in crate::atomic) fn assign<Job>(
196 &mut self,
197 worker: &WorkerAttempt,
198 payload: Job,
199 ) -> Option<(CompletionCorrelation, Assignment<Job>)> {
200 let assignment = self.issue()?;
201 let token = Arc::new(());
202 let correlation = CompletionCorrelation {
203 assignment,
204 worker: worker.clone(),
205 token: Arc::clone(&token),
206 };
207 let authority = CompletionAuthority {
208 assignment,
209 worker: worker.clone(),
210 token,
211 };
212 Some((correlation, Assignment::issued(payload, authority)))
213 }
214}
215
216pub(in crate::atomic) struct CompletionAuthority {
217 assignment: AssignmentId,
218 worker: WorkerAttempt,
219 token: Arc<()>,
220}
221
222#[derive(Clone)]
223pub(in crate::atomic) struct CompletionCorrelation {
224 assignment: AssignmentId,
225 worker: WorkerAttempt,
226 token: Arc<()>,
227}
228
229impl CompletionCorrelation {
230 fn compare(&self, authority: &CompletionAuthority) -> CorrelationMatch {
231 match (
232 self.assignment == authority.assignment,
233 self.worker == authority.worker,
234 Arc::ptr_eq(&self.token, &authority.token),
235 ) {
236 (true, true, true) => CorrelationMatch::Exact,
237 _ => CorrelationMatch::Foreign,
238 }
239 }
240}
241
242pub(in crate::atomic) enum CorrelationMatch {
243 Exact,
244 Foreign,
245}
246
247pub struct AssignmentReceipt {
249 assignment: AssignmentId,
250 worker: WorkerAttempt,
251}
252
253impl AssignmentReceipt {
254 fn issued(correlation: &CompletionCorrelation) -> Self {
255 Self {
256 assignment: correlation.assignment,
257 worker: correlation.worker.clone(),
258 }
259 }
260
261 fn compare(&self, correlation: &CompletionCorrelation) -> CorrelationMatch {
262 match (
263 self.assignment == correlation.assignment,
264 self.worker == correlation.worker,
265 ) {
266 (true, true) => CorrelationMatch::Exact,
267 _ => CorrelationMatch::Foreign,
268 }
269 }
270}
271
272#[must_use = "worker assignment delivery must settle or remain in lifecycle custody"]
274pub struct AssignWorker<P, Job>
275where
276 P: Protocol<Msg = Assignment<Job>>,
277 P::Addr: RecipientAddress,
278{
279 target: EstablishedRecipient<P>,
280 assignment: Assignment<Job>,
281 receipt: AssignmentReceipt,
282}
283
284impl<P, Job> AssignWorker<P, Job>
285where
286 P: Protocol<Msg = Assignment<Job>>,
287 P::Addr: RecipientAddress,
288{
289 pub(in crate::atomic) fn new(
290 target: EstablishedRecipient<P>,
291 correlation: &CompletionCorrelation,
292 assignment: Assignment<Job>,
293 ) -> Self {
294 Self {
295 target,
296 assignment,
297 receipt: AssignmentReceipt::issued(correlation),
298 }
299 }
300
301 #[must_use]
304 pub fn target(&self) -> EstablishedRecipient<P> {
305 self.target.clone()
306 }
307
308 #[must_use]
310 pub(in crate::atomic) fn receipt(&self) -> AssignmentReceipt {
311 AssignmentReceipt {
312 assignment: self.receipt.assignment,
313 worker: self.receipt.worker.clone(),
314 }
315 }
316
317 #[must_use]
319 pub(in crate::atomic) fn into_parts(
320 self,
321 ) -> (EstablishedRecipient<P>, Assignment<Job>, AssignmentReceipt) {
322 (self.target, self.assignment, self.receipt)
323 }
324
325 pub(in crate::atomic) fn returned(
326 target: EstablishedRecipient<P>,
327 assignment: Assignment<Job>,
328 receipt: AssignmentReceipt,
329 ) -> Self {
330 Self {
331 target,
332 assignment,
333 receipt,
334 }
335 }
336
337 pub async fn settle<Host, RootEvent, Path>(
340 progress: &mut Option<
341 InterpretationProgress<
342 Self,
343 <Self as ActionItem>::Custody,
344 ItemSettlement<Self, AssignmentReceipt, ExactDeliveryReason, Never>,
345 >,
346 >,
347 host: &mut Host,
348 ) where
349 Host: InterpretItem<EstablishedDelivery<P>, RootEvent, Path>,
350 <P::Addr as RecipientAddress>::Established<P>: Send,
351 Job: Send,
352 {
353 <Self as ActionItem>::prepare_interpretation(progress);
354 if let Some(InterpretationProgress::Interpreting(custody)) = progress {
355 if let Some((input, received)) = <Self as ActionItem>::interpretation_input(custody) {
356 <Host as InterpretItem<EstablishedDelivery<P>, RootEvent, Path>>::interpret_item(
357 host, input, received,
358 )
359 .await;
360 }
361 }
362 <Self as ActionItem>::finish_interpretation(progress);
363 }
364}
365
366impl<P, Job> ActionItem for AssignWorker<P, Job>
367where
368 P: Protocol<Msg = Assignment<Job>>,
369 P::Addr: RecipientAddress,
370 <P::Addr as RecipientAddress>::Established<P>: Send,
371 Job: Send,
372{
373 type Custody = (
374 AssignmentReceipt,
375 (Option<EstablishedDelivery<P>>, Option<Self::Reply>),
376 );
377 type Input<'a>
378 = &'a mut Option<EstablishedDelivery<P>>
379 where
380 Self: 'a;
381 type Reply = ItemSettlement<EstablishedDelivery<P>, (), ExactDeliveryReason, Never>;
382
383 fn prepare_interpretation(
384 progress: &mut Option<
385 InterpretationProgress<
386 Self,
387 Self::Custody,
388 ItemSettlement<Self, Self::Accepted, Self::Rejection, Self::Prerequisite>,
389 >,
390 >,
391 ) {
392 if !matches!(progress, Some(InterpretationProgress::Original(_))) {
393 return;
394 }
395 match progress.take() {
396 Some(InterpretationProgress::Original(Self {
397 target,
398 assignment,
399 receipt,
400 })) => {
401 *progress = Some(InterpretationProgress::Interpreting((
402 receipt,
403 (Some(EstablishedDelivery::new(target, assignment)), None),
404 )));
405 }
406 retained => *progress = retained,
407 }
408 }
409 fn interpretation_input<'a>(
410 custody: &'a mut Self::Custody,
411 ) -> Option<(Self::Input<'a>, &'a mut Option<Self::Reply>)>
412 where
413 Self: 'a,
414 {
415 let (_, (input, received)) = custody;
416 if input.is_some() && received.is_none() {
417 Some((input, received))
418 } else {
419 None
420 }
421 }
422 fn finish_interpretation(
423 progress: &mut Option<
424 InterpretationProgress<
425 Self,
426 Self::Custody,
427 ItemSettlement<Self, Self::Accepted, Self::Rejection, Self::Prerequisite>,
428 >,
429 >,
430 ) {
431 if !matches!(
432 progress,
433 Some(InterpretationProgress::Interpreting((_, (None, Some(_)))))
434 ) {
435 return;
436 }
437 match progress.take() {
438 Some(InterpretationProgress::Interpreting((receipt, (None, Some(reply))))) => {
439 let received = match reply {
440 ItemSettlement::Accepted(()) => ItemSettlement::Accepted(receipt),
441 ItemSettlement::Rejected {
442 item: EstablishedDelivery { to, message },
443 reason,
444 } => ItemSettlement::Rejected {
445 item: Self {
446 target: to,
447 assignment: message,
448 receipt,
449 },
450 reason,
451 },
452 ItemSettlement::Corrupt {
453 item: EstablishedDelivery { to, message },
454 fault,
455 } => ItemSettlement::Corrupt {
456 item: Self {
457 target: to,
458 assignment: message,
459 receipt,
460 },
461 fault,
462 },
463 ItemSettlement::Blocked { prerequisite, .. } => match prerequisite {},
464 };
465 let interpretation = match received {
466 received @ ItemSettlement::Corrupt { .. } => Interpretation::Corrupt(received),
467 received => Interpretation::Complete(received),
468 };
469 *progress = Some(InterpretationProgress::Completed(interpretation));
470 }
471 retained => *progress = retained,
472 }
473 }
474
475 type Accepted = AssignmentReceipt;
476 type Rejection = ExactDeliveryReason;
477 type Prerequisite = Never;
478}
479
480impl<P, Job> SourceAction for AssignWorker<P, Job>
481where
482 P: Protocol<Msg = Assignment<Job>>,
483 P::Addr: RecipientAddress,
484 <P::Addr as RecipientAddress>::Established<P>: Send,
485 Job: Send,
486{
487 type Source = Self;
488}
489
490pub(in crate::atomic) struct CustomerJob<Job, Customer> {
491 pub(in crate::atomic) id: JobId,
492 pub(in crate::atomic) admitted: AdmissionOrdinal,
493 pub(in crate::atomic) payload: Job,
494 pub(in crate::atomic) customer: Customer,
495}
496
497enum AssignmentDelivery<WorkerResult, A>
498where
499 A: behavior::Address,
500{
501 AwaitingReceipt,
502 Accepted,
503 CompletionHasPriority {
504 completion: Completion<WorkerResult>,
505 later_exit: Option<crate::ChildStopped<A>>,
506 },
507 WorkerExitHasPriority {
508 stopped: crate::ChildStopped<A>,
509 later_completion: Option<Completion<WorkerResult>>,
510 },
511}
512
513pub(in crate::atomic) struct AssignmentShutdown<Job, Customer, WorkerResult, A>
514where
515 A: behavior::Address,
516{
517 pub(in crate::atomic) customer: CustomerJob<Job, Customer>,
518 pub(in crate::atomic) worker: WorkerAttempt,
519 pub(in crate::atomic) stopped: Option<crate::ChildStopped<A>>,
520 pub(in crate::atomic) completion: Option<Completion<WorkerResult>>,
521}
522
523pub(in crate::atomic) struct AssignedJob<Job, Customer, WorkerResult, A>
524where
525 A: behavior::Address,
526{
527 pub(in crate::atomic) customer: CustomerJob<Job, Customer>,
528 correlation: CompletionCorrelation,
529 delivery: AssignmentDelivery<WorkerResult, A>,
530}
531
532pub(in crate::atomic) enum AssignmentReceiptOutcome<Job, Customer, WorkerResult, A>
533where
534 A: behavior::Address,
535{
536 AwaitingCompletion(AssignedJob<Job, Customer, WorkerResult, A>),
537 JobCompleted {
538 customer: CustomerJob<Job, Customer>,
539 result: WorkerResult,
540 stopped: Option<crate::ChildStopped<A>>,
541 },
542 JobInterrupted {
543 customer: CustomerJob<Job, Customer>,
544 stopped: crate::ChildStopped<A>,
545 late_completion: Option<Completion<WorkerResult>>,
546 },
547}
548
549pub(in crate::atomic) enum WorkerCompletionOutcome<Job, Customer, WorkerResult, A>
550where
551 A: behavior::Address,
552{
553 AwaitingReceipt(AssignedJob<Job, Customer, WorkerResult, A>),
554 JobCompleted {
555 customer: CustomerJob<Job, Customer>,
556 result: WorkerResult,
557 stopped: Option<crate::ChildStopped<A>>,
558 },
559}
560
561pub(in crate::atomic) enum AssignmentRejectionOutcome<Job, Customer, WorkerResult, A>
562where
563 A: behavior::Address,
564{
565 JobReturned {
566 customer: CustomerJob<Job, Customer>,
567 stopped: Option<crate::ChildStopped<A>>,
568 },
569 ConflictingCompletion {
570 customer: CustomerJob<Job, Customer>,
571 assignment: Assignment<Job>,
572 completion: Completion<WorkerResult>,
573 stopped: Option<crate::ChildStopped<A>>,
574 },
575}
576
577pub(in crate::atomic) enum WorkerExitOutcome<Job, Customer, WorkerResult, A>
578where
579 A: behavior::Address,
580{
581 AwaitingReceipt(AssignedJob<Job, Customer, WorkerResult, A>),
582 JobInterrupted {
583 customer: CustomerJob<Job, Customer>,
584 stopped: crate::ChildStopped<A>,
585 late_completion: Option<Completion<WorkerResult>>,
586 },
587}
588
589impl<Job, Customer, WorkerResult, A> AssignedJob<Job, Customer, WorkerResult, A>
590where
591 A: behavior::Address,
592{
593 pub(in crate::atomic) fn new(
594 customer: CustomerJob<Job, Customer>,
595 correlation: CompletionCorrelation,
596 ) -> Self {
597 Self {
598 customer,
599 correlation,
600 delivery: AssignmentDelivery::AwaitingReceipt,
601 }
602 }
603
604 pub(in crate::atomic) fn compare_receipt(
605 &self,
606 receipt: &AssignmentReceipt,
607 ) -> CorrelationMatch {
608 receipt.compare(&self.correlation)
609 }
610
611 pub(in crate::atomic) fn compare_completion(
612 &self,
613 completion: &Completion<WorkerResult>,
614 ) -> CorrelationMatch {
615 self.correlation.compare(completion.authority())
616 }
617
618 pub(in crate::atomic) fn shutdown(self) -> AssignmentShutdown<Job, Customer, WorkerResult, A> {
619 let Self {
620 customer,
621 correlation,
622 delivery,
623 } = self;
624 let worker = correlation.worker;
625 match delivery {
626 AssignmentDelivery::AwaitingReceipt | AssignmentDelivery::Accepted => {
627 AssignmentShutdown {
628 customer,
629 worker,
630 stopped: None,
631 completion: None,
632 }
633 }
634 AssignmentDelivery::CompletionHasPriority {
635 completion,
636 later_exit,
637 } => AssignmentShutdown {
638 customer,
639 worker,
640 stopped: later_exit,
641 completion: Some(completion),
642 },
643 AssignmentDelivery::WorkerExitHasPriority {
644 stopped,
645 later_completion,
646 } => AssignmentShutdown {
647 customer,
648 worker,
649 stopped: Some(stopped),
650 completion: later_completion,
651 },
652 }
653 }
654
655 pub(in crate::atomic) fn accept_receipt(
656 self,
657 receipt: AssignmentReceipt,
658 ) -> Result<AssignmentReceiptOutcome<Job, Customer, WorkerResult, A>, (Self, AssignmentReceipt)>
659 {
660 if let CorrelationMatch::Foreign = receipt.compare(&self.correlation) {
661 return Err((self, receipt));
662 }
663 let Self {
664 customer,
665 correlation,
666 delivery,
667 } = self;
668 Ok(match delivery {
669 AssignmentDelivery::AwaitingReceipt => {
670 AssignmentReceiptOutcome::AwaitingCompletion(Self {
671 customer,
672 correlation,
673 delivery: AssignmentDelivery::Accepted,
674 })
675 }
676 AssignmentDelivery::CompletionHasPriority {
677 completion,
678 later_exit,
679 } => {
680 let (result, _) = completion.into_parts();
681 AssignmentReceiptOutcome::JobCompleted {
682 customer,
683 result,
684 stopped: later_exit,
685 }
686 }
687 AssignmentDelivery::WorkerExitHasPriority {
688 stopped,
689 later_completion,
690 } => AssignmentReceiptOutcome::JobInterrupted {
691 customer,
692 stopped,
693 late_completion: later_completion,
694 },
695 AssignmentDelivery::Accepted => {
696 return Err((
697 Self {
698 customer,
699 correlation,
700 delivery: AssignmentDelivery::Accepted,
701 },
702 receipt,
703 ));
704 }
705 })
706 }
707
708 pub(in crate::atomic) fn accept_rejection(
709 self,
710 assignment: Assignment<Job>,
711 ) -> Result<AssignmentRejectionOutcome<Job, Customer, WorkerResult, A>, (Self, Assignment<Job>)>
712 {
713 let (payload, authority) = assignment.into_parts();
714 match self.correlation.compare(&authority) {
715 CorrelationMatch::Foreign => {
716 return Err((self, Assignment::issued(payload, authority)));
717 }
718 CorrelationMatch::Exact => {}
719 }
720 let assignment = Assignment::issued(payload, authority);
721 Ok(match self.delivery {
722 AssignmentDelivery::AwaitingReceipt => AssignmentRejectionOutcome::JobReturned {
723 customer: self.customer,
724 stopped: None,
725 },
726 AssignmentDelivery::CompletionHasPriority {
727 completion,
728 later_exit,
729 } => AssignmentRejectionOutcome::ConflictingCompletion {
730 customer: self.customer,
731 assignment,
732 completion,
733 stopped: later_exit,
734 },
735 AssignmentDelivery::WorkerExitHasPriority {
736 stopped,
737 later_completion: None,
738 } => AssignmentRejectionOutcome::JobReturned {
739 customer: self.customer,
740 stopped: Some(stopped),
741 },
742 AssignmentDelivery::WorkerExitHasPriority {
743 stopped,
744 later_completion: Some(completion),
745 } => AssignmentRejectionOutcome::ConflictingCompletion {
746 customer: self.customer,
747 assignment,
748 completion,
749 stopped: Some(stopped),
750 },
751 AssignmentDelivery::Accepted => {
752 return Err((
753 Self {
754 customer: self.customer,
755 correlation: self.correlation,
756 delivery: AssignmentDelivery::Accepted,
757 },
758 assignment,
759 ));
760 }
761 })
762 }
763
764 pub(in crate::atomic) fn accept_completion(
765 self,
766 completion: Completion<WorkerResult>,
767 ) -> Result<
768 WorkerCompletionOutcome<Job, Customer, WorkerResult, A>,
769 (Self, Completion<WorkerResult>),
770 > {
771 match self.correlation.compare(completion.authority()) {
772 CorrelationMatch::Foreign => return Err((self, completion)),
773 CorrelationMatch::Exact => {}
774 }
775 let Self {
776 customer,
777 correlation,
778 delivery,
779 } = self;
780 Ok(match delivery {
781 AssignmentDelivery::AwaitingReceipt => WorkerCompletionOutcome::AwaitingReceipt(Self {
782 customer,
783 correlation,
784 delivery: AssignmentDelivery::CompletionHasPriority {
785 completion,
786 later_exit: None,
787 },
788 }),
789 AssignmentDelivery::Accepted => {
790 let (result, _) = completion.into_parts();
791 WorkerCompletionOutcome::JobCompleted {
792 customer,
793 result,
794 stopped: None,
795 }
796 }
797 AssignmentDelivery::WorkerExitHasPriority {
798 stopped,
799 later_completion: None,
800 } => WorkerCompletionOutcome::AwaitingReceipt(Self {
801 customer,
802 correlation,
803 delivery: AssignmentDelivery::WorkerExitHasPriority {
804 stopped,
805 later_completion: Some(completion),
806 },
807 }),
808 delivery @ (AssignmentDelivery::CompletionHasPriority { .. }
809 | AssignmentDelivery::WorkerExitHasPriority {
810 later_completion: Some(_),
811 ..
812 }) => {
813 return Err((
814 Self {
815 customer,
816 correlation,
817 delivery,
818 },
819 completion,
820 ));
821 }
822 })
823 }
824
825 pub(in crate::atomic) fn accept_worker_exit(
826 self,
827 stopped: crate::ChildStopped<A>,
828 ) -> Result<WorkerExitOutcome<Job, Customer, WorkerResult, A>, (Self, crate::ChildStopped<A>)>
829 {
830 if stopped.child != self.correlation.worker.creation() {
831 return Err((self, stopped));
832 }
833 let Self {
834 customer,
835 correlation,
836 delivery,
837 } = self;
838 Ok(match delivery {
839 AssignmentDelivery::AwaitingReceipt => WorkerExitOutcome::AwaitingReceipt(Self {
840 customer,
841 correlation,
842 delivery: AssignmentDelivery::WorkerExitHasPriority {
843 stopped,
844 later_completion: None,
845 },
846 }),
847 AssignmentDelivery::Accepted => WorkerExitOutcome::JobInterrupted {
848 customer,
849 stopped,
850 late_completion: None,
851 },
852 AssignmentDelivery::CompletionHasPriority {
853 completion,
854 later_exit: None,
855 } => WorkerExitOutcome::AwaitingReceipt(Self {
856 customer,
857 correlation,
858 delivery: AssignmentDelivery::CompletionHasPriority {
859 completion,
860 later_exit: Some(stopped),
861 },
862 }),
863 delivery @ (AssignmentDelivery::WorkerExitHasPriority { .. }
864 | AssignmentDelivery::CompletionHasPriority {
865 later_exit: Some(_),
866 ..
867 }) => {
868 return Err((
869 Self {
870 customer,
871 correlation,
872 delivery,
873 },
874 stopped,
875 ));
876 }
877 })
878 }
879
880 pub(in crate::atomic) fn observed_stop(&self) -> Option<&crate::ChildStopped<A>> {
881 match &self.delivery {
882 AssignmentDelivery::CompletionHasPriority {
883 later_exit: Some(stopped),
884 ..
885 }
886 | AssignmentDelivery::WorkerExitHasPriority { stopped, .. } => Some(stopped),
887 AssignmentDelivery::AwaitingReceipt
888 | AssignmentDelivery::Accepted
889 | AssignmentDelivery::CompletionHasPriority {
890 later_exit: None, ..
891 } => None,
892 }
893 }
894}
895
896#[cfg(test)]
897mod tests {
898 use core::future::Future;
899 use std::collections::VecDeque;
900 use std::sync::Arc;
901 use std::time::Instant;
902
903 use behavior::{
904 ActionItem, Address, CreationSequence, EstablishedDelivery, EstablishedRecipient,
905 ExactDeliveryReason, Here, InterpretItem, InterpretationProgress, ItemSettlement, MailAddr,
906 MessageProtocol, Protocol, RecipientAddress,
907 };
908
909 use super::{
910 AcceptedJobSequence, AssignWorker, AssignedJob, Assignment, AssignmentReceipt,
911 AssignmentReceiptOutcome, AssignmentRejectionOutcome, AssignmentSequence,
912 CompletionAuthority, CorrelationMatch, CustomerJob, WorkerCompletionOutcome,
913 WorkerExitOutcome,
914 };
915 use crate::atomic::WorkerAttempt;
916 use crate::{ChildStopped, Crash};
917
918 fn worker() -> WorkerAttempt {
919 let mut creations = CreationSequence::new();
920 let creation = creations
921 .issue()
922 .unwrap_or_else(|| panic!("test creation correlation is available"));
923 WorkerAttempt::issued(creation)
924 }
925
926 fn customer() -> CustomerJob<u8, u8> {
927 let mut jobs = AcceptedJobSequence::new();
928 let (id, admitted) = jobs
929 .issue()
930 .unwrap_or_else(|| panic!("test job correlation is available"));
931 CustomerJob {
932 id,
933 admitted,
934 payload: 7,
935 customer: 9,
936 }
937 }
938
939 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
940 struct DeliveryAddr(u64);
941
942 impl Address for DeliveryAddr {
943 type Nonce = u64;
944 }
945
946 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
947 struct DeliveryEndpoint(u64);
948
949 impl RecipientAddress for DeliveryAddr {
950 type Established<P>
951 = DeliveryEndpoint
952 where
953 P: Protocol<Addr = Self>;
954 }
955
956 struct MoveJob(Box<str>);
957
958 type AssignmentProtocol = MessageProtocol<DeliveryAddr, Assignment<MoveJob>>;
959 type ExactAssignment = EstablishedDelivery<AssignmentProtocol>;
960
961 enum DeliveryAdmission {
962 Accept,
963 Reject,
964 }
965
966 struct AssignmentDeliveryHost {
967 decisions: VecDeque<DeliveryAdmission>,
968 observed_payloads: Vec<usize>,
969 }
970
971 impl InterpretItem<ExactAssignment, (), Here> for AssignmentDeliveryHost {
972 fn interpret_item<'a>(
973 &'a mut self,
974 input: &'a mut Option<ExactAssignment>,
975 received: &'a mut Option<<ExactAssignment as ActionItem>::Reply>,
976 ) -> impl Future<Output = ()> + Send + 'a
977 where
978 ExactAssignment: 'a,
979 {
980 async move {
981 if received.is_some() {
982 return;
983 }
984 let Some(delivery) = input.take() else {
985 return;
986 };
987 *received = Some({
988 self.observed_payloads
989 .push(delivery.message.payload().0.as_ptr() as usize);
990 match self
991 .decisions
992 .pop_front()
993 .expect("one decision per delivery")
994 {
995 DeliveryAdmission::Accept => ItemSettlement::Accepted(()),
996 DeliveryAdmission::Reject => ItemSettlement::Rejected {
997 item: delivery,
998 reason: ExactDeliveryReason::ClosedRecipient,
999 },
1000 }
1001 });
1002 }
1003 }
1004 }
1005
1006 #[tokio::test]
1007 async fn exact_assignment_settlement_preserves_original_correlation_after_transfer() {
1008 let worker = worker();
1009 let mut assignments = AssignmentSequence::new();
1010 let (first_correlation, first_assignment) = assignments
1011 .assign(&worker, MoveJob(Box::from("first")))
1012 .expect("first assignment correlation");
1013 let (second_correlation, second_assignment) = assignments
1014 .assign(&worker, MoveJob(Box::from("second")))
1015 .expect("second assignment correlation");
1016 let first_payload = first_assignment.payload().0.as_ptr();
1017 let second_payload = second_assignment.payload().0.as_ptr();
1018 let first_target = EstablishedRecipient::issued(DeliveryEndpoint(41));
1019 let second_target = EstablishedRecipient::issued(DeliveryEndpoint(42));
1020 let first = AssignWorker::<AssignmentProtocol, _>::new(
1021 first_target.clone(),
1022 &first_correlation,
1023 first_assignment,
1024 );
1025 let second = AssignWorker::<AssignmentProtocol, _>::new(
1026 second_target,
1027 &second_correlation,
1028 second_assignment,
1029 );
1030 let mut host = AssignmentDeliveryHost {
1031 decisions: VecDeque::from([DeliveryAdmission::Accept, DeliveryAdmission::Reject]),
1032 observed_payloads: Vec::new(),
1033 };
1034
1035 let ItemSettlement::Accepted(second_receipt) = ({
1036 let mut progress = Some(InterpretationProgress::Original(second));
1037 AssignWorker::<AssignmentProtocol, MoveJob>::settle::<_, (), Here>(
1038 &mut progress,
1039 &mut host,
1040 )
1041 .await;
1042 let Some(InterpretationProgress::Completed(settlement)) = progress else {
1043 panic!("the exact host returns its complete original settlement");
1044 };
1045 settlement.into_settlement()
1046 }) else {
1047 panic!("second assignment admission returns its held receipt");
1048 };
1049 assert!(matches!(
1050 second_receipt.compare(&second_correlation),
1051 CorrelationMatch::Exact
1052 ));
1053 let ItemSettlement::Rejected {
1054 item: first_returned,
1055 reason: ExactDeliveryReason::ClosedRecipient,
1056 } = ({
1057 let mut progress = Some(InterpretationProgress::Original(first));
1058 AssignWorker::<AssignmentProtocol, MoveJob>::settle::<_, (), Here>(
1059 &mut progress,
1060 &mut host,
1061 )
1062 .await;
1063 let Some(InterpretationProgress::Completed(settlement)) = progress else {
1064 panic!("the exact host returns its complete original settlement");
1065 };
1066 settlement.into_settlement()
1067 })
1068 else {
1069 panic!("first assignment returns its actual rejected delivery");
1070 };
1071 assert_eq!(first_returned.target(), first_target);
1072 let (_, returned_assignment, first_receipt) = first_returned.into_parts();
1073 assert_eq!(returned_assignment.payload().0.as_ptr(), first_payload);
1074 assert_eq!(&*returned_assignment.payload().0, "first");
1075 assert!(matches!(
1076 first_receipt.compare(&first_correlation),
1077 CorrelationMatch::Exact
1078 ));
1079 assert_eq!(
1080 host.observed_payloads,
1081 [second_payload as usize, first_payload as usize]
1082 );
1083 assert!(host.decisions.is_empty());
1084 }
1085
1086 #[test]
1087 fn completion_waits_for_delivery_acceptance() {
1088 let worker = worker();
1089 let mut assignments = AssignmentSequence::new();
1090 let (correlation, execution) = assignments
1091 .assign(&worker, 7)
1092 .unwrap_or_else(|| panic!("test assignment correlation is available"));
1093 let receipt = AssignmentReceipt::issued(&correlation);
1094 let completion = execution.complete(21).into_inner();
1095 let assigned: AssignedJob<_, _, _, MailAddr> = AssignedJob::new(customer(), correlation);
1096
1097 let waiting = assigned
1098 .accept_completion(completion)
1099 .unwrap_or_else(|_| panic!("exact completion is admitted"));
1100 let WorkerCompletionOutcome::AwaitingReceipt(waiting) = waiting else {
1101 panic!("completion alone cannot resolve customer custody")
1102 };
1103 let completed = waiting
1104 .accept_receipt(receipt)
1105 .unwrap_or_else(|_| panic!("exact delivery receipt is admitted"));
1106 let AssignmentReceiptOutcome::JobCompleted {
1107 customer,
1108 result,
1109 stopped,
1110 } = completed
1111 else {
1112 panic!("accepted delivery releases the retained completion")
1113 };
1114 assert_eq!(customer.payload, 7);
1115 assert_eq!(result, 21);
1116 assert_eq!(stopped, None);
1117 }
1118
1119 #[test]
1120 fn stop_waits_for_delivery_acceptance_and_wins_terminal_order() {
1121 let worker = worker();
1122 let mut assignments = AssignmentSequence::new();
1123 let (correlation, _) = assignments
1124 .assign(&worker, 7)
1125 .unwrap_or_else(|| panic!("test assignment correlation is available"));
1126 let receipt = AssignmentReceipt::issued(&correlation);
1127 let stop = ChildStopped::new(worker.creation(), Err(Crash::Failed), Instant::now());
1128 let assigned: AssignedJob<u8, u8, u8, MailAddr> = AssignedJob::new(customer(), correlation);
1129
1130 let waiting = assigned
1131 .accept_worker_exit(stop)
1132 .unwrap_or_else(|_| panic!("exact stop is admitted"));
1133 let WorkerExitOutcome::AwaitingReceipt(waiting) = waiting else {
1134 panic!("stop alone cannot resolve customer custody")
1135 };
1136 let interrupted = waiting
1137 .accept_receipt(receipt)
1138 .unwrap_or_else(|_| panic!("exact delivery receipt is admitted"));
1139 let AssignmentReceiptOutcome::JobInterrupted {
1140 customer,
1141 stopped,
1142 late_completion,
1143 } = interrupted
1144 else {
1145 panic!("accepted delivery releases the retained stop")
1146 };
1147 assert_eq!(customer.payload, 7);
1148 assert_eq!(stopped.child, worker.creation());
1149 match late_completion {
1150 None => {}
1151 Some(_) => panic!("no later completion was observed"),
1152 }
1153 }
1154
1155 #[test]
1156 fn rejected_delivery_returns_the_execution_and_customer_obligation() {
1157 let worker = worker();
1158 let mut assignments = AssignmentSequence::new();
1159 let (correlation, execution) = assignments
1160 .assign(&worker, 7)
1161 .unwrap_or_else(|| panic!("test assignment correlation is available"));
1162 let assigned: AssignedJob<u8, u8, u8, MailAddr> = AssignedJob::new(customer(), correlation);
1163
1164 let returned = assigned
1165 .accept_rejection(execution)
1166 .unwrap_or_else(|_| panic!("exact returned execution reunites authority"));
1167 let AssignmentRejectionOutcome::JobReturned { customer, stopped } = returned else {
1168 panic!("rejected delivery reunites authority and returns customer custody")
1169 };
1170 assert_eq!(customer.payload, 7);
1171 assert_eq!(stopped, None);
1172 }
1173
1174 #[test]
1175 fn foreign_delivery_rejection_preserves_both_assignments() {
1176 let worker = worker();
1177 let mut assignments = AssignmentSequence::new();
1178 let (correlation, _) = assignments
1179 .assign(&worker, 7)
1180 .unwrap_or_else(|| panic!("first assignment correlation is available"));
1181 let receipt = AssignmentReceipt::issued(&correlation);
1182 let (_, foreign) = assignments
1183 .assign(&worker, 11)
1184 .unwrap_or_else(|| panic!("second assignment correlation is available"));
1185 let assigned: AssignedJob<u8, u8, u8, MailAddr> = AssignedJob::new(customer(), correlation);
1186
1187 let (assigned, foreign) = match assigned.accept_rejection(foreign) {
1188 Err(returned) => returned,
1189 Ok(_) => panic!("a foreign assignment cannot reunite authority"),
1190 };
1191 assert_eq!(foreign.payload(), &11);
1192 let accepted = assigned
1193 .accept_receipt(receipt)
1194 .unwrap_or_else(|_| panic!("the original assignment remains current"));
1195 assert!(matches!(
1196 accepted,
1197 AssignmentReceiptOutcome::AwaitingCompletion(_)
1198 ));
1199 }
1200
1201 #[test]
1202 fn returned_authority_after_completion_is_a_contradiction() {
1203 let worker = worker();
1204 let mut assignments = AssignmentSequence::new();
1205 let (correlation, execution) = assignments
1206 .assign(&worker, 7)
1207 .unwrap_or_else(|| panic!("test assignment correlation is available"));
1208 let returned = crate::atomic::Assignment::issued(
1209 7,
1210 CompletionAuthority {
1211 assignment: correlation.assignment,
1212 worker: correlation.worker.clone(),
1213 token: Arc::clone(&correlation.token),
1214 },
1215 );
1216 let completion = execution.complete(21).into_inner();
1217 let assigned: AssignedJob<u8, u8, u8, MailAddr> = AssignedJob::new(customer(), correlation);
1218 let waiting = assigned
1219 .accept_completion(completion)
1220 .unwrap_or_else(|_| panic!("exact completion is admitted"));
1221 let WorkerCompletionOutcome::AwaitingReceipt(waiting) = waiting else {
1222 panic!("completion must remain pending before delivery settlement")
1223 };
1224
1225 let contradiction = waiting
1226 .accept_rejection(returned)
1227 .unwrap_or_else(|_| panic!("the returned authority has exact correlation"));
1228 let AssignmentRejectionOutcome::ConflictingCompletion {
1229 customer,
1230 assignment,
1231 completion,
1232 stopped,
1233 } = contradiction
1234 else {
1235 panic!("returned authority contradicts its retained completion")
1236 };
1237 assert_eq!(customer.payload, 7);
1238 assert_eq!(assignment.payload(), &7);
1239 let (result, _) = completion.into_parts();
1240 assert_eq!(result, 21);
1241 assert_eq!(stopped, None);
1242 }
1243}