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