1use crate::log::{Event, Logger};
12use crate::message::{Message, Request};
13use anyhow::{Context as _, bail};
14use async_trait::async_trait;
15use serde::{Serialize, de::DeserializeOwned};
16use std::collections::{HashMap, HashSet};
17use std::ops::ControlFlow::{self, Break, Continue};
18use tokio::sync::mpsc::{UnboundedReceiver, UnboundedSender};
19use tokio::sync::oneshot;
20use tokio_util::sync::CancellationToken;
21use uuid::Uuid;
22
23pub type ActorId = String;
25
26pub type Builder<B> =
28 Box<dyn FnOnce(Context<<B as Behavior>::Message, <B as Behavior>::Log>) -> B + Send>;
29
30pub struct ActorInit<B: Behavior> {
33 pub behavior: Builder<B>,
36 pub can_send_to: HashSet<ActorId>,
38 pub can_shut_down: HashSet<ActorId>,
40 pub has_logger: bool,
42}
43
44pub(crate) struct Actor<B: Behavior> {
47 pub(crate) behavior: B,
50 pub(crate) ready: Option<oneshot::Sender<()>>,
53 pub(crate) start: oneshot::Receiver<()>,
55 pub(crate) mailbox: UnboundedReceiver<Envelope<B::Message>>,
57}
58
59impl<B: Behavior> Actor<B> {
60 pub(crate) async fn run(mut self) -> anyhow::Result<()> {
77 let shutdown = self.behavior.context().shutdown.mine.clone();
78 let Some(initialized) = run_unless_stopped(&shutdown, self.behavior.initialize()).await
79 else {
80 return Ok(());
81 };
82 initialized?;
83 if let Some(ready) = self.ready.take() {
85 let _ = ready.send(());
86 }
87 let mut start_consumed = false;
88 loop {
89 let flow = tokio::select! {
96 biased;
97 _ = shutdown.cancelled() => Break(()),
98 signal = &mut self.start, if !start_consumed => {
99 start_consumed = true;
100 self.open(signal).await?
101 }
102 envelope = self.mailbox.recv() => {
103 self.deliver(envelope.expect("mailbox closed")).await?
106 }
107 };
108 if flow.is_break() {
109 break;
110 }
111 }
112 self.behavior.clean_up().await?;
113 Ok(())
114 }
115
116 async fn open(
120 &mut self,
121 signal: Result<(), oneshot::error::RecvError>,
122 ) -> anyhow::Result<ControlFlow<()>> {
123 if signal.is_err() {
124 return Ok(Break(()));
125 }
126 let shutdown = self.behavior.context().shutdown.mine.clone();
127 let Some(opened) = run_unless_stopped(&shutdown, self.behavior.start()).await else {
128 return Ok(Break(()));
129 };
130 opened?;
131 Ok(Continue(()))
132 }
133
134 async fn deliver(&mut self, envelope: Envelope<B::Message>) -> anyhow::Result<ControlFlow<()>> {
140 let shutdown = self.behavior.context().shutdown.mine.clone();
141 match envelope {
142 Envelope::Statement(message) => {
143 let Some(received) =
144 run_unless_stopped(&shutdown, self.behavior.receive(&message)).await
145 else {
146 return Ok(Break(()));
147 };
148 received?;
149 }
150 Envelope::Request(request) => {
151 let Some(answered) =
152 run_unless_stopped(&shutdown, self.behavior.answer(request.message())).await
153 else {
154 return Ok(Break(()));
155 };
156 let _ = request.reply(answered?);
159 }
160 }
161 Ok(Continue(()))
162 }
163}
164
165async fn run_unless_stopped<T>(
168 shutdown: &CancellationToken,
169 step: impl Future<Output = T>,
170) -> Option<T> {
171 tokio::select! {
172 biased;
173 _ = shutdown.cancelled() => None,
174 outcome = step => Some(outcome),
175 }
176}
177
178pub struct Context<M: Message, L = M> {
186 pub id: ActorId,
188 pub episode: Uuid,
190 pub(crate) mailboxes: HashMap<ActorId, UnboundedSender<Envelope<M>>>,
193 pub(crate) shutdown: Shutdown,
195 pub(crate) log: Option<Logger<L>>,
197}
198
199impl<M: Message, L> Context<M, L> {
200 pub fn log(&self, payload: L) {
204 if let Some(log) = &self.log {
205 let _ = log.send(Event::now(self.episode, payload));
206 }
207 }
208
209 pub fn send(&self, message: M, to: HashSet<ActorId>) -> anyhow::Result<()> {
218 let senders = to
219 .iter()
220 .map(|id| self.mailbox_of(id))
221 .collect::<anyhow::Result<Vec<_>>>()?;
222 for sender in senders {
223 let _ = sender.send(Envelope::Statement(message.clone()));
225 }
226 Ok(())
227 }
228
229 pub async fn request(
239 &self,
240 message: M,
241 to: HashSet<ActorId>,
242 ) -> anyhow::Result<HashMap<ActorId, Vec<M>>> {
243 if to.contains(&self.id) {
244 bail!(
245 "{} cannot request from itself: it would wait forever",
246 self.id
247 );
248 }
249 let senders = to
250 .iter()
251 .map(|id| self.mailbox_of(id).map(|sender| (id, sender)))
252 .collect::<anyhow::Result<Vec<_>>>()?;
253 let mut pending = Vec::new();
254 for (id, sender) in senders {
255 let (request, reply) = Request::new(message.clone());
256 if sender.send(Envelope::Request(request)).is_ok() {
257 pending.push((id.clone(), reply));
258 }
259 }
260 let mut replies = HashMap::new();
261 for (id, reply) in pending {
262 if let Ok(reply) = reply.await {
264 replies.insert(id, reply.messages);
265 }
266 }
267 Ok(replies)
268 }
269
270 fn mailbox_of(&self, id: &ActorId) -> anyhow::Result<&UnboundedSender<Envelope<M>>> {
277 self.mailboxes.get(id).with_context(|| {
278 format!(
279 "{} cannot send to {id}: not an actor it may send to",
280 self.id
281 )
282 })
283 }
284
285 pub fn stop(&self, who: &ActorId) -> anyhow::Result<()> {
291 match self.shutdown.others.get(who) {
292 Some(token) => {
293 token.cancel();
294 Ok(())
295 }
296 None => bail!(
297 "{} cannot stop {who}: not an actor it may shut down",
298 self.id
299 ),
300 }
301 }
302
303 pub fn shutdown(&self) {
307 self.shutdown.mine.cancel();
308 }
309}
310
311#[async_trait]
363pub trait Behavior: Send {
364 type Message: Message;
366 type Log: Serialize + DeserializeOwned + Send + 'static;
368 fn context(&self) -> &Context<Self::Message, Self::Log>;
370 async fn initialize(&mut self) -> anyhow::Result<()> {
373 Ok(())
374 }
375 async fn receive(&mut self, _message: &Self::Message) -> anyhow::Result<()> {
377 Ok(())
378 }
379 async fn answer(&mut self, _message: &Self::Message) -> anyhow::Result<Vec<Self::Message>> {
382 Ok(vec![])
383 }
384 async fn start(&mut self) -> anyhow::Result<()> {
387 Ok(())
388 }
389 async fn clean_up(&mut self) -> anyhow::Result<()> {
392 Ok(())
393 }
394}
395
396pub(crate) struct Shutdown {
399 pub(crate) mine: CancellationToken,
401 pub(crate) others: HashMap<ActorId, CancellationToken>,
403}
404
405#[derive(Debug)]
407pub(crate) enum Envelope<M: Message> {
408 Statement(M),
410 Request(Request<M>),
412}
413
414#[cfg(test)]
415mod tests {
416 use super::*;
417 use serde::{Deserialize, Serialize};
418 use std::sync::Arc;
419 use std::time::Duration;
420 use tokio::sync::Semaphore;
421 use tokio::sync::mpsc::unbounded_channel;
422 use tokio::time::{sleep, timeout};
423
424 #[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
425 struct Note(String);
426 impl Message for Note {}
427
428 fn note(text: &str) -> Note {
429 Note(text.to_string())
430 }
431
432 struct Echo {
437 context: Context<Note>,
438 gate: Option<Arc<Semaphore>>,
439 slow_start: bool,
440 broken: bool,
441 }
442 #[async_trait]
443 impl Behavior for Echo {
444 type Message = Note;
445 type Log = Note;
446 fn context(&self) -> &Context<Note> {
447 &self.context
448 }
449 async fn receive(&mut self, _message: &Note) -> anyhow::Result<()> {
450 self.step().await
451 }
452 async fn answer(&mut self, message: &Note) -> anyhow::Result<Vec<Note>> {
453 self.step().await?;
454 Ok(vec![message.clone()])
455 }
456 async fn start(&mut self) -> anyhow::Result<()> {
457 if self.slow_start
458 && let Some(gate) = &self.gate
459 {
460 gate.acquire().await?.forget();
461 }
462 Ok(())
463 }
464 }
465
466 impl Echo {
467 async fn step(&self) -> anyhow::Result<()> {
469 if self.broken {
470 bail!("echo is broken");
471 }
472 if let Some(gate) = &self.gate {
473 gate.acquire().await?.forget();
474 }
475 Ok(())
476 }
477 }
478
479 struct Mute(Context<Note>);
482 #[async_trait]
483 impl Behavior for Mute {
484 type Message = Note;
485 type Log = Note;
486 fn context(&self) -> &Context<Note> {
487 &self.0
488 }
489 }
490
491 struct Tally {
494 context: Context<Note>,
495 seen: usize,
496 }
497 #[async_trait]
498 impl Behavior for Tally {
499 type Message = Note;
500 type Log = Note;
501 fn context(&self) -> &Context<Note> {
502 &self.context
503 }
504 async fn receive(&mut self, _message: &Note) -> anyhow::Result<()> {
505 self.seen += 1;
506 Ok(())
507 }
508 async fn answer(&mut self, _message: &Note) -> anyhow::Result<Vec<Note>> {
509 Ok(vec![note("seen"); self.seen])
510 }
511 }
512
513 struct Rig<B: Behavior = Echo> {
516 actor: Actor<B>,
517 sender: UnboundedSender<Envelope<Note>>,
518 start: oneshot::Sender<()>,
519 stop: CancellationToken,
520 }
521
522 impl<B: Behavior<Message = Note, Log = Note>> Rig<B> {
523 fn context(&self) -> &Context<Note> {
524 self.actor.behavior.context()
525 }
526 }
527
528 fn rig(
530 name: &str,
531 others: HashMap<ActorId, UnboundedSender<Envelope<Note>>>,
532 can_stop: HashMap<ActorId, CancellationToken>,
533 ) -> Rig {
534 rig_with(name, others, can_stop, |context| Echo {
535 context,
536 gate: None,
537 slow_start: false,
538 broken: false,
539 })
540 }
541
542 fn rig_with<B: Behavior<Message = Note, Log = Note>>(
545 name: &str,
546 others: HashMap<ActorId, UnboundedSender<Envelope<Note>>>,
547 can_stop: HashMap<ActorId, CancellationToken>,
548 build: impl FnOnce(Context<Note>) -> B,
549 ) -> Rig<B> {
550 let (sender, mailbox) = unbounded_channel();
551 let (ready, _) = oneshot::channel();
552 let (start, started) = oneshot::channel();
553 let stop = CancellationToken::new();
554 let mut mailboxes = others;
555 mailboxes.insert(id(name), sender.clone());
556 let context = Context {
557 id: id(name),
558 episode: Uuid::new_v4(),
559 mailboxes,
560 shutdown: Shutdown {
561 mine: stop.clone(),
562 others: can_stop,
563 },
564 log: None,
565 };
566 let actor = Actor {
567 behavior: build(context),
568 ready: Some(ready),
569 start: started,
570 mailbox,
571 };
572 Rig {
573 actor,
574 sender,
575 start,
576 stop,
577 }
578 }
579
580 fn id(name: &str) -> ActorId {
581 name.to_string()
582 }
583
584 #[tokio::test]
585 async fn send_reaches_each_actor_it_is_addressed_to() {
586 let (bob, mut bob_mailbox) = unbounded_channel();
587 let (cat, mut cat_mailbox) = unbounded_channel();
588 let (dan, mut dan_mailbox) = unbounded_channel();
589 let ann = rig(
590 "ann",
591 HashMap::from([(id("bob"), bob), (id("cat"), cat), (id("dan"), dan)]),
592 HashMap::new(),
593 );
594
595 ann.context()
596 .send(note("hello"), HashSet::from([id("bob"), id("cat")]))
597 .unwrap();
598
599 for mailbox in [&mut bob_mailbox, &mut cat_mailbox] {
600 let heard = mailbox.recv().await.unwrap();
601 assert!(
602 matches!(&heard, Envelope::Statement(Note(text)) if text == "hello"),
603 "{heard:?}"
604 );
605 }
606 assert!(dan_mailbox.try_recv().is_err(), "nothing should reach dan");
607 }
608
609 #[tokio::test]
610 async fn send_fails_without_sending_if_a_recipient_may_not_be_addressed() {
611 let (bob, mut bob_mailbox) = unbounded_channel();
612 let ann = rig("ann", HashMap::from([(id("bob"), bob)]), HashMap::new());
613
614 let error = ann
615 .context()
616 .send(note("psst"), HashSet::from([id("bob"), id("zed")]))
617 .unwrap_err();
618
619 assert!(error.to_string().contains("zed"), "{error}");
620 assert!(bob_mailbox.try_recv().is_err(), "nothing should reach bob");
621 }
622
623 #[tokio::test]
624 async fn send_to_itself_lands_in_the_actors_own_mailbox() {
625 let mut ann = rig("ann", HashMap::new(), HashMap::new());
626
627 ann.context()
628 .send(note("remember this"), HashSet::from([id("ann")]))
629 .unwrap();
630
631 let heard = ann.actor.mailbox.recv().await.unwrap();
632 assert!(
633 matches!(&heard, Envelope::Statement(Note(text)) if text == "remember this"),
634 "{heard:?}"
635 );
636 }
637
638 #[tokio::test]
639 async fn request_collects_a_reply_from_every_recipient() {
640 let bob = rig("bob", HashMap::new(), HashMap::new());
641 let cat = rig("cat", HashMap::new(), HashMap::new());
642 let ann = rig(
643 "ann",
644 HashMap::from([
645 (id("bob"), bob.sender.clone()),
646 (id("cat"), cat.sender.clone()),
647 ]),
648 HashMap::new(),
649 );
650 let bob_running = tokio::spawn(bob.actor.run());
651 let cat_running = tokio::spawn(cat.actor.run());
652 bob.start.send(()).unwrap();
653 cat.start.send(()).unwrap();
654
655 let replies = ann
656 .context()
657 .request(note("who's there?"), HashSet::from([id("bob"), id("cat")]))
658 .await
659 .unwrap();
660
661 let expected = HashMap::from([
662 (id("bob"), vec![note("who's there?")]),
663 (id("cat"), vec![note("who's there?")]),
664 ]);
665 assert_eq!(replies, expected);
666 bob.stop.cancel();
667 cat.stop.cancel();
668 bob_running.await.unwrap().unwrap();
669 cat_running.await.unwrap().unwrap();
670 }
671
672 #[tokio::test]
673 async fn a_behavior_keeps_its_state_between_steps() {
674 let bob = rig_with("bob", HashMap::new(), HashMap::new(), |context| Tally {
675 context,
676 seen: 0,
677 });
678 let ann = rig(
679 "ann",
680 HashMap::from([(id("bob"), bob.sender.clone())]),
681 HashMap::new(),
682 );
683 let running = tokio::spawn(bob.actor.run());
684 bob.start.send(()).unwrap();
685 for _ in 0..2 {
686 ann.context()
687 .send(note("one more"), HashSet::from([id("bob")]))
688 .unwrap();
689 }
690
691 let replies = ann
692 .context()
693 .request(note("how many?"), HashSet::from([id("bob")]))
694 .await
695 .unwrap();
696
697 let expected = HashMap::from([(id("bob"), vec![note("seen"), note("seen")])]);
698 assert_eq!(replies, expected);
699 bob.stop.cancel();
700 running.await.unwrap().unwrap();
701 }
702
703 #[tokio::test]
704 async fn the_default_receive_ignores_the_statement() {
705 let bob = rig_with("bob", HashMap::new(), HashMap::new(), Mute);
706 let ann = rig(
707 "ann",
708 HashMap::from([(id("bob"), bob.sender.clone())]),
709 HashMap::new(),
710 );
711 let running = tokio::spawn(bob.actor.run());
712 bob.start.send(()).unwrap();
713
714 ann.context()
715 .send(note("whatever"), HashSet::from([id("bob")]))
716 .unwrap();
717
718 let replies = ann
720 .context()
721 .request(note("still there?"), HashSet::from([id("bob")]))
722 .await
723 .unwrap();
724 assert_eq!(replies, HashMap::from([(id("bob"), vec![])]));
725 bob.stop.cancel();
726 running.await.unwrap().unwrap();
727 }
728
729 #[tokio::test]
730 async fn the_default_answer_is_nothing() {
731 let bob = rig_with("bob", HashMap::new(), HashMap::new(), Mute);
732 let ann = rig(
733 "ann",
734 HashMap::from([(id("bob"), bob.sender.clone())]),
735 HashMap::new(),
736 );
737 let running = tokio::spawn(bob.actor.run());
738 bob.start.send(()).unwrap();
739
740 let replies = ann
741 .context()
742 .request(note("anything?"), HashSet::from([id("bob")]))
743 .await
744 .unwrap();
745
746 assert_eq!(replies, HashMap::from([(id("bob"), vec![])]));
747 bob.stop.cancel();
748 running.await.unwrap().unwrap();
749 }
750
751 #[tokio::test]
752 async fn request_fails_without_sending_if_a_recipient_may_not_be_addressed() {
753 let (bob, mut bob_mailbox) = unbounded_channel();
754 let ann = rig("ann", HashMap::from([(id("bob"), bob)]), HashMap::new());
755
756 let error = ann
757 .context()
758 .request(note("psst"), HashSet::from([id("bob"), id("zed")]))
759 .await
760 .unwrap_err();
761
762 assert!(error.to_string().contains("zed"), "{error}");
763 assert!(bob_mailbox.try_recv().is_err(), "nothing should reach bob");
764 }
765
766 #[tokio::test]
767 async fn request_refuses_to_ask_the_actor_itself() {
768 let ann = rig("ann", HashMap::new(), HashMap::new());
769
770 let error = ann
771 .context()
772 .request(note("hello me"), HashSet::from([id("ann")]))
773 .await
774 .unwrap_err();
775
776 assert!(error.to_string().contains("itself"), "{error}");
777 }
778
779 #[tokio::test]
780 async fn request_leaves_out_a_recipient_that_has_stopped() {
781 let (bob, bob_mailbox) = unbounded_channel();
782 let ann = rig("ann", HashMap::from([(id("bob"), bob)]), HashMap::new());
783 drop(bob_mailbox);
784
785 let replies = ann
786 .context()
787 .request(note("anyone?"), HashSet::from([id("bob")]))
788 .await
789 .unwrap();
790
791 assert!(replies.is_empty(), "{replies:?}");
792 }
793
794 #[tokio::test]
795 async fn stop_cancels_an_actor_it_may_shut_down_and_refuses_others() {
796 let bob = rig("bob", HashMap::new(), HashMap::new());
797 let ann = rig(
798 "ann",
799 HashMap::new(),
800 HashMap::from([(id("bob"), bob.stop.clone())]),
801 );
802
803 ann.context().stop(&id("bob")).unwrap();
804 assert!(bob.stop.is_cancelled());
805
806 let error = ann.context().stop(&id("zed")).unwrap_err();
807 assert!(error.to_string().contains("zed"), "{error}");
808 }
809
810 #[tokio::test]
811 async fn log_sends_a_stamped_event_down_the_logger_if_there_is_one() {
812 let (logger, mut events) = unbounded_channel();
813 let mut ann = rig("ann", HashMap::new(), HashMap::new());
814 ann.actor.behavior.context.log = Some(logger);
815
816 ann.context().log(note("for the record"));
817
818 let event = events.recv().await.unwrap();
819 assert_eq!(event.payload, note("for the record"));
820 }
821
822 #[tokio::test]
823 async fn log_does_nothing_without_a_logger() {
824 let ann = rig("ann", HashMap::new(), HashMap::new());
825
826 ann.context().log(note("into the void"));
827 }
828
829 #[tokio::test]
830 async fn shutdown_cancels_the_actors_own_token() {
831 let ann = rig("ann", HashMap::new(), HashMap::new());
832
833 ann.context().shutdown();
834
835 assert!(ann.stop.is_cancelled());
836 }
837
838 fn gated_rig(name: &str) -> (Rig, Arc<Semaphore>) {
840 let gate = Arc::new(Semaphore::new(0));
841 let mut rig = rig(name, HashMap::new(), HashMap::new());
842 rig.actor.behavior.gate = Some(Arc::clone(&gate));
843 (rig, gate)
844 }
845
846 #[tokio::test(start_paused = true)]
847 async fn a_kill_stops_an_actor_in_the_middle_of_a_step() {
848 let (bob, _gate) = gated_rig("bob");
849 let running = tokio::spawn(bob.actor.run());
850 bob.start.send(()).unwrap();
851 bob.sender
852 .send(Envelope::Statement(note("take your time")))
853 .unwrap();
854 sleep(Duration::from_secs(1)).await;
856
857 bob.stop.cancel();
858
859 let stopped = timeout(Duration::from_secs(5), running).await;
860 stopped.expect("bob should stop").unwrap().unwrap();
861 }
862
863 #[tokio::test(start_paused = true)]
864 async fn a_kill_stops_an_actor_in_the_middle_of_starting() {
865 let (mut bob, _gate) = gated_rig("bob");
866 bob.actor.behavior.slow_start = true;
867 let running = tokio::spawn(bob.actor.run());
868 bob.start.send(()).unwrap();
869 sleep(Duration::from_secs(1)).await;
871
872 bob.stop.cancel();
873
874 let stopped = timeout(Duration::from_secs(5), running).await;
875 stopped.expect("bob should stop").unwrap().unwrap();
876 }
877
878 #[tokio::test(start_paused = true)]
879 async fn a_behavior_with_no_opening_move_starts_and_waits() {
880 let bob = rig_with("bob", HashMap::new(), HashMap::new(), Mute);
881 let running = tokio::spawn(bob.actor.run());
882 bob.start.send(()).unwrap();
883 sleep(Duration::from_secs(1)).await;
885
886 bob.stop.cancel();
887
888 running.await.unwrap().unwrap();
889 }
890
891 #[tokio::test]
892 async fn an_actor_stops_when_the_episode_is_gone_before_it_starts() {
893 let bob = rig("bob", HashMap::new(), HashMap::new());
894 let running = tokio::spawn(bob.actor.run());
895
896 drop(bob.start);
897
898 running.await.unwrap().unwrap();
899 }
900
901 #[tokio::test]
902 async fn a_behavior_that_fails_to_receive_a_statement_fails_the_actor() {
903 let mut bob = rig("bob", HashMap::new(), HashMap::new());
904 bob.actor.behavior.broken = true;
905 let running = tokio::spawn(bob.actor.run());
906 bob.start.send(()).unwrap();
907
908 bob.sender.send(Envelope::Statement(note("hello"))).unwrap();
909
910 let error = running.await.unwrap().unwrap_err();
911 assert!(error.to_string().contains("broken"), "{error}");
912 }
913
914 #[tokio::test]
915 async fn a_behavior_that_fails_to_answer_a_request_fails_the_actor() {
916 let mut bob = rig("bob", HashMap::new(), HashMap::new());
917 bob.actor.behavior.broken = true;
918 let running = tokio::spawn(bob.actor.run());
919 bob.start.send(()).unwrap();
920 let (request, _reply) = Request::new(note("well?"));
921
922 bob.sender.send(Envelope::Request(request)).unwrap();
923
924 let error = running.await.unwrap().unwrap_err();
925 assert!(error.to_string().contains("broken"), "{error}");
926 }
927
928 #[tokio::test(start_paused = true)]
929 async fn a_recipient_killed_mid_step_is_left_out_of_the_replies() {
930 let (bob, _gate) = gated_rig("bob");
931 let ann = rig(
932 "ann",
933 HashMap::from([(id("bob"), bob.sender.clone())]),
934 HashMap::new(),
935 );
936 let bob_running = tokio::spawn(bob.actor.run());
937 bob.start.send(()).unwrap();
938 let asking = tokio::spawn(async move {
939 ann.context()
940 .request(note("well?"), HashSet::from([id("bob")]))
941 .await
942 });
943 sleep(Duration::from_secs(1)).await;
944
945 bob.stop.cancel();
946
947 let replies = asking.await.unwrap().unwrap();
948 assert!(replies.is_empty(), "{replies:?}");
949 bob_running.await.unwrap().unwrap();
950 }
951
952 #[tokio::test(start_paused = true)]
953 async fn a_responder_whose_asker_gave_up_carries_on() {
954 let (bob, gate) = gated_rig("bob");
955 let ann = rig(
956 "ann",
957 HashMap::from([(id("bob"), bob.sender.clone())]),
958 HashMap::new(),
959 );
960 let bob_running = tokio::spawn(bob.actor.run());
961 bob.start.send(()).unwrap();
962
963 let asking = ann
966 .context()
967 .request(note("well?"), HashSet::from([id("bob")]));
968 let gave_up = tokio::time::timeout(Duration::from_secs(1), asking).await;
969 assert!(gave_up.is_err(), "{gave_up:?}");
970 gate.add_permits(1);
971 sleep(Duration::from_secs(1)).await;
972
973 bob.stop.cancel();
974 bob_running.await.unwrap().unwrap();
975 }
976}